From 9efd3fe177eac7136285893fa4cfb3bcfc0e961e Mon Sep 17 00:00:00 2001 From: Brain <1a02c72794dcd0f07058a353bc3a81f4028b8c77c92c87fce6d5c8b85970a20b@buzz.block.builderlab.xyz> Date: Sat, 12 Sep 2026 08:42:05 -0600 Subject: [PATCH 01/20] docs(workflows): establish capability and UI handoff contract Signed-off-by: Brain <1a02c72794dcd0f07058a353bc3a81f4028b8c77c92c87fce6d5c8b85970a20b@buzz.block.builderlab.xyz> --- docs/workflows.md | 85 +++++++++++++++++++++ package.json | 3 +- pnpm-lock.yaml | 34 ++++++--- src/features/workflows/types.ts | 130 ++++++++++++++++++++++++++++++++ 4 files changed, 240 insertions(+), 12 deletions(-) create mode 100644 docs/workflows.md create mode 100644 src/features/workflows/types.ts diff --git a/docs/workflows.md b/docs/workflows.md new file mode 100644 index 00000000..6d94e7b4 --- /dev/null +++ b/docs/workflows.md @@ -0,0 +1,85 @@ +# Workflows capability and bundled UI handoff + +Status: implementation contract, not a shipped or live-validated feature. +Wes approved session FOUNDATION wiring and workflow-only relay save/delete repair +on 2026-09-12 (Buzz event `768f982eb3295e1bcbc69614d61b38deb3dc608e62864a5c58df9d8453f7fac9`). + +## Ownership and base + +App baseline: `17f90c18fff6b86bc029e710401fb2b60bc385ea`. +Brain owns `src/features/workflows/**`, relay transport/outbox/session integration, +`dev/` host adapters, catalogs, dependencies, this document, and the separate +legacy relay repair. Pinky owns `src/bundled/workflows/**` and adjacent UI/helper +tests in a separate worktree. No shared live-tree mutations. + +The type contract is [types.ts](../src/features/workflows/types.ts). +UI imports that capability by type and receives the captured session's +`workflows` property once integration lands; build/test UI compositions against +explicit fixture capabilities meanwhile. Do not implement an alternate host in +bundled code. Use the existing `pages` + `relay` injection and +`useRelayConnection`, not new plugins/author API or a router. + +## First complete UI slice + +- Channel-scoped **Saved configurations** list and raw YAML detail. Definitions + expose canonical author/channel/UUID, signed event revision, timestamp and raw + text. They do not fabricate runtime status or execution authority. `partial` + signals the bounded query limit, not lifecycle verification. +- New drafts explicitly disabled. Message/reaction triggers; Send Message/Delay. + Reuse pure legacy YAML/form/duration/schedule/condition/template helpers and + their tests selectively. Use shared design-system components and read its + stewardship instructions before composition. Preserve unsupported YAML and + incomplete header edits; no rewrite merely on opening or changing selection. +- Save with existing definition for compare-and-swap; failed/conflicted/unknown + writes retain drafts. Show accepted delivery separately from domain success. + Delete requires confirmation and host availability; an old relay's generic + accepted kind-5 receipt does not prove deletion. +- Manual run and bounded real run/trace next, before advanced editor polish. + Only returned run ID correlates a run; never choose newest run as recovery. + Approval rows are read-only; their hash is not an approval token. +- Form and YAML share restrictions. Webhook create/transition cannot bypass the + one-time-secret capability. Never write secret into drafts, logs, ordinary + journal, messages or clipboard automatically. Secret reveal is optional until + its display lifecycle is tested; otherwise keep those saves unavailable. + +## Read and write lifetime + +Each read view starts idle; the UI subscribes and calls refresh on interest, +then disposes on unmount/selection change. It owns no background poll unless an +active-run detail is visible; pause on hidden and terminal state. Runs return one +20-row page with an opaque exact `(before,beforeId)` pair; dispose the prior page +before opening another. Never reconstruct the cursor from second-granularity rows. +Failures show retry and do not become empty, deleted or permission-denied guesses. + +Views and operations purge before access-change callbacks. Remount by captured +scope + generation. Draft keys use stable community/viewer/coordinate, never +just channel/UUID; generation is not a durable key. Disable/unmount releases UI +interest but neither disables server workflows nor discards accepted intent. + +`save`, `delete`, `trigger` return local signed-intent IDs synchronously; subscribe +to operations for outcome. The shared outbox journals intent/signature before +send. Event echo cannot cancel the only result-bearing receipt. Restored signed +intent never auto-runs; retry uses exactly that event ID and does not promise +recovery of lost secret/run receipts. Unknown outcome is actionable information, +not permission to automatically submit a new trigger. Bounded ephemeral result +state is separate from ordinary delivery persistence. + +## Relay compatibility and unresolved historical state + +The approved forward repair includes workflow save runtime/event transaction +consistency and timestamp-ordered atomic deletion using coordinate serialization. +It does not authorize generic command refactoring, blind old-delete replay, +destructive historical reconciliation, or a new lifecycle endpoint. Historical +configuration rows remain unverified; authorized runs reads prove presence only +at the read. A repaired forward-delete capability must be positively identified +before Delete is enabled. The exact host compatibility signal is implemented and +tested with that relay change; do not infer it from version strings or kind lists. + +## Acceptance + +Production-seam receipt ordering/duplicate/unknown tests; revision and permission +checks; access purge and A->B->A fencing; persistence and secret isolation; real +DB rollback/stale/concurrent deletion tests in the relay; UI keyboard/focus, +dirty-close/conflict drafts, narrow layouts and YAML ownership tests. Fixture +feedback can precede final package gates. Live identity/signing/destructive +workflow trials require a separate consented test, not this implementation approval. diff --git a/package.json b/package.json index ec19fb25..42cb1655 100644 --- a/package.json +++ b/package.json @@ -56,7 +56,8 @@ "react-markdown": "10.1.0", "remark-breaks": "4.0.0", "remark-gfm": "4.0.1", - "virtua": "0.51.0" + "virtua": "0.51.0", + "yaml": "2.8.3" }, "devDependencies": { "@biomejs/biome": "2.5.12", diff --git a/pnpm-lock.yaml b/pnpm-lock.yaml index bad39e68..b6ca9e40 100644 --- a/pnpm-lock.yaml +++ b/pnpm-lock.yaml @@ -74,6 +74,9 @@ importers: virtua: specifier: 0.51.0 version: 0.51.0(react-dom@19.2.8(react@19.2.8))(react@19.2.8) + yaml: + specifier: 2.8.3 + version: 2.8.3 devDependencies: '@biomejs/biome': specifier: 2.5.12 @@ -98,7 +101,7 @@ importers: version: 19.2.7(@types/react@19.2.18) '@vitejs/plugin-react': specifier: 6.1.1 - version: 6.1.1(vite@8.2.2(@types/node@24.13.3)(jiti@2.7.0)) + version: 6.1.1(vite@8.2.2(@types/node@24.13.3)(jiti@2.7.0)(yaml@2.8.3)) postcss: specifier: 8.5.28 version: 8.5.28 @@ -113,10 +116,10 @@ importers: version: 7.16.0 vite: specifier: 8.2.2 - version: 8.2.2(@types/node@24.13.3)(jiti@2.7.0) + version: 8.2.2(@types/node@24.13.3)(jiti@2.7.0)(yaml@2.8.3) vitest: specifier: 4.1.11 - version: 4.1.11(@types/node@24.13.3)(vite@8.2.2(@types/node@24.13.3)(jiti@2.7.0)) + version: 4.1.11(@types/node@24.13.3)(vite@8.2.2(@types/node@24.13.3)(jiti@2.7.0)(yaml@2.8.3)) packages: @@ -184,6 +187,7 @@ packages: engines: {node: '>=14.21.3'} cpu: [arm64] os: [linux] + libc: [glibc] '@biomejs/cli-linux-x64-musl@2.5.12': resolution: {integrity: sha512-8A0oDW58/w9f/PQNYuq0sGUZtGtGrkNF4Z6n0PUoXpLCshi85vtKTv1XSznQawhdE4MXJ8ufpzHXyLFe87M/+w==} @@ -1587,6 +1591,11 @@ packages: engines: {node: '>=8'} hasBin: true + yaml@2.8.3: + resolution: {integrity: sha512-AvbaCLOO2Otw/lW5bmh9d/WEdcDFdQp2Z2ZUH3pX9U2ihyUY0nvLv7J6TrWowklRGPYbB/IuIMfYgxaCPg5Bpg==} + engines: {node: '>= 14.6'} + hasBin: true + zwitch@2.0.4: resolution: {integrity: sha512-bXE4cR/kVZhKZX/RjPEflHaKVhUVl85noU3v6b8apfQEc1x4A+zBxjZ4lN8LqGd6WZ3dl98pY4o717VFmoPp+A==} @@ -2036,10 +2045,10 @@ snapshots: '@ungap/structured-clone@1.4.0': {} - '@vitejs/plugin-react@6.1.1(vite@8.2.2(@types/node@24.13.3)(jiti@2.7.0))': + '@vitejs/plugin-react@6.1.1(vite@8.2.2(@types/node@24.13.3)(jiti@2.7.0)(yaml@2.8.3))': dependencies: '@rolldown/pluginutils': 1.0.1 - vite: 8.2.2(@types/node@24.13.3)(jiti@2.7.0) + vite: 8.2.2(@types/node@24.13.3)(jiti@2.7.0)(yaml@2.8.3) '@vitest/expect@4.1.11': dependencies: @@ -2050,13 +2059,13 @@ snapshots: chai: 6.2.2 tinyrainbow: 3.1.1 - '@vitest/mocker@4.1.11(vite@8.2.2(@types/node@24.13.3)(jiti@2.7.0))': + '@vitest/mocker@4.1.11(vite@8.2.2(@types/node@24.13.3)(jiti@2.7.0)(yaml@2.8.3))': dependencies: '@vitest/spy': 4.1.11 estree-walker: 3.0.3 magic-string: 0.30.21 optionalDependencies: - vite: 8.2.2(@types/node@24.13.3)(jiti@2.7.0) + vite: 8.2.2(@types/node@24.13.3)(jiti@2.7.0)(yaml@2.8.3) '@vitest/pretty-format@4.1.11': dependencies: @@ -2951,7 +2960,7 @@ snapshots: react: 19.2.8 react-dom: 19.2.8(react@19.2.8) - vite@8.2.2(@types/node@24.13.3)(jiti@2.7.0): + vite@8.2.2(@types/node@24.13.3)(jiti@2.7.0)(yaml@2.8.3): dependencies: lightningcss: 1.33.0 picomatch: 4.0.7 @@ -2962,11 +2971,12 @@ snapshots: '@types/node': 24.13.3 fsevents: 2.3.3 jiti: 2.7.0 + yaml: 2.8.3 - vitest@4.1.11(@types/node@24.13.3)(vite@8.2.2(@types/node@24.13.3)(jiti@2.7.0)): + vitest@4.1.11(@types/node@24.13.3)(vite@8.2.2(@types/node@24.13.3)(jiti@2.7.0)(yaml@2.8.3)): dependencies: '@vitest/expect': 4.1.11 - '@vitest/mocker': 4.1.11(vite@8.2.2(@types/node@24.13.3)(jiti@2.7.0)) + '@vitest/mocker': 4.1.11(vite@8.2.2(@types/node@24.13.3)(jiti@2.7.0)(yaml@2.8.3)) '@vitest/pretty-format': 4.1.11 '@vitest/runner': 4.1.11 '@vitest/snapshot': 4.1.11 @@ -2983,7 +2993,7 @@ snapshots: tinyexec: 1.3.1 tinyglobby: 0.2.17 tinyrainbow: 3.1.1 - vite: 8.2.2(@types/node@24.13.3)(jiti@2.7.0) + vite: 8.2.2(@types/node@24.13.3)(jiti@2.7.0)(yaml@2.8.3) why-is-node-running: 2.3.0 optionalDependencies: '@types/node': 24.13.3 @@ -2995,4 +3005,6 @@ snapshots: siginfo: 2.0.0 stackback: 0.0.2 + yaml@2.8.3: {} + zwitch@2.0.4: {} diff --git a/src/features/workflows/types.ts b/src/features/workflows/types.ts new file mode 100644 index 00000000..37cf0711 --- /dev/null +++ b/src/features/workflows/types.ts @@ -0,0 +1,130 @@ +import type { Delivery } from "../relay/outbox"; + +/** Canonical signed coordinate, bound to the owning community session. */ +export type WorkflowReference = Readonly<{ + id: string; + owner: string; + channelId: string; +}>; + +/** Configuration intent, not proof of runtime presence, enabled state or authority. */ +export type WorkflowDefinition = WorkflowReference & + Readonly<{ + revision: string; + createdAt: number; + yaml: string; + }>; + +export type WorkflowDefinitions = Readonly<{ + items: readonly WorkflowDefinition[]; + /** A bounded configuration snapshot is not a complete runtime inventory. */ + partial: boolean; +}>; + +/** Views belong to a captured session; dispose only releases this read interest. */ +export interface WorkflowView { + snapshot(): Readonly<{ + status: "idle" | "loading" | "ready" | "error" | "unavailable"; + data: T; + error?: string; + }>; + subscribe(listener: () => void): () => void; + refresh(): Promise; + dispose(): void; +} + +/** Preserve the relay cursor pair verbatim, including fractional timestamp precision. */ +export type WorkflowRunCursor = Readonly<{ before: string; beforeId: string }>; +export type WorkflowRun = Readonly<{ + id: string; + workflowId: string; + status: + | "pending" + | "running" + | "waiting_approval" + | "completed" + | "failed" + | "cancelled"; + currentStep: number; + trace: readonly unknown[]; + startedAt: number | null; + completedAt: number | null; + createdAt: number; + errorCode: string | null; + errorMessage: string | null; +}>; +export type WorkflowRunPage = Readonly<{ + runs: readonly WorkflowRun[]; + next: WorkflowRunCursor | null; +}>; +export type WorkflowApproval = Readonly<{ + /** Hashed reference, NEVER an actionable approval token. */ + reference: string; + runId: string; + stepId: string; + status: "pending" | "granted" | "denied" | "expired"; + note: string | null; + createdAt: number; +}>; + +/** Delivery evidence and domain outcome are deliberately separate. No secret here. */ +export type WorkflowOperation = Readonly<{ + eventId: string; + workflow: WorkflowReference; + action: "save" | "delete" | "trigger"; + delivery: Delivery; + outcome: "pending" | "succeeded" | "rejected" | "unknown"; + error?: string; + runId?: string; + /** A secret can be consumed once, never journaled or automatically copied. */ + secretAvailable: boolean; +}>; + +/** Host availability, NOT per-row permission; the relay remains authoritative. */ +export type WorkflowAvailability = Readonly<{ + definitions: boolean; + history: boolean; + save: boolean; + trigger: boolean; + /** Requires the repaired relay lifecycle contract, not merely kind-5 support. */ + delete: boolean; + webhookSecrets: boolean; +}>; + +/** Bundled UI contract. No socket, signer, arbitrary HTTP, scheduler or approval writes. */ +export interface WorkflowCapability { + readonly availability: WorkflowAvailability; + definitions(channelId: string): WorkflowView; + runs( + workflow: WorkflowReference, + cursor?: WorkflowRunCursor, + ): WorkflowView; + approvals( + workflow: WorkflowReference, + runId: string, + ): WorkflowView; + /** Synchronous local intent ID; follow operations for delivery and domain completion. + * Existing definitions preserve author/channel/id and use their signed revision. + * New definitions get a new UUID. YAML mode uses the same host restrictions. + */ + save( + input: Readonly<{ + channelId: string; + yaml: string; + existing?: WorkflowDefinition; + }>, + ): string; + delete(workflow: WorkflowDefinition): string; + trigger(workflow: WorkflowDefinition): string; + operations: Readonly<{ + snapshot(): readonly WorkflowOperation[]; + subscribe(listener: () => void): () => void; + /** Exact signed replay only; never creates a new event or recovers a lost receipt. */ + retry(eventId: string): void; + dismiss(eventId: string): Promise; + }>; + /** Consume in response to explicit reveal; caller must clear display on scope/access loss. + * Returns undefined after consumption, clear-cache, access revocation or disposal. + */ + takeWebhookSecret(eventId: string): string | undefined; +} From ced79b9b9416b8475dc94219c33e4505a6fc9376 Mon Sep 17 00:00:00 2001 From: Brain <1a02c72794dcd0f07058a353bc3a81f4028b8c77c92c87fce6d5c8b85970a20b@buzz.block.builderlab.xyz> Date: Sat, 12 Sep 2026 08:57:00 -0600 Subject: [PATCH 02/20] feat(workflows): add session history and receipt-safe command foundation Signed-off-by: Brain <1a02c72794dcd0f07058a353bc3a81f4028b8c77c92c87fce6d5c8b85970a20b@buzz.block.builderlab.xyz> --- dev/relay-broker.mjs | 73 +++- dev/workflow-broker.test.mjs | 223 ++++++++++ docs/workflows.md | 26 ++ src/features/relay/outbox-receipts.test.ts | 231 +++++++++++ src/features/relay/outbox.ts | 114 +++-- src/features/relay/receipt.ts | 21 + src/features/relay/session.ts | 35 +- src/features/relay/transport.ts | 75 +++- src/features/workflows/capability.test.ts | 208 ++++++++++ src/features/workflows/capability.ts | 460 +++++++++++++++++++++ src/features/workflows/host.ts | 13 + src/features/workflows/http.test.ts | 94 +++++ src/features/workflows/http.ts | 97 +++++ src/features/workflows/protocol.test.ts | 126 ++++++ src/features/workflows/protocol.ts | 286 +++++++++++++ src/features/workflows/session.test.ts | 195 +++++++++ 16 files changed, 2220 insertions(+), 57 deletions(-) create mode 100644 dev/workflow-broker.test.mjs create mode 100644 src/features/relay/outbox-receipts.test.ts create mode 100644 src/features/relay/receipt.ts create mode 100644 src/features/workflows/capability.test.ts create mode 100644 src/features/workflows/capability.ts create mode 100644 src/features/workflows/host.ts create mode 100644 src/features/workflows/http.test.ts create mode 100644 src/features/workflows/http.ts create mode 100644 src/features/workflows/protocol.test.ts create mode 100644 src/features/workflows/protocol.ts create mode 100644 src/features/workflows/session.test.ts diff --git a/dev/relay-broker.mjs b/dev/relay-broker.mjs index 9ec90c45..f2d15596 100644 --- a/dev/relay-broker.mjs +++ b/dev/relay-broker.mjs @@ -1,3 +1,8 @@ +import { + workflowReadPath, + workflowReadText, +} from "../src/features/workflows/http.ts"; +import { readReceiptText } from "../src/features/relay/receipt.ts"; import { decodeReadState, signReadState, @@ -507,6 +512,7 @@ export function relayBrokerPlugin({ ...(await getAuthority(relay)), relayUrl: relay, writeKinds: [9], + workflowReads: true, sidebarPreferences: true, readState: true, agentLibrary: true, @@ -697,6 +703,8 @@ export function relayBrokerPlugin({ "/api/relay/claim", "/api/relay/accept-policy", "/api/relay/gifs", + "/api/relay/workflow-runs", + "/api/relay/workflow-approvals", ].includes(route) || req.method !== "POST" ) @@ -713,6 +721,25 @@ export function relayBrokerPlugin({ } catch { return json(res, 400, { error: "Filter body is not JSON" }); } + let workflowPath; + if ( + [ + "/api/relay/workflow-runs", + "/api/relay/workflow-approvals", + ].includes(route) + ) { + try { + workflowPath = workflowReadPath( + route.slice("/api/relay/".length), + filters, + ); + } catch { + return json(res, 400, { + error: "Invalid workflow read", + sent: false, + }); + } + } const profile = route === "/api/relay/profile"; const claim = route === "/api/relay/claim"; const policy = route === "/api/relay/accept-policy"; @@ -836,6 +863,7 @@ export function relayBrokerPlugin({ !claim && !policy && !gifs && + !workflowPath && !readPublishing && !snapshot && !validFilters(filters) @@ -844,15 +872,18 @@ export function relayBrokerPlugin({ const gifSearchPath = gifs ? await getGifSearchPath(relay) : null; if (gifs && !gifSearchPath) return json(res, 404, { error: "GIF search is unavailable" }); - const upstreamPath = gifs - ? gifSearchPath - : profile || publishing || readPublishing - ? "/events" - : claim - ? "/api/invites/claim" - : policy - ? "/api/invites/accept-policy" - : "/query"; + const upstreamPath = + workflowPath ?? + (gifs + ? gifSearchPath + : profile || publishing || readPublishing + ? "/events" + : claim + ? "/api/invites/claim" + : policy + ? "/api/invites/accept-policy" + : "/query"); + const method = workflowPath ? "GET" : "POST"; if (inflight >= MAX_INFLIGHT) return json(res, 429, { error: "Query concurrency limit", @@ -861,7 +892,7 @@ export function relayBrokerPlugin({ inflight++; try { const lane = admissions(relay, viewer).api; - const body = JSON.stringify(filters); + const body = workflowPath ? undefined : JSON.stringify(filters); // A browser that gave up (the client's ten-second deadline) must also release // this upstream request, or hung requests exhaust the inflight budget. const cancel = new AbortController(); @@ -891,11 +922,15 @@ export function relayBrokerPlugin({ content: "", tags: [ ["u", `${relay}${upstreamPath}`], - ["method", "POST"], - [ - "payload", - createHash("sha256").update(body).digest("hex"), - ], + ["method", method], + ...(body === undefined + ? [] + : [ + [ + "payload", + createHash("sha256").update(body).digest("hex"), + ], + ]), ["nonce", randomBytes(16).toString("hex")], ], }, @@ -907,7 +942,7 @@ export function relayBrokerPlugin({ connectsBefore = upstream.connects(); upstreamStart = performance.now(); return fetchUpstream(`${relay}${upstreamPath}`, { - method: "POST", + method, headers: { "Content-Type": "application/json", Authorization: @@ -933,7 +968,11 @@ export function relayBrokerPlugin({ const text = snapshot && response.ok ? await readSnapshotText(response) - : await response.text(); + : workflowPath && response.ok + ? await workflowReadText(response) + : publishing && response.ok + ? await readReceiptText(response) + : await response.text(); // The relay's own service time separates server work from network time. const relayMs = Number( response.headers.get("x-envoy-upstream-service-time"), diff --git a/dev/workflow-broker.test.mjs b/dev/workflow-broker.test.mjs new file mode 100644 index 00000000..702dbb1f --- /dev/null +++ b/dev/workflow-broker.test.mjs @@ -0,0 +1,223 @@ +import { createServer } from "node:http"; +import { setTimeout as delay } from "node:timers/promises"; +import { expect, it } from "vitest"; +import { getPublicKey, verifyEvent } from "nostr-tools"; +import { relayBrokerPlugin } from "./relay-broker.mjs"; +import { connectBrokerTransport } from "../src/features/relay/transport.ts"; +import { WORKFLOW_READ_BYTES } from "../src/features/workflows/http.ts"; + +const id = "11111111-1111-4111-8111-111111111111"; +const runId = "22222222-2222-4222-8222-222222222222"; +const cursor = { before: "2026-09-12T14:44:19.123456+00:00", beforeId: runId }; +async function harness( + respond = () => Response.json({ runs: [], next: null }), +) { + const key = new Uint8Array(32); + key[31] = 8; + const viewer = getPublicKey(key), + calls = [], + logs = []; + let handler; + const server = createServer((req, res) => { + if (!req.headers.origin) req.headers.origin = `http://${req.headers.host}`; + handler(req, res); + }); + await relayBrokerPlugin({ + relayUrl: "https://a.workflow.test", + communityAliases: JSON.stringify({ secondary: "https://b.workflow.test" }), + identity: () => key, + authority: async () => ({ relayAuthor: viewer }), + upstreamFetch: async (url, init) => { + const auth = JSON.parse( + Buffer.from(init.headers.Authorization.slice(6), "base64").toString(), + ); + expect(verifyEvent(auth)).toBe(true); + expect(auth.pubkey).toBe(viewer); + const call = { url: String(url), init, auth }; + calls.push(call); + return respond(call); + }, + }).configureServer({ + httpServer: server, + config: { + logger: { + info() {}, + error(text) { + logs.push(text); + }, + }, + }, + middlewares: { + use(fn) { + handler = fn; + }, + }, + }); + await new Promise((resolve) => server.listen(0, "127.0.0.1", resolve)); + const base = `http://127.0.0.1:${server.address().port}`; + return { + base, + calls, + logs, + post: (route, body, headers = {}) => + fetch(`${base}/api/relay/${route}`, { + method: "POST", + headers: { "Content-Type": "application/json", ...headers }, + body: JSON.stringify(body), + }), + async close() { + server.closeAllConnections(); + await new Promise((resolve) => server.close(resolve)); + }, + }; +} +const signal = () => new AbortController().signal; +it("real broker scoped history signs exact GET path/cursor and captured principal without startup reads or workflow writes", async () => { + const h = await harness(({ url }) => + Response.json( + url.endsWith("approvals") ? { approvals: [] } : { runs: [], next: null }, + ), + ); + try { + const first = await connectBrokerTransport(h.base); + const other = await connectBrokerTransport(h.base, undefined, "secondary"); + expect(h.calls).toHaveLength(0); + expect(first.writer.kinds).toEqual([9]); + expect(first.workflows.lifecycleVersion).toBeUndefined(); + await first.workflows.runs(id, cursor, signal()); + await other.workflows.approvals(id, runId, signal()); + await first.workflows.runs(id, undefined, signal()); + expect(h.calls.map((call) => call.url)).toEqual([ + `https://a.workflow.test/workflows/${id}/runs?limit=20&before=2026-09-12T14%3A44%3A19.123456%2B00%3A00&before_id=${runId}`, + `https://b.workflow.test/workflows/${id}/runs/${runId}/approvals`, + `https://a.workflow.test/workflows/${id}/runs?limit=20`, + ]); + for (const { url, init, auth } of h.calls) { + expect(init.method).toBe("GET"); + expect(init.body).toBeUndefined(); + expect(init.redirect).toBe("error"); + expect(auth.tags).toContainEqual(["u", url]); + expect(auth.tags).toContainEqual(["method", "GET"]); + expect(auth.tags.some(([name]) => name === "payload")).toBe(false); + expect(auth.tags.some(([name]) => name === "nonce")).toBe(true); + } + expect(new Set(h.calls.map(({ auth }) => auth.id)).size).toBe(3); + } finally { + await h.close(); + } +}); +it("broker rejects arbitrary targets, cursors, limits and untrusted origin before upstream dispatch", async () => { + const h = await harness(); + try { + for (const body of [ + null, + { id: "../secret" }, + { id, limit: 100 }, + { id, url: "https://evil.test" }, + { id, cursor: { before: cursor.before } }, + { id, cursor: { ...cursor, before: "not a date" } }, + { id, cursor: { ...cursor, beforeId: "../" } }, + { id, cursor: null }, + ]) { + expect((await h.post("workflow-runs", body)).status).toBe(400); + } + expect( + (await h.post("workflow-approvals", { id, runId: "../" })).status, + ).toBe(400); + expect( + (await h.post("workflow-runs", { id }, { Origin: "https://evil.test" })) + .status, + ).toBe(403); + expect(h.calls).toHaveLength(0); + } finally { + await h.close(); + } +}); +it("workflow quota gates the existing query lane while another community remains independent", async () => { + const h = await harness(({ url }) => + url.startsWith("https://a.") + ? Response.json( + { error: "rate-limited: quota exceeded; retry in 0s" }, + { status: 429 }, + ) + : Response.json([]), + ); + try { + const transport = await connectBrokerTransport(h.base); + await expect( + transport.workflows.runs(id, undefined, signal()), + ).rejects.toMatchObject({ status: 429, retryAfterMs: 1000 }); + await expect( + transport.query([{ kinds: [0], limit: 1 }]), + ).rejects.toMatchObject({ status: 429 }); + expect(h.calls).toHaveLength(1); + const other = await connectBrokerTransport(h.base, undefined, "secondary"); + await other.query([{ kinds: [0], limit: 1 }]); + expect(h.calls).toHaveLength(2); + } finally { + await h.close(); + } +}); +it("workflow 404 remains unavailable, over-budget body is cancelled without leaking bytes", async () => { + let large = false, + cancelled = false; + const h = await harness(() => + large + ? new Response( + new ReadableStream({ + start(controller) { + controller.enqueue( + new TextEncoder().encode( + "PRIVATE".repeat(Math.ceil(WORKFLOW_READ_BYTES / 7) + 1), + ), + ); + }, + cancel() { + cancelled = true; + }, + }), + ) + : Response.json({ error: "PRIVATE missing workflow" }, { status: 404 }), + ); + try { + const transport = await connectBrokerTransport(h.base); + await expect( + transport.workflows.runs(id, undefined, signal()), + ).rejects.toMatchObject({ status: 404 }); + large = true; + await expect( + transport.workflows.runs(id, undefined, signal()), + ).rejects.toThrow(); + expect(cancelled).toBe(true); + expect(h.logs.join(" ")).not.toContain("PRIVATE"); + } finally { + await h.close(); + } +}); +it("closing workflow interest aborts the actual broker upstream request", async () => { + const h = await harness( + ({ init }) => + new Promise((_resolve, reject) => { + init.signal.addEventListener( + "abort", + () => reject(init.signal.reason), + { once: true }, + ); + }), + ); + try { + const transport = await connectBrokerTransport(h.base), + cancel = new AbortController(); + const pending = transport.workflows.runs(id, undefined, cancel.signal); + const rejection = expect(pending).rejects.toThrow(); + for (let i = 0; i < 100 && !h.calls.length; i++) await delay(5); + expect(h.calls).toHaveLength(1); + cancel.abort(); + await rejection; + for (let i = 0; i < 100 && !h.calls[0].init.signal.aborted; i++) + await delay(5); + expect(h.calls[0].init.signal.aborted).toBe(true); + } finally { + await h.close(); + } +}); diff --git a/docs/workflows.md b/docs/workflows.md index 6d94e7b4..0463a817 100644 --- a/docs/workflows.md +++ b/docs/workflows.md @@ -83,3 +83,29 @@ DB rollback/stale/concurrent deletion tests in the relay; UI keyboard/focus, dirty-close/conflict drafts, narrow layouts and YAML ownership tests. Fixture feedback can precede final package gates. Live identity/signing/destructive workflow trials require a separate consented test, not this implementation approval. + +## Host checkpoint (2026-09-12) + +The app implements lazy structured history reads in both signed and dev-broker +hosts. The broker exposes only `workflow-runs` / `workflow-approvals` POST inputs, +constructs fixed upstream GET routes, preserves the exact cursor pair, and shares +the captured principal's API admission. History capability means the adapter +exists, not that an older relay serves the endpoint: failures remain explicit. +Responses are stream-bounded to 1 MiB before parsing; command receipts to 16 KiB. + +All workflow writes remain unavailable in real host connections at this +checkpoint. Fixtures may supply `WorkflowHost.lifecycleVersion = 1` to exercise +commands. No real host advertises that evidence until the forward relay repair +and compatibility handshake are implemented and reviewed. Receipt tests do not +prove a deployed database transaction. + +Revocation puts existing and newly opened denied views in `unavailable`, purges +all data before callbacks, and cancels late results. A regrant requires explicit +fresh interest; it never revives an old snapshot. UI must not reopen a recovered +private draft from an unavailable view. Secret reveal remains disabled. + +A save receipt does not populate a read view. Refresh explicitly and match both +`operation.workflow` and `operation.eventId === definition.revision` before +permitting resave. A coordinate-only old/concurrent head is not save readback. +Generic kind-5 deletes retain their previous behavior: only workflow-coordinate +kind-5 operations use workflow validation and receipt semantics. diff --git a/src/features/relay/outbox-receipts.test.ts b/src/features/relay/outbox-receipts.test.ts new file mode 100644 index 00000000..a463a10e --- /dev/null +++ b/src/features/relay/outbox-receipts.test.ts @@ -0,0 +1,231 @@ +import { assert, afterEach, expect, it, vi } from "vitest"; +import { createOutbox, PublishRejected, type OutgoingEvent } from "./outbox"; +import type { RelayEvent } from "./events"; +import { flush, keypair, signed } from "./testing"; +import { readReceiptText } from "./receipt"; + +const owners: ReturnType[] = []; +afterEach(() => { + for (const owner of owners.splice(0)) owner.dispose(); + vi.useRealTimers(); +}); +function setup(timeoutMs = 10000) { + const key = keypair(); + let saved: readonly OutgoingEvent[] = []; + let settle: ((message: string) => void) | undefined; + let reject: ((error: Error) => void) | undefined; + let published: RelayEvent | undefined; + let signal: AbortSignal | undefined; + const onReceipt = vi.fn(); + const sign = vi.fn(async (template) => + signed(key, structuredClone(template)), + ); + const publish = vi.fn((event: RelayEvent, abort: AbortSignal) => { + expect(saved.find((row) => row.signed?.id === event.id)).toBeDefined(); + published = event; + signal = abort; + return new Promise((resolve, fail) => { + settle = resolve; + reject = fail; + }); + }); + const storage = { + load: () => saved, + save: (rows: readonly OutgoingEvent[]) => { + saved = structuredClone(rows); + }, + }; + const owner = createOutbox(key.pubkey, { sign, publish }, storage, { + needsReceipt: (event) => [30620, 46020].includes(event.kind), + onReceipt, + timeoutMs, + }); + owners.push(owner); + return { + ...owner, + key, + storage, + sign, + publish, + onReceipt, + saved: () => saved, + published: () => { + assert.exists(published); + return published; + }, + signal: () => { + assert.exists(signal); + return signal; + }, + settle: (message: string) => { + assert.exists(settle); + settle(message); + }, + reject: (error: Error) => { + assert.exists(reject); + reject(error); + }, + send: (kind = 30620) => + owner.outbox.send({ + kind, + content: "disabled workflow", + tags: [ + ["h", "channel"], + ["d", "workflow"], + ], + }), + }; +} +it("preserves command receipt after echo, without persisting secret or aborting publication", async () => { + const h = setup(); + const id = h.send(); + await flush(); + h.observe([h.published()]); + expect(h.signal().aborted).toBe(false); + expect(h.outbox.snapshot()[0]?.delivery).toBe("seen"); + expect(h.onReceipt).not.toHaveBeenCalled(); + h.settle('response:{"webhook_secret":"one-time-fixture"}'); + await flush(); + expect(h.onReceipt).toHaveBeenCalledExactlyOnceWith( + h.published(), + 'response:{"webhook_secret":"one-time-fixture"}', + ); + expect(h.outbox.snapshot()).toEqual([]); + expect(h.local.snapshot()[0]).toMatchObject({ + event: { id }, + delivery: "seen", + }); + expect(JSON.stringify(h.saved())).not.toContain("one-time-fixture"); +}); +it("receipt before echo is retained once and message echo keeps its existing cancellation behavior", async () => { + const h = setup(); + h.send(); + await flush(); + h.settle("response:{}"); + await flush(); + expect(h.outbox.snapshot()[0]?.delivery).toBe("accepted"); + h.observe([h.published()]); + expect(h.onReceipt).toHaveBeenCalledTimes(1); + const message = setup(); + message.send(9); + await flush(); + message.observe([message.published()]); + expect(message.signal().aborted).toBe(true); + message.settle("irrelevant"); + await flush(); + expect(message.onReceipt).not.toHaveBeenCalled(); +}); +it("echo plus lost receipt remains seen but settles command result as unavailable", async () => { + const h = setup(); + h.send(); + await flush(); + h.observe([h.published()]); + h.reject(new Error("connection lost")); + await flush(); + expect(h.local.snapshot()[0]?.delivery).toBe("seen"); + expect(h.onReceipt).toHaveBeenCalledExactlyOnceWith(h.published(), undefined); +}); +it("receipt waits are bounded and late results after disposal never publish", async () => { + vi.useFakeTimers(); + const h = setup(100); + h.send(); + await vi.advanceTimersByTimeAsync(0); + h.observe([h.published()]); + await vi.advanceTimersByTimeAsync(101); + expect(h.signal().aborted).toBe(true); + expect(h.onReceipt).toHaveBeenCalledExactlyOnceWith(h.published(), undefined); + const late = setup(); + late.send(); + await vi.advanceTimersByTimeAsync(0); + late.dispose(); + late.settle("secret"); + await vi.advanceTimersByTimeAsync(0); + expect(late.onReceipt).not.toHaveBeenCalled(); +}); +it("restores signed command intent without sending, exact retry receives only duplicate outcome", async () => { + const h = setup(); + const id = h.send(46020); + await flush(); + h.dispose(); + const receipt = vi.fn(); + const publish = vi.fn( + async (_event: RelayEvent) => "duplicate: already processed", + ); + const sign = vi.fn(async () => { + throw new Error("must not sign again"); + }); + const restored = createOutbox(h.key.pubkey, { sign, publish }, h.storage, { + needsReceipt: (event) => event.kind === 46020, + onReceipt: receipt, + }); + owners.push(restored); + await restored.ready; + expect(publish).not.toHaveBeenCalled(); + restored.outbox.retry(id); + await flush(); + expect(sign).not.toHaveBeenCalled(); + expect(publish.mock.calls[0]?.[0]).toEqual(h.published()); + expect(receipt).toHaveBeenCalledExactlyOnceWith( + h.published(), + "duplicate: already processed", + ); +}); +it("bounds receipt bytes while streaming before JSON decoding", async () => { + const cancel = vi.fn(); + const body = new ReadableStream({ + start(c) { + c.enqueue(new Uint8Array(17000)); + }, + cancel, + }); + await expect(readReceiptText(new Response(body))).rejects.toThrow( + "size limit", + ); + expect(cancel).toHaveBeenCalled(); + expect(await readReceiptText(new Response("response:{}"))).toBe( + "response:{}", + ); +}); + +it("seen commands can retry the exact signed event and dismiss without losing observation evidence", async () => { + const h = setup(); + const id = h.send(); + await flush(); + h.observe([h.published()]); + h.reject(new Error("lost")); + await flush(); + const original = h.published(); + expect(h.outbox.snapshot()).toEqual([]); + h.outbox.retry(id); + await flush(); + expect(h.sign).toHaveBeenCalledTimes(1); + expect(h.publish).toHaveBeenCalledTimes(2); + expect(h.published()).toEqual(original); + h.reject(new PublishRejected("response:{secret:PRIVATE}")); + await flush(); + expect(h.local.snapshot()[0]?.delivery).toBe("seen"); + expect(JSON.stringify(h.saved())).not.toContain("PRIVATE"); + await h.outbox.dismiss(id); + expect(h.local.snapshot()).toEqual([]); + expect(h.saved()).toEqual([]); +}); +it("rejection text never journals command secrets; successful seen retry stays seen", async () => { + const h = setup(); + const id = h.send(); + await flush(); + h.reject(new PublishRejected("PRIVATE")); + await flush(); + expect(h.outbox.snapshot()[0]?.delivery).toBe("failed"); + expect(JSON.stringify(h.saved())).not.toContain("PRIVATE"); + h.outbox.retry(id); + await flush(); + h.observe([h.published()]); + h.settle("response:{}"); + await flush(); + h.outbox.retry(id); + await flush(); + h.settle("duplicate:"); + await flush(); + expect(h.local.snapshot()[0]?.delivery).toBe("seen"); + expect(h.outbox.snapshot()).toEqual([]); +}); diff --git a/src/features/relay/outbox.ts b/src/features/relay/outbox.ts index 24ff813e..90a40b50 100644 --- a/src/features/relay/outbox.ts +++ b/src/features/relay/outbox.ts @@ -45,8 +45,13 @@ export function createOutbox( profiling = createRelayProfiler(), notifyListener = (listener: () => void) => listener(), preparePublish, + needsReceipt = () => false, + onReceipt = (_event: EventData, _message: string | undefined) => {}, }: { timeoutMs?: number; + /** Commands await their receipt even after a verified echo. Never persisted. */ + needsReceipt?: (event: EventData) => boolean; + onReceipt?: (event: EventData, message: string | undefined) => void; onAccepted?: (event: RelayEvent) => void; profiling?: RelayProfiler; notifyListener?: (listener: () => void) => void; @@ -58,6 +63,7 @@ export function createOutbox( ) => Promise<(() => void) | undefined>; } = {}, ) { + const awaitsReceipt = needsReceipt; let snapshot: readonly OutgoingEvent[] = Object.freeze([]); let visible: readonly OutgoingEvent[] = snapshot; let finalSnapshot: readonly OutgoingEvent[] | undefined; @@ -223,17 +229,20 @@ export function createOutbox( clearTimeout(attempt.timer); attempts.delete(id); const item = find(id); - if (!closed && item) + if (!closed && item) { + if (awaitsReceipt(item.event)) onReceipt(item.event, undefined); saveStatus({ ...item, delivery: failedDelivery(attempt), error: error.message, }); + } } } function failedDelivery(attempt: Attempt): Delivery { return attempt.previousDelivery === "unknown" || - attempt.previousDelivery === "accepted" + attempt.previousDelivery === "accepted" || + attempt.previousDelivery === "seen" ? attempt.previousDelivery : "failed"; } @@ -308,38 +317,55 @@ export function createOutbox( // transport publisher crosses that boundary, including for signed retries. if (closed || !find(id)) return; signal.throwIfAborted(); - await profiling.measureAsync("send.publish", id, () => { + const receipt = await profiling.measureAsync("send.publish", id, () => { check?.(); publishing = true; return Promise.race([writer.publish(signed, signal), aborted]); }); if (closed || signal.aborted) return; + if (awaitsReceipt(signed)) + onReceipt(signed, typeof receipt === "string" ? receipt : undefined); const latest = find(id); if (latest) - saveStatus({ ...latest, delivery: "accepted", error: undefined }); + saveStatus({ + ...latest, + delivery: + latest.delivery === "seen" || attempt.previousDelivery === "seen" + ? "seen" + : "accepted", + error: undefined, + }); onAccepted(signed); } catch (error) { const latest = find(id); // A verified observation ends the attempt even if its HTTP ACK never arrives. total(!closed && !latest ? "ok" : "error"); if (closed) return; + if (latest && awaitsReceipt(latest.event)) + onReceipt(latest.signed ?? latest.event, undefined); if (latest) saveStatus({ ...latest, delivery: - attempt.previousDelivery === "accepted" - ? "accepted" - : publishing && !(error instanceof PublishRejected) - ? "unknown" - : failedDelivery(attempt), - error: `${ - attempt.previousDelivery === "unknown" || - attempt.previousDelivery === "accepted" - ? error instanceof PublishRejected - ? "Retry blocked: " - : "Retry failed: " - : "" - }${error instanceof Error ? error.message : String(error)}`, + latest.delivery === "seen" || attempt.previousDelivery === "seen" + ? "seen" + : attempt.previousDelivery === "accepted" + ? "accepted" + : publishing && !(error instanceof PublishRejected) + ? "unknown" + : failedDelivery(attempt), + error: awaitsReceipt(latest.event) + ? publishing && !(error instanceof PublishRejected) + ? "Workflow delivery could not be confirmed; retain this operation to retry." + : "Workflow command rejected; retain the draft and refresh before retrying." + : `${ + attempt.previousDelivery === "unknown" || + attempt.previousDelivery === "accepted" + ? error instanceof PublishRejected + ? "Retry blocked: " + : "Retry failed: " + : "" + }${error instanceof Error ? error.message : String(error)}`, }); if (publishing && latest?.signed && !(error instanceof PublishRejected)) onAccepted(latest.signed); @@ -347,6 +373,15 @@ export function createOutbox( total(); clearTimeout(attempt.timer); if (attempts.get(id) === attempt) attempts.delete(id); + const observed = find(id); + if (!closed && observed?.delivery === "seen") { + completed.set(id, observed); + snapshot = Object.freeze( + snapshot.filter((item) => item.event.id !== id), + ); + notify(); + void persist(id).catch(() => {}); + } if (!closed) for (const queued of snapshot) if (queued.delivery === "sending") void deliver(queued.event.id); @@ -414,15 +449,28 @@ export function createOutbox( return event.id; }, retry(id: string) { - const item = find(id); + const retained = completed.peek(id); + const item = + find(id) ?? + (retained && awaitsReceipt(retained.event) ? retained : undefined); if (!closed && item && !attempts.has(id)) { + if (!find(id)) { + if (snapshot.length >= MAX_PENDING) + throw new Error("Too many outstanding operations"); + completed.delete(id); + snapshot = Object.freeze([...snapshot, item]); + } replace({ ...item, delivery: "sending", error: undefined }); schedule(id, undefined, item.delivery); } }, async dismiss(id: string) { if (closed || attempts.has(id)) return; - const previous = find(id); + const retained = completed.peek(id); + const previous = + find(id) ?? + (retained && awaitsReceipt(retained.event) ? retained : undefined); + if (previous && awaitsReceipt(previous.event)) completed.delete(id); snapshot = Object.freeze(snapshot.filter((item) => item.event.id !== id)); notify(); try { @@ -467,15 +515,33 @@ export function createOutbox( const [first] = confirmed; if (!first) return; for (const event of confirmed) { - completed.set( - event.id, - Object.freeze({ event, signed: event, delivery: "seen" }), - ); + if (awaitsReceipt(event) && attempts.has(event.id)) { + snapshot = Object.freeze( + snapshot.map((item) => + item.event.id === event.id + ? Object.freeze({ + ...item, + signed: event, + delivery: "seen" as const, + }) + : item, + ), + ); + } else + completed.set( + event.id, + Object.freeze({ event, signed: event, delivery: "seen" }), + ); } snapshot = Object.freeze( - snapshot.filter((item) => !byId.has(item.event.id)), + snapshot.filter( + (item) => + !byId.has(item.event.id) || + (awaitsReceipt(item.event) && attempts.has(item.event.id)), + ), ); for (const event of confirmed) { + if (awaitsReceipt(event) && attempts.has(event.id)) continue; const attempt = attempts.get(event.id); attempt?.controller?.abort(); clearTimeout(attempt?.timer); diff --git a/src/features/relay/receipt.ts b/src/features/relay/receipt.ts new file mode 100644 index 00000000..35dd5896 --- /dev/null +++ b/src/features/relay/receipt.ts @@ -0,0 +1,21 @@ +/** Bound command receipt bytes before parsing or retaining secret-bearing text. */ +export async function readReceiptText(response: Response): Promise { + if (!response.body) throw new Error("Relay delivery receipt body missing"); + const reader = response.body.getReader(); + const decoder = new TextDecoder("utf-8", { fatal: true }); + let bytes = 0; + let text = ""; + try { + while (true) { + const { value, done } = await reader.read(); + if (done) return text + decoder.decode(); + bytes += value.byteLength; + if (bytes > 16 * 1024) + throw new Error("Relay delivery receipt exceeds the size limit"); + text += decoder.decode(value, { stream: true }); + } + } finally { + await reader.cancel().catch(() => {}); + reader.releaseLock(); + } +} diff --git a/src/features/relay/session.ts b/src/features/relay/session.ts index abbac089..0411cea2 100644 --- a/src/features/relay/session.ts +++ b/src/features/relay/session.ts @@ -1,4 +1,6 @@ // FOUNDATION: One relay session owns reads, local intent, delivery and shared views. +import { createWorkflows } from "../workflows/capability"; +import { isWorkflowOperation } from "../workflows/protocol"; import { createRelayReader, type ReadOptions, @@ -99,6 +101,11 @@ export function createRelaySession( ...writer, async sign(template, signal) { validateMentionEvent(template); + workflows.validate({ + ...template, + id: "", + pubkey: transport.viewer, + }); return writer.sign(template, signal); }, }, @@ -113,7 +120,19 @@ export function createRelaySession( profiling, notifyListener: notify, onAccepted: (event) => confirm(event), - preparePublish: prepareMentionPublication, + needsReceipt: isWorkflowOperation, + onReceipt: (event, message) => workflows.receipt(event, message), + preparePublish: async (event, signal) => { + workflows.validate(event); + const checkMentions = await prepareMentionPublication( + event, + signal, + ); + return () => { + workflows.validate(event); + checkMentions?.(); + }; + }, }, ) : undefined; @@ -171,6 +190,7 @@ export function createRelaySession( emoji.clear(); agentLibrary.clear(); archives.clear(); + workflows.clear(); for (const purge of views.values()) purge(); commit(); unread.purge(); @@ -315,6 +335,15 @@ export function createRelaySession( }, ); canAccess = channels.canAccess; + const workflows = createWorkflows({ + reader: transport ? verified : undefined, + viewer: transport?.viewer ?? "", + outbox: writes?.outbox, + local: localViews, + host: transport?.workflows, + canAccess: (channelId) => canAccess(channelId), + notify, + }); const readScope = `${transport?.scope ?? transport?.relayAuthor ?? "offline"}:${transport?.viewer ?? ""}`; const reads = createReadState({ viewer: transport?.viewer ?? "", @@ -614,6 +643,7 @@ export function createRelaySession( profiles: profiles.queries, emoji: emoji.queries, agentLibrary: agentLibrary.queries, + workflows: workflows.capability, archives: archives.queries, media: (url: string) => transport?.media(url), /** A plugin may request writes from this same interface when the host supports them. */ @@ -895,6 +925,7 @@ export function createRelaySession( requests.invalidate(); agentLibrary.clear(); archives.clear(); + workflows.clear(); channels.staleHeads(); unread.stale(); } @@ -959,6 +990,7 @@ export function createRelaySession( emoji.clear(); agentLibrary.clear(); archives.clear(); + workflows.clear(); await channels.clearCache(); publishLive(); }, @@ -978,6 +1010,7 @@ export function createRelaySession( channels.dispose(); profiles.dispose(); emoji.dispose(); + workflows.dispose(); agentLibrary.dispose(); archives.dispose(); }, diff --git a/src/features/relay/transport.ts b/src/features/relay/transport.ts index d7b0738a..5868ef96 100644 --- a/src/features/relay/transport.ts +++ b/src/features/relay/transport.ts @@ -1,3 +1,6 @@ +import { workflowHost, workflowReadPath } from "../workflows/http"; +import type { WorkflowHost } from "../workflows/host"; +import { readReceiptText } from "./receipt"; import type { ReadStateHost, ReadStateSigning } from "./read-state-host"; import { parseReadSnapshot, @@ -31,9 +34,14 @@ import { eventDto, type ReadFilter, type RelayEvent } from "./events"; export interface RelayWriter { readonly kinds?: readonly number[]; sign(event: EventTemplate, signal: AbortSignal): Promise; - publish(event: RelayEvent, signal: AbortSignal): Promise; + /** Accepted receipt text is ephemeral; callers must never journal it. */ + publish( + event: RelayEvent, + signal: AbortSignal, + ): Promise | Promise; } export interface ReadTransport { + readonly workflows?: WorkflowHost; /** Host-projected local library; display only, never relay authority. */ readonly readAgentLibrary?: AgentLibraryReader; /** Host-only decoder of the viewer's two signed sidebar preference coordinates. */ @@ -141,6 +149,7 @@ export async function connectBrokerTransport( relayAuthor?: unknown; archiveAuthority?: unknown; writeKinds?: number[]; + workflowReads?: boolean; relayUrl?: string; live?: boolean; sidebarPreferences?: boolean; @@ -174,6 +183,19 @@ export async function connectBrokerTransport( ...(typeof session.archiveAuthority === "string" ? { archiveAuthority: session.archiveAuthority } : {}), + ...(session.workflowReads === true + ? { + workflows: workflowHost((route, body, signal) => + fetch(`${endpoint}/${route}`, { + method: "POST", + credentials: "same-origin", + headers: { "Content-Type": "application/json" }, + body: JSON.stringify(body), + signal, + }), + ), + } + : {}), ...(session.agentLibrary ? { readAgentLibrary: async (signal: AbortSignal) => { @@ -312,7 +334,7 @@ export async function connectBrokerTransport( signal, }); recordServerTiming(result, profiling, event.id); - await acceptPublish(result, event.id); + return acceptPublish(result, event.id); }, }, } @@ -398,6 +420,19 @@ export async function connectSignedTransport( }, }; }, + workflows: workflowHost((route, body, signal) => + signedRequest( + signer, + `${httpOrigin}${workflowReadPath(route, body)}`, + undefined, + signal, + profiling, + route, + principal().api, + "foreground", + "GET", + ), + ), scope: httpOrigin, viewer, relayAuthor, @@ -405,8 +440,8 @@ export async function connectSignedTransport( writer: { sign: (event) => signer.signEvent(event), async publish(event, signal) { - await acceptPublish( - await signedPost( + return acceptPublish( + await signedRequest( signer, `${httpOrigin}/events`, event, @@ -424,7 +459,7 @@ export async function connectSignedTransport( }, }, async query(filters, signal, requestId = "read", priority = "foreground") { - const result = await signedPost( + const result = await signedRequest( signer, `${httpOrigin}/query`, filters, @@ -453,7 +488,7 @@ export async function connectSignedTransport( }; } -async function signedPost( +async function signedRequest( signer: Signer, url: string, value: unknown, @@ -462,13 +497,20 @@ async function signedPost( id: string, admission: Parameters[0], priority: "foreground" | "background" = "foreground", + method: "POST" | "GET" = "POST", ) { signal?.throwIfAborted(); return admission.prepare(async () => { - const body = JSON.stringify(value); - const payload = hex( - await crypto.subtle.digest("SHA-256", new TextEncoder().encode(body)), - ); + const body = method === "POST" ? JSON.stringify(value) : undefined; + const payload = + body === undefined + ? undefined + : hex( + await crypto.subtle.digest( + "SHA-256", + new TextEncoder().encode(body), + ), + ); if (signal?.aborted) throw signal.reason; const auth = await profiling.measureAsync("http.auth", id, () => signer.signEvent({ @@ -477,8 +519,8 @@ async function signedPost( content: "", tags: [ ["u", url], - ["method", "POST"], - ["payload", payload], + ["method", method], + ...(payload === undefined ? [] : [["payload", payload]]), ["nonce", crypto.randomUUID()], ], }), @@ -499,12 +541,13 @@ async function signedPost( ); return profiling.measureAsync("http.fetch", id, () => fetch(url, { - method: "POST", + method, + redirect: "error", headers: { Authorization: `Nostr ${btoa(JSON.stringify(auth))}`, "Content-Type": "application/json", }, - body, + ...(body === undefined ? {} : { body }), signal: signal ?? null, }), ); @@ -533,7 +576,8 @@ async function acceptPublish(response: Response, id: string) { `Relay delivery could not be confirmed (${response.status})`, ); } - const result = (await response.json()) as { + const text = await readReceiptText(response); + const result = JSON.parse(text) as { accepted?: unknown; event_id?: unknown; message?: unknown; @@ -546,6 +590,7 @@ async function acceptPublish(response: Response, id: string) { ? result.message : "Relay rejected the message", ); + return typeof result.message === "string" ? result.message : ""; } function recordServerTiming( diff --git a/src/features/workflows/capability.test.ts b/src/features/workflows/capability.test.ts new file mode 100644 index 00000000..ab3a317e --- /dev/null +++ b/src/features/workflows/capability.test.ts @@ -0,0 +1,208 @@ +import { afterEach, expect, it, vi } from "vitest"; +import { createWorkflows } from "./capability"; +import { + createOutbox, + PublishRejected, + type OutgoingEvent, +} from "../relay/outbox"; +import { keypair, signed, flush } from "../relay/testing"; +import type { RelayEvent } from "../relay/events"; +const id = "11111111-1111-4111-8111-111111111111", + channelId = "22222222-2222-4222-8222-222222222222", + runId = "33333333-3333-4333-8333-333333333333"; +const yaml = + "name: Fixture\nenabled: false\ntrigger:\n on: message_posted\nsteps:\n - id: send\n action: send_message\n text: Hi\n"; +const disposers: (() => void)[] = []; +afterEach(() => { + for (const dispose of disposers.splice(0)) dispose(); +}); +function setup() { + const key = keypair(); + let saved: readonly OutgoingEvent[] = [], + allowed = true; + let settle!: (value: string) => void, reject!: (error: Error) => void; + const sign = vi.fn(async (template: Parameters[1]) => + signed(key, template), + ); + const publish = vi.fn( + (_event: RelayEvent, _signal: AbortSignal) => + new Promise((resolve, fail) => { + settle = resolve; + reject = fail; + }), + ); + const outbox = createOutbox( + key.pubkey, + { kinds: [30620, 46020, 5], sign, publish }, + { + load: () => [], + save: (rows) => { + saved = structuredClone(rows); + }, + }, + { + needsReceipt: (event) => [30620, 46020, 5].includes(event.kind), + onReceipt: (event, message) => workflows.receipt(event, message), + }, + ); + const read = vi.fn(async () => [] as RelayEvent[]); + const workflows = createWorkflows({ + viewer: key.pubkey, + reader: { read }, + outbox: outbox.outbox, + local: outbox.local, + host: { + lifecycleVersion: 1, + runs: async () => ({ runs: [], next: null }), + approvals: async () => ({ approvals: [] }), + }, + canAccess: () => allowed, + }); + disposers.push(() => { + workflows.dispose(); + outbox.dispose(); + }); + const definition = { + id, + channelId, + owner: key.pubkey, + revision: "a".repeat(64), + createdAt: 1, + yaml, + }; + return { + ...workflows, + outbox, + read, + sign, + publish, + definition, + saved: () => saved, + settle: (value: string) => settle(value), + reject: (error: Error) => reject(error), + revoke: () => { + allowed = false; + workflows.clear(); + }, + }; +} +it("save preserves exact YAML/coordinate/revision, receipt success is distinct from signed head readback", async () => { + const h = setup(); + const view = h.capability.definitions(channelId); + const operation = h.capability.save({ + channelId, + yaml, + existing: h.definition, + }); + await flush(); + const event = h.publish.mock.calls[0]?.[0]; + expect(event).toBeDefined(); + expect(event?.content).toBe(yaml); + expect(event?.tags).toContainEqual([ + "expected-revision", + h.definition.revision, + ]); + expect(event?.tags).toContainEqual(["d", id]); + h.settle(`response:${JSON.stringify({ workflow_id: id })}`); + await flush(); + expect(h.capability.operations.snapshot()[0]).toMatchObject({ + eventId: operation, + outcome: "succeeded", + }); + // Local echo and receipt never masquerade as a freshly read committed head. + expect(view.snapshot()).toMatchObject({ + status: "idle", + data: { items: [] }, + }); + await view.refresh(); + expect(view.snapshot().data.items).toEqual([]); +}); +it.each([ + [JSON.stringify({ run_id: runId }), "succeeded", runId], + [JSON.stringify({ workflow_id: id, run_id: runId }), "succeeded", runId], + [JSON.stringify({ workflow_id: runId, run_id: runId }), "unknown", undefined], + [JSON.stringify({ run_id: "bad" }), "unknown", undefined], + [ + JSON.stringify({ run_id: runId, webhook_secret: "PRIVATE" }), + "unknown", + undefined, + ], + ["{PRIVATE malformed", "unknown", undefined], +])( + "manual run only correlates validated returned run ID: %s", + async (payload, outcome, expectedRun) => { + const h = setup(); + h.capability.trigger(h.definition); + await flush(); + h.settle(`response:${payload}`); + await flush(); + expect(h.capability.operations.snapshot()[0]?.outcome).toBe(outcome); + expect(h.capability.operations.snapshot()[0]?.runId).toBe(expectedRun); + expect(h.read).not.toHaveBeenCalled(); + expect(JSON.stringify(h.saved())).not.toContain("PRIVATE"); + expect(JSON.stringify(h.capability.operations.snapshot())).not.toContain( + "PRIVATE", + ); + }, +); +it("lost receipt plus echo stays unknown; same signed retry and dismiss preserve operation identity", async () => { + const h = setup(); + const operation = h.capability.trigger(h.definition); + await flush(); + const event = h.publish.mock.calls[0]?.[0]; + if (!event) throw new Error("missing publication"); + h.outbox.observe([event]); + h.reject(new Error("PRIVATE lost")); + await flush(); + expect(h.capability.operations.snapshot()[0]).toMatchObject({ + eventId: operation, + delivery: "seen", + outcome: "unknown", + }); + h.capability.operations.retry(operation); + await flush(); + expect(h.publish).toHaveBeenCalledTimes(2); + expect(h.sign).toHaveBeenCalledTimes(1); + expect(h.publish.mock.calls[1]?.[0]).toEqual(event); + h.settle("duplicate: already processed"); + await flush(); + expect(h.capability.operations.snapshot()[0]?.outcome).toBe("unknown"); + await h.capability.operations.dismiss(operation); + expect(h.capability.operations.snapshot()).toEqual([]); +}); +it("explicit rejection is rejected, not unknown; revocation fences late receipts without discarding durable intent", async () => { + const h = setup(); + h.capability.trigger(h.definition); + await flush(); + h.reject(new PublishRejected("PRIVATE rejection")); + await flush(); + expect(h.capability.operations.snapshot()[0]?.outcome).toBe("rejected"); + const id = h.capability.trigger(h.definition); + await flush(); + h.revoke(); + h.settle(`response:${JSON.stringify({ run_id: runId })}`); + await flush(); + expect(h.capability.operations.snapshot()).toEqual([]); + expect(h.saved().some((row) => row.event.id === id)).toBe(true); + expect(JSON.stringify(h.saved())).not.toContain("PRIVATE"); +}); +it("webhook saves are blocked through raw YAML; stale/legacy deletion receipt never proves deletion", async () => { + const h = setup(); + expect(() => + h.capability.save({ + channelId, + yaml: yaml.replace("message_posted", "webhook"), + }), + ).toThrow("secret"); + expect(h.publish).not.toHaveBeenCalled(); + h.capability.delete(h.definition); + await flush(); + h.settle(""); + await flush(); + expect(h.capability.operations.snapshot()[0]?.outcome).toBe("unknown"); + h.capability.delete(h.definition); + await flush(); + h.settle(`response:${JSON.stringify({ workflow_id: id, deleted: true })}`); + await flush(); + expect(h.capability.operations.snapshot()[1]?.outcome).toBe("succeeded"); +}); diff --git a/src/features/workflows/capability.ts b/src/features/workflows/capability.ts new file mode 100644 index 00000000..19ad608b --- /dev/null +++ b/src/features/workflows/capability.ts @@ -0,0 +1,460 @@ +import type { EventData } from "../relay/events"; +import type { RelayReader } from "../relay/reader"; +import type { Outbox, LocalEvents } from "../relay/outbox"; +import type { WorkflowHost } from "./host"; +import type { + WorkflowCapability, + WorkflowOperation, + WorkflowReference, + WorkflowView, + WorkflowDefinition, +} from "./types"; +import { + definition, + isWorkflowOperation, + parseApprovals, + parseRuns, + record, + validateReference, + validateWorkflowEvent, + workflowReference, + UUID, +} from "./protocol"; + +/** Session-owned configuration snapshots and bounded result state; never an engine. */ +export function createWorkflows({ + reader, + viewer, + outbox, + local, + host, + canAccess, + notify = (listener: () => void) => listener(), +}: { + reader: RelayReader | undefined; + viewer: string; + outbox: Outbox | undefined; + local: LocalEvents | undefined; + host: WorkflowHost | undefined; + canAccess(channel: string): boolean; + notify?: (listener: () => void) => void; +}) { + let closed = false; + const views = new Set<{ clear(): void; emit(): void; dispose(): void }>(); + const listeners = new Set<() => void>(); + type Result = { + outcome: WorkflowOperation["outcome"]; + runId?: string; + error?: string; + }; + const results = new Map(); + const receiptInterest = new Set(); + // Webhook save/reveal stays unavailable until the explicit UI secret lifetime is integrated. + const availability = Object.freeze({ + definitions: !!reader, + history: !!host, + save: host?.lifecycleVersion === 1 && !!outbox?.supports(30620), + trigger: host?.lifecycleVersion === 1 && !!outbox?.supports(46020), + delete: host?.lifecycleVersion === 1 && !!outbox?.supports(5), + webhookSecrets: false, + }); + let operations: readonly WorkflowOperation[] = Object.freeze([]); + function rebuild() { + operations = closed + ? Object.freeze([]) + : Object.freeze( + (local?.snapshot() ?? []) + .flatMap((item): WorkflowOperation[] => { + if (!isWorkflowOperation(item.event)) return []; + let workflow: WorkflowReference; + try { + workflow = workflowReference(item.event); + } catch { + return []; + } + if (workflow.owner !== viewer || !canAccess(workflow.channelId)) + return []; + const result = results.get(item.event.id); + const outcome = + result?.outcome ?? + (item.delivery === "sending" + ? "pending" + : item.delivery === "failed" + ? "rejected" + : "unknown"); + return [ + Object.freeze({ + eventId: item.event.id, + workflow, + action: + item.event.kind === 30620 + ? "save" + : item.event.kind === 5 + ? "delete" + : "trigger", + delivery: item.delivery, + outcome, + secretAvailable: false, + ...(result?.runId ? { runId: result.runId } : {}), + ...((result?.error ?? item.error) !== undefined + ? { error: (result?.error ?? item.error) as string } + : {}), + }), + ]; + }) + .slice(-256), + ); + const active = new Set(operations.map((op) => op.eventId)); + for (const id of results.keys()) if (!active.has(id)) results.delete(id); + for (const listener of listeners) notify(listener); + } + const stop = local?.subscribe(rebuild); + rebuild(); + function assertAccess(reference: WorkflowReference) { + validateReference(reference); + if (closed || !canAccess(reference.channelId)) + throw new Error( + "Workflow access unavailable; refresh channel membership", + ); + } + function view( + channelId: string, + available: boolean, + empty: T, + load: (signal: AbortSignal) => Promise, + ): WorkflowView { + if (!UUID.test(channelId)) throw new Error("Invalid workflow channel"); + if (views.size >= 16) + throw new Error("Too many workflow views; close another detail first"); + let disposed = false; + let controller: AbortController | undefined; + let pending: Promise | undefined; + const subscribers = new Set<() => void>(); + type Snapshot = ReturnType["snapshot"]>; + let snapshot: Snapshot = Object.freeze({ + status: + available && !closed && canAccess(channelId) ? "idle" : "unavailable", + data: empty, + }); + const emit = () => { + for (const listener of subscribers) notify(listener); + }; + function clear() { + controller?.abort(); + controller = undefined; + pending = undefined; + snapshot = Object.freeze({ + status: + available && !closed && !disposed && canAccess(channelId) + ? "idle" + : "unavailable", + data: empty, + }); + } + const owner = { + clear, + emit, + dispose() { + disposed = true; + clear(); + subscribers.clear(); + views.delete(owner); + }, + }; + views.add(owner); + return Object.freeze({ + snapshot: () => snapshot, + subscribe(listener: () => void) { + if (closed || disposed) return () => {}; + subscribers.add(listener); + return () => { + subscribers.delete(listener); + }; + }, + refresh() { + if (closed || disposed || !available) return Promise.resolve(); + if (!canAccess(channelId)) { + clear(); + emit(); + return Promise.resolve(); + } + if (pending) return pending; + const owned = new AbortController(); + controller = owned; + const signal = AbortSignal.any([ + owned.signal, + AbortSignal.timeout(10000), + ]); + snapshot = Object.freeze({ status: "loading", data: snapshot.data }); + pending = Promise.resolve() + .then(() => { + signal.throwIfAborted(); + if (!canAccess(channelId)) + throw new Error("Workflow channel access unavailable"); + return load(signal); + }) + .then((data) => { + if ( + closed || + disposed || + controller !== owned || + signal.aborted || + !canAccess(channelId) + ) + return; + snapshot = Object.freeze({ status: "ready", data }); + emit(); + }) + .catch(() => { + if ( + closed || + disposed || + controller !== owned || + owned.signal.aborted + ) + return; + snapshot = Object.freeze({ + status: "error", + data: empty, + error: + "Workflow read unavailable. Retry; this is not proof of deletion.", + }); + emit(); + }) + .finally(() => { + if (controller === owned) { + controller = undefined; + pending = undefined; + } + }); + const started = pending; + emit(); + return started; + }, + dispose: owner.dispose, + }); + } + function assertOperation(kind: number) { + const enabled = + kind === 30620 + ? availability.save + : kind === 46020 + ? availability.trigger + : availability.delete; + if (!enabled) + throw new Error("Reliable workflow writes are unavailable on this relay"); + } + function send( + kind: 30620 | 46020 | 5, + workflow: WorkflowReference, + yaml = "", + revision?: string, + ) { + assertAccess(workflow); + assertOperation(kind); + const tags = [ + ["h", workflow.channelId], + ...(kind === 5 + ? [["a", `30620:${workflow.owner}:${workflow.id}`]] + : [["d", workflow.id]]), + ...(revision ? [["expected-revision", revision]] : []), + ]; + const input = { kind, content: yaml, tags }; + validateWorkflowEvent( + { ...input, pubkey: viewer, id: "", created_at: 0 }, + viewer, + availability, + ); + if (!outbox) throw new Error("Workflow publishing unavailable"); + if (receiptInterest.size >= 256) + throw new Error("Too many unresolved workflow commands"); + const id = outbox.send(input); + receiptInterest.add(id); + return id; + } + const capability = Object.freeze({ + availability, + definitions(channelId) { + return view( + channelId, + !!reader, + Object.freeze({ + items: Object.freeze([] as WorkflowDefinition[]), + partial: false, + }), + async (signal) => { + if (!reader) throw new Error("Workflow definitions unavailable"); + const events = await reader.read( + [{ kinds: [30620], "#h": [channelId], limit: 100 }], + { signal, fresh: true }, + ); + const coordinates = new Map(); + for (const event of events) { + const row = definition(event); + if (row.channelId !== channelId) + throw new Error("Mismatched workflow channel"); + const key = `${row.owner}:${row.id}`; + const old = coordinates.get(key); + if ( + !old || + row.createdAt > old.createdAt || + (row.createdAt === old.createdAt && row.revision < old.revision) + ) + coordinates.set(key, row); + } + return Object.freeze({ + items: Object.freeze([...coordinates.values()]), + partial: events.length >= 100, + }); + }, + ); + }, + runs(workflow, cursor) { + assertAccess(workflow); + return view( + workflow.channelId, + !!host, + Object.freeze({ runs: Object.freeze([]), next: null }), + async (signal) => { + if (!host) throw new Error("Workflow history unavailable"); + return parseRuns( + await host.runs(workflow.id, cursor, signal), + workflow.id, + ); + }, + ); + }, + approvals(workflow, runId) { + assertAccess(workflow); + return view( + workflow.channelId, + !!host, + Object.freeze([]), + async (signal) => { + if (!host) throw new Error("Workflow history unavailable"); + return parseApprovals( + await host.approvals(workflow.id, runId, signal), + workflow.id, + runId, + ); + }, + ); + }, + save({ channelId, yaml, existing }) { + if (existing && existing.channelId !== channelId) + throw new Error("Workflow channel cannot change"); + return send( + 30620, + existing ?? { channelId, owner: viewer, id: crypto.randomUUID() }, + yaml, + existing?.revision, + ); + }, + delete(workflow) { + return send(5, workflow); + }, + trigger(workflow) { + return send(46020, workflow); + }, + operations: Object.freeze({ + snapshot: () => operations, + subscribe(listener) { + listeners.add(listener); + return () => { + listeners.delete(listener); + }; + }, + retry(id) { + const op = operations.find((row) => row.eventId === id); + if (!op || op.outcome === "succeeded" || op.delivery === "sending") + return; + assertAccess(op.workflow); + receiptInterest.add(id); + results.delete(id); + outbox?.retry(id); + }, + async dismiss(id) { + await outbox?.dismiss(id); + results.delete(id); + receiptInterest.delete(id); + rebuild(); + }, + }), + takeWebhookSecret() { + return undefined; + }, + }); + return { + capability, + validate(event: EventData) { + if (!isWorkflowOperation(event)) return; + assertOperation(event.kind); + const reference = validateWorkflowEvent(event, viewer, availability); + assertAccess(reference); + }, + receipt(event: EventData, message: string | undefined) { + if (closed || !receiptInterest.delete(event.id)) return; + if (message === undefined) { + results.delete(event.id); + rebuild(); + return; + } + let result: Result = { + outcome: "unknown", + error: + "Delivery may have succeeded, but its result is unavailable. Do not submit a new command to retry.", + }; + try { + assertAccess(workflowReference(event)); + if (message?.startsWith("response:")) { + const value: unknown = JSON.parse(message.slice(9)); + const reference = workflowReference(event); + if ( + record(value) && + (event.kind === 46020 + ? value.workflow_id === undefined || + value.workflow_id === reference.id + : value.workflow_id === reference.id) && + value.webhook_secret === undefined + ) { + if (event.kind === 30620) result = { outcome: "succeeded" }; + else if ( + event.kind === 46020 && + typeof value.run_id === "string" && + UUID.test(value.run_id) + ) + result = { outcome: "succeeded", runId: value.run_id }; + else if ( + event.kind === 5 && + host?.lifecycleVersion === 1 && + value.deleted === true + ) + result = { outcome: "succeeded" }; + } + } + } catch { + /* Do not leak receipt text into errors/journal. */ + } + results.set(event.id, result); + rebuild(); + }, + clear() { + results.clear(); + receiptInterest.clear(); + for (const owned of views) owned.clear(); + rebuild(); + for (const owned of views) owned.emit(); + // Individual view snapshots were cleared before any view callbacks. + // Views refresh on explicit UI interest; no startup/background fanout. + }, + dispose() { + closed = true; + for (const owned of [...views]) owned.dispose(); + results.clear(); + receiptInterest.clear(); + stop?.(); + rebuild(); + listeners.clear(); + }, + }; +} diff --git a/src/features/workflows/host.ts b/src/features/workflows/host.ts new file mode 100644 index 00000000..5f486596 --- /dev/null +++ b/src/features/workflows/host.ts @@ -0,0 +1,13 @@ +import type { WorkflowRunCursor } from "./types"; + +/** Host-owned authenticated reads on the captured relay principal/admission lane. */ +export interface WorkflowHost { + /** Positive forward lifecycle contract evidence, never inferred from kind support. */ + readonly lifecycleVersion?: 1; + runs( + id: string, + cursor: WorkflowRunCursor | undefined, + signal: AbortSignal, + ): Promise; + approvals(id: string, runId: string, signal: AbortSignal): Promise; +} diff --git a/src/features/workflows/http.test.ts b/src/features/workflows/http.test.ts new file mode 100644 index 00000000..4cfb4e10 --- /dev/null +++ b/src/features/workflows/http.test.ts @@ -0,0 +1,94 @@ +import { afterEach, expect, it, vi } from "vitest"; +import { verifyEvent } from "nostr-tools"; +import { connectSignedTransport } from "../relay/transport"; +import { keypair, signed } from "../relay/testing"; +import { WORKFLOW_READ_BYTES, workflowReadText } from "./http"; +const id = "11111111-1111-4111-8111-111111111111"; +const runId = "22222222-2222-4222-8222-222222222222"; +afterEach(() => vi.unstubAllGlobals()); +it("direct signed transport history uses exact GET URL, no payload and the existing principal quota lane", async () => { + const key = keypair(); + const signer = { + getPublicKey: async () => key.pubkey, + signEvent: async (template: Parameters[1]) => + signed(key, template), + }; + const fetcher = vi.fn(async (url: string, init?: RequestInit) => { + const auth = JSON.parse( + atob(new Headers(init?.headers).get("Authorization")?.slice(6) ?? ""), + ); + expect(verifyEvent(auth)).toBe(true); + expect(auth.pubkey).toBe(key.pubkey); + expect(auth.tags).toContainEqual(["u", url]); + expect(auth.tags).toContainEqual(["method", "GET"]); + expect(auth.tags.some(([name]: string[]) => name === "payload")).toBe( + false, + ); + expect(init?.body).toBeUndefined(); + expect(init?.redirect).toBe("error"); + return Response.json( + { error: "rate-limited: quota exceeded; retry in 0s" }, + { status: 429 }, + ); + }); + vi.stubGlobal("fetch", fetcher); + const transport = await connectSignedTransport( + signer, + "https://workflow-direct.test", + key.pubkey, + ); + expect(fetcher).not.toHaveBeenCalled(); + const cursor = { + before: "2026-09-12T14:44:19.123456+00:00", + beforeId: runId, + }; + await expect( + transport.workflows?.runs(id, cursor, new AbortController().signal), + ).rejects.toMatchObject({ status: 429 }); + expect(fetcher.mock.calls[0]?.[0]).toBe( + `https://workflow-direct.test/workflows/${id}/runs?limit=20&before=2026-09-12T14%3A44%3A19.123456%2B00%3A00&before_id=${runId}`, + ); + await expect( + transport.query([{ kinds: [0], limit: 1 }]), + ).rejects.toMatchObject({ status: 429 }); + expect(fetcher).toHaveBeenCalledTimes(1); +}); +it("structured body budget counts stream bytes, cancels overflow, rejects invalid UTF8", async () => { + const cancel = vi.fn(); + const response = new Response( + new ReadableStream({ + start(controller) { + controller.enqueue(new Uint8Array(WORKFLOW_READ_BYTES + 1)); + }, + cancel, + }), + ); + await expect(workflowReadText(response)).rejects.toThrow("size limit"); + expect(cancel).toHaveBeenCalledTimes(1); + await expect( + workflowReadText(new Response(new Uint8Array([0xff]))), + ).rejects.toThrow(); +}); +it("direct workflow cancellation/invalid arguments never sign or dispatch", async () => { + const key = keypair(), + signEvent = vi.fn(async (template: Parameters[1]) => + signed(key, template), + ); + const fetcher = vi.fn(); + vi.stubGlobal("fetch", fetcher); + const transport = await connectSignedTransport( + { getPublicKey: async () => key.pubkey, signEvent }, + "https://workflow-cancel.test", + key.pubkey, + ); + const cancel = new AbortController(); + cancel.abort(); + await expect( + transport.workflows?.runs(id, undefined, cancel.signal), + ).rejects.toThrow(); + await expect( + transport.workflows?.approvals(id, "../", new AbortController().signal), + ).rejects.toThrow(); + expect(signEvent).not.toHaveBeenCalled(); + expect(fetcher).not.toHaveBeenCalled(); +}); diff --git a/src/features/workflows/http.ts b/src/features/workflows/http.ts new file mode 100644 index 00000000..c9c20326 --- /dev/null +++ b/src/features/workflows/http.ts @@ -0,0 +1,97 @@ +import { ReadError } from "../relay/errors"; +import { readApiFailure } from "../relay/http-admission"; +import type { WorkflowHost } from "./host"; +import { approvalsPath, record, runsPath } from "./protocol"; + +export const WORKFLOW_READ_BYTES = 1024 * 1024; + +/** Fixed routes only: browser input can never select an upstream URL or page size. */ +export function workflowReadPath(route: string, body: unknown): string { + if (!record(body) || typeof body.id !== "string") + throw new Error("Invalid workflow read"); + if (route === "workflow-runs") { + if (Object.keys(body).some((key) => !["id", "cursor"].includes(key))) + throw new Error("Invalid workflow read fields"); + const cursor = body.cursor; + if ( + cursor !== undefined && + (!record(cursor) || + typeof cursor.before !== "string" || + typeof cursor.beforeId !== "string" || + Object.keys(cursor).some( + (key) => !["before", "beforeId"].includes(key), + )) + ) + throw new Error("Invalid workflow cursor"); + return runsPath(body.id, cursor as Parameters[1]); + } + if ( + route === "workflow-approvals" && + typeof body.runId === "string" && + Object.keys(body).every((key) => ["id", "runId"].includes(key)) + ) + return approvalsPath(body.id, body.runId); + throw new Error("Invalid workflow read route"); +} + +/** Stream-bound before parsing. Never include upstream text in an error/log. */ +export async function workflowReadText(response: Response): Promise { + if (!response.body) throw new Error("Workflow response body missing"); + const reader = response.body.getReader(); + const decoder = new TextDecoder("utf-8", { fatal: true }); + let bytes = 0, + text = ""; + try { + while (true) { + const { value, done } = await reader.read(); + if (done) return text + decoder.decode(); + bytes += value.byteLength; + if (bytes > WORKFLOW_READ_BYTES) + throw new Error("Workflow response exceeds the size limit"); + text += decoder.decode(value, { stream: true }); + } + } finally { + await reader.cancel().catch(() => {}); + reader.releaseLock(); + } +} + +/** Adapters supply authentication/admission; the capability validates domain rows. */ +export function workflowHost( + request: ( + route: string, + body: unknown, + signal: AbortSignal, + ) => Promise, +): WorkflowHost { + async function read(route: string, body: unknown, signal: AbortSignal) { + workflowReadPath(route, body); + signal = AbortSignal.any([signal, AbortSignal.timeout(10000)]); + signal.throwIfAborted(); + const response = await request(route, body, signal); + if (!response.ok) { + const failure = await readApiFailure(response); + throw new ReadError( + response.status === 401 || response.status === 403 + ? "denied" + : "unavailable", + failure.error, + response.status, + failure.retryAfterMs, + ); + } + const text = await workflowReadText(response); + signal.throwIfAborted(); + try { + return JSON.parse(text) as unknown; + } catch { + throw new Error("Invalid workflow response"); + } + } + return Object.freeze({ + runs: (id, cursor, signal) => + read("workflow-runs", { id, ...(cursor ? { cursor } : {}) }, signal), + approvals: (id, runId, signal) => + read("workflow-approvals", { id, runId }, signal), + } satisfies WorkflowHost); +} diff --git a/src/features/workflows/protocol.test.ts b/src/features/workflows/protocol.test.ts new file mode 100644 index 00000000..cf0ad582 --- /dev/null +++ b/src/features/workflows/protocol.test.ts @@ -0,0 +1,126 @@ +import { expect, it } from "vitest"; +import { + isWorkflowOperation, + validateWorkflowEvent, + parseRuns, + parseApprovals, + workflowYaml, +} from "./protocol"; +const owner = "a".repeat(64), + id = "11111111-1111-4111-8111-111111111111", + runId = "22222222-2222-4222-8222-222222222222"; +const yaml = + "name: Fixture\nenabled: false\ntrigger:\n on: message_posted\nsteps:\n - id: 1_send\n action: send_message\n text: Hi\n"; +const base = { + id: "b".repeat(64), + pubkey: owner, + created_at: 1, + kind: 30620, + tags: [ + ["h", id], + ["d", id], + ], + content: yaml, +}; +it("strict authoring boundary rejects malformed coordinates/tags and webhook raw bypass, without taking generic deletions", () => { + expect( + validateWorkflowEvent(base, owner, { delete: true, webhookSecrets: false }), + ).toEqual({ id, owner, channelId: id }); + for (const event of [ + { ...base, pubkey: "c".repeat(64) }, + { + ...base, + tags: [ + ["h", "../"], + ["d", id], + ], + }, + { ...base, tags: [...base.tags, ["d", id]] }, + { ...base, tags: [...base.tags, ["client-id", "x"], ["client-id", "y"]] }, + { ...base, tags: [...base.tags, ["expected-revision", "bad"]] }, + { ...base, content: yaml.replace("message_posted", "webhook") }, + { ...base, kind: 46020 }, + { ...base, created_at: Infinity }, + ]) + expect(() => + validateWorkflowEvent(event, owner, { + delete: true, + webhookSecrets: false, + }), + ).toThrow(); + expect(isWorkflowOperation({ kind: 5, tags: [["e", "b".repeat(64)]] })).toBe( + false, + ); + expect( + isWorkflowOperation({ kind: 5, tags: [["a", `30620:${owner}:${id}`]] }), + ).toBe(true); + expect(() => + validateWorkflowEvent( + { + ...base, + kind: 5, + content: "", + tags: [ + ["h", id], + ["a", `30620:${owner}:${id}`], + ], + }, + owner, + { delete: false, webhookSecrets: false }, + ), + ).toThrow("deletion"); +}); +it("YAML remains unchanged and malformed input/duplicate steps are rejected", () => { + expect(workflowYaml(yaml)).toMatchObject({ enabled: false }); + for (const text of [ + "[]", + yaml.replace("enabled: false", ""), + yaml.replace("1_send", "my-step"), + `${yaml} - id: 1_send\n action: delay\n duration: 1s\n`, + "x".repeat(24001), + "name: [bad", + ]) { + expect(() => workflowYaml(text)).toThrow(); + } +}); +it("run/approval rows require matching identities and preserve exact cursor precision", () => { + const row = { + id: runId, + workflow_id: id, + status: "completed", + current_step: 1, + execution_trace: [], + started_at: 1, + completed_at: 2, + created_at: 1, + error_code: null, + error_message: null, + }; + const before = "2026-09-12T14:44:19.123456+00:00"; + const raw = { runs: [row], next: { before, before_id: runId } }; + expect(parseRuns(raw, id).next).toEqual({ before, beforeId: runId }); + for (const value of [ + { ...raw, runs: [row, row] }, + { ...raw, runs: [{ ...row, workflow_id: runId }] }, + { ...raw, runs: [{ ...row, status: "imaginary" }] }, + { ...raw, next: { before } }, + { ...raw, runs: Array(21).fill(row) }, + ]) { + expect(() => parseRuns(value, id)).toThrow(); + } + const approval = { + workflow_id: id, + run_id: runId, + approval_ref: owner, + step_id: "a", + status: "pending", + created_at: 1, + note: null, + }; + expect( + parseApprovals({ approvals: [approval] }, id, runId)[0]?.reference, + ).toBe(owner); + expect(() => + parseApprovals({ approvals: [{ ...approval, run_id: id }] }, id, runId), + ).toThrow(); +}); diff --git a/src/features/workflows/protocol.ts b/src/features/workflows/protocol.ts new file mode 100644 index 00000000..4f0fcff2 --- /dev/null +++ b/src/features/workflows/protocol.ts @@ -0,0 +1,286 @@ +import { parseDocument } from "yaml"; +import type { EventData } from "../relay/events"; +import type { + WorkflowDefinition, + WorkflowReference, + WorkflowRunCursor, + WorkflowRunPage, + WorkflowApproval, +} from "./types"; + +export const WORKFLOW_KINDS = [30620, 46020, 5] as const; +export function isWorkflowOperation( + event: Pick, +): boolean { + return ( + event.kind === 30620 || + event.kind === 46020 || + (event.kind === 5 && + event.tags.some( + ([name, value]) => name === "a" && value?.startsWith("30620:"), + )) + ); +} +export const UUID = + /^[0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12}$/; +const HEX = /^[0-9a-f]{64}$/; +export function record(value: unknown): value is Record { + return typeof value === "object" && value !== null && !Array.isArray(value); +} +export function validateReference(value: WorkflowReference) { + if ( + !UUID.test(value.id) || + !UUID.test(value.channelId) || + !HEX.test(value.owner) + ) + throw new Error("Invalid workflow coordinate"); +} +function one(event: Pick, name: string): string { + const tags = event.tags.filter((tag) => tag[0] === name); + if (tags.length !== 1 || tags[0]?.length !== 2 || !tags[0][1]) + throw new Error(`Invalid workflow ${name} tag`); + return tags[0][1]; +} +export function workflowReference(event: EventData): WorkflowReference { + const channelId = one(event, "h"); + let id: string, + owner = event.pubkey; + if (event.kind === 5) { + const parts = one(event, "a").split(":"); + if (parts.length !== 3 || parts[0] !== "30620") + throw new Error("Invalid workflow deletion coordinate"); + owner = parts[1] ?? ""; + id = parts[2] ?? ""; + } else id = one(event, "d"); + const reference = { id, owner, channelId }; + validateReference(reference); + return reference; +} +export function definition(event: EventData): WorkflowDefinition { + if (event.kind !== 30620 || !HEX.test(event.id)) + throw new Error("Invalid workflow definition"); + return Object.freeze({ + ...workflowReference(event), + revision: event.id, + createdAt: event.created_at, + yaml: event.content, + }); +} +/** Structural authoring checks, not a replacement for relay language/role validation. */ +export function workflowYaml(text: string) { + if (new TextEncoder().encode(text).byteLength > 24000) + throw new Error("Workflow YAML exceeds 24 KB"); + const doc = parseDocument(text); + if (doc.errors.length) throw new Error("Workflow YAML is invalid"); + const value: unknown = doc.toJS({ maxAliasCount: 50 }); + if ( + !record(value) || + typeof value.name !== "string" || + !value.name.trim() || + typeof value.enabled !== "boolean" || + !record(value.trigger) || + typeof value.trigger.on !== "string" || + !Array.isArray(value.steps) || + !value.steps.length || + value.steps.length > 100 + ) + throw new Error( + "Workflow needs a name, explicit enabled state, trigger and 1–100 steps", + ); + const ids = new Set(); + for (const step of value.steps) { + if ( + !record(step) || + typeof step.id !== "string" || + !/^[A-Za-z0-9_]{1,64}$/.test(step.id) || + ids.has(step.id) || + typeof step.action !== "string" + ) + throw new Error("Workflow steps need unique identifiers and actions"); + ids.add(step.id); + } + return { + webhook: value.trigger.on === "webhook", + enabled: value.enabled, + name: value.name, + }; +} +/** Shared session/broker signing boundary. Only canonical workflow operations, never generic kind 5. */ +export function validateWorkflowEvent( + event: EventData, + viewer: string, + options: { delete: boolean; webhookSecrets: boolean }, +) { + if ( + typeof event.content !== "string" || + !Number.isSafeInteger(event.created_at) || + event.created_at < 0 || + !Array.isArray(event.tags) || + event.tags.length > 8 || + event.tags.some( + (tag) => + !Array.isArray(tag) || + tag.length !== 2 || + tag.some((value) => typeof value !== "string" || value.length > 256), + ) + ) + throw new Error("Malformed workflow command"); + if (!WORKFLOW_KINDS.includes(event.kind as 30620 | 46020 | 5)) + throw new Error("Unsupported workflow operation"); + const reference = workflowReference(event); + if (reference.owner !== viewer || event.pubkey !== viewer) + throw new Error("Only the workflow author can manage it"); + const allowed = + event.kind === 5 + ? ["h", "a", "client-id"] + : ["h", "d", "expected-revision", "client-id"]; + if ( + new Set(event.tags.map(([name]) => name)).size !== event.tags.length || + event.tags.some((tag) => !allowed.includes(tag[0] ?? "")) + ) + throw new Error("Unsupported workflow command tag"); + if (event.kind === 5 && !options.delete) + throw new Error("Reliable workflow deletion is unavailable on this relay"); + if (event.kind !== 30620 && event.content !== "") + throw new Error("Workflow command content must be empty"); + if (event.kind === 30620) { + const revisions = event.tags.filter( + ([name]) => name === "expected-revision", + ); + if (revisions.length && !HEX.test(one(event, "expected-revision"))) + throw new Error("Invalid expected workflow revision"); + if (workflowYaml(event.content).webhook && !options.webhookSecrets) + throw new Error("Webhook saves require secure one-time-secret handling"); + } else if (event.tags.some(([name]) => name === "expected-revision")) + throw new Error("Unexpected workflow revision tag"); + return reference; +} +export function runsPath(id: string, cursor?: WorkflowRunCursor) { + if (!UUID.test(id)) throw new Error("Invalid workflow ID"); + const query = new URLSearchParams({ limit: "20" }); + if (cursor) { + validateCursor(cursor); + query.set("before", cursor.before); + query.set("before_id", cursor.beforeId); + } + return `/workflows/${id}/runs?${query}`; +} +export function approvalsPath(id: string, runId: string) { + if (!UUID.test(id) || !UUID.test(runId)) + throw new Error("Invalid workflow/run ID"); + return `/workflows/${id}/runs/${runId}/approvals`; +} +function validateCursor(cursor: WorkflowRunCursor) { + if ( + !UUID.test(cursor.beforeId) || + cursor.before.length > 40 || + !/^\d{4}-\d{2}-\d{2}T/.test(cursor.before) || + !Number.isFinite(Date.parse(cursor.before)) + ) + throw new Error("Invalid workflow run cursor"); +} +const number = (v: unknown): v is number => + Number.isSafeInteger(v) && (v as number) >= 0; +const nullableNumber = (v: unknown): v is number | null => + v === null || number(v); +const nullableText = (v: unknown): v is string | null => + v === null || (typeof v === "string" && v.length <= 16000); +export function parseRuns(raw: unknown, workflowId: string): WorkflowRunPage { + if (!record(raw) || !Array.isArray(raw.runs) || raw.runs.length > 20) + throw new Error("Invalid workflow run response"); + const ids = new Set(); + const runs = raw.runs.map((run) => { + if ( + !record(run) || + typeof run.id !== "string" || + !UUID.test(run.id) || + ids.has(run.id) || + run.workflow_id !== workflowId || + ![ + "pending", + "running", + "waiting_approval", + "completed", + "failed", + "cancelled", + ].includes(String(run.status)) || + !number(run.current_step) || + !number(run.created_at) || + !nullableNumber(run.started_at) || + !nullableNumber(run.completed_at) || + !Array.isArray(run.execution_trace) || + run.execution_trace.length > 1000 || + !nullableText(run.error_code) || + !nullableText(run.error_message) + ) + throw new Error("Invalid workflow run row"); + ids.add(run.id); + return Object.freeze({ + id: run.id, + workflowId, + status: run.status as WorkflowRunPage["runs"][number]["status"], + currentStep: run.current_step, + trace: Object.freeze(run.execution_trace), + startedAt: run.started_at, + completedAt: run.completed_at, + createdAt: run.created_at, + errorCode: run.error_code, + errorMessage: run.error_message, + }); + }); + let next: WorkflowRunCursor | null = null; + if (raw.next !== null) { + if ( + !record(raw.next) || + typeof raw.next.before !== "string" || + typeof raw.next.before_id !== "string" || + !runs.length + ) + throw new Error("Invalid workflow run cursor"); + next = Object.freeze({ + before: raw.next.before, + beforeId: raw.next.before_id, + }); + validateCursor(next); + } + return Object.freeze({ runs: Object.freeze(runs), next }); +} +export function parseApprovals( + raw: unknown, + workflowId: string, + runId: string, +): readonly WorkflowApproval[] { + if ( + !record(raw) || + !Array.isArray(raw.approvals) || + raw.approvals.length > 1000 + ) + throw new Error("Invalid workflow approvals response"); + return Object.freeze( + raw.approvals.map((row) => { + if ( + !record(row) || + row.workflow_id !== workflowId || + row.run_id !== runId || + typeof row.approval_ref !== "string" || + !HEX.test(row.approval_ref) || + typeof row.step_id !== "string" || + row.step_id.length > 256 || + !["pending", "granted", "denied", "expired"].includes( + String(row.status), + ) || + !nullableText(row.note) || + !number(row.created_at) + ) + throw new Error("Invalid workflow approval row"); + return Object.freeze({ + reference: row.approval_ref, + runId, + stepId: row.step_id, + status: row.status as WorkflowApproval["status"], + note: row.note, + createdAt: row.created_at, + }); + }), + ); +} diff --git a/src/features/workflows/session.test.ts b/src/features/workflows/session.test.ts new file mode 100644 index 00000000..428a9eac --- /dev/null +++ b/src/features/workflows/session.test.ts @@ -0,0 +1,195 @@ +import { afterEach, expect, it, vi } from "vitest"; +import { createRelaySession } from "../relay/session"; +import type { RelayEvent } from "../relay/events"; +import { keypair, roster, signed, scriptedTransport } from "../relay/testing"; +const channelId = "11111111-1111-4111-8111-111111111111"; +const id = "22222222-2222-4222-8222-222222222222"; +const relay = keypair(), + viewer = keypair(); +const definition = signed(viewer, { + kind: 30620, + created_at: 10, + content: "PRIVATE yaml", + tags: [ + ["h", channelId], + ["d", id], + ], +}); +const reference = { id, channelId, owner: viewer.pubkey }; +const owners: ReturnType[] = []; +afterEach(() => { + for (const owner of owners.splice(0)) owner.dispose(); +}); +function setup() { + const wire = scriptedTransport(viewer.pubkey, relay.pubkey); + let incoming!: (events: readonly RelayEvent[]) => void; + let resolveRuns!: (value: unknown) => void; + const runs = vi.fn( + (_id: string, _cursor: unknown, _signal: AbortSignal) => + new Promise((resolve) => { + resolveRuns = resolve; + }), + ); + const owner = createRelaySession({ + ...wire.transport, + workflows: { runs, approvals: async () => ({ approvals: [] }) }, + subscribe(callbacks) { + incoming = callbacks.receive; + return { update() {}, retry() {}, dispose() {} }; + }, + }); + owners.push(owner); + return { + ...wire, + ...owner, + runs, + emit: (events: readonly RelayEvent[]) => incoming(events), + resolveRuns: (value: unknown) => resolveRuns(value), + }; +} +it("workflow views start lazy and purge to unavailable before any channel/operation/view observer runs", async () => { + const h = setup(); + h.emit([roster(relay, channelId, [viewer.pubkey], 1)]); + const definitions = h.session.workflows.definitions(channelId), + history = h.session.workflows.runs(reference); + expect(h.pending).toHaveLength(0); + expect(h.runs).not.toHaveBeenCalled(); + const loading = definitions.refresh(); + await vi.waitFor(() => expect(h.pending).toHaveLength(1)); + const request = h.next(); + expect(request.filters).toEqual([ + { kinds: [30620], "#h": [channelId], limit: 100 }, + ]); + // A second refresh shares its pending read. + expect(definitions.refresh()).toBe(loading); + request.respond([definition]); + await loading; + expect(definitions.snapshot().status).toBe("ready"); + expect(history.snapshot().status).toBe("idle"); +}); +it("authoritative revocation clears all saved and structured data before callbacks, rejects late results and denies fresh views", async () => { + const h = setup(); + h.emit([roster(relay, channelId, [viewer.pubkey], 1)]); + const definitions = h.session.workflows.definitions(channelId), + history = h.session.workflows.runs(reference); + const loading = definitions.refresh(); + await vi.waitFor(() => expect(h.pending).toHaveLength(1)); + h.next().respond([definition]); + await loading; + expect(definitions.snapshot().data.items[0]?.revision).toBe(definition.id); + const runRead = history.refresh(); + await vi.waitFor(() => expect(h.runs).toHaveBeenCalledTimes(1)); + const checked = vi.fn(() => { + expect(definitions.snapshot()).toMatchObject({ + status: "unavailable", + data: { items: [] }, + }); + expect(history.snapshot()).toMatchObject({ + status: "unavailable", + data: { runs: [] }, + }); + expect(h.session.workflows.operations.snapshot()).toEqual([]); + }); + definitions.subscribe(checked); + history.subscribe(checked); + h.session.workflows.operations.subscribe(checked); + h.session.channels.subscribeList(checked); + h.emit([roster(relay, channelId, [], 2)]); + expect(checked).toHaveBeenCalled(); + expect(h.runs.mock.calls[0]?.[2].aborted).toBe(true); + const denied = h.session.workflows.definitions(channelId); + expect(denied.snapshot().status).toBe("unavailable"); + await denied.refresh(); + expect(h.pending).toHaveLength(0); + h.resolveRuns({ runs: [], next: null }); + await runRead; + expect(history.snapshot().status).toBe("unavailable"); +}); +it("regrant cannot resurrect stale history; clear-cache and dispose cancel interest", async () => { + const h = setup(); + h.emit([roster(relay, channelId, [viewer.pubkey], 1)]); + const history = h.session.workflows.runs(reference); + const first = history.refresh(); + await vi.waitFor(() => expect(h.runs).toHaveBeenCalledTimes(1)); + h.emit([roster(relay, channelId, [], 2)]); + h.emit([roster(relay, channelId, [viewer.pubkey], 3)]); + h.resolveRuns({ runs: [], next: null }); + await first; + expect(history.snapshot().status).toBe("unavailable"); + const fresh = history.refresh(); + await vi.waitFor(() => expect(h.runs).toHaveBeenCalledTimes(2)); + h.resolveRuns({ runs: [], next: null }); + await fresh; + expect(history.snapshot().status).toBe("ready"); + const next = history.refresh(); + await vi.waitFor(() => expect(h.runs).toHaveBeenCalledTimes(3)); + await h.clearCache(); + expect(h.runs.mock.calls[2]?.[2].aborted).toBe(true); + expect(history.snapshot()).toMatchObject({ + status: "idle", + data: { runs: [] }, + }); + h.resolveRuns({ runs: [], next: null }); + await next; + expect(history.snapshot().status).toBe("idle"); + h.dispose(); + expect(history.snapshot().status).toBe("unavailable"); +}); +it("a loading observer can revoke without leaving a wedged pending read", async () => { + const h = setup(); + h.emit([roster(relay, channelId, [viewer.pubkey], 1)]); + const history = h.session.workflows.runs(reference); + const stop = history.subscribe(() => { + if (history.snapshot().status === "loading") + h.emit([roster(relay, channelId, [], 2)]); + }); + await history.refresh(); + expect(h.runs).not.toHaveBeenCalled(); + expect(history.snapshot().status).toBe("unavailable"); + stop(); + h.emit([roster(relay, channelId, [viewer.pubkey], 3)]); + const fresh = history.refresh(); + await vi.waitFor(() => expect(h.runs).toHaveBeenCalledTimes(1)); + h.resolveRuns({ runs: [], next: null }); + await fresh; + expect(history.snapshot().status).toBe("ready"); +}); +it("old host keeps every workflow command unavailable, even through direct session outbox", async () => { + const wire = scriptedTransport(viewer.pubkey, relay.pubkey), + sign = vi.fn(async (template: Parameters[1]) => + signed(viewer, template), + ); + const owner = createRelaySession( + { ...wire.transport, writer: { sign, publish: async () => "" } }, + { + outboxStorage: { load: async () => [], save: async () => {} }, + }, + ); + owners.push(owner); + expect(owner.session.workflows.availability).toMatchObject({ + save: false, + delete: false, + trigger: false, + }); + const workflow = { + ...reference, + yaml: definition.content, + revision: definition.id, + createdAt: 10, + }; + expect(() => owner.session.workflows.trigger(workflow)).toThrow( + "unavailable", + ); + owner.session.outbox?.send({ + kind: 46020, + tags: [ + ["h", channelId], + ["d", id], + ], + content: "", + }); + await vi.waitFor(() => + expect(owner.session.outbox?.snapshot()[0]?.delivery).toBe("failed"), + ); + expect(sign).not.toHaveBeenCalled(); +}); From 23dda527fc6f2dc3e43b4909b4126ea565a5dc32 Mon Sep 17 00:00:00 2001 From: Brain <1a02c72794dcd0f07058a353bc3a81f4028b8c77c92c87fce6d5c8b85970a20b@buzz.block.builderlab.xyz> Date: Sat, 12 Sep 2026 08:58:56 -0600 Subject: [PATCH 03/20] test(workflows): route editor journey through browser lanes Signed-off-by: Brain <1a02c72794dcd0f07058a353bc3a81f4028b8c77c92c87fce6d5c8b85970a20b@buzz.block.builderlab.xyz> --- tests/browser/workflows.spec.mjs | 1 + 1 file changed, 1 insertion(+) create mode 100644 tests/browser/workflows.spec.mjs diff --git a/tests/browser/workflows.spec.mjs b/tests/browser/workflows.spec.mjs new file mode 100644 index 00000000..af5a0e4e --- /dev/null +++ b/tests/browser/workflows.spec.mjs @@ -0,0 +1 @@ +import "../../src/bundled/workflows/workflows.journey.mjs"; From 7dcb4db41015bab36d46d39beda3072be23e9c23 Mon Sep 17 00:00:00 2001 From: Brain <1a02c72794dcd0f07058a353bc3a81f4028b8c77c92c87fce6d5c8b85970a20b@buzz.block.builderlab.xyz> Date: Sat, 12 Sep 2026 09:00:47 -0600 Subject: [PATCH 04/20] fix(workflows): preserve legacy enabled default in signing checks Signed-off-by: Brain <1a02c72794dcd0f07058a353bc3a81f4028b8c77c92c87fce6d5c8b85970a20b@buzz.block.builderlab.xyz> --- src/features/workflows/protocol.test.ts | 22 +++++++++++++++++++++- src/features/workflows/protocol.ts | 7 ++++--- 2 files changed, 25 insertions(+), 4 deletions(-) diff --git a/src/features/workflows/protocol.test.ts b/src/features/workflows/protocol.test.ts index cf0ad582..aabc83d6 100644 --- a/src/features/workflows/protocol.test.ts +++ b/src/features/workflows/protocol.test.ts @@ -74,7 +74,9 @@ it("YAML remains unchanged and malformed input/duplicate steps are rejected", () expect(workflowYaml(yaml)).toMatchObject({ enabled: false }); for (const text of [ "[]", - yaml.replace("enabled: false", ""), + yaml.replace("enabled: false", "enabled: null"), + yaml.replace("enabled: false", 'enabled: "false"'), + yaml.replace("enabled: false", "enabled: 0"), yaml.replace("1_send", "my-step"), `${yaml} - id: 1_send\n action: delay\n duration: 1s\n`, "x".repeat(24001), @@ -83,6 +85,24 @@ it("YAML remains unchanged and malformed input/duplicate steps are rejected", () expect(() => workflowYaml(text)).toThrow(); } }); +it("legacy omitted enabled and explicit toggles cross the real event boundary unchanged", () => { + for (const [line, enabled] of [ + ["", true], + ["enabled: true\n", true], + ["enabled: false\n", false], + ] as const) { + const content = yaml.replace("enabled: false\n", line); + const event = { ...base, content }; + expect(workflowYaml(content).enabled).toBe(enabled); + expect( + validateWorkflowEvent(event, owner, { + delete: false, + webhookSecrets: false, + }), + ).toEqual({ id, owner, channelId: id }); + expect(event.content).toBe(content); + } +}); it("run/approval rows require matching identities and preserve exact cursor precision", () => { const row = { id: runId, diff --git a/src/features/workflows/protocol.ts b/src/features/workflows/protocol.ts index 4f0fcff2..f666d6d3 100644 --- a/src/features/workflows/protocol.ts +++ b/src/features/workflows/protocol.ts @@ -77,7 +77,7 @@ export function workflowYaml(text: string) { !record(value) || typeof value.name !== "string" || !value.name.trim() || - typeof value.enabled !== "boolean" || + (value.enabled !== undefined && typeof value.enabled !== "boolean") || !record(value.trigger) || typeof value.trigger.on !== "string" || !Array.isArray(value.steps) || @@ -85,7 +85,7 @@ export function workflowYaml(text: string) { value.steps.length > 100 ) throw new Error( - "Workflow needs a name, explicit enabled state, trigger and 1–100 steps", + "Workflow needs a name, boolean enabled state if present, trigger and 1–100 steps", ); const ids = new Set(); for (const step of value.steps) { @@ -101,7 +101,8 @@ export function workflowYaml(text: string) { } return { webhook: value.trigger.on === "webhook", - enabled: value.enabled, + // Match legacy WorkflowDef: omission means enabled; do not rewrite YAML. + enabled: value.enabled !== false, name: value.name, }; } From 90b2becfa3c9452f429bb26f4d5a464ad143e678 Mon Sep 17 00:00:00 2001 From: Brain <1a02c72794dcd0f07058a353bc3a81f4028b8c77c92c87fce6d5c8b85970a20b@buzz.block.builderlab.xyz> Date: Sat, 12 Sep 2026 09:09:27 -0600 Subject: [PATCH 05/20] fix(workflows): make broker dependencies native-loader compatible Signed-off-by: Brain <1a02c72794dcd0f07058a353bc3a81f4028b8c77c92c87fce6d5c8b85970a20b@buzz.block.builderlab.xyz> --- src/features/workflows/host.ts | 2 +- src/features/workflows/http.ts | 8 ++++---- src/features/workflows/protocol.ts | 4 ++-- 3 files changed, 7 insertions(+), 7 deletions(-) diff --git a/src/features/workflows/host.ts b/src/features/workflows/host.ts index 5f486596..308dae53 100644 --- a/src/features/workflows/host.ts +++ b/src/features/workflows/host.ts @@ -1,4 +1,4 @@ -import type { WorkflowRunCursor } from "./types"; +import type { WorkflowRunCursor } from "./types.ts"; /** Host-owned authenticated reads on the captured relay principal/admission lane. */ export interface WorkflowHost { diff --git a/src/features/workflows/http.ts b/src/features/workflows/http.ts index c9c20326..0cc46320 100644 --- a/src/features/workflows/http.ts +++ b/src/features/workflows/http.ts @@ -1,7 +1,7 @@ -import { ReadError } from "../relay/errors"; -import { readApiFailure } from "../relay/http-admission"; -import type { WorkflowHost } from "./host"; -import { approvalsPath, record, runsPath } from "./protocol"; +import { ReadError } from "../relay/errors.ts"; +import { readApiFailure } from "../relay/http-admission.ts"; +import type { WorkflowHost } from "./host.ts"; +import { approvalsPath, record, runsPath } from "./protocol.ts"; export const WORKFLOW_READ_BYTES = 1024 * 1024; diff --git a/src/features/workflows/protocol.ts b/src/features/workflows/protocol.ts index f666d6d3..4154dea3 100644 --- a/src/features/workflows/protocol.ts +++ b/src/features/workflows/protocol.ts @@ -1,12 +1,12 @@ import { parseDocument } from "yaml"; -import type { EventData } from "../relay/events"; +import type { EventData } from "../relay/events.ts"; import type { WorkflowDefinition, WorkflowReference, WorkflowRunCursor, WorkflowRunPage, WorkflowApproval, -} from "./types"; +} from "./types.ts"; export const WORKFLOW_KINDS = [30620, 46020, 5] as const; export function isWorkflowOperation( From 6167504f20ec20825ea80cea9d42e133e4429f36 Mon Sep 17 00:00:00 2001 From: Pinky <5f5ab050ec58ae208332edd544ebf705221e24c1b86d82a6ca07038a7a8f6ac9@buzz.block.builderlab.xyz> Date: Sat, 12 Sep 2026 08:59:06 -0600 Subject: [PATCH 06/20] feat(workflows): add guarded configuration editor and run history UI Signed-off-by: Pinky <5f5ab050ec58ae208332edd544ebf705221e24c1b86d82a6ca07038a7a8f6ac9@buzz.block.builderlab.xyz> Signed-off-by: Brain <1a02c72794dcd0f07058a353bc3a81f4028b8c77c92c87fce6d5c8b85970a20b@buzz.block.builderlab.xyz> --- src/bundled/workflows/ConfirmAction.tsx | 46 ++ src/bundled/workflows/WorkflowChannel.tsx | 441 ++++++++++++ src/bundled/workflows/WorkflowEditor.tsx | 191 +++++ src/bundled/workflows/WorkflowForm.tsx | 217 ++++++ src/bundled/workflows/WorkflowOperations.tsx | 73 ++ src/bundled/workflows/WorkflowRuns.tsx | 190 +++++ src/bundled/workflows/WorkflowsPage.tsx | 138 ++++ src/bundled/workflows/cronExpression.test.mjs | 54 ++ src/bundled/workflows/cronExpression.ts | 157 +++++ src/bundled/workflows/editor-model.test.ts | 45 ++ src/bundled/workflows/editor-model.ts | 119 ++++ src/bundled/workflows/fixture.html | 1 + src/bundled/workflows/fixture.tsx | 63 ++ src/bundled/workflows/fixtures.ts | 203 ++++++ src/bundled/workflows/index.tsx | 12 + src/bundled/workflows/manifest.json | 1 + src/bundled/workflows/useWorkflowView.ts | 16 + .../workflowActivationWarning.test.mjs | 107 +++ .../workflows/workflowActivationWarning.ts | 161 +++++ .../workflows/workflowDuration.test.mjs | 55 ++ src/bundled/workflows/workflowDuration.ts | 129 ++++ .../workflows/workflowFormTypes.test.mjs | 322 +++++++++ src/bundled/workflows/workflowFormTypes.ts | 657 ++++++++++++++++++ .../workflows/workflowYamlDocument.test.mjs | 152 ++++ src/bundled/workflows/workflowYamlDocument.ts | 115 +++ src/bundled/workflows/workflows.css | 129 ++++ src/bundled/workflows/workflows.journey.mjs | 122 ++++ .../workflows/workflows.playwright.config.mjs | 24 + 28 files changed, 3940 insertions(+) create mode 100644 src/bundled/workflows/ConfirmAction.tsx create mode 100644 src/bundled/workflows/WorkflowChannel.tsx create mode 100644 src/bundled/workflows/WorkflowEditor.tsx create mode 100644 src/bundled/workflows/WorkflowForm.tsx create mode 100644 src/bundled/workflows/WorkflowOperations.tsx create mode 100644 src/bundled/workflows/WorkflowRuns.tsx create mode 100644 src/bundled/workflows/WorkflowsPage.tsx create mode 100644 src/bundled/workflows/cronExpression.test.mjs create mode 100644 src/bundled/workflows/cronExpression.ts create mode 100644 src/bundled/workflows/editor-model.test.ts create mode 100644 src/bundled/workflows/editor-model.ts create mode 100644 src/bundled/workflows/fixture.html create mode 100644 src/bundled/workflows/fixture.tsx create mode 100644 src/bundled/workflows/fixtures.ts create mode 100644 src/bundled/workflows/index.tsx create mode 100644 src/bundled/workflows/manifest.json create mode 100644 src/bundled/workflows/useWorkflowView.ts create mode 100644 src/bundled/workflows/workflowActivationWarning.test.mjs create mode 100644 src/bundled/workflows/workflowActivationWarning.ts create mode 100644 src/bundled/workflows/workflowDuration.test.mjs create mode 100644 src/bundled/workflows/workflowDuration.ts create mode 100644 src/bundled/workflows/workflowFormTypes.test.mjs create mode 100644 src/bundled/workflows/workflowFormTypes.ts create mode 100644 src/bundled/workflows/workflowYamlDocument.test.mjs create mode 100644 src/bundled/workflows/workflowYamlDocument.ts create mode 100644 src/bundled/workflows/workflows.css create mode 100644 src/bundled/workflows/workflows.journey.mjs create mode 100644 src/bundled/workflows/workflows.playwright.config.mjs diff --git a/src/bundled/workflows/ConfirmAction.tsx b/src/bundled/workflows/ConfirmAction.tsx new file mode 100644 index 00000000..206857e1 --- /dev/null +++ b/src/bundled/workflows/ConfirmAction.tsx @@ -0,0 +1,46 @@ +import { AlertDialog } from "@base-ui/react/alert-dialog"; +import { Button } from "../../shared/design-system/ui/Button"; + +export function ConfirmAction({ + title, + description, + action, + onConfirm, + onCancel, +}: { + title: string; + description: string; + action: string; + onConfirm: () => void; + onCancel: () => void; +}) { + return ( + { + if (!open) onCancel(); + }} + > + + + + + {title} + + + {description} + +
+ + +
+
+
+
+ ); +} diff --git a/src/bundled/workflows/WorkflowChannel.tsx b/src/bundled/workflows/WorkflowChannel.tsx new file mode 100644 index 00000000..953acfc6 --- /dev/null +++ b/src/bundled/workflows/WorkflowChannel.tsx @@ -0,0 +1,441 @@ +import { + useEffect, + useMemo, + useRef, + useState, + useSyncExternalStore, +} from "react"; +import type { + WorkflowCapability, + WorkflowDefinition, +} from "../../features/workflows/types"; +import { Button } from "../../shared/design-system/ui/Button"; +import { ConfirmAction } from "./ConfirmAction"; +import { WorkflowEditor } from "./WorkflowEditor"; +import { WorkflowOperations } from "./WorkflowOperations"; +import { WorkflowRuns } from "./WorkflowRuns"; +import { exactSaveReadback } from "./editor-model"; +import { DEFAULT_FORM_STATE, formStateToYaml } from "./workflowFormTypes"; +import { readWorkflowDocumentFields } from "./workflowYamlDocument"; +import { useWorkflowView } from "./useWorkflowView"; + +type Draft = { + original: WorkflowDefinition | undefined; + yaml: string; + initial: string; + operationId?: string; +}; + +export function WorkflowChannel({ + capability, + channelId, + channelName, + viewer, + onDraftRiskChange, +}: { + capability: WorkflowCapability; + channelId: string; + channelName: string; + viewer: string; + onDraftRiskChange?: (atRisk: boolean) => void; +}) { + const view = useMemo( + () => capability.definitions(channelId), + [capability, channelId], + ); + const snapshot = useWorkflowView(view); + const operations = useSyncExternalStore( + capability.operations.subscribe, + capability.operations.snapshot, + capability.operations.snapshot, + ); + const submission = useRef(null); + const [draft, setDraft] = useState(null); + const [pendingSelection, setPendingSelection] = useState< + WorkflowDefinition | "new" | "close" | null + >(null); + const [confirmDelete, setConfirmDelete] = useState(false); + const [error, setError] = useState(null); + const [readRuns, setReadRuns] = useState(false); + const operation = draft?.operationId + ? operations.find((item) => item.eventId === draft.operationId) + : undefined; + const ownOperations = operations.filter( + (item) => item.workflow.channelId === channelId, + ); + const busy = + !!draft?.operationId && (!operation || operation.outcome === "pending"); + const readonly = !!draft?.original && draft.original.owner !== viewer; + const dirty = !!draft && draft.yaml !== draft.initial; + const atRisk = dirty || !!draft?.operationId; + const unresolvedWrite = ownOperations.some( + (item) => + (item.outcome === "pending" || item.outcome === "unknown") && + (draft?.original + ? item.workflow.id === draft.original.id && + item.workflow.owner === draft.original.owner + : item.action === "save"), + ); + useEffect(() => { + onDraftRiskChange?.(atRisk); + return () => onDraftRiskChange?.(false); + }, [atRisk, onDraftRiskChange]); + useEffect(() => { + if (!atRisk) return; + const warn = (event: BeforeUnloadEvent) => event.preventDefault(); + window.addEventListener("beforeunload", warn); + return () => window.removeEventListener("beforeunload", warn); + }, [atRisk]); + const open = (next: WorkflowDefinition | "new" | "close") => { + const yaml = + next === "new" + ? formStateToYaml({ ...DEFAULT_FORM_STATE, name: "Untitled workflow" }) + : next === "close" + ? "" + : next.yaml; + setDraft( + next === "close" + ? null + : { + original: typeof next === "string" ? undefined : next, + yaml, + initial: yaml, + }, + ); + submission.current = null; + setError(null); + setReadRuns(false); + setPendingSelection(null); + setConfirmDelete(false); + }; + const select = (next: WorkflowDefinition | "new" | "close") => { + if (atRisk) setPendingSelection(next); + else open(next); + }; + useEffect(() => { + if (operation?.eventId && operation.outcome === "succeeded") + void view.refresh(); + }, [operation?.eventId, operation?.outcome, view]); + useEffect(() => { + if (!operation || !draft || snapshot.status !== "ready") return; + const saved = exactSaveReadback(operation, snapshot.data.items); + if (saved) { + submission.current = null; + setDraft({ original: saved, yaml: saved.yaml, initial: saved.yaml }); + } + }, [operation, snapshot, draft]); + // A cleared/unavailable view withdraws the saved private definition from display. + // Unsaved user-authored drafts never become a second retained definition cache. + useEffect(() => { + if (snapshot.status === "unavailable" || snapshot.status === "idle") { + submission.current = null; + setDraft(null); + setPendingSelection(null); + setConfirmDelete(false); + setReadRuns(false); + setError(null); + } + }, [snapshot.status]); + const save = () => { + if ( + !draft || + submission.current || + busy || + unresolvedWrite || + readonly || + !capability.availability.save || + draft.operationId + ) + return; + try { + submission.current = "submitting"; + const operationId = capability.save({ + channelId, + yaml: draft.yaml, + ...(draft.original ? { existing: draft.original } : {}), + }); + submission.current = operationId; + setDraft({ ...draft, operationId }); + setError(null); + } catch (cause) { + submission.current = null; + setError( + cause instanceof Error + ? cause.message + : "The workflow could not be submitted. Your draft is retained.", + ); + } + }; + const remove = () => { + if ( + !draft?.original || + submission.current || + readonly || + unresolvedWrite || + !capability.availability.delete || + draft.operationId + ) + return; + try { + submission.current = "submitting"; + const operationId = capability.delete(draft.original); + submission.current = operationId; + setDraft({ ...draft, operationId }); + setConfirmDelete(false); + setError(null); + } catch (cause) { + submission.current = null; + setError( + cause instanceof Error + ? cause.message + : "Deletion could not be submitted. Your draft is retained.", + ); + } + }; + const trigger = () => { + if ( + !draft?.original || + submission.current || + dirty || + unresolvedWrite || + draft.operationId || + readonly || + !capability.availability.trigger + ) + return; + try { + submission.current = "submitting"; + submission.current = capability.trigger(draft.original); + setError(null); + } catch (cause) { + submission.current = null; + setError( + cause instanceof Error ? cause.message : "Run could not be submitted.", + ); + } + }; + useEffect(() => { + const active = operations.find( + (item) => item.eventId === submission.current, + ); + if ( + active?.action === "trigger" && + (active.outcome === "succeeded" || active.outcome === "rejected") + ) + submission.current = null; + }, [operations]); + let blocked: string | undefined; + if (!capability.availability.save) + blocked = "Saving is unavailable from this host."; + else if (unresolvedWrite && !draft?.operationId) + blocked = + "An operation for this workflow is unresolved. Review its retained identity below; do not submit a replacement."; + else if (draft?.operationId) + blocked = + operation?.outcome === "succeeded" + ? operation.action === "delete" + ? "Deletion completed. Close this draft; the configuration list is being refreshed." + : "Save completed; waiting for a readback of this exact signed revision. A different head must be reviewed before editing again." + : operation?.outcome === "rejected" + ? "Save rejected. Your draft is retained; review the error before retrying." + : "This operation has not been resolved. Your draft and operation identity are retained."; + return ( +
+
+

Saved configurations

+ + +
+

+ Configured state is not runtime health. Historical configurations may no + longer have a runtime workflow. +

+ {snapshot.status === "loading" && ( +

Reading configurations…

+ )} + {snapshot.status === "idle" && ( +

Configurations cleared. Refresh to read again.

+ )} + {snapshot.status === "unavailable" && ( +

Workflow definitions are unavailable.

+ )} + {snapshot.status === "error" && ( +

+ {snapshot.error ?? + "Configurations could not be read. Use Refresh configurations to retry."} +

+ )} + {snapshot.data.partial && ( +

This is a bounded, partial list.

+ )} + {snapshot.status === "ready" && !snapshot.data.items.length && ( +

+ No saved configurations returned for this channel. Create a disabled + draft to start. +

+ )} +
    + {snapshot.data.items.map((definition) => { + const header = readWorkflowDocumentFields(definition.yaml); + return ( +
  • + + + {header.editable + ? header.enabled === false + ? "Configured disabled" + : "Configured enabled" + : "Unreadable configuration"} + {definition.owner !== viewer ? " · Read-only" : ""} + +
  • + ); + })} +
+ {draft && + snapshot.status !== "unavailable" && + snapshot.status !== "idle" && ( +
+
+

+ {draft.original ? "Workflow details" : "New workflow"} +

+ +
+ {draft.original && ( + <> +

+ Owner: {draft.original.owner} +

+

+ Revision: {draft.original.revision} +

+ + )} + {readonly && ( +

+ This definition belongs to another identity. Only its author can + manage it here. +

+ )} + setDraft({ ...draft, yaml })} + onSave={save} + readOnly={readonly} + busy={busy} + locked={!!draft.operationId} + blocked={blocked} + /> +

+ Drafts stay in this editor only. Leaving the Workflows page or + reloading discards unsaved text, but does not cancel submitted + operations. +

+ {error && ( +

+ {error} +

+ )} + {operation?.outcome === "rejected" && ( + + )} + {draft.original && ( +
+ + {!readonly && ( + <> + + + + )} +
+ )} + {draft.original && !capability.availability.delete && !readonly && ( +

+ Delete is unavailable until the relay proves support for + consistent workflow deletion. +

+ )} + {draft.original && readRuns && ( + + )} + {confirmDelete && ( + setConfirmDelete(false)} + /> + )} +
+ )} + + {pendingSelection && + snapshot.status !== "idle" && + snapshot.status !== "unavailable" && ( + open(pendingSelection)} + onCancel={() => setPendingSelection(null)} + /> + )} +
+ ); +} diff --git a/src/bundled/workflows/WorkflowEditor.tsx b/src/bundled/workflows/WorkflowEditor.tsx new file mode 100644 index 00000000..00907ffc --- /dev/null +++ b/src/bundled/workflows/WorkflowEditor.tsx @@ -0,0 +1,191 @@ +import { Input } from "@base-ui/react/input"; +import { useId, useState } from "react"; +import { Button } from "../../shared/design-system/ui/Button"; +import { Switch } from "../../shared/design-system/ui/Switch"; +import { Tabs } from "../../shared/design-system/ui/Tabs"; +import { ConfirmAction } from "./ConfirmAction"; +import { WorkflowForm } from "./WorkflowForm"; +import { draftError, hasWebhookTrigger, visualForm } from "./editor-model"; +import { getWorkflowActivationWarning } from "./workflowActivationWarning"; +import { formStateToYaml, type WorkflowFormState } from "./workflowFormTypes"; +import { + readWorkflowDocumentFields, + yamlWithWorkflowEnabled, + yamlWithWorkflowName, +} from "./workflowYamlDocument"; + +export function WorkflowEditor({ + yaml, + onChange, + onSave, + readOnly = false, + blocked, + busy = false, + locked = false, +}: { + yaml: string; + onChange: (yaml: string) => void; + onSave: () => void; + readOnly?: boolean; + blocked?: string | undefined; + busy?: boolean; + locked?: boolean; +}) { + const id = useId(); + const [mode, setMode] = useState<"form" | "yaml">(() => + visualForm(yaml).ok ? "form" : "yaml", + ); + const [formDraft, setFormDraft] = useState(null); + const [formYaml, setFormYaml] = useState(yaml); + const [activating, setActivating] = useState(false); + const [modeError, setModeError] = useState(null); + const fields = readWorkflowDocumentFields(yaml); + const parsed = visualForm(yaml); + // An incomplete form is still an editable draft. An external YAML change must + // be reparsed instead of reviving stale form state. + const form = + formYaml === yaml && formDraft + ? formDraft + : parsed.ok + ? parsed.state + : null; + const error = draftError(yaml); + const secretGate = hasWebhookTrigger(yaml) + ? "Webhook-trigger saves are unavailable until secure one-time-secret display is supported." + : undefined; + const unavailable = blocked || secretGate; + const mutateForm = (state: WorkflowFormState) => { + const next = formStateToYaml(state); + setFormDraft(state); + setFormYaml(next); + onChange(next); + }; + const changeHeader = ( + next: string | null, + patch: Partial, + ) => { + if (next === null) return; + if (form) { + setFormDraft({ ...form, ...patch }); + setFormYaml(next); + } + onChange(next); + }; + const warning = getWorkflowActivationWarning(yaml); + const submit = () => { + if (readOnly || busy || locked || unavailable || error) return; + if (fields.enabled !== false) setActivating(true); + else onSave(); + }; + return ( +
+
+ + + changeHeader(yamlWithWorkflowEnabled(yaml, enabled), { enabled }) + } + /> +
+ { + if (next === "form" && !form) { + setModeError(parsed.ok ? null : parsed.error); + return; + } + setModeError(null); + setMode(next); + }} + /> + {modeError && ( +

+ {modeError} +

+ )} + {mode === "form" && form ? ( + + ) : ( +
+ +