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:
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:
@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.
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:
- Read the failed run: its inputs, outputs, logs, and report.
- Check the inputs. A Pandera report shows whether the data changed. If you have no data contract yet, add one.
- Compare against history to see whether the score is unusual.
- Re-run it. If a retry succeeds, the failure was transient.
- 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:
- Open or update a ticket in Linear, Jira, or ClickUp.
- Comment on a pull request or publish a check run in GitHub.
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
- Event-driven automation for what starts a run.
- Run with notifications for declarative notifications.
- Slack integration for messages, approvals, and the receiver.