# (c) cavaliba.com - tests / helper_pipeline

import yaml

from app_data.loader import load_broker


def run_pipeline(pipeline_name, schema_names, aaa=None, dryrun=False):
    from app_data.models import DataTask
    from app_data.tasks import submit_pipeline

    if aaa is None:
        aaa = {"perms": ["p_data_admin", "p_pipeline_run"]}

    handle, err = submit_pipeline(
        pipeline_name=pipeline_name,
        schema_names=schema_names,
        dryrun=dryrun,
        aaa=aaa,
        owner_type="test",
        owner_id="system",
        sync=True,
    )
    if not handle:
        return 0, 0, [err or "submit_pipeline failed"]

    dt = DataTask.objects.get(handle=handle)
    output = dt.output or {}

    total_ok = output.get("total_ok", 0)
    total_discarded = output.get("total_discarded", 0)
    errors = []
    for entry in output.get("results", []):
        errors.extend(entry.get("errors", []))

    return total_ok, total_discarded, errors


def add_pipeline_noop():

    datalist = yaml.safe_load("""
        - classname: _pipeline
          keyname: pipeline_noop
          displayname: pipeline_noop
          description: pipeline_noop
          is_enabled: True
          content: |
                csv_delimiter: '|'
                classname: test1
                keyfield: keyname
                tasks:
                - field_noop
        """)
    aaa = {"perms": ["p_pipeline_create"]}
    load_broker(datalist=datalist, aaa=aaa)
