# -*- coding: utf-8 -*-
# vim: sw=4:ts=4:expandtab
# Pipe {{ pipe_name }} generated by riko

from riko import Context
from riko.bado import run
{%- if use_collection %}
from riko.collections import AsyncCollection
{%- endif %}
from riko.modules._subpipe import mark_subpipe
{%- for module in uniq_modules %}
{%- if module.is_sub_pipe %}
from riko.pypipelines.{{ module.name }} import {{ module.pipe_name }}
{%- elif module.name != "output" %}
from riko.modules.{{ module.name }} import async_pipe as {{ module.alias }}
{%- endif %}
{%- endfor %}
{%- if raw_confs %}
from riko.types.modules import (
{%- for raw_conf in raw_confs %}
    {{ raw_conf }}{% if not loop.last %},{% endif %}
{%- endfor %}
)
{%- endif %}


async def {{ pipe_name }}(item=None, context: Context | None = None, **_):
    if context and context.describe_input:
        {{ last_module }} = {{ inputs }}
    elif context and context.describe_dependencies:
        {{ last_module }} = {{ dependencies }}
    else:
{%- for module in modules %}
    {%- if module.splits %}
        splits = {{ module.expr }}
    {%- for i in range(module.splits) %}
        {{ module.id }}_{{ i }} = next(splits)
    {%- endfor %}
    {%- else %}
        {{ module.id }} = {% if module.name == "output" %}{{ module.expr }}{% else %}await {{ module.expr }}{% endif %}
    {%- endif %}
{%- endfor %}

    return {{ last_module }}


mark_subpipe({{ pipe_name }}, subtype="{{ subtype }}")


async def _main():
    for i in await {{ pipe_name }}():
        print(i)


if __name__ == "__main__":
    run(_main)
