diff --git a/v2/integrations/flyte-plugins/clickup/clickup_tasks.py b/v2/integrations/flyte-plugins/clickup/clickup_tasks.py index 15bf0588..31268e3b 100644 --- a/v2/integrations/flyte-plugins/clickup/clickup_tasks.py +++ b/v2/integrations/flyte-plugins/clickup/clickup_tasks.py @@ -8,13 +8,11 @@ # main = "replay_sample_delivery" # params = "" # /// -"""The task a ClickUp webhook launches, over ClickUp's REST API. +"""Tasks launched by the ClickUp webhook receiver. -ClickUp ships no official Python SDK, and its API is a handful of REST calls — -so `httpx` directly beats a thin third-party wrapper, and the plugin's job stops -at the webhook. +`close_ticket` sets a task's status through ClickUp's REST API, using `httpx`. -`replay_sample_delivery` is the entrypoint, and needs no ClickUp workspace: +`replay_sample_delivery` runs without a ClickUp workspace: flyte run --local clickup_tasks.py replay_sample_delivery """ @@ -34,12 +32,11 @@ @env.task async def close_ticket(task_id: str) -> str: - """Read a ticket, then close it only if it is not closed already. + """Close a ticket unless it's already complete. - The pre-check is the point. ClickUp accepts a redundant status write, so - without it a redelivered webhook would produce a second, misleading - audit-log entry on the ticket — `run_once` keeps duplicate *runs* away, but - an idempotent task is what keeps a re-run from lying. + `run_once` prevents duplicate runs, not duplicate writes within a run. + ClickUp accepts a redundant status write and logs it, so the task checks + the current status first. """ import os @@ -57,8 +54,8 @@ async def close_ticket(task_id: str) -> str: # {{docs-fragment replay}} -# Its own environment, with no secrets: the replay needs none, and a task -# environment that declares a secret cannot start until that secret exists. +# A separate environment with no secrets, so the replay runs before any +# secret is created. replay_env = flyte.TaskEnvironment( name="clickup-replay", image=flyte.Image.from_debian_base().with_pip_packages("flyteplugins-clickup"), @@ -68,7 +65,7 @@ async def close_ticket(task_id: str) -> str: @replay_env.task async def replay_sample_delivery() -> dict[str, str]: - """Verify and parse the real delivery the plugin ships. No workspace needed.""" + """Verify and parse the sample delivery bundled with the plugin.""" import flyteplugins.clickup as plugin secret = "a-test-signing-secret" @@ -78,25 +75,23 @@ async def replay_sample_delivery() -> dict[str, str]: assert plugin.verify(body, headers, secret), "a correctly signed delivery must verify" assert not plugin.verify(body, headers, "wrong-secret"), "a bad signature must not" - # The wire contract, which the round trip above cannot check: `verify` and - # `SAMPLE_DELIVERY` agree with each other whatever the header is called, so - # a wrong name passes conformance and then rejects every real delivery. - # ClickUp signs with `X-Signature` -- not `X-Clickup-Signature`, the - # name it looks like it should have and the name that shipped broken. + # Check the header name too. The sample's headers come from the plugin, so + # the round trip above passes whatever the header is called. ClickUp sends + # `X-Signature`, not `X-Clickup-Signature`. assert list(headers) == ["X-Signature"], f"unexpected signature header: {list(headers)}" assert not plugin.verify(body, {"X-Clickup-Signature": headers["X-Signature"]}, secret), ( - "the old, wrong header name must not verify" + "X-Clickup-Signature must not verify" ) event = plugin.parse(headers, body) return { - # `taskCreated` — one string, because ClickUp sends no separate action. + # `taskCreated`: ClickUp sends a single event name with no action. "qualified_type": event.qualified_type, "action_is_none": str(event.action is None), "scope": event.scope or "", "title": event.title or "", "dedupe_key": event.dedupe_key(), - # The header a real delivery carries the signature in. + # The header that carries the signature. "signature_header": next(iter(headers)), } # {{/docs-fragment replay}} diff --git a/v2/integrations/flyte-plugins/clickup/clickup_webhooks.py b/v2/integrations/flyte-plugins/clickup/clickup_webhooks.py index 0d096a77..ab959f9a 100644 --- a/v2/integrations/flyte-plugins/clickup/clickup_webhooks.py +++ b/v2/integrations/flyte-plugins/clickup/clickup_webhooks.py @@ -5,9 +5,9 @@ # "flyteplugins-clickup[app]>=2.10.7", # ] # /// -"""The ClickUp webhook receiver. +"""ClickUp webhook receiver. -Deploy it, then paste the payload URL the dashboard shows into +Deploy the app, then enter the payload URL from its dashboard in ClickUp under Space Settings -> Integrations -> Webhooks: python clickup_webhooks.py @@ -20,11 +20,10 @@ from flyte.extras.webhooks import WebhookAppEnvironment, WebhookEvent, run_once from flyteplugins.clickup import ClickUpProvider, events -# CLICKUP_WEBHOOK_SECRET is mounted from the provider's `default_secret_env`. +# CLICKUP_WEBHOOK_SECRET is mounted automatically. # -# `scopes` matches the ClickUp list id. The provider reads it from the top level -# on list-scoped events and from the nested task on task-scoped ones, so one -# allowlist attributes both. +# `scopes` lists ClickUp list IDs. The provider reads the list ID from both +# list events and task events. app_env = WebhookAppEnvironment( name="clickup-webhooks", providers=[ClickUpProvider()], @@ -38,10 +37,9 @@ # {{docs-fragment handler}} @app_env.on_event(events.Task.STATUS_UPDATED) async def on_status_updated(event: WebhookEvent) -> dict: - """Launch a run when a ticket changes status, once per change. + """Launch a run once per status change. - ClickUp does not split type and action: the event name is one camelCase - string, so `qualified_type` is `taskStatusUpdated` and `action` is None. + `qualified_type` is `taskStatusUpdated`, and `action` is None. """ import flyte.remote as remote diff --git a/v2/integrations/flyte-plugins/github/github_tasks.py b/v2/integrations/flyte-plugins/github/github_tasks.py index 60ff1f93..4e4dfbb4 100644 --- a/v2/integrations/flyte-plugins/github/github_tasks.py +++ b/v2/integrations/flyte-plugins/github/github_tasks.py @@ -8,22 +8,14 @@ # main = "replay_sample_delivery" # params = "" # /// -"""The tasks a GitHub webhook launches — plus the one Flyte beats PyGithub at. +"""Tasks launched by the GitHub webhook receiver. -Three things live here: +- `triage_pr` labels a pull request by size, calling PyGithub directly. +- `gated_merge` uses `review_pr` to pause the run until a person reviews the + pull request in the Flyte UI. +- `clone_at_head` mints a short-lived GitHub App installation token. -- `triage_pr` calls PyGithub directly. There is deliberately no Flyte wrapper - around the GitHub API: PyGithub is maintained by people who get deprecation - notices first, and a task is just a function, so a wrapper would only add a - surface to keep in sync with someone else's release calendar. -- `gated_merge` uses `review_pr`, which *is* the plugin's job: it parks the run - on a `flyte.new_condition` until a human answers. The condition is the part - only Flyte can do. -- `clone_at_head` mints a short-lived GitHub App token instead of holding a - personal access token. - -`replay_sample_delivery` is the entrypoint, and runs with no GitHub account, no -webhook, and no credentials: +`replay_sample_delivery` runs without a GitHub account, webhook, or credentials: flyte run --local github_tasks.py replay_sample_delivery """ @@ -34,8 +26,8 @@ env = flyte.TaskEnvironment( name="github-triage", image=flyte.Image.from_debian_base().with_pip_packages("PyGithub"), - # The PR-reading token. Separate from the webhook signing secret, which - # belongs to the receiver app and never reaches a task. + # Token for reading the pull request. The webhook secret is mounted on + # the receiver app, not here. secrets=[flyte.Secret(key="github-token", as_env_var="GITHUB_TOKEN")], resources=flyte.Resources(cpu=1, memory="512Mi"), ) @@ -43,7 +35,7 @@ @env.task async def triage_pr(repo: str, number: int) -> str: - """Label a pull request by size. Plain PyGithub, called from a task.""" + """Label a pull request by size.""" import os from github import Auth, Github @@ -68,12 +60,11 @@ async def triage_pr(repo: str, number: int) -> str: @review_env.task async def gated_merge(repo: str, number: int) -> str: - """Park the run until a human answers, then branch on a typed decision. + """Wait for a review decision in the Flyte UI, then act on it. - `review_pr` collects the pull request's metadata, raises a condition - carrying it as JSON, and waits. The reviewer answers in the Flyte UI. The run - survives restarts while it waits, because the condition is durable state on - the backend rather than a held-open process. + `review_pr` creates a condition carrying the pull request's metadata and + waits for a reviewer to answer it. The condition is stored on the backend, + so the run survives restarts while it waits. """ from flyteplugins.github import review_pr @@ -89,8 +80,8 @@ async def gated_merge(repo: str, number: int) -> str: agent_env = flyte.TaskEnvironment( name="github-agent", image=flyte.Image.from_debian_base().with_pip_packages("flyteplugins-github[auth]"), - # A kebab-case secret key upper-cases into the environment variable - # `mint_installation_token` reads, so no `as_env_var=` is needed. + # Each key maps to the environment variable `mint_installation_token` + # reads (github-app-id -> GITHUB_APP_ID), so `as_env_var=` isn't needed. secrets=[ flyte.Secret(key="github-app-id"), flyte.Secret(key="github-app-installation-id"), @@ -102,16 +93,16 @@ async def gated_merge(repo: str, number: int) -> str: @agent_env.task async def clone_at_head(repo: str) -> str: - """Mint a one-hour App token per operation instead of storing a PAT. + """Build an authenticated clone URL with a GitHub App installation token. - A token this short-lived is plenty for a clone or a `gh pr create`, and - useless to anyone who later finds it in a log. + The token expires after one hour. """ import asyncio from flyteplugins.github import clone_url, mint_installation_token - # Synchronous, one HTTPS round trip — keep it off the event loop. + # mint_installation_token is synchronous; run it off the event loop. + # It returns None if the App credentials are missing. token = await asyncio.to_thread(mint_installation_token) url = clone_url(repo, token) # Never return or log the token itself. @@ -120,8 +111,8 @@ async def clone_at_head(repo: str) -> str: # {{docs-fragment replay}} -# Its own environment, with no secrets: the replay needs none, and a task -# environment that declares a secret cannot start until that secret exists. +# A separate environment with no secrets, so the replay runs before any +# secret is created. replay_env = flyte.TaskEnvironment( name="github-replay", image=flyte.Image.from_debian_base().with_pip_packages("flyteplugins-github"), @@ -131,12 +122,10 @@ async def clone_at_head(repo: str) -> str: @replay_env.task async def replay_sample_delivery() -> dict[str, str]: - """Verify and parse the real delivery the plugin ships. No account needed. + """Verify and parse the sample delivery bundled with the plugin. - Every provider plugin exports a `SAMPLE_DELIVERY`: a trimmed but real - payload, plus a function that signs it. It is what the plugin's own - conformance test replays, which is how `verify` and `parse` are checked - against something GitHub actually sent rather than against each other. + `SAMPLE_DELIVERY` is a recorded GitHub payload and a function that signs + it. The body is real; the signature header is generated by the plugin. """ import flyteplugins.github as plugin @@ -149,13 +138,13 @@ async def replay_sample_delivery() -> dict[str, str]: event = plugin.parse(headers, body) return { - # What `on_event` matches on: `events.PullRequest.OPENED` is this string. + # Matched by `on_event`. `events.PullRequest.OPENED` equals this string. "qualified_type": event.qualified_type, - # What `scopes` filters on. + # Matched against `scopes`. "scope": event.scope or "", "title": event.title or "", "actor": event.actor or "", - # What `run_once` dedupes on. + # The key `run_once` deduplicates on. "dedupe_key": event.dedupe_key(), } # {{/docs-fragment replay}} diff --git a/v2/integrations/flyte-plugins/github/github_webhooks.py b/v2/integrations/flyte-plugins/github/github_webhooks.py index 3ebfb150..8b59c2ef 100644 --- a/v2/integrations/flyte-plugins/github/github_webhooks.py +++ b/v2/integrations/flyte-plugins/github/github_webhooks.py @@ -5,14 +5,12 @@ # "flyteplugins-github[app]>=2.10.6", # ] # /// -"""The GitHub webhook receiver: verify a delivery, launch a run, return. +"""GitHub webhook receiver. -This is an app, not a task. It does no work itself — it authenticates GitHub's -HMAC, normalizes the delivery into a `WebhookEvent`, and launches the tasks in -`github_tasks.py`. Keeping the two apart is what lets those tasks be run, -tested, and retried on their own. +An app that verifies GitHub deliveries, parses them into `WebhookEvent`s, and +launches the tasks in `github_tasks.py`. Deploy those tasks first. -Deploy it, then point GitHub at the payload URL the dashboard shows: +Deploy the app, then set the payload URL from its dashboard in GitHub: python github_webhooks.py """ @@ -24,15 +22,11 @@ from flyte.extras.webhooks import WebhookAppEnvironment, WebhookEvent, run_once from flyteplugins.github import GitHubProvider, events -# One provider, one route at /webhook/github, one dashboard at /. +# Serves the receiver at /webhook/github and a setup dashboard at /. +# GITHUB_WEBHOOK_SECRET is mounted automatically. # -# GITHUB_WEBHOOK_SECRET is mounted for you from the provider's -# `default_secret_env`, so it does not need naming again in `secrets=`. -# -# `scopes` is an allowlist of repositories. A delivery from anywhere else is -# acknowledged — so GitHub stops retrying it — but never dispatched. So is a -# delivery carrying no repository at all: an allowlist cannot vouch for an -# event it cannot attribute. +# `scopes` lists the repositories to act on. Deliveries from other +# repositories, or with no repository, are acknowledged but not dispatched. app_env = WebhookAppEnvironment( name="github-webhooks", providers=[GitHubProvider()], @@ -48,15 +42,10 @@ async def on_pull_request_opened(event: WebhookEvent) -> dict: """Launch triage once per pull request. - `run_once` is what makes this safe to call repeatedly. GitHub retries any - non-2xx delivery and an operator may re-send one by hand; both arrive with - the same `dedupe_key()`, and only the first launches a run. A *later* change - to the same pull request gets its own key, because the key folds in the - provider's own timestamp. - - Handlers must `await run_once.aio(...)` rather than call the blocking form: - the blocking form stalls the app's event loop, and GitHub times a delivery - out in ten seconds. + Retried and resent deliveries have the same `dedupe_key()`, so `run_once` + launches only one run for them. Use `await run_once.aio(...)`: the blocking + form stalls the app's event loop, and GitHub times out a delivery after + ten seconds. """ import flyte.remote as remote @@ -68,7 +57,7 @@ async def on_pull_request_opened(event: WebhookEvent) -> dict: number=event.payload["pull_request"]["number"], ) if not result.created: - # An earlier delivery of this same event already launched it. + # An earlier delivery of this event already launched a run. return {"skipped": result.run.name, "url": result.run.url} return {"run": result.run.name, "url": result.run.url} # {{/docs-fragment handler}} @@ -78,8 +67,7 @@ async def on_pull_request_opened(event: WebhookEvent) -> dict: if __name__ == "__main__": flyte.init_from_config(root_dir=pathlib.Path(__file__).parent) deployment = flyte.serve(app_env) - # The dashboard lists every provider's payload URL, whether its secret is - # mounted, and how it is verified — paste the GitHub row into + # The dashboard shows the payload URL to enter in GitHub under # Settings -> Webhooks -> Add webhook. print(f"Setup dashboard: {deployment.url}") # {{/docs-fragment serve}} diff --git a/v2/integrations/flyte-plugins/jira/jira_tasks.py b/v2/integrations/flyte-plugins/jira/jira_tasks.py index 59c375ea..2085c7b3 100644 --- a/v2/integrations/flyte-plugins/jira/jira_tasks.py +++ b/v2/integrations/flyte-plugins/jira/jira_tasks.py @@ -8,9 +8,11 @@ # main = "replay_sample_delivery" # params = "" # /// -"""The task a Jira webhook launches, over the `jira` client. +"""Tasks launched by the Jira webhook receiver. -`replay_sample_delivery` is the entrypoint, and needs no Jira site: +`triage_issue` comments on an issue and transitions it, using the `jira` client. + +`replay_sample_delivery` runs without a Jira site: flyte run --local jira_tasks.py replay_sample_delivery """ @@ -32,11 +34,10 @@ @env.task async def triage_issue(issue_key: str) -> str: - """Comment on an issue and move it to In Progress. Plain `jira` client. + """Comment on an issue and move it to In Progress. - Transitions are named per workflow, not globally, so this resolves the name - to an id rather than hard-coding one — a hard-coded id breaks the first time - somebody edits the project's workflow. + Transition IDs differ between workflows, so this looks up the transition + by name instead of hard-coding an ID. """ import asyncio import os @@ -56,14 +57,14 @@ def _work() -> str: return transition["name"] return "no matching transition" - # The `jira` client is synchronous; keep it off the event loop. + # The `jira` client is synchronous; run it off the event loop. return await asyncio.to_thread(_work) # {{/docs-fragment task}} # {{docs-fragment replay}} -# Its own environment, with no secrets: the replay needs none, and a task -# environment that declares a secret cannot start until that secret exists. +# A separate environment with no secrets, so the replay runs before any +# secret is created. replay_env = flyte.TaskEnvironment( name="jira-replay", image=flyte.Image.from_debian_base().with_pip_packages("flyteplugins-jira"), @@ -73,11 +74,10 @@ def _work() -> str: @replay_env.task async def replay_sample_delivery() -> dict[str, str]: - """Verify and parse the real delivery the plugin ships. No Jira site needed. + """Verify and parse the sample delivery bundled with the plugin. - Note what `verify` does here: it compares a shared token in constant time, - rather than checking a signature. The sample's "sign" function only sets the - header, because there is nothing to sign. + Jira doesn't sign deliveries, so `verify` compares a shared token in the + `X-Webhook-Token` header. The sample's sign function only sets that header. """ import flyteplugins.jira as plugin @@ -90,13 +90,14 @@ async def replay_sample_delivery() -> dict[str, str]: event = plugin.parse(headers, body) return { - # `jira:issue_created` — Jira namespaces some, but not all, event names. + # For example, `jira:issue_created`. Some Jira event names have the + # `jira:` prefix and some don't. "qualified_type": event.qualified_type, "scope": event.scope or "", "resource_id": event.resource_id or "", "title": event.title or "", "dedupe_key": event.dedupe_key(), - # False — and the setup dashboard says so. + # False for Jira. The setup dashboard shows this too. "provider_signs_deliveries": str(plugin.JiraProvider().signed), } # {{/docs-fragment replay}} diff --git a/v2/integrations/flyte-plugins/jira/jira_webhooks.py b/v2/integrations/flyte-plugins/jira/jira_webhooks.py index 81e2aa63..59e47d52 100644 --- a/v2/integrations/flyte-plugins/jira/jira_webhooks.py +++ b/v2/integrations/flyte-plugins/jira/jira_webhooks.py @@ -5,17 +5,14 @@ # "flyteplugins-jira[app]>=2.10.6", # ] # /// -"""The Jira webhook receiver — the one provider that does not sign. +"""Jira webhook receiver. -Jira Cloud sends no signature, so there is no HMAC to check. `JiraProvider` -authenticates with a shared token in an `X-Webhook-Token` header instead, and -reports `signed=False` so the dashboard says so plainly rather than implying a -guarantee that is absent. +Jira Cloud doesn't sign deliveries. `JiraProvider` checks a shared token in the +`X-Webhook-Token` header instead, and reports `signed=False` on the dashboard. -Jira cannot send custom headers itself, so something in front of this app has to -inject that header — an API gateway, an ingress rule, or a Jira Automation rule -using *Send web request*, which can. Read the guide's authentication section -before exposing this route. +Jira webhooks can't set custom headers, so an API gateway, an ingress rule, or +a Jira Automation rule using Send web request must add the header. See the +Authentication section of the Jira integration guide before exposing this route. python jira_webhooks.py """ @@ -27,11 +24,9 @@ from flyte.extras.webhooks import WebhookAppEnvironment, WebhookEvent, run_once from flyteplugins.jira import JiraProvider, events -# JIRA_WEBHOOK_TOKEN is mounted from the provider's `default_secret_env`. It is a -# shared token, not a signing secret: anything holding it can post a delivery, -# and the token travels on every request rather than signing one. So treat the -# proxy in front and the `scopes` allowlist below as part of the auth story, not -# as extras. +# JIRA_WEBHOOK_TOKEN is mounted automatically. It's a shared token, sent with +# every request: anyone who has it can post a delivery. Restrict `scopes` to +# the projects the app should act on. app_env = WebhookAppEnvironment( name="jira-webhooks", providers=[JiraProvider()], @@ -48,8 +43,7 @@ async def on_issue_created(event: WebhookEvent) -> dict: """Launch triage once per new issue. - `event.resource_id` is the issue key (`PROJ-1`). That is the stable handle - the Jira API takes, and unlike the numeric id it is also what a human reads. + `event.resource_id` is the issue key, such as `PROJ-1`. """ import flyte.remote as remote diff --git a/v2/integrations/flyte-plugins/linear/linear_tasks.py b/v2/integrations/flyte-plugins/linear/linear_tasks.py index 225f20a4..3613007f 100644 --- a/v2/integrations/flyte-plugins/linear/linear_tasks.py +++ b/v2/integrations/flyte-plugins/linear/linear_tasks.py @@ -8,13 +8,11 @@ # main = "replay_sample_delivery" # params = "" # /// -"""The task a Linear webhook launches, over Linear's GraphQL API. +"""Tasks launched by the Linear webhook receiver. -Linear ships no official Python SDK and does not need one: its API is a single -GraphQL endpoint, so `gql` is the maintained client and the task below calls it -directly. The plugin's job stops at the webhook. +`triage_issue` comments on an issue through Linear's GraphQL API, using `gql`. -`replay_sample_delivery` is the entrypoint, and needs no Linear workspace: +`replay_sample_delivery` runs without a Linear workspace: flyte run --local linear_tasks.py replay_sample_delivery """ @@ -32,7 +30,7 @@ @env.task async def triage_issue(issue_id: str, title: str) -> str: - """Comment on an issue. Plain `gql` against Linear's one GraphQL endpoint.""" + """Comment on an issue.""" import os from gql import Client, gql @@ -40,7 +38,7 @@ async def triage_issue(issue_id: str, title: str) -> str: transport = HTTPXAsyncTransport( url="https://api.linear.app/graphql", - # Linear takes the API key raw, with no "Bearer " prefix. + # Linear expects the API key with no "Bearer " prefix. headers={"Authorization": os.environ["LINEAR_API_KEY"]}, ) mutation = gql( @@ -60,8 +58,8 @@ async def triage_issue(issue_id: str, title: str) -> str: # {{docs-fragment replay}} -# Its own environment, with no secrets: the replay needs none, and a task -# environment that declares a secret cannot start until that secret exists. +# A separate environment with no secrets, so the replay runs before any +# secret is created. replay_env = flyte.TaskEnvironment( name="linear-replay", image=flyte.Image.from_debian_base().with_pip_packages("flyteplugins-linear"), @@ -71,7 +69,7 @@ async def triage_issue(issue_id: str, title: str) -> str: @replay_env.task async def replay_sample_delivery() -> dict[str, str]: - """Verify and parse the real delivery the plugin ships. No workspace needed.""" + """Verify and parse the sample delivery bundled with the plugin.""" import flyteplugins.linear as plugin secret = "a-test-signing-secret" @@ -81,25 +79,23 @@ async def replay_sample_delivery() -> dict[str, str]: assert plugin.verify(body, headers, secret), "a correctly signed delivery must verify" assert not plugin.verify(body, headers, "wrong-secret"), "a bad signature must not" - # The wire contract, which the round trip above cannot check: `verify` and - # `SAMPLE_DELIVERY` agree with each other whatever the header is called, so - # a wrong name passes conformance and then rejects every real delivery. - # Linear signs with `Linear-Signature` -- note the missing `X-` prefix, - # which looks like a typo and is not one. + # Check the header name too. The sample's headers come from the plugin, so + # the round trip above passes whatever the header is called. Linear sends + # `Linear-Signature`, with no `X-` prefix. assert list(headers) == ["Linear-Signature"], f"unexpected signature header: {list(headers)}" assert not plugin.verify(body, {"X-Linear-Signature": headers["Linear-Signature"]}, secret), ( - "the old, wrong header name must not verify" + "X-Linear-Signature must not verify" ) event = plugin.parse(headers, body) return { - # `Issue.create` — Linear is one of the providers that splits the two. + # `Issue.create`: Linear sends the type and action separately. "qualified_type": event.qualified_type, "scope": event.scope or "", "title": event.title or "", "url": event.url or "", "dedupe_key": event.dedupe_key(), - # The header a real delivery carries the signature in. + # The header that carries the signature. "signature_header": next(iter(headers)), } # {{/docs-fragment replay}} diff --git a/v2/integrations/flyte-plugins/linear/linear_webhooks.py b/v2/integrations/flyte-plugins/linear/linear_webhooks.py index ac27fbd1..6109fc8a 100644 --- a/v2/integrations/flyte-plugins/linear/linear_webhooks.py +++ b/v2/integrations/flyte-plugins/linear/linear_webhooks.py @@ -5,10 +5,10 @@ # "flyteplugins-linear[app]>=2.10.7", # ] # /// -"""The Linear webhook receiver. +"""Linear webhook receiver. -Deploy it, then paste the payload URL the dashboard shows into -Linear Settings -> API -> Webhooks: +Deploy the app, then enter the payload URL from its dashboard in Linear under +Settings -> API -> Webhooks: python linear_webhooks.py """ @@ -20,12 +20,10 @@ from flyte.extras.webhooks import WebhookAppEnvironment, WebhookEvent, run_once from flyteplugins.linear import LinearProvider, events -# LINEAR_WEBHOOK_SECRET is mounted from the provider's `default_secret_env`. +# LINEAR_WEBHOOK_SECRET is mounted automatically. # -# `scopes` matches Linear's team id, which is what the provider puts in -# `WebhookEvent.scope` — including on Comment and Reaction payloads, where the -# team id is nested on the issue rather than sent at the top level. Without that -# fallback an allowlist would drop every non-Issue event as unattributable. +# `scopes` lists Linear team IDs. For Comment and Reaction events, the provider +# reads the team ID from the related issue. app_env = WebhookAppEnvironment( name="linear-webhooks", providers=[LinearProvider()], @@ -41,8 +39,7 @@ async def on_issue_created(event: WebhookEvent) -> dict: """Launch triage once per new issue. - Linear splits type and action, so the constant is `Issue.CREATE` and the - normalized `qualified_type` reads `Issue.create`. + The constant `Issue.CREATE` matches the `qualified_type` `Issue.create`. """ import flyte.remote as remote diff --git a/v2/integrations/flyte-plugins/slack/slack_tasks.py b/v2/integrations/flyte-plugins/slack/slack_tasks.py index 2353e6f2..d5d188c1 100644 --- a/v2/integrations/flyte-plugins/slack/slack_tasks.py +++ b/v2/integrations/flyte-plugins/slack/slack_tasks.py @@ -7,18 +7,14 @@ # main = "replay_sample_delivery" # params = "" # /// -"""Sending to Slack from tasks, and gating a run on a button click. +"""Tasks that post to Slack and wait for approval button clicks. -Receiving is the receiver's job; sending is a task's. This plugin is one of two -that carry more than a webhook provider, because two things here are not -reshaped JSON: +- `answer` posts a threaded reply with `notify.post`, then edits it with + `notify.update`. +- `deploy_with_approval` uses `approval.request` to pause the run until + someone clicks a button. -- `notify` replaces the fifty lines of `requests` and `ok`-checking every - integration ends up hand-rolling. -- `approval` is the round trip — post buttons, park the run on a condition, - resume when someone clicks. The condition is the part only Flyte can do. - -`replay_sample_delivery` is the entrypoint, and needs no Slack workspace: +`replay_sample_delivery` runs without a Slack workspace: flyte run --local slack_tasks.py replay_sample_delivery """ @@ -30,8 +26,8 @@ env = flyte.TaskEnvironment( name="slack-bot", image=flyte.Image.from_debian_base().with_pip_packages("flyteplugins-slack"), - # The `xoxb-` credential from OAuth & Permissions. Posting needs the - # `chat:write` scope, and the bot needs a `/invite` into the channel. + # 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"), ) @@ -39,11 +35,9 @@ @env.task async def answer(channel: str, text: str, thread_ts: str) -> str: - """Post a threaded reply, then edit it in place when the work finishes. + """Post a threaded reply, then edit it when the work finishes. - `post` returns the message's `ts`, which is both the thread anchor and the - address `update` edits — so a progress counter is two calls, not a second - API surface to learn. + `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}") @@ -54,15 +48,12 @@ async def answer(channel: str, text: str, thread_ts: str) -> str: # {{docs-fragment approval}} @env.task async def deploy_with_approval(release: str, channel: str = "C0DEPLOYS") -> str: - """Ask Slack for a decision, and block until somebody clicks. - - `approval.request` posts Block Kit buttons and parks the run on a - `flyte.new_condition`, then replaces the buttons with a "decided by" line so - nobody clicks twice. + """Post approval buttons and wait for a click. - The same condition is answerable from the Flyte UI, so an approval nobody - clicks in Slack is not stuck: the run shows the same prompt, and either path - resolves it. + `approval.request` posts Block Kit buttons and waits on a + `flyte.new_condition`. The handler added by `approval.register` resolves + the condition when someone clicks. The condition can also be resolved + from the Flyte UI. """ decision = await approval.request.aio( channel, @@ -77,8 +68,8 @@ async def deploy_with_approval(release: str, channel: str = "C0DEPLOYS") -> str: # {{docs-fragment replay}} -# Its own environment, with no secrets: the replay needs none, and a task -# environment that declares a secret cannot start until that secret exists. +# A separate environment with no secrets, so the replay runs before any +# secret is created. replay_env = flyte.TaskEnvironment( name="slack-replay", image=flyte.Image.from_debian_base().with_pip_packages("flyteplugins-slack"), @@ -88,11 +79,10 @@ async def deploy_with_approval(release: str, channel: str = "C0DEPLOYS") -> str: @replay_env.task async def replay_sample_delivery() -> dict[str, str]: - """Verify and parse the real delivery the plugin ships. No workspace needed. + """Verify and parse the sample delivery bundled with the plugin. - Slack's sample signs at call time rather than carrying a fixed signature, - because its verifier enforces a five-minute replay window — a delivery - signed at a hard-coded timestamp would start failing the moment it aged out. + The provider rejects deliveries more than five minutes old, so + `SAMPLE_DELIVERY` signs the payload with the current time when called. """ import flyteplugins.slack as plugin @@ -110,7 +100,7 @@ async def replay_sample_delivery() -> dict[str, str]: "title": event.title or "", "actor": event.actor or "", "dedupe_key": event.dedupe_key(), - # The replay window, in seconds. Slack is the one provider that has one. + # Maximum delivery age, in seconds. "max_request_age": str(plugin.MAX_REQUEST_AGE_SECONDS), } # {{/docs-fragment replay}} diff --git a/v2/integrations/flyte-plugins/slack/slack_webhooks.py b/v2/integrations/flyte-plugins/slack/slack_webhooks.py index aeac13bd..36647890 100644 --- a/v2/integrations/flyte-plugins/slack/slack_webhooks.py +++ b/v2/integrations/flyte-plugins/slack/slack_webhooks.py @@ -5,16 +5,14 @@ # "flyteplugins-slack[app]>=2.10.6", # ] # /// -"""The Slack webhook receiver: Events API, interactivity, and slash commands. +"""Slack webhook receiver for Events API callbacks, interactivity, and slash commands. -Slack is the broadest provider in this family because it delivers three -different shapes to the same route — event callbacks as JSON, interactivity -payloads and slash commands as form bodies. One `SlackProvider()` verifies and -normalizes all three; `on_event` is what tells them apart. +`SlackProvider` verifies and parses all three delivery types on one route; +`on_event` selects between them. -Deploy it, then paste the payload URL the dashboard shows into all three fields -at api.slack.com/apps (Event Subscriptions, Interactivity, and each slash -command): +Deploy the app, then enter the payload URL from its dashboard at +api.slack.com/apps under Event Subscriptions, Interactivity & Shortcuts, and +each slash command: python slack_webhooks.py """ @@ -26,9 +24,8 @@ from flyte.extras.webhooks import WebhookAppEnvironment, WebhookEvent, run_once from flyteplugins.slack import SlackProvider, approval, events, notify -# SLACK_SIGNING_SECRET is mounted from the provider's `default_secret_env`. That -# is the *signing secret* under Basic Information — not the `xoxb-` bot token -# that `notify` sends with, which belongs on a task environment instead. +# SLACK_SIGNING_SECRET is mounted automatically. It's the signing secret from +# Basic Information, not the bot token, which goes on the task environment. app_env = WebhookAppEnvironment( name="slack-webhooks", providers=[SlackProvider()], @@ -36,10 +33,9 @@ resources=flyte.Resources(cpu=1, memory="512Mi"), ) -# One line, and every approval button posted by `approval.request` is answered -# from here on: the handler reads the run, action, and condition names off the -# button's own `value`, looks the condition up, and signals it. No configuration, -# because the button carries everything needed to answer it. +# Adds a handler that resolves the condition behind each button posted by +# `approval.request`. Each button's value carries the run, action, and +# condition names, so no other configuration is needed. approval.register(app_env) # {{/docs-fragment app}} @@ -47,10 +43,10 @@ # {{docs-fragment handler}} @app_env.on_event(events.AppMention.ANY) async def on_mention(event: WebhookEvent) -> dict: - """Answer an @-mention by launching a run, once per message. + """Launch a run for each @-mention. - Slack's dedupe key is per message. To collapse a whole thread onto one run, - build your own key from `thread_ts` and pass that to `run_once` instead. + 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 @@ -68,12 +64,11 @@ async def on_mention(event: WebhookEvent) -> dict: @app_env.on_event(events.Command, action="/deploy") async def on_deploy_command(event: WebhookEvent) -> dict: - """A slash command. `respond` needs no token at all. + """Acknowledge the /deploy slash command. - It posts to the `response_url` every interaction and slash command carries, - which makes it the zero-setup way to answer the click that launched you. - Slack only shows a synchronous reply if it arrives within three seconds, so - acknowledge here and let the launched run post the real answer. + `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}