Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
37 changes: 16 additions & 21 deletions v2/integrations/flyte-plugins/clickup/clickup_tasks.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
"""
Expand All @@ -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

Expand All @@ -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"),
Expand All @@ -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"
Expand All @@ -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}}
Expand Down
16 changes: 7 additions & 9 deletions v2/integrations/flyte-plugins/clickup/clickup_webhooks.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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()],
Expand All @@ -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

Expand Down
65 changes: 27 additions & 38 deletions v2/integrations/flyte-plugins/github/github_tasks.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
"""
Expand All @@ -34,16 +26,16 @@
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"),
)


@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
Expand All @@ -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

Expand All @@ -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"),
Expand All @@ -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.
Expand All @@ -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"),
Expand All @@ -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

Expand All @@ -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}}
Expand Down
40 changes: 14 additions & 26 deletions v2/integrations/flyte-plugins/github/github_webhooks.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
"""
Expand All @@ -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()],
Expand All @@ -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

Expand All @@ -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}}
Expand All @@ -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}}
Loading
Loading