diff --git a/packages/plugins/paperclip-plugin-github-mirror/README.md b/packages/plugins/paperclip-plugin-github-mirror/README.md index 2188413ccbfb..32a704ab2163 100644 --- a/packages/plugins/paperclip-plugin-github-mirror/README.md +++ b/packages/plugins/paperclip-plugin-github-mirror/README.md @@ -35,12 +35,54 @@ plugin state. ## Behaviour under failure GitHub being down, rate-limiting, or rejecting the token must not stop Paperclip from -processing events. Every handler is wrapped: failures are logged with a `retryable` flag -(429, rate-limited 403, and 5xx are retryable; other 4xx are not) and the worker keeps +processing events. Every handler is wrapped: failures are logged and the worker keeps running. -Mirroring is idempotent — the GitHub issue number is stored in plugin state, so repeated -events never create duplicates. +**Retry.** Each write is attempted up to three times, with a jittered 1s/4s backoff, when +GitHub asked us to back off (429, rate-limited 403), fell over (5xx), or the worker→host +call timed out without an answer. Everything else — 401, 404, 422 — fails on the first +attempt, because it will fail identically on the second. If GitHub named a wait via +`retry-after` or `x-ratelimit-reset`, that wait is used instead of the backoff. + +**Timeouts.** The plugin does not set one, and cannot: `ctx.http.fetch` serializes only +method, headers and body, so an `AbortSignal` never reaches the host. It does not need to. +Each call is already bounded at 30s twice over — by the SDK's worker→host call timer and by +the host's own `AbortController` — and the retry budget is capped at that same 30s so three +attempts cannot occupy a handler for a minute and a half. + +**Duplicate creates.** Mirroring is idempotent: the GitHub issue number is stored in plugin +state, so repeated events never create duplicates. Creating the issue is the one step that +cannot simply be repeated, so it is written down first — see below. + +## The create outbox + +Event delivery is fire-and-forget. The host pushes events as a JSON-RPC notification and +drops them outright when the worker is down; there is no replay. So if a create is +interrupted between the POST and the write that records its number, nothing would ever +mention it again — and the next event for that task would happily create a second GitHub +issue. + +Before posting, the plugin upserts a `mirror-create` entity keyed `:` +with status `pending`. On success it stores the number in plugin state and flips the record +to `done`. A `pending` record therefore means exactly one thing: a create was attempted and +we do not know how it ended. + +A scheduled job (`*/5 * * * *`) resolves those: + +| Record | Becomes | Why | +|---|---|---| +| `pending`, number known (in the record or in plugin state) | `done` | The create demonstrably succeeded; only the closing write was lost | +| `pending`, no number, started under 2 minutes ago | unchanged | Could still be in flight | +| `pending`, no number, older than that | `uncertain` | Logged once, and never attempted again | + +`uncertain` is terminal on purpose. The issue may or may not exist on GitHub, and finding +out would mean reading GitHub back — which this plugin does not do, at all, by design. So +it records the ambiguity where a human can see it instead of guessing. That trades a silent +duplicate for a task that is visibly not mirrored, which is the lesser of the two. + +The outbox lives in `ctx.entities` rather than `ctx.state` because entities can be +enumerated. Plugin state cannot: the host's state store has a `list`, but it is not exposed +over the worker→host RPC, so a queue kept there could never find its own pending work. ## Development diff --git a/packages/plugins/paperclip-plugin-github-mirror/src/constants.ts b/packages/plugins/paperclip-plugin-github-mirror/src/constants.ts index 6ccf8fee37cb..1d0055c8d3ed 100644 --- a/packages/plugins/paperclip-plugin-github-mirror/src/constants.ts +++ b/packages/plugins/paperclip-plugin-github-mirror/src/constants.ts @@ -9,6 +9,54 @@ export const STATE_KEYS = { lastStatus: "last-mirrored-status", } as const; +/** + * Entity type for the create outbox — one record per Paperclip issue we have + * tried to mirror, written before the POST so an interrupted create leaves a + * trace instead of nothing. + * + * It lives in `ctx.entities` rather than `ctx.state` for one reason: entities + * can be enumerated (`ctx.entities.list`), and plugin state cannot. The state + * store does have a `list`, but it is not exposed over the worker→host RPC, so + * an outbox kept there could never find its own pending work. + */ +export const OUTBOX_ENTITY_TYPE = "mirror-create"; + +/** Job key for the outbox drain, declared in the manifest. */ +export const OUTBOX_DRAIN_JOB = "drain-mirror-outbox"; + +/** + * Statuses a create record moves through. + * + * `uncertain` is terminal and deliberately final: it means a create may or may + * not have reached GitHub and we have no way to find out — the mirror is + * write-only, so it will not go and look. Refusing forever turns what used to + * be a silent duplicate issue into one recorded fact a human can act on. + * + * `failed` is the case `uncertain` must not swallow. When GitHub answers with a + * status code — a 500, or a 429 that outlived the retry budget — the create + * definitively did not happen, so there is nothing to be uncertain about and no + * duplicate to fear. Those records stay retryable: a later event creates the + * issue. Collapsing them into `uncertain` would refuse forever on the strength + * of an answer that said "no", which is worse than the behaviour this file + * replaced, where a failed create was simply retried on the next event. + * + * The distinction is exactly whether GitHub replied. It did — `failed`. We + * never found out — `pending`, and `uncertain` once the grace window passes. + */ +export const OUTBOX_STATUS = { + pending: "pending", + done: "done", + uncertain: "uncertain", + failed: "failed", +} as const; + +/** + * How long a `pending` record is left alone before the drain will call it + * `uncertain`. Comfortably past the 30s a single call can take, so the drain + * never condemns a create that is still in flight. + */ +export const OUTBOX_PENDING_GRACE_MS = 120_000; + /** * Emitted by the escalation plugin when a task is handed to a human. Subscribed to * rather than reimplemented, so escalation policy stays in one place. diff --git a/packages/plugins/paperclip-plugin-github-mirror/src/github.ts b/packages/plugins/paperclip-plugin-github-mirror/src/github.ts index a30995b78833..d7c3ee2d4c12 100644 --- a/packages/plugins/paperclip-plugin-github-mirror/src/github.ts +++ b/packages/plugins/paperclip-plugin-github-mirror/src/github.ts @@ -3,9 +3,31 @@ * * Deliberately write-only: the mirror never reads GitHub state back into * Paperclip, so there is no `get`/`list` here. Paperclip stays the store of - * record; GitHub is a viewing surface. + * record; GitHub is a viewing surface. Retry does not change that: a failed + * write is tried again, never read back to find out what happened. + * + * ## The 30s budget + * + * A plugin cannot cancel its own request. `ctx.http.fetch` is a JSON-RPC call + * whose `init` is serialized down to method, headers and body, so an + * `AbortSignal` passed here would be silently dropped and a `Promise.race` + * around the await would only shorten the plugin's wait while the host socket + * kept running. + * + * It does not need one. Every call is already bounded twice at 30 seconds: + * + * - worker side, the SDK's `callHost` timer (`DEFAULT_RPC_TIMEOUT_MS` in + * `worker-rpc-host.ts`; `runWorker` never passes `rpcTimeoutMs`), which + * rejects with a JSON-RPC timeout; + * - host side, an `AbortController` armed with `PLUGIN_FETCH_TIMEOUT_MS` + * (`plugin-host-services.ts`), which aborts the socket itself. + * + * So the ceiling is the runtime's, not ours, and the only thing this file owes + * it is that retrying stays inside it — see `RETRY_BUDGET_MS` in `retry.ts`. */ +import { withRetry, type RetryDeps } from "./retry.js"; + export interface GithubClientOptions { /** `owner/repo`. */ repository: string; @@ -13,6 +35,8 @@ export interface GithubClientOptions { token: string; fetchImpl: (url: string, init?: RequestInit) => Promise; apiBaseUrl?: string; + /** Retry timing seams. Tests inject them; production uses the defaults. */ + retryDeps?: Partial; } export interface GithubIssueRef { @@ -29,12 +53,42 @@ export class GithubApiError extends Error { readonly status: number, /** True when GitHub asked us to back off rather than rejecting the request outright. */ readonly retryable: boolean, + /** + * How long GitHub asked us to wait, in ms, when it said so via + * `retry-after` or `x-ratelimit-reset`. `null` when it did not. + */ + readonly retryAfterMs: number | null = null, ) { super(message); this.name = "GithubApiError"; } } +/** + * GitHub names a wait in one of two ways: `retry-after` (seconds, on secondary + * rate limits and abuse detection) or `x-ratelimit-reset` (epoch seconds, on + * primary rate limits). Anything absent, unparseable, or in the past yields + * `null`, and the caller falls back to its own backoff. + */ +function parseRetryAfterMs(headers: Headers, nowMs: number): number | null { + const retryAfter = headers.get("retry-after"); + if (retryAfter) { + const seconds = Number(retryAfter); + if (Number.isFinite(seconds) && seconds >= 0) return Math.round(seconds * 1000); + } + + const reset = headers.get("x-ratelimit-reset"); + if (reset) { + const resetSeconds = Number(reset); + if (Number.isFinite(resetSeconds)) { + const waitMs = resetSeconds * 1000 - nowMs; + if (waitMs > 0) return Math.round(waitMs); + } + } + + return null; +} + /** `owner/repo` → validated parts. Throws on anything else so a typo fails loudly at setup. */ export function parseRepository(repository: string): { owner: string; repo: string } { const match = /^([A-Za-z0-9._-]+)\/([A-Za-z0-9._-]+)$/.exec(repository.trim()); @@ -56,7 +110,17 @@ export class GithubClient { this.apiBaseUrl = options.apiBaseUrl ?? DEFAULT_API_BASE; } + /** + * One write, retried on the failures that can succeed on a second try. The + * retry lives here rather than in the handlers so every call gets it, and so + * a handler that ends up in `guard` has genuinely exhausted its options + * rather than given up on the first 502. + */ private async request(path: string, init: RequestInit): Promise { + return withRetry(() => this.attempt(path, init), this.options.retryDeps); + } + + private async attempt(path: string, init: RequestInit): Promise { const response = await this.options.fetchImpl(`${this.apiBaseUrl}${path}`, { ...init, headers: { @@ -83,6 +147,7 @@ export class GithubClient { `GitHub ${init.method ?? "GET"} ${path} failed: ${response.status} ${detail.slice(0, 300)}`, response.status, rateLimited || response.status >= 500, + parseRetryAfterMs(response.headers, Date.now()), ); } diff --git a/packages/plugins/paperclip-plugin-github-mirror/src/manifest.ts b/packages/plugins/paperclip-plugin-github-mirror/src/manifest.ts index a24c734312ce..597837c2c018 100644 --- a/packages/plugins/paperclip-plugin-github-mirror/src/manifest.ts +++ b/packages/plugins/paperclip-plugin-github-mirror/src/manifest.ts @@ -1,5 +1,5 @@ import type { PaperclipPluginManifestV1 } from "@paperclipai/plugin-sdk"; -import { PLUGIN_ID, PLUGIN_VERSION } from "./constants.js"; +import { OUTBOX_DRAIN_JOB, PLUGIN_ID, PLUGIN_VERSION } from "./constants.js"; const manifest: PaperclipPluginManifestV1 = { id: PLUGIN_ID, @@ -17,10 +17,20 @@ const manifest: PaperclipPluginManifestV1 = { "plugin.state.write", "http.outbound", "secrets.read-ref", + "jobs.schedule", ], entrypoints: { worker: "./dist/worker.js", }, + jobs: [ + { + jobKey: OUTBOX_DRAIN_JOB, + displayName: "Resolve interrupted mirrors", + description: + "Closes create records left open by an interrupted mirror, and marks the ones that cannot be confirmed. Reads plugin state only — it never calls GitHub.", + schedule: "*/5 * * * *", + }, + ], instanceConfigSchema: { type: "object", properties: { diff --git a/packages/plugins/paperclip-plugin-github-mirror/src/outbox.ts b/packages/plugins/paperclip-plugin-github-mirror/src/outbox.ts new file mode 100644 index 000000000000..7b5d61c03954 --- /dev/null +++ b/packages/plugins/paperclip-plugin-github-mirror/src/outbox.ts @@ -0,0 +1,223 @@ +/** + * The create outbox. + * + * ## Why an outbox at all + * + * Creating the mirrored GitHub issue is the one step that cannot be repeated + * safely: a second POST makes a second issue. Everything else the mirror does + * — patching, commenting — is keyed on a number we already hold. + * + * Event delivery gives no help here. The host pushes events as a fire-and-forget + * JSON-RPC notification and drops it outright when the worker is down + * (`plugin-worker-manager.ts`, `notify()` returns early unless the status is + * `running`). There is no replay: an event that arrives while the worker is + * restarting is simply gone. So the plugin cannot wait to be told again — it + * has to write down what it is about to do, before it does it. + * + * ## Why entities and not state + * + * `ctx.state` is `get`/`set`/`delete` only, with a last-write-wins upsert on + * the host side — no compare-and-set, no conditional write. Worse for an + * outbox: the host's state store has a `list`, but it is not exposed over the + * worker→host RPC, so a queue kept in state could never enumerate its own + * pending work. `ctx.entities` has `upsert(externalId)` plus `list(query)` and + * is enumerable, so the outbox lives there. + * + * ## What this buys, and what it does not + * + * Intent is recorded before the POST, so a create that dies mid-flight leaves a + * `pending` record. What it cannot do is find out what actually happened: the + * mirror is write-only and stays that way, so a `pending` record with no number + * becomes `uncertain` — recorded, logged, visible, and never retried. That is a + * downgrade from a silent duplicate to a known unknown, not a cure. + */ + +import type { PluginContext, PluginEntityRecord, ScopeKey } from "@paperclipai/plugin-sdk"; +import { + OUTBOX_ENTITY_TYPE, + OUTBOX_PENDING_GRACE_MS, + OUTBOX_STATUS, + STATE_KEYS, +} from "./constants.js"; + +/** Records examined per drain run, and the page size used to walk them. */ +const PAGE_SIZE = 200; +const MAX_RECORDS_PER_RUN = 2_000; + +export interface OutboxData { + companyId: string; + issueId: string; + /** When the create was first attempted — the clock the grace window uses. */ + startedAt: string; + /** Set once the POST has returned a number. */ + mirroredNumber?: number; +} + +export function issueScope(issueId: string): { scopeKind: "issue"; scopeId: string } { + return { scopeKind: "issue", scopeId: issueId }; +} + +export async function readMirroredNumber( + ctx: PluginContext, + issueId: string, +): Promise { + const key: ScopeKey = { ...issueScope(issueId), stateKey: STATE_KEYS.mirroredNumber }; + const stored = await ctx.state.get(key); + return typeof stored === "number" ? stored : null; +} + +/** One record per Paperclip issue, per company. */ +function externalId(companyId: string, issueId: string): string { + return `${companyId}:${issueId}`; +} + +function readData(record: PluginEntityRecord): OutboxData | null { + const data = record.data as Partial | undefined; + if (!data || typeof data.issueId !== "string" || typeof data.companyId !== "string") return null; + return { + companyId: data.companyId, + issueId: data.issueId, + startedAt: typeof data.startedAt === "string" ? data.startedAt : "", + mirroredNumber: typeof data.mirroredNumber === "number" ? data.mirroredNumber : undefined, + }; +} + +async function write( + ctx: PluginContext, + data: OutboxData, + status: string, + title: string | undefined, +): Promise { + return ctx.entities.upsert({ + entityType: OUTBOX_ENTITY_TYPE, + scopeKind: "issue", + scopeId: data.issueId, + externalId: externalId(data.companyId, data.issueId), + title, + status, + data: { ...data }, + }); +} + +export async function findCreateRecord( + ctx: PluginContext, + companyId: string, + issueId: string, +): Promise { + const [found] = await ctx.entities.list({ + entityType: OUTBOX_ENTITY_TYPE, + externalId: externalId(companyId, issueId), + limit: 1, + }); + return found ?? null; +} + +/** + * Written before the POST, so an interrupted create is not invisible. Returns + * the record so the closing write can keep its original `startedAt`. + */ +export async function recordIntent( + ctx: PluginContext, + companyId: string, + issueId: string, + title: string, +): Promise { + return write( + ctx, + { companyId, issueId, startedAt: new Date().toISOString() }, + OUTBOX_STATUS.pending, + title, + ); +} + +/** Written after the number has been persisted to state, closing the record. */ +export async function recordMirrored( + ctx: PluginContext, + record: PluginEntityRecord, + mirroredNumber: number, +): Promise { + const data = readData(record); + if (!data) return; + await write(ctx, { ...data, mirroredNumber }, OUTBOX_STATUS.done, record.title ?? undefined); +} + +/** + * GitHub answered with a status code, so the create definitively did not + * happen. Terminal for this attempt but not for the task: a later event is + * free to create the issue, because there is no duplicate to fear. + */ +export async function recordFailed( + ctx: PluginContext, + record: PluginEntityRecord, +): Promise { + const data = readData(record); + if (!data) return; + await write(ctx, data, OUTBOX_STATUS.failed, record.title ?? undefined); +} + +/** + * Terminal. The create may or may not exist on GitHub; finding out would mean + * reading GitHub back, which this plugin does not do. + */ +export async function recordUncertain( + ctx: PluginContext, + record: PluginEntityRecord, +): Promise { + const data = readData(record); + if (!data) return; + await write(ctx, data, OUTBOX_STATUS.uncertain, record.title ?? undefined); +} + +function isPastGrace(startedAt: string, now: number): boolean { + const started = Date.parse(startedAt); + // An unreadable timestamp is treated as old: leaving it pending forever would + // hide it, and hiding it is the failure mode this whole file exists to fix. + if (Number.isNaN(started)) return true; + return now - started >= OUTBOX_PENDING_GRACE_MS; +} + +/** + * Resolves every `pending` record left behind by an interrupted create. + * + * Reads plugin state and plugin entities. Never GitHub — no lookup, no search, + * no verification call. A record either has a number we already stored, or it + * does not and becomes `uncertain`. + */ +export async function drainOutbox(ctx: PluginContext): Promise { + const now = Date.now(); + + for (let offset = 0; offset < MAX_RECORDS_PER_RUN; offset += PAGE_SIZE) { + const page = await ctx.entities.list({ + entityType: OUTBOX_ENTITY_TYPE, + limit: PAGE_SIZE, + offset, + }); + + for (const record of page) { + if (record.status !== OUTBOX_STATUS.pending) continue; + const data = readData(record); + if (!data) continue; + + // The number may have been persisted to state even though the record was + // never closed — the two writes are not atomic and cannot be. + const mirroredNumber = data.mirroredNumber ?? (await readMirroredNumber(ctx, data.issueId)); + if (mirroredNumber !== null && mirroredNumber !== undefined) { + await recordMirrored(ctx, record, mirroredNumber); + continue; + } + + // Still inside the window where a create could legitimately be in flight. + if (!isPastGrace(data.startedAt, now)) continue; + + await recordUncertain(ctx, record); + // Logged exactly once: the record leaves `pending` on this pass and the + // drain never looks at it again. + ctx.logger.error( + "GitHub mirror: a create was interrupted and cannot be confirmed — no further attempts will be made", + { issueId: data.issueId, companyId: data.companyId, startedAt: data.startedAt }, + ); + } + + if (page.length < PAGE_SIZE) break; + } +} diff --git a/packages/plugins/paperclip-plugin-github-mirror/src/retry.ts b/packages/plugins/paperclip-plugin-github-mirror/src/retry.ts new file mode 100644 index 000000000000..965dbed91cba --- /dev/null +++ b/packages/plugins/paperclip-plugin-github-mirror/src/retry.ts @@ -0,0 +1,104 @@ +/** + * Retry for the mirror's GitHub writes. + * + * A plugin cannot bound its own HTTP call: `ctx.http.fetch` serializes only + * method, headers and body, so an `AbortSignal` never reaches the host. The + * bound comes from the runtime instead, and it is the reason this file has a + * budget at all — see the note above `RETRY_BUDGET_MS`. + * + * Only two failures are worth trying again: GitHub asked us to back off (or + * fell over), and the worker→host call timed out without an answer. Everything + * else — a 401, a 404, a 422 — will fail identically on the second attempt, so + * it is raised immediately. + */ + +import { GithubApiError } from "./github.js"; + +/** Total attempts for one write, the first try included. */ +export const RETRY_ATTEMPTS = 3; + +/** Backoff before attempts 2 and 3, before jitter. */ +export const RETRY_BACKOFF_MS = [1_000, 4_000] as const; + +/** + * Wall clock the whole retry sequence may consume. + * + * Each attempt is already bounded twice at 30s: the SDK's `callHost` timer + * (`worker-rpc-host.ts`, `DEFAULT_RPC_TIMEOUT_MS`, which `runWorker` never + * overrides) and the host's own `AbortController` + * (`plugin-host-services.ts`, `PLUGIN_FETCH_TIMEOUT_MS`). Retrying multiplies + * that: three attempts could otherwise occupy a handler for a minute and a + * half. This budget caps the sequence at one call's worth of time, so a + * first attempt that burns the full 30s leaves nothing to retry with — which + * is the right answer, not a bug. + */ +export const RETRY_BUDGET_MS = 30_000; + +/** Jitter spread, as a fraction of the base delay. */ +const JITTER_RATIO = 0.25; + +/** JSON-RPC code the SDK raises when a worker→host call times out. */ +const RPC_TIMEOUT_CODE = -32003; + +export interface RetryDeps { + now(): number; + sleep(ms: number): Promise; + /** `[0, 1)`. Injected so tests get a deterministic delay. */ + random(): number; +} + +const DEFAULT_DEPS: RetryDeps = { + now: () => Date.now(), + sleep: (ms) => new Promise((resolve) => setTimeout(resolve, ms)), + random: () => Math.random(), +}; + +/** + * A timed-out worker→host call. The plugin never learns whether the request + * reached GitHub, which is exactly why the create path records its intent + * before posting. + */ +function isRpcTimeout(error: unknown): boolean { + if (!error || typeof error !== "object") return false; + if ((error as { code?: unknown }).code === RPC_TIMEOUT_CODE) return true; + const message = (error as { message?: unknown }).message; + return typeof message === "string" && /timed out after \d+ *ms/i.test(message); +} + +export function isRetryable(error: unknown): boolean { + if (error instanceof GithubApiError) return error.retryable; + return isRpcTimeout(error); +} + +/** `retry-after` / `x-ratelimit-reset` if GitHub named a time, else jittered backoff. */ +function delayFor(error: unknown, attempt: number, deps: RetryDeps): number { + if (error instanceof GithubApiError && error.retryAfterMs !== null) { + return error.retryAfterMs; + } + const base = RETRY_BACKOFF_MS[Math.min(attempt - 1, RETRY_BACKOFF_MS.length - 1)]; + // Spread retries of simultaneous events so they do not re-collide on the + // same rate limit. + return Math.round(base * (1 + (deps.random() * 2 - 1) * JITTER_RATIO)); +} + +export async function withRetry( + run: () => Promise, + overrides: Partial = {}, +): Promise { + const deps: RetryDeps = { ...DEFAULT_DEPS, ...overrides }; + const startedAt = deps.now(); + + for (let attempt = 1; ; attempt += 1) { + try { + return await run(); + } catch (error) { + if (attempt >= RETRY_ATTEMPTS || !isRetryable(error)) throw error; + + const delay = delayFor(error, attempt, deps); + // Waiting past the budget would leave no room for the attempt the wait + // is for, so stop and let the caller record the failure now. + if (deps.now() - startedAt + delay >= RETRY_BUDGET_MS) throw error; + await deps.sleep(delay); + } + } +} diff --git a/packages/plugins/paperclip-plugin-github-mirror/src/worker.ts b/packages/plugins/paperclip-plugin-github-mirror/src/worker.ts index e6e4d4150406..6f150310c37d 100644 --- a/packages/plugins/paperclip-plugin-github-mirror/src/worker.ts +++ b/packages/plugins/paperclip-plugin-github-mirror/src/worker.ts @@ -7,7 +7,23 @@ import { } from "@paperclipai/plugin-sdk"; import type { IssueStatus } from "@paperclipai/shared"; import { GithubApiError, GithubClient } from "./github.js"; -import { ESCALATION_EVENT, STATE_KEYS, TELEGRAM_UNDELIVERED_EVENT } from "./constants.js"; +import { + ESCALATION_EVENT, + OUTBOX_DRAIN_JOB, + OUTBOX_STATUS, + STATE_KEYS, + TELEGRAM_UNDELIVERED_EVENT, +} from "./constants.js"; +import { + drainOutbox, + findCreateRecord, + issueScope, + readMirroredNumber, + recordIntent, + recordFailed, + recordMirrored, + recordUncertain, +} from "./outbox.js"; import { formatBody, formatBudgetComment, @@ -56,15 +72,6 @@ async function clientFor( }); } -function issueScope(issueId: string) { - return { scopeKind: "issue" as const, scopeId: issueId }; -} - -async function readMirroredNumber(ctx: PluginContext, issueId: string): Promise { - const stored = await ctx.state.get({ ...issueScope(issueId), stateKey: STATE_KEYS.mirroredNumber }); - return typeof stored === "number" ? stored : null; -} - /** * Handlers must never take the worker down: a GitHub outage or a revoked token * is an observability problem, not a reason to stop processing Paperclip events. @@ -100,7 +107,12 @@ const plugin = definePlugin({ } }; - /** Creates the GitHub issue once and remembers its number. Idempotent. */ + /** + * Creates the GitHub issue once and remembers its number. Idempotent, and + * where it cannot be idempotent it refuses rather than guesses — see + * `outbox.ts` for why a create is the one step that must be written down + * before it is attempted. + */ const ensureMirrored = async (event: PluginEvent): Promise => { const issueId = event.entityId; if (!issueId) return; @@ -109,16 +121,46 @@ const plugin = definePlugin({ if (!config) return; if (await readMirroredNumber(ctx, issueId)) return; + // No number in state, but a record of an earlier attempt: that attempt + // may already have created the issue. Posting again would duplicate it, + // and the mirror will not read GitHub back to find out which it is. + const attempted = await findCreateRecord(ctx, event.companyId, issueId); + if (attempted && attempted.status !== OUTBOX_STATUS.failed) { + if (attempted.status === OUTBOX_STATUS.pending) { + await recordUncertain(ctx, attempted); + ctx.logger.error( + "GitHub mirror: an earlier create was never confirmed — refusing to create a second issue", + { issueId, companyId: event.companyId }, + ); + } + return; + } + // A `failed` record falls through on purpose: GitHub answered, so the + // issue does not exist and creating it now cannot duplicate anything. + const issue = await ctx.issues.get(issueId, event.companyId); if (!issue) return; + const intent = await recordIntent(ctx, event.companyId, issueId, formatTitle(issue)); + const github = await clientFor(ctx, event.companyId, config); - const created = await github.createIssue({ - title: formatTitle(issue), - body: formatBody(issue), - labels: [statusLabel(issue.status)], - }); + let created; + try { + created = await github.createIssue({ + title: formatTitle(issue), + body: formatBody(issue), + labels: [statusLabel(issue.status)], + }); + } catch (error) { + // Only a reply from GitHub proves nothing was created. Anything else — + // a timeout, a dead socket — leaves the outcome unknown, so the record + // stays `pending` for the drain to condemn. + if (error instanceof GithubApiError) await recordFailed(ctx, intent); + throw error; + } + // State first: it is what every later handler reads. The outbox record + // is closed after, and the drain repairs the gap if we die in between. await ctx.state.set( { ...issueScope(issueId), stateKey: STATE_KEYS.mirroredNumber }, created.number, @@ -127,6 +169,8 @@ const plugin = definePlugin({ { ...issueScope(issueId), stateKey: STATE_KEYS.lastStatus }, issue.status, ); + await recordMirrored(ctx, intent, created.number); + ctx.logger.info("Mirrored Paperclip issue to GitHub", { issueId, githubIssue: created.number, @@ -283,6 +327,12 @@ const plugin = definePlugin({ }), ); + // Resolves creates that were interrupted between the POST and the state + // write. Plugin-local only: it reads state and entities, never GitHub. + ctx.jobs.register(OUTBOX_DRAIN_JOB, () => + guard(ctx, OUTBOX_DRAIN_JOB, () => drainOutbox(ctx)), + ); + ctx.logger.info("GitHub mirror ready"); }, diff --git a/packages/plugins/paperclip-plugin-github-mirror/tests/plugin.spec.ts b/packages/plugins/paperclip-plugin-github-mirror/tests/plugin.spec.ts index 91fc13ecbf10..391ff999feed 100644 --- a/packages/plugins/paperclip-plugin-github-mirror/tests/plugin.spec.ts +++ b/packages/plugins/paperclip-plugin-github-mirror/tests/plugin.spec.ts @@ -1,11 +1,16 @@ -import { describe, expect, it } from "vitest"; -import { createTestHarness } from "@paperclipai/plugin-sdk/testing"; +import { describe, expect, it, vi } from "vitest"; +import { createTestHarness, type TestHarness } from "@paperclipai/plugin-sdk/testing"; import type { Issue } from "@paperclipai/shared"; import manifest from "../src/manifest.js"; import plugin from "../src/worker.js"; import { GithubApiError, GithubClient, parseRepository } from "../src/github.js"; import { formatBody, formatTitle, githubStateFor, statusLabel } from "../src/mirror.js"; -import { STATE_KEYS } from "../src/constants.js"; +import { + OUTBOX_DRAIN_JOB, + OUTBOX_ENTITY_TYPE, + OUTBOX_STATUS, + STATE_KEYS, +} from "../src/constants.js"; const COMPANY_ID = "c_1"; @@ -42,8 +47,23 @@ function makeIssue(overrides: Partial = {}): Issue { type HttpFetch = (url: string, init?: RequestInit) => Promise; -/** Records every outbound call and answers with a canned GitHub response. */ -function stubFetch() { +function okResponse() { + return new Response(JSON.stringify({ number: 77, html_url: "https://github.com/o/r/issues/77" }), { + status: 200, + headers: { "content-type": "application/json" }, + }); +} + +/** A canned failure, as GitHub would send it. */ +function errorResponse(status: number, headers: Record = {}) { + return () => new Response("nope", { status, headers }); +} + +/** + * Records every outbound call. Answers from `queue` while it lasts, then with a + * successful create — so a test only has to spell out the failures it cares about. + */ +function stubFetch(queue: Array<() => Response> = []) { const calls: Array<{ url: string; method: string; body: unknown; headers: Record }> = []; const impl = async (url: string, init?: RequestInit): Promise => { const headers = (init?.headers ?? {}) as Record; @@ -53,15 +73,56 @@ function stubFetch() { body: init?.body ? JSON.parse(String(init.body)) : undefined, headers, }); - return new Response(JSON.stringify({ number: 77, html_url: "https://github.com/o/r/issues/77" }), { - status: 200, - headers: { "content-type": "application/json" }, - }); + const next = queue.shift(); + return next ? next() : okResponse(); }; return { calls, impl }; } -async function setupHarness(config: Record = {}) { +/** + * Drives a promise that sleeps between retries without spending the backoff in + * real time. Only `setTimeout` is faked: the retry budget reads `Date.now()`, + * and faking that too would make every attempt look instantaneous. + */ +async function withoutBackoffDelays(start: () => Promise): Promise { + vi.useFakeTimers({ toFake: ["setTimeout"] }); + try { + const promise = start(); + let settled = false; + // Marks completion without producing a second promise: a rejecting one + // would be reported as unhandled before the loop below reaches its `await`. + void promise.then( + () => { + settled = true; + }, + () => { + settled = true; + }, + ); + // Each pass flushes microtasks and fires any backoff timer that is now due. + for (let i = 0; i < 20 && !settled; i += 1) { + await vi.advanceTimersByTimeAsync(10_000); + } + return await promise; + } finally { + vi.useRealTimers(); + } +} + +/** The create record the outbox keeps for a Paperclip issue, if any. */ +async function outboxRecord(harness: TestHarness, issueId = "iss_1") { + const [record] = await harness.ctx.entities.list({ + entityType: OUTBOX_ENTITY_TYPE, + externalId: `${COMPANY_ID}:${issueId}`, + limit: 1, + }); + return record ?? null; +} + +async function setupHarness( + config: Record = {}, + responses: Array<() => Response> = [], +) { const harness = createTestHarness({ manifest, capabilities: [...manifest.capabilities, "events.emit"], @@ -72,7 +133,7 @@ async function setupHarness(config: Record = {}) { ...config, }, }); - const fetcher = stubFetch(); + const fetcher = stubFetch(responses); // The harness performs a real network fetch by default — replace it so the // suite stays offline and can assert on the exact requests. harness.ctx.http.fetch = fetcher.impl as HttpFetch; @@ -146,13 +207,54 @@ describe("github client", () => { fetchImpl: async () => new Response("nope", { status, headers }), }); - await expect(make(429).addComment(1, "x")).rejects.toMatchObject({ retryable: true }); await expect( - make(403, { "x-ratelimit-remaining": "0" }).addComment(1, "x"), + withoutBackoffDelays(() => make(429).addComment(1, "x")), + ).rejects.toMatchObject({ retryable: true }); + await expect( + withoutBackoffDelays(() => make(403, { "x-ratelimit-remaining": "0" }).addComment(1, "x")), ).rejects.toMatchObject({ retryable: true }); await expect(make(422).addComment(1, "x")).rejects.toBeInstanceOf(GithubApiError); await expect(make(422).addComment(1, "x")).rejects.toMatchObject({ retryable: false }); }); + + it("reads the wait GitHub asks for, and ignores one it cannot use", async () => { + const make = (status: number, headers: Record) => + new GithubClient({ + repository: "acme/content", + token: "t", + fetchImpl: async () => new Response("nope", { status, headers }), + }); + + // `retry-after` is in seconds. + await expect( + withoutBackoffDelays(() => make(429, { "retry-after": "2" }).addComment(1, "x")), + ).rejects.toMatchObject({ retryAfterMs: 2000 }); + // `x-ratelimit-reset` is an epoch timestamp; one in the past says nothing. + await expect( + withoutBackoffDelays(() => make(429, { "x-ratelimit-reset": "1" }).addComment(1, "x")), + ).rejects.toMatchObject({ retryAfterMs: null }); + }); + + it("retries a worker→host call that timed out without an answer", async () => { + // The SDK drops `signal` when it serializes `init`, so the plugin cannot + // bound the call itself; a timeout arrives as a rejection like this one. + let attempts = 0; + const client = new GithubClient({ + repository: "acme/content", + token: "t", + fetchImpl: async () => { + attempts += 1; + if (attempts === 1) { + throw new Error('Worker→host call "http.fetch" timed out after 30000ms'); + } + return okResponse(); + }, + }); + + await withoutBackoffDelays(() => client.addComment(1, "x")); + + expect(attempts).toBe(2); + }); }); describe("mirror behaviour", () => { @@ -246,7 +348,9 @@ describe("mirror behaviour", () => { await plugin.definition.setup(harness.ctx); await expect( - harness.emit("issue.created", {}, { entityId: "iss_1", companyId: COMPANY_ID }), + withoutBackoffDelays(() => + harness.emit("issue.created", {}, { entityId: "iss_1", companyId: COMPANY_ID }), + ), ).resolves.not.toThrow(); expect(harness.logs.some((l) => l.level === "error")).toBe(true); @@ -330,6 +434,197 @@ describe("mirror behaviour", () => { expect(fetcher.calls.length).toBe(before); }); + it("retries a 502 and still ends up with exactly one GitHub issue", async () => { + const { harness, fetcher } = await setupHarness({}, [errorResponse(502)]); + + await withoutBackoffDelays(() => + harness.emit("issue.created", {}, { entityId: "iss_1", companyId: COMPANY_ID }), + ); + + expect(fetcher.calls.filter((c) => c.method === "POST")).toHaveLength(2); + expect( + harness.getState({ scopeKind: "issue", scopeId: "iss_1", stateKey: STATE_KEYS.mirroredNumber }), + ).toBe(77); + expect((await outboxRecord(harness))?.status).toBe(OUTBOX_STATUS.done); + }); + + it("does not retry a request GitHub rejected outright", async () => { + // A 422 fails identically on the second attempt; retrying only wastes the + // handler's budget. + const { harness, fetcher } = await setupHarness({}, [ + errorResponse(422), + errorResponse(422), + errorResponse(422), + ]); + + await withoutBackoffDelays(() => + harness.emit("issue.created", {}, { entityId: "iss_1", companyId: COMPANY_ID }), + ); + + expect(fetcher.calls.filter((c) => c.method === "POST")).toHaveLength(1); + }); + + it("gives up after three attempts and leaves the create recorded", async () => { + const { harness, fetcher } = await setupHarness({}, [ + errorResponse(429), + errorResponse(429), + errorResponse(429), + errorResponse(429), + ]); + + await withoutBackoffDelays(() => + harness.emit("issue.created", {}, { entityId: "iss_1", companyId: COMPANY_ID }), + ); + + expect(fetcher.calls.filter((c) => c.method === "POST")).toHaveLength(3); + // GitHub answered every time, so nothing was created and we know it. The + // record says `failed`, not `pending`: there is nothing to be uncertain + // about, and the drain must not condemn it. + expect((await outboxRecord(harness))?.status).toBe(OUTBOX_STATUS.failed); + expect( + harness.getState({ scopeKind: "issue", scopeId: "iss_1", stateKey: STATE_KEYS.mirroredNumber }), + ).toBeFalsy(); + }); + + it("creates the issue on a later event after GitHub refused the first attempt", async () => { + // The regression this guards: treating an exhausted retry the same as a lost + // outcome would refuse forever, and the task would never be mirrored at all — + // strictly worse than the behaviour before the outbox existed, where a failed + // create was simply retried on the next event. + const { harness, fetcher } = await setupHarness({}, [ + errorResponse(429), + errorResponse(429), + errorResponse(429), + ]); + + await withoutBackoffDelays(() => + harness.emit("issue.created", {}, { entityId: "iss_1", companyId: COMPANY_ID }), + ); + expect((await outboxRecord(harness))?.status).toBe(OUTBOX_STATUS.failed); + + await harness.emit("issue.created", {}, { entityId: "iss_1", companyId: COMPANY_ID }); + + expect(fetcher.calls.filter((c) => c.method === "POST")).toHaveLength(4); + expect((await outboxRecord(harness))?.status).toBe(OUTBOX_STATUS.done); + expect( + harness.getState({ scopeKind: "issue", scopeId: "iss_1", stateKey: STATE_KEYS.mirroredNumber }), + ).toBe(77); + }); + + it("leaves a lost outcome pending even though a refusal would be failed", async () => { + // Same handler, opposite evidence: no reply from GitHub means the create may + // have landed, so this one must stay `pending` for the drain. + const { harness } = await setupHarness({}, []); + harness.ctx.http.fetch = (async () => { + throw new Error("worker→host call timed out after 30000ms"); + }) as typeof harness.ctx.http.fetch; + + await withoutBackoffDelays(() => + harness.emit("issue.created", {}, { entityId: "iss_1", companyId: COMPANY_ID }), + ); + + expect((await outboxRecord(harness))?.status).toBe(OUTBOX_STATUS.pending); + }); + + it("refuses a second create when the first one's outcome was lost", async () => { + // The crash window: GitHub accepted the issue, then the state write died. + // Nothing on this side knows the number, and the mirror will not go and ask. + const { harness, fetcher } = await setupHarness(); + const realSet = harness.ctx.state.set.bind(harness.ctx.state); + let failed = false; + harness.ctx.state.set = async (input, value) => { + if (!failed) { + failed = true; + throw new Error("state write lost"); + } + return realSet(input, value); + }; + + await harness.emit("issue.created", {}, { entityId: "iss_1", companyId: COMPANY_ID }); + await harness.emit("issue.created", {}, { entityId: "iss_1", companyId: COMPANY_ID }); + + expect(fetcher.calls.filter((c) => c.method === "POST")).toHaveLength(1); + expect((await outboxRecord(harness))?.status).toBe(OUTBOX_STATUS.uncertain); + expect( + harness.logs.some((l) => l.level === "error" && l.message.includes("refusing to create")), + ).toBe(true); + + // The drain does not try to resolve it either: the mirror is write-only, so + // asking GitHub what happened is not an option it has. + const callsBeforeDrain = fetcher.calls.length; + await harness.runJob(OUTBOX_DRAIN_JOB); + + expect(fetcher.calls).toHaveLength(callsBeforeDrain); + expect((await outboxRecord(harness))?.status).toBe(OUTBOX_STATUS.uncertain); + }); + + it("closes a create whose record was never closed, without calling GitHub", async () => { + // The other half of the same window: the number reached state, the record + // did not. That one is knowable from plugin-local state alone. + const { harness, fetcher } = await setupHarness(); + const realUpsert = harness.ctx.entities.upsert.bind(harness.ctx.entities); + let upserts = 0; + harness.ctx.entities.upsert = async (input) => { + upserts += 1; + if (upserts === 2) throw new Error("entity write lost"); + return realUpsert(input); + }; + + await harness.emit("issue.created", {}, { entityId: "iss_1", companyId: COMPANY_ID }); + + expect( + harness.getState({ scopeKind: "issue", scopeId: "iss_1", stateKey: STATE_KEYS.mirroredNumber }), + ).toBe(77); + expect((await outboxRecord(harness))?.status).toBe(OUTBOX_STATUS.pending); + + const callsBeforeDrain = fetcher.calls.length; + await harness.runJob(OUTBOX_DRAIN_JOB); + + expect(fetcher.calls).toHaveLength(callsBeforeDrain); + const record = await outboxRecord(harness); + expect(record?.status).toBe(OUTBOX_STATUS.done); + expect(record?.data).toMatchObject({ mirroredNumber: 77 }); + }); + + it("records a stale unfinished create as uncertain rather than retrying it", async () => { + const { harness, fetcher } = await setupHarness(); + const stale = new Date(Date.now() - 10 * 60 * 1000).toISOString(); + await harness.ctx.entities.upsert({ + entityType: OUTBOX_ENTITY_TYPE, + scopeKind: "issue", + scopeId: "iss_9", + externalId: `${COMPANY_ID}:iss_9`, + status: OUTBOX_STATUS.pending, + data: { companyId: COMPANY_ID, issueId: "iss_9", startedAt: stale }, + }); + + await harness.runJob(OUTBOX_DRAIN_JOB); + + expect((await outboxRecord(harness, "iss_9"))?.status).toBe(OUTBOX_STATUS.uncertain); + expect(fetcher.calls).toHaveLength(0); + expect( + harness.logs.some((l) => l.level === "error" && l.message.includes("cannot be confirmed")), + ).toBe(true); + }); + + it("leaves a create that is still in flight alone", async () => { + // The drain runs on a schedule and must not condemn a create that simply + // has not come back yet. + const { harness } = await setupHarness(); + await harness.ctx.entities.upsert({ + entityType: OUTBOX_ENTITY_TYPE, + scopeKind: "issue", + scopeId: "iss_9", + externalId: `${COMPANY_ID}:iss_9`, + status: OUTBOX_STATUS.pending, + data: { companyId: COMPANY_ID, issueId: "iss_9", startedAt: new Date().toISOString() }, + }); + + await harness.runJob(OUTBOX_DRAIN_JOB); + + expect((await outboxRecord(harness, "iss_9"))?.status).toBe(OUTBOX_STATUS.pending); + }); + it("reports a budget stop as a pause, not a failure", async () => { const { harness, fetcher } = await setupHarness(); await harness.emit("issue.created", {}, { entityId: "iss_1", companyId: COMPANY_ID });