Get a run's outcome to the people who need it, and see what production is doing without reading logs.

Notifications and observability

Notifications tell a person that a specific run finished or failed. Observability lets someone investigate what the system is doing across runs. Keep them separate: if every run posts to a channel, people mute the channel and miss the message that mattered.

Notification Observability
Answers “This run finished, or failed” “What is the system doing”
Audience Whoever must act Whoever is investigating
Destination Slack, email, or a webhook Traces and metrics in your own backend
Volume Low High, and queried rather than read

Notify on a terminal state

Start with the built-in notifications, which need no plugin. Run with notifications sends a message when a run reaches a terminal state, through email, Slack, Teams, or a generic webhook. Trigger notifications does the same for triggered runs.

They are declarative, need no code, and fire even if the task itself crashed.

Post from inside a run

When the message needs results the run computed, such as which model won or which rows failed, post from inside the run with the Slack integration:

slack_tasks.py
env = flyte.TaskEnvironment(
    name="slack-bot",
    image=flyte.Image.from_debian_base().with_pip_packages("flyteplugins-slack"),
    # The bot token (xoxb-...) from OAuth & Permissions. Posting requires the
    # `chat:write` scope, and the bot must be invited to the channel.
    secrets=[flyte.Secret(key="SLACK_BOT_TOKEN", as_env_var="SLACK_BOT_TOKEN")],
    resources=flyte.Resources(cpu=1, memory="512Mi"),
)

@env.task
async def answer(channel: str, text: str, thread_ts: str) -> str:
    """Post a threaded reply, then edit it when the work finishes.

    `post` returns the message's `ts`, which `update` uses to edit it.
    """
    ts = await notify.post(channel, f"Working on: {text}", thread_ts=thread_ts)
    await notify.update(channel, ts, f"Done: {text}")
    return ts

post returns the message’s ts, which serves as both the thread anchor and the ID that update edits. To report progress, post once when work starts and edit the message as it proceeds, instead of posting a new message at each step.

To limit which tasks hold the bot token, deploy the ready-made task environment in notify once, with flyte.deploy(notify.env). Only that environment mounts the token. Other runs post by calling its send task, without mounting the secret.

Reply to the person who asked

When a person started the run, with a slash command, a button, or a mention, reply where they are:

slack_webhooks.py
@app_env.on_event(events.AppMention.ANY)
async def on_mention(event: WebhookEvent) -> dict:
    """Launch a run for each @-mention.

    The dedupe key identifies one message. For one run per thread, build a
    key from `thread_ts` and pass it to `run_once` instead.
    """
    import flyte.remote as remote

    task = remote.Task.get(name="slack-bot.answer", auto_version="latest")
    slack_event = event.payload["event"]
    result = await run_once.aio(
        task,
        key=event.dedupe_key(),
        channel=event.scope,
        text=event.title or "",
        thread_ts=slack_event.get("thread_ts") or slack_event["ts"],
    )
    return {"run": result.run.name, "created": result.created}

@app_env.on_event(events.Command, action="/deploy")
async def on_deploy_command(event: WebhookEvent) -> dict:
    """Acknowledge the /deploy slash command.

    `respond` posts to the command's `response_url` and needs no bot token.
    Slack expects a reply within three seconds, so acknowledge here and do
    longer work in a launched run.
    """
    await notify.respond(event.payload["response_url"], "Deploy queued.")
    return {"ok": True}

respond needs no token. It posts to the response_url that every interaction and slash command carries, which Slack accepts for 30 minutes and up to five times.

Acknowledge from the handler and let the run post the answer. Slack shows a synchronous reply only if it arrives within 3 seconds.

Observe across runs

For a single failed run, the run itself is usually enough: it records inputs, outputs, logs, and task reports.

For questions that span runs and services, export telemetry. OpenTelemetry sends task and agent telemetry to any OTLP backend, and Grafana Agent sends it to Grafana Cloud. Flyte spans then join your application’s traces, so a request that calls into a Flyte run stays one trace. This is most useful for agents, where the question is often about the whole trajectory, such as why an answer took nine turns.

Track what produced an output

Every run records its inputs, outputs, code version, and image. Deploy with --version <commit-sha>, as described in CI/CD deployments, to tie each run to a commit.

To trace an output to the data and model behind it, register datasets and models as artifacts. Artifacts record lineage across runs: which model version came from which training run, on which dataset version.

MLflow and Weights & Biases keep metric history, so you can tell whether a score is unusual compared with previous weeks.

Investigate a failure

Work through these in order, cheapest first:

  1. Read the failed run: its inputs, outputs, logs, and report.
  2. Check the inputs. A Pandera report shows whether the data changed. If you have no data contract yet, add one.
  3. Compare against history to see whether the score is unusual.
  4. Re-run it. If a retry succeeds, the failure was transient.
  5. Open the trace when the question spans services.

Two things set up in advance make this faster: typed task signatures, so a malformed input fails at the task boundary, and error handling that separates recoverable failures from ones that should stop the run.

File the result in your tracker

A run that finds something to act on can file it:

Make these tasks idempotent, so a retried run doesn’t file a duplicate ticket. Launch with run_once, and have the task check whether the ticket exists before creating it. The ClickUp example shows this.

See also