from prefect import flow, serve
from prefect.events import emit_event
from prefect.events.schemas import DeploymentTrigger


@flow
def upstream_foo():
    emit_event(
        event="foo",
        resource={
            "prefect.resource.id": str(id(upstream_foo)),
        },
    )


@flow
def upstream_bar():
    emit_event(
        event="bar",
        resource={
            "prefect.resource.id": str(id(upstream_bar)),
        },
    )


@flow
def upstream_baz():
    emit_event(
        event="baz",
        resource={
            "prefect.resource.id": str(id(upstream_baz)),
        },
    )


@flow(log_prints=True)
def downstream():
    print("i was triggered!")


trigger_on_upstream = DeploymentTrigger(
    name="Wait for N upstream deployments",
    enabled=True,
    match_related={
        "prefect.resource.name": "upstream-*",
        "prefect.resource.role": "flow",
    },
    after={"foo", "bar", "baz"},  # the events to wait for
    within=60,
)


if __name__ == "__main__":
    deployments = [
        (
            f.to_deployment(f.name)
            if f.name != "downstream"
            else f.to_deployment(f.name, triggers=[trigger_on_upstream])
        )
        for f in [upstream_foo, upstream_bar, upstream_baz, downstream]
    ]

    serve(*deployments)
	