import { env } from "@code/env/convex"; import type { Doc } from "./_generated/dataModel"; import { mutation, query } from "./_generated/server"; import type { MutationCtx, QueryCtx } from "./_generated/server"; import { v } from "convex/values"; const FLUE_SCHEMA_VERSION = "4"; const DURABILITY_DEFAULT_MAX_ATTEMPTS = 10; const DURABILITY_DEFAULT_TIMEOUT_MS = 3_600_000; const LEASE_DURATION_MS = 30_000; type SubmissionStatus = "queued" | "running" | "terminalizing" | "settled"; type RunStatus = "active" | "completed" | "errored"; type SettledOutcome = "completed" | "failed" | "aborted"; type SubmissionDoc = Doc<"flueSubmissions">; type AttemptMarkerDoc = Doc<"flueAttemptMarkers">; type ConversationStreamDoc = Doc<"flueConversationStreams">; type ConversationBatchDoc = Doc<"flueConversationBatches">; type EventEntryDoc = Doc<"flueEventEntries">; type RunDoc = Doc<"flueRuns">; type AttachmentDoc = Doc<"flueAttachments">; type SubmissionRow = { readonly sequence: number; readonly submissionId: string; readonly sessionKey: string; readonly kind: "dispatch" | "direct"; readonly inputJson: string; readonly chunksJson: string | null; readonly status: SubmissionStatus; readonly acceptedAt: number; readonly canonicalReadyAt: number | null; readonly attemptId: string | null; readonly inputAppliedAt: number | null; readonly recoveryRequestedAt: number | null; readonly abortRequestedAt: number | null; readonly startedAt: number | null; readonly error: string | null; readonly attemptCount: number; readonly maxRetry: number; readonly timeoutAt: number; readonly ownerId: string | null; readonly leaseExpiresAt: number; readonly traceCarrierJson: string | null; }; type SettlementObligationRow = { readonly submissionId: string; readonly sessionKey: string; readonly attemptId: string; readonly recordId: string; readonly recordJson: string; }; type AttemptMarkerRow = { readonly submissionId: string; readonly attemptId: string; readonly createdAt: number; }; type ConversationStreamIdentity = { readonly agentName: string; readonly instanceId: string; }; type ConversationProducerClaim = { readonly producerId: string; readonly producerEpoch: number; readonly incarnation: string; readonly nextProducerSequence: number; readonly offset: string; }; type ConversationStreamRow = { readonly identity: ConversationStreamIdentity; readonly incarnation: string; readonly nextOffset: number; readonly producerId: string | null; readonly producerEpoch: number; readonly nextProducerSequence: number; }; type ConversationBatchRow = { readonly offset: number; readonly recordsJson: string; }; type EventEntryRow = { readonly offset: number; readonly dataJson: string; }; type RunRow = { readonly runId: string; readonly workflowName: string; readonly status: RunStatus; readonly startedAt: string; readonly endedAt: string | null; readonly isError: boolean | null; readonly durationMs: number | null; readonly inputJson: string | null; readonly resultJson: string | null; readonly errorJson: string | null; readonly traceCarrierJson: string | null; }; type RunPointerRow = { readonly runId: string; readonly workflowName: string; readonly status: RunStatus; readonly startedAt: string; readonly endedAt: string | null; readonly durationMs: number | null; readonly isError: boolean | null; }; type AttachmentRefWire = { readonly id: string; readonly mimeType: string; readonly size: number; readonly digest: string; readonly filename?: string; }; type AttachmentWire = { readonly attachment: AttachmentRefWire; readonly bytes: ArrayBuffer; }; type AdmitSubmissionResponse = { readonly kind: "submission" | "retained_receipt" | "conflict"; readonly submission?: SubmissionRow; readonly receipt?: { readonly submissionId: string; readonly acceptedAt: number }; }; type AppendResult = { readonly offset: number; readonly appended: boolean }; type ListRunsCursor = { readonly startedAt: string; readonly runId: string }; type OwnedConversationRecord = { readonly id?: string; readonly type?: string; readonly submissionId?: string; readonly attemptId?: string; }; const submissionKind = v.union(v.literal("dispatch"), v.literal("direct")); const runStatus = v.union( v.literal("active"), v.literal("completed"), v.literal("errored"), ); const attachmentRef = v.object({ id: v.string(), mimeType: v.string(), size: v.number(), digest: v.string(), filename: v.optional(v.string()), }); const conversationIdentity = v.object({ agentName: v.string(), instanceId: v.string(), }); const tokenArgs = { token: v.string() }; const assertToken = (token: string): void => { if (token !== env.FLUE_DB_TOKEN) { throw new Error("[flue] Unauthorized persistence request."); } }; const safeJsonParse = (text: string): T => JSON.parse(text) as T; const toSubmissionRow = (doc: SubmissionDoc): SubmissionRow => ({ sequence: doc.sequence, submissionId: doc.submissionId, sessionKey: doc.sessionKey, kind: doc.kind, inputJson: doc.inputJson, chunksJson: doc.chunksJson, status: doc.status, acceptedAt: doc.acceptedAt, canonicalReadyAt: doc.canonicalReadyAt ?? null, attemptId: doc.attemptId ?? null, inputAppliedAt: doc.inputAppliedAt ?? null, recoveryRequestedAt: doc.recoveryRequestedAt ?? null, abortRequestedAt: doc.abortRequestedAt ?? null, startedAt: doc.startedAt ?? null, error: doc.error ?? null, attemptCount: doc.attemptCount, maxRetry: doc.maxRetry, timeoutAt: doc.timeoutAt, ownerId: doc.ownerId ?? null, leaseExpiresAt: doc.leaseExpiresAt, traceCarrierJson: doc.traceCarrierJson ?? null, }); const toSettlementObligationRow = ( doc: SubmissionDoc, ): SettlementObligationRow | null => { if ( doc.attemptId === undefined || doc.settlementRecordId === undefined || doc.settlementRecordJson === undefined ) { return null; } return { submissionId: doc.submissionId, sessionKey: doc.sessionKey, attemptId: doc.attemptId, recordId: doc.settlementRecordId, recordJson: doc.settlementRecordJson, }; }; const toAttemptMarkerRow = (doc: AttemptMarkerDoc): AttemptMarkerRow => ({ submissionId: doc.submissionId, attemptId: doc.attemptId, createdAt: doc.createdAt, }); const toConversationStreamRow = ( doc: ConversationStreamDoc, ): ConversationStreamRow => ({ identity: safeJsonParse(doc.identityJson), incarnation: doc.incarnation, nextOffset: doc.nextOffset - 1, producerId: doc.producerId ?? null, producerEpoch: doc.producerEpoch, nextProducerSequence: doc.nextProducerSequence, }); const toConversationBatchRow = (doc: ConversationBatchDoc): ConversationBatchRow => ({ offset: doc.seq, recordsJson: doc.recordsJson, }); const toEventEntryRow = (doc: EventEntryDoc): EventEntryRow => ({ offset: doc.seq, dataJson: doc.dataJson, }); const toRunRow = (doc: RunDoc): RunRow => ({ runId: doc.runId, workflowName: doc.workflowName, status: doc.status, startedAt: doc.startedAt, endedAt: doc.endedAt ?? null, isError: doc.isError ?? null, durationMs: doc.durationMs ?? null, inputJson: doc.inputJson ?? null, resultJson: doc.resultJson ?? null, errorJson: doc.errorJson ?? null, traceCarrierJson: doc.traceCarrierJson ?? null, }); const toRunPointerRow = (doc: RunDoc): RunPointerRow => ({ runId: doc.runId, workflowName: doc.workflowName, status: doc.status, startedAt: doc.startedAt, endedAt: doc.endedAt ?? null, durationMs: doc.durationMs ?? null, isError: doc.isError ?? null, }); const toAttachmentWire = (doc: AttachmentDoc): AttachmentWire => ({ attachment: { id: doc.attachmentId, mimeType: doc.mimeType, size: doc.size, digest: doc.digest, ...(doc.filename === undefined ? {} : { filename: doc.filename }), }, bytes: doc.bytes, }); const compareBuffers = (left: ArrayBuffer, right: ArrayBuffer): boolean => { const leftBytes = new Uint8Array(left); const rightBytes = new Uint8Array(right); if (leftBytes.byteLength !== rightBytes.byteLength) { return false; } for (let i = 0; i < leftBytes.byteLength; i += 1) { if (leftBytes[i] !== rightBytes[i]) { return false; } } return true; }; const sameAttachment = ( existing: AttachmentDoc, input: { readonly conversationId: string; readonly attachment: AttachmentRefWire; readonly bytes: ArrayBuffer; }, ): boolean => existing.conversationId === input.conversationId && existing.attachmentId === input.attachment.id && existing.mimeType === input.attachment.mimeType && existing.size === input.attachment.size && existing.digest === input.attachment.digest && existing.filename === input.attachment.filename && compareBuffers(existing.bytes, input.bytes); const compareRunPointerDesc = (left: RunPointerRow, right: RunPointerRow): number => { if (left.startedAt !== right.startedAt) { return left.startedAt < right.startedAt ? 1 : -1; } if (left.runId !== right.runId) { return left.runId < right.runId ? 1 : -1; } return 0; }; const isUnsettled = (status: SubmissionStatus): boolean => status !== "settled"; const parseSessionInstance = (sessionKey: string): string | undefined => { if (!sessionKey.startsWith("agent-session:")) { return undefined; } try { const parsed = JSON.parse(sessionKey.slice("agent-session:".length)) as unknown; if ( Array.isArray(parsed) && typeof parsed[0] === "string" ) { return parsed[0]; } } catch {} return undefined; }; const parseOwnedRecord = (record: unknown): OwnedConversationRecord | null => { if (typeof record !== "object" || record === null) { return null; } const value = record as Record; return { ...(typeof value.id === "string" ? { id: value.id } : {}), ...(typeof value.type === "string" ? { type: value.type } : {}), ...(typeof value.submissionId === "string" ? { submissionId: value.submissionId } : {}), ...(typeof value.attemptId === "string" ? { attemptId: value.attemptId } : {}), }; }; const parseSettledOutcome = (recordJson: string | undefined): SettledOutcome | undefined => { if (recordJson === undefined) { return undefined; } try { const parsed = JSON.parse(recordJson) as unknown; if (typeof parsed !== "object" || parsed === null) { return undefined; } const value = parsed as Record; return value.outcome === "completed" || value.outcome === "failed" || value.outcome === "aborted" ? value.outcome : undefined; } catch { return undefined; } }; type ReadCtx = QueryCtx | MutationCtx; const getSubmissionDoc = async ( ctx: ReadCtx, submissionId: string, ): Promise => ctx.db .query("flueSubmissions") .withIndex("by_submissionId", (q) => q.eq("submissionId", submissionId)) .unique(); const getConversationStreamDoc = async ( ctx: ReadCtx, path: string, ): Promise => ctx.db .query("flueConversationStreams") .withIndex("by_path", (q) => q.eq("path", path)) .unique(); const assertSubmissionAuthorization = async ( ctx: ReadCtx, path: string, submission: | { readonly submissionId: string; readonly attemptId: string; } | undefined, recordsJson: string, ): Promise => { const records = safeJsonParse(recordsJson); const owned = records .map(parseOwnedRecord) .filter( (record): record is OwnedConversationRecord => record !== null && (record.submissionId !== undefined || record.attemptId !== undefined), ); if (submission === undefined) { if (owned.length > 0) { throw new Error( `[flue] Conversation stream \"${path}\" received submission-owned records without authorization.`, ); } return; } if ( owned.some( (record) => record.submissionId !== submission.submissionId || record.attemptId !== submission.attemptId, ) ) { throw new Error( `[flue] Conversation stream \"${path}\" record ownership does not match the authorized submission attempt.`, ); } const stream = await getConversationStreamDoc(ctx, path); const stored = await getSubmissionDoc(ctx, submission.submissionId); if (stream === null || stored === null) { throw new Error( `[flue] Submission attempt no longer owns work for agent instance \"${path}\".`, ); } const streamIdentity = safeJsonParse(stream.identityJson); const terminalizingSettlement = stored.status === "terminalizing" && stored.attemptId === submission.attemptId && owned.length === 1 && owned[0]?.type === "submission_settled" && owned[0]?.id === stored.settlementRecordId && JSON.stringify(records[0]) === stored.settlementRecordJson; const validInstance = parseSessionInstance(stored.sessionKey) === streamIdentity.instanceId; if ( !terminalizingSettlement && !( stored.status === "running" && stored.attemptId === submission.attemptId && validInstance ) ) { throw new Error( `[flue] Submission attempt no longer owns work for agent instance \"${path}\".`, ); } }; export const checkSchemaVersion = query({ args: tokenArgs, handler: async (ctx, args) => { assertToken(args.token); const meta = await ctx.db .query("flueMeta") .withIndex("by_key", (q) => q.eq("key", "schema_version")) .unique(); return meta?.value ?? FLUE_SCHEMA_VERSION; }, }); export const getSubmission = query({ args: { ...tokenArgs, submissionId: v.string() }, handler: async (ctx, args) => { assertToken(args.token); const submission = await ctx.db .query("flueSubmissions") .withIndex("by_submissionId", (q) => q.eq("submissionId", args.submissionId)) .unique(); return submission === null ? null : toSubmissionRow(submission); }, }); export const hasUnsettledSubmissions = query({ args: tokenArgs, handler: async (ctx, args) => { assertToken(args.token); const rows = await ctx.db.query("flueSubmissions").collect(); return rows.some((row) => isUnsettled(row.status)); }, }); export const listRunnableSubmissions = query({ args: tokenArgs, handler: async (ctx, args) => { assertToken(args.token); const rows = (await ctx.db.query("flueSubmissions").collect()).sort( (left, right) => left.sequence - right.sequence, ); return rows .filter( (row) => row.status === "queued" && row.canonicalReadyAt !== undefined && !rows.some( (candidate) => candidate.sessionKey === row.sessionKey && candidate.sequence < row.sequence && isUnsettled(candidate.status), ), ) .map(toSubmissionRow); }, }); export const listUnreadySubmissions = query({ args: tokenArgs, handler: async (ctx, args) => { assertToken(args.token); return (await ctx.db.query("flueSubmissions").collect()) .filter( (row) => row.status === "queued" && row.canonicalReadyAt === undefined, ) .sort((left, right) => left.sequence - right.sequence) .map(toSubmissionRow); }, }); export const listRunningSubmissions = query({ args: tokenArgs, handler: async (ctx, args) => { assertToken(args.token); return (await ctx.db.query("flueSubmissions").collect()) .filter((row) => row.status === "running") .sort((left, right) => left.sequence - right.sequence) .map(toSubmissionRow); }, }); export const listPendingSubmissionSettlements = query({ args: tokenArgs, handler: async (ctx, args) => { assertToken(args.token); return (await ctx.db.query("flueSubmissions").collect()) .filter((row) => row.status === "terminalizing") .sort((left, right) => left.sequence - right.sequence) .map(toSettlementObligationRow) .filter((row): row is SettlementObligationRow => row !== null); }, }); export const replaceSubmissionAttempt = mutation({ args: { ...tokenArgs, submissionId: v.string(), attemptId: v.string(), nextAttemptId: v.string(), ownerId: v.optional(v.string()), leaseExpiresAt: v.optional(v.number()), }, handler: async (ctx, args) => { assertToken(args.token); const submission = await ctx.db .query("flueSubmissions") .withIndex("by_submissionId", (q) => q.eq("submissionId", args.submissionId)) .unique(); if ( submission === null || submission.status !== "running" || submission.attemptId !== args.attemptId ) { return null; } const now = Date.now(); await ctx.db.patch(submission._id, { attemptId: args.nextAttemptId, recoveryRequestedAt: undefined, startedAt: now, attemptCount: submission.attemptCount + 1, ...(args.ownerId === undefined ? { ownerId: undefined } : { ownerId: args.ownerId }), ...(args.leaseExpiresAt === undefined ? { leaseExpiresAt: submission.leaseExpiresAt } : { leaseExpiresAt: args.leaseExpiresAt }), updatedAt: now, }); const updated = await ctx.db.get(submission._id); return updated === null ? null : toSubmissionRow(updated); }, }); export const admitSubmission = mutation({ args: { ...tokenArgs, input: v.object({ kind: submissionKind, submissionId: v.string(), sessionKey: v.string(), acceptedAt: v.number(), inputJson: v.string(), chunksJson: v.optional(v.string()), traceCarrierJson: v.optional(v.string()), }), }, handler: async (ctx, args): Promise => { assertToken(args.token); const existing = await ctx.db .query("flueSubmissions") .withIndex("by_submissionId", (q) => q.eq("submissionId", args.input.submissionId)) .unique(); if (existing !== null) { const exactMatch = existing.kind === args.input.kind && existing.sessionKey === args.input.sessionKey && existing.acceptedAt === args.input.acceptedAt && existing.inputJson === args.input.inputJson && existing.chunksJson === (args.input.chunksJson ?? "[]") && (existing.traceCarrierJson ?? undefined) === args.input.traceCarrierJson; if (!exactMatch) { return { kind: "conflict" }; } if (existing.kind === "dispatch") { return { kind: "retained_receipt", receipt: { submissionId: existing.submissionId, acceptedAt: existing.acceptedAt, }, }; } return { kind: "submission", submission: toSubmissionRow(existing) }; } const all = await ctx.db.query("flueSubmissions").collect(); const nextSequence = all.reduce( (max, row) => (row.sequence > max ? row.sequence : max), -1, ) + 1; const now = Date.now(); const rowId = await ctx.db.insert("flueSubmissions", { submissionId: args.input.submissionId, sessionKey: args.input.sessionKey, sequence: nextSequence, kind: args.input.kind, inputJson: args.input.inputJson, chunksJson: args.input.chunksJson ?? "[]", traceCarrierJson: args.input.traceCarrierJson, status: "queued", acceptedAt: args.input.acceptedAt, attemptCount: 0, maxRetry: DURABILITY_DEFAULT_MAX_ATTEMPTS, timeoutAt: 0, leaseExpiresAt: 0, updatedAt: now, }); const created = await ctx.db.get(rowId); if (created === null) { throw new Error("[flue] Failed to create submission row."); } return { kind: "submission", submission: toSubmissionRow(created) }; }, }); export const markSubmissionCanonicalReady = mutation({ args: { ...tokenArgs, submissionId: v.string() }, handler: async (ctx, args) => { assertToken(args.token); const submission = await ctx.db .query("flueSubmissions") .withIndex("by_submissionId", (q) => q.eq("submissionId", args.submissionId)) .unique(); if (submission === null || submission.status !== "queued") { return null; } if (submission.canonicalReadyAt === undefined) { await ctx.db.patch(submission._id, { canonicalReadyAt: Date.now(), updatedAt: Date.now(), }); } const updated = await ctx.db.get(submission._id); return updated === null ? null : toSubmissionRow(updated); }, }); export const claimSubmission = mutation({ args: { ...tokenArgs, submissionId: v.string(), attemptId: v.string(), ownerId: v.string(), leaseExpiresAt: v.number(), }, handler: async (ctx, args) => { assertToken(args.token); const rows = (await ctx.db.query("flueSubmissions").collect()).sort( (left, right) => left.sequence - right.sequence, ); const submission = rows.find((row) => row.submissionId === args.submissionId) ?? null; if ( submission === null || submission.status !== "queued" || submission.canonicalReadyAt === undefined ) { return null; } const hasEarlierUnsettled = rows.some( (row) => row.sessionKey === submission.sessionKey && row.sequence < submission.sequence && isUnsettled(row.status), ); if (hasEarlierUnsettled) { return null; } const now = Date.now(); await ctx.db.patch(submission._id, { status: "running", attemptId: args.attemptId, startedAt: now, attemptCount: submission.attemptCount + 1, ownerId: args.ownerId, leaseExpiresAt: args.leaseExpiresAt, timeoutAt: submission.timeoutAt > 0 ? submission.timeoutAt : now + DURABILITY_DEFAULT_TIMEOUT_MS, maxRetry: submission.maxRetry > 0 ? submission.maxRetry : DURABILITY_DEFAULT_MAX_ATTEMPTS, updatedAt: now, }); const updated = await ctx.db.get(submission._id); return updated === null ? null : toSubmissionRow(updated); }, }); export const markSubmissionInputApplied = mutation({ args: { ...tokenArgs, submissionId: v.string(), attemptId: v.string(), maxRetry: v.optional(v.number()), timeoutAt: v.optional(v.number()), }, handler: async (ctx, args) => { assertToken(args.token); const submission = await ctx.db .query("flueSubmissions") .withIndex("by_submissionId", (q) => q.eq("submissionId", args.submissionId)) .unique(); if ( submission === null || submission.status !== "running" || submission.attemptId !== args.attemptId ) { return false; } const patch: Partial = {}; if (submission.inputAppliedAt === undefined) { patch.inputAppliedAt = Date.now(); if (args.maxRetry !== undefined) { patch.maxRetry = args.maxRetry; } if (args.timeoutAt !== undefined) { patch.timeoutAt = args.timeoutAt; } } patch.updatedAt = Date.now(); await ctx.db.patch(submission._id, patch); return true; }, }); export const requestSubmissionRecovery = mutation({ args: { ...tokenArgs, submissionId: v.string(), attemptId: v.string() }, handler: async (ctx, args) => { assertToken(args.token); const submission = await ctx.db .query("flueSubmissions") .withIndex("by_submissionId", (q) => q.eq("submissionId", args.submissionId)) .unique(); if ( submission === null || submission.status !== "running" || submission.attemptId !== args.attemptId ) { return false; } await ctx.db.patch(submission._id, { recoveryRequestedAt: submission.recoveryRequestedAt ?? Date.now(), updatedAt: Date.now(), }); return true; }, }); export const requestSessionAbort = mutation({ args: { ...tokenArgs, sessionKey: v.string() }, handler: async (ctx, args) => { assertToken(args.token); const rows = (await ctx.db.query("flueSubmissions").collect()).filter( (row) => row.sessionKey === args.sessionKey && (row.status === "queued" || row.status === "running"), ); const now = Date.now(); for (const row of rows) { if (row.abortRequestedAt === undefined) { await ctx.db.patch(row._id, { abortRequestedAt: now, updatedAt: now }); } } return rows.map((row) => row.submissionId); }, }); export const requeueSubmissionBeforeInputApplied = mutation({ args: { ...tokenArgs, submissionId: v.string(), attemptId: v.string() }, handler: async (ctx, args) => { assertToken(args.token); const submission = await ctx.db .query("flueSubmissions") .withIndex("by_submissionId", (q) => q.eq("submissionId", args.submissionId)) .unique(); if ( submission === null || submission.status !== "running" || submission.attemptId !== args.attemptId || submission.inputAppliedAt !== undefined ) { return false; } await ctx.db.patch(submission._id, { status: "queued", attemptId: undefined, recoveryRequestedAt: undefined, startedAt: undefined, ownerId: undefined, leaseExpiresAt: 0, updatedAt: Date.now(), }); return true; }, }); export const reserveSubmissionSettlement = mutation({ args: { ...tokenArgs, submissionId: v.string(), attemptId: v.string(), recordId: v.string(), recordJson: v.string(), }, handler: async (ctx, args) => { assertToken(args.token); const submission = await ctx.db .query("flueSubmissions") .withIndex("by_submissionId", (q) => q.eq("submissionId", args.submissionId)) .unique(); if (submission === null) { return null; } if ( submission.status === "terminalizing" && submission.attemptId === args.attemptId && submission.settlementRecordId === args.recordId && submission.settlementRecordJson === args.recordJson ) { return toSettlementObligationRow(submission); } if ( submission.status !== "running" || submission.attemptId !== args.attemptId ) { return null; } await ctx.db.patch(submission._id, { status: "terminalizing", settlementRecordId: args.recordId, settlementRecordJson: args.recordJson, updatedAt: Date.now(), }); const updated = await ctx.db.get(submission._id); return updated === null ? null : toSettlementObligationRow(updated); }, }); export const finalizeSubmissionSettlement = mutation({ args: { ...tokenArgs, submissionId: v.string(), attemptId: v.string(), recordId: v.string() }, handler: async (ctx, args) => { assertToken(args.token); const submission = await ctx.db .query("flueSubmissions") .withIndex("by_submissionId", (q) => q.eq("submissionId", args.submissionId)) .unique(); if ( submission === null || submission.status !== "terminalizing" || submission.attemptId !== args.attemptId || submission.settlementRecordId !== args.recordId ) { return false; } await ctx.db.patch(submission._id, { status: "settled", settledOutcome: parseSettledOutcome(submission.settlementRecordJson) ?? "completed", ownerId: undefined, leaseExpiresAt: 0, updatedAt: Date.now(), }); return true; }, }); export const settleSubmission = mutation({ args: { ...tokenArgs, submissionId: v.string(), attemptId: v.string(), outcome: v.union(v.literal("completed"), v.literal("failed")), errorJson: v.optional(v.string()), }, handler: async (ctx, args) => { assertToken(args.token); const submission = await ctx.db .query("flueSubmissions") .withIndex("by_submissionId", (q) => q.eq("submissionId", args.submissionId)) .unique(); if ( submission === null || submission.status !== "running" || submission.attemptId !== args.attemptId ) { return false; } await ctx.db.patch(submission._id, { status: "settled", settledOutcome: args.outcome, error: args.errorJson, ownerId: undefined, leaseExpiresAt: 0, updatedAt: Date.now(), }); return true; }, }); export const insertAttemptMarker = mutation({ args: { ...tokenArgs, submissionId: v.string(), attemptId: v.string() }, handler: async (ctx, args) => { assertToken(args.token); const existing = await ctx.db .query("flueAttemptMarkers") .withIndex("by_submissionId_and_attemptId", (q) => q.eq("submissionId", args.submissionId).eq("attemptId", args.attemptId), ) .unique(); if (existing !== null) { return toAttemptMarkerRow(existing); } const rowId = await ctx.db.insert("flueAttemptMarkers", { submissionId: args.submissionId, attemptId: args.attemptId, createdAt: Date.now(), }); const created = await ctx.db.get(rowId); if (created === null) { throw new Error("[flue] Failed to create attempt marker."); } return toAttemptMarkerRow(created); }, }); export const deleteAttemptMarker = mutation({ args: { ...tokenArgs, submissionId: v.string(), attemptId: v.string() }, handler: async (ctx, args) => { assertToken(args.token); const existing = await ctx.db .query("flueAttemptMarkers") .withIndex("by_submissionId_and_attemptId", (q) => q.eq("submissionId", args.submissionId).eq("attemptId", args.attemptId), ) .collect(); for (const row of existing) { await ctx.db.delete(row._id); } }, }); export const listAttemptMarkers = query({ args: tokenArgs, handler: async (ctx, args) => { assertToken(args.token); return (await ctx.db.query("flueAttemptMarkers").collect()) .sort((left, right) => left.createdAt - right.createdAt) .map(toAttemptMarkerRow); }, }); export const renewLeases = mutation({ args: { ...tokenArgs, ownerId: v.string(), submissionIds: v.array(v.string()) }, handler: async (ctx, args) => { assertToken(args.token); const wanted = new Set(args.submissionIds); const now = Date.now(); const nextLease = now + LEASE_DURATION_MS; const rows = await ctx.db.query("flueSubmissions").collect(); for (const row of rows) { if ( row.status === "running" && row.ownerId === args.ownerId && wanted.has(row.submissionId) ) { await ctx.db.patch(row._id, { leaseExpiresAt: nextLease, updatedAt: now, }); } } }, }); export const listExpiredSubmissions = query({ args: tokenArgs, handler: async (ctx, args) => { assertToken(args.token); const now = Date.now(); return (await ctx.db.query("flueSubmissions").collect()) .filter( (row) => row.status === "running" && row.leaseExpiresAt > 0 && row.leaseExpiresAt < now, ) .sort((left, right) => left.sequence - right.sequence) .map(toSubmissionRow); }, }); export const createConversationStream = mutation({ args: { ...tokenArgs, path: v.string(), identity: conversationIdentity }, handler: async (ctx, args) => { assertToken(args.token); const existing = await ctx.db .query("flueConversationStreams") .withIndex("by_path", (q) => q.eq("path", args.path)) .unique(); const identityJson = JSON.stringify(args.identity); if (existing !== null) { if (existing.identityJson !== identityJson) { throw new Error( `[flue] Conversation stream \"${args.path}\" identity conflicts with the existing stream.`, ); } return; } await ctx.db.insert("flueConversationStreams", { path: args.path, identityJson, incarnation: crypto.randomUUID(), producerEpoch: 0, nextProducerSequence: 0, nextOffset: 0, closed: false, createdAt: Date.now(), }); }, }); export const acquireConversationProducer = mutation({ args: { ...tokenArgs, path: v.string(), producerId: v.string() }, handler: async (ctx, args): Promise => { assertToken(args.token); const stream = await ctx.db .query("flueConversationStreams") .withIndex("by_path", (q) => q.eq("path", args.path)) .unique(); if (stream === null) { throw new Error(`[flue] Conversation stream \"${args.path}\" does not exist.`); } const producerEpoch = stream.producerEpoch + 1; await ctx.db.patch(stream._id, { producerId: args.producerId, producerEpoch, nextProducerSequence: 0, }); return { producerId: args.producerId, producerEpoch, incarnation: stream.incarnation, nextProducerSequence: 0, offset: `${stream.nextOffset - 1}`, }; }, }); export const appendConversationBatch = mutation({ args: { ...tokenArgs, path: v.string(), producerId: v.string(), producerEpoch: v.number(), incarnation: v.string(), producerSequence: v.number(), submission: v.optional( v.object({ submissionId: v.string(), attemptId: v.string() }), ), recordsJson: v.string(), }, handler: async (ctx, args): Promise => { assertToken(args.token); const records = safeJsonParse(args.recordsJson); if (records.length === 0) { throw new Error( `[flue] Conversation stream \"${args.path}\" cannot append an empty canonical batch.`, ); } const stream = await ctx.db .query("flueConversationStreams") .withIndex("by_path", (q) => q.eq("path", args.path)) .unique(); if (stream === null) { throw new Error(`[flue] Conversation stream \"${args.path}\" does not exist.`); } if ( stream.producerId !== args.producerId || stream.producerEpoch !== args.producerEpoch || stream.incarnation !== args.incarnation ) { throw new Error( `[flue] Conversation stream \"${args.path}\" producer ownership is stale.`, ); } const existing = await ctx.db .query("flueConversationBatches") .withIndex("by_path_producer_epoch_producerSequence", (q) => q .eq("path", args.path) .eq("producerId", args.producerId) .eq("producerEpoch", args.producerEpoch) .eq("producerSequence", args.producerSequence), ) .unique(); if (existing !== null) { const sameSubmission = existing.submissionId === args.submission?.submissionId && existing.attemptId === args.submission?.attemptId; if (!sameSubmission || existing.recordsJson !== args.recordsJson) { throw new Error( `[flue] Conversation stream \"${args.path}\" producer sequence has conflicting content.`, ); } return { offset: existing.seq, appended: false }; } if (stream.nextProducerSequence !== args.producerSequence) { throw new Error( `[flue] Conversation stream \"${args.path}\" producer sequence is not the next expected value.`, ); } await assertSubmissionAuthorization(ctx, args.path, args.submission, args.recordsJson); const seq = stream.nextOffset; await ctx.db.insert("flueConversationBatches", { path: args.path, seq, producerId: args.producerId, producerEpoch: args.producerEpoch, producerSequence: args.producerSequence, recordsJson: args.recordsJson, submissionId: args.submission?.submissionId, attemptId: args.submission?.attemptId, appendedAt: Date.now(), }); await ctx.db.patch(stream._id, { nextOffset: seq + 1, nextProducerSequence: stream.nextProducerSequence + 1, }); return { offset: seq, appended: true }; }, }); export const readConversationBatches = query({ args: { ...tokenArgs, path: v.string(), afterOffset: v.number(), limit: v.number() }, handler: async (ctx, args) => { assertToken(args.token); const stream = await ctx.db .query("flueConversationStreams") .withIndex("by_path", (q) => q.eq("path", args.path)) .unique(); if (stream === null) { return { batches: [] as ConversationBatchRow[], nextOffset: -1, upToDate: true, }; } const rows = (await ctx.db .query("flueConversationBatches") .withIndex("by_path_and_seq", (q) => q.eq("path", args.path)) .collect()) .filter((row) => row.seq > args.afterOffset) .sort((left, right) => left.seq - right.seq); const page = rows.slice(0, args.limit); return { batches: page.map(toConversationBatchRow), nextOffset: page.length > 0 ? page[page.length - 1]!.seq : args.afterOffset, upToDate: rows.length <= args.limit, }; }, }); export const getConversationStreamMeta = query({ args: { ...tokenArgs, path: v.string() }, handler: async (ctx, args) => { assertToken(args.token); const stream = await ctx.db .query("flueConversationStreams") .withIndex("by_path", (q) => q.eq("path", args.path)) .unique(); return stream === null ? null : toConversationStreamRow(stream); }, }); export const deleteConversationStream = mutation({ args: { ...tokenArgs, path: v.string() }, handler: async (ctx, args) => { assertToken(args.token); const batches = await ctx.db .query("flueConversationBatches") .withIndex("by_path_and_seq", (q) => q.eq("path", args.path)) .collect(); for (const batch of batches) { await ctx.db.delete(batch._id); } const stream = await ctx.db .query("flueConversationStreams") .withIndex("by_path", (q) => q.eq("path", args.path)) .unique(); if (stream !== null) { await ctx.db.delete(stream._id); } }, }); export const createEventStream = mutation({ args: { ...tokenArgs, path: v.string() }, handler: async (ctx, args) => { assertToken(args.token); const existing = await ctx.db .query("flueEventStreams") .withIndex("by_path", (q) => q.eq("path", args.path)) .unique(); if (existing !== null) { return; } await ctx.db.insert("flueEventStreams", { path: args.path, nextSeq: 0, closed: false, createdAt: Date.now(), }); }, }); export const appendEvent = mutation({ args: { ...tokenArgs, path: v.string(), dataJson: v.string() }, handler: async (ctx, args): Promise => { assertToken(args.token); const stream = await ctx.db .query("flueEventStreams") .withIndex("by_path", (q) => q.eq("path", args.path)) .unique(); if (stream === null) { throw new Error(`[flue] Event stream \"${args.path}\" does not exist.`); } if (stream.closed) { throw new Error(`[flue] Event stream \"${args.path}\" is closed.`); } const seq = stream.nextSeq; await ctx.db.insert("flueEventEntries", { path: args.path, seq, dataJson: args.dataJson, appendedAt: Date.now(), }); await ctx.db.patch(stream._id, { nextSeq: seq + 1 }); return { offset: seq, appended: true }; }, }); export const appendEventOnce = mutation({ args: { ...tokenArgs, path: v.string(), key: v.string(), dataJson: v.string() }, handler: async (ctx, args): Promise => { assertToken(args.token); const stream = await ctx.db .query("flueEventStreams") .withIndex("by_path", (q) => q.eq("path", args.path)) .unique(); if (stream === null) { throw new Error(`[flue] Event stream \"${args.path}\" does not exist.`); } if (stream.closed) { throw new Error(`[flue] Event stream \"${args.path}\" is closed.`); } const existing = await ctx.db .query("flueEventEntries") .withIndex("by_path_and_onceKey", (q) => q.eq("path", args.path).eq("onceKey", args.key), ) .unique(); if (existing !== null) { if (existing.dataJson !== args.dataJson) { throw new Error( `[flue] Event key \"${args.key}\" already has a conflicting payload.`, ); } return { offset: existing.seq, appended: false }; } const seq = stream.nextSeq; await ctx.db.insert("flueEventEntries", { path: args.path, seq, onceKey: args.key, dataJson: args.dataJson, appendedAt: Date.now(), }); await ctx.db.patch(stream._id, { nextSeq: seq + 1 }); return { offset: seq, appended: true }; }, }); export const readEvents = query({ args: { ...tokenArgs, path: v.string(), afterOffset: v.number(), limit: v.number() }, handler: async (ctx, args) => { assertToken(args.token); const stream = await ctx.db .query("flueEventStreams") .withIndex("by_path", (q) => q.eq("path", args.path)) .unique(); if (stream === null) { return { events: [] as EventEntryRow[], nextOffset: -1, upToDate: true, closed: false, }; } const rows = (await ctx.db .query("flueEventEntries") .withIndex("by_path_and_seq", (q) => q.eq("path", args.path)) .collect()) .filter((row) => row.seq > args.afterOffset) .sort((left, right) => left.seq - right.seq); const page = rows.slice(0, args.limit); return { events: page.map(toEventEntryRow), nextOffset: page.length > 0 ? page[page.length - 1]!.seq : args.afterOffset, upToDate: rows.length <= args.limit, closed: stream.closed, }; }, }); export const closeEventStream = mutation({ args: { ...tokenArgs, path: v.string() }, handler: async (ctx, args) => { assertToken(args.token); const stream = await ctx.db .query("flueEventStreams") .withIndex("by_path", (q) => q.eq("path", args.path)) .unique(); if (stream !== null && !stream.closed) { await ctx.db.patch(stream._id, { closed: true }); } }, }); export const getEventStreamMeta = query({ args: { ...tokenArgs, path: v.string() }, handler: async (ctx, args) => { assertToken(args.token); const stream = await ctx.db .query("flueEventStreams") .withIndex("by_path", (q) => q.eq("path", args.path)) .unique(); return stream === null ? null : { nextOffset: stream.nextSeq - 1, closed: stream.closed }; }, }); export const createRun = mutation({ args: { ...tokenArgs, runId: v.string(), workflowName: v.string(), startedAt: v.string(), inputJson: v.string(), traceCarrierJson: v.optional(v.string()), }, handler: async (ctx, args) => { assertToken(args.token); const existing = await ctx.db .query("flueRuns") .withIndex("by_runId", (q) => q.eq("runId", args.runId)) .unique(); if (existing !== null) { return; } await ctx.db.insert("flueRuns", { runId: args.runId, workflowName: args.workflowName, status: "active", startedAt: args.startedAt, inputJson: args.inputJson, traceCarrierJson: args.traceCarrierJson, }); }, }); export const endRun = mutation({ args: { ...tokenArgs, runId: v.string(), endedAt: v.string(), isError: v.boolean(), durationMs: v.number(), resultJson: v.optional(v.string()), errorJson: v.optional(v.string()), }, handler: async (ctx, args) => { assertToken(args.token); const run = await ctx.db .query("flueRuns") .withIndex("by_runId", (q) => q.eq("runId", args.runId)) .unique(); if (run === null) { return; } await ctx.db.patch(run._id, { status: args.isError ? "errored" : "completed", endedAt: args.endedAt, isError: args.isError, durationMs: args.durationMs, resultJson: args.resultJson, errorJson: args.errorJson, }); }, }); export const getRun = query({ args: { ...tokenArgs, runId: v.string() }, handler: async (ctx, args) => { assertToken(args.token); const run = await ctx.db .query("flueRuns") .withIndex("by_runId", (q) => q.eq("runId", args.runId)) .unique(); return run === null ? null : toRunRow(run); }, }); export const lookupRun = query({ args: { ...tokenArgs, runId: v.string() }, handler: async (ctx, args) => { assertToken(args.token); const run = await ctx.db .query("flueRuns") .withIndex("by_runId", (q) => q.eq("runId", args.runId)) .unique(); return run === null ? null : { runId: run.runId, workflowName: run.workflowName }; }, }); export const listRuns = query({ args: { ...tokenArgs, status: v.optional(runStatus), workflowName: v.optional(v.string()), limit: v.number(), cursor: v.optional(v.object({ startedAt: v.string(), runId: v.string() })), }, handler: async (ctx, args) => { assertToken(args.token); const cursor: ListRunsCursor | undefined = args.cursor; const rows = (await ctx.db.query("flueRuns").collect()) .map(toRunPointerRow) .filter( (row) => (args.status === undefined || row.status === args.status) && (args.workflowName === undefined || row.workflowName === args.workflowName), ) .sort(compareRunPointerDesc) .filter((row) => { if (cursor === undefined) { return true; } if (row.startedAt !== cursor.startedAt) { return row.startedAt < cursor.startedAt; } return row.runId < cursor.runId; }); const page = rows.slice(0, args.limit + 1); return { rows: page.slice(0, args.limit), hasMore: page.length > args.limit, }; }, }); export const putAttachment = mutation({ args: { ...tokenArgs, streamPath: v.string(), conversationId: v.string(), attachment: attachmentRef, bytes: v.bytes(), }, handler: async (ctx, args) => { assertToken(args.token); const existing = await ctx.db .query("flueAttachments") .withIndex("by_streamPath_and_attachmentId", (q) => q.eq("streamPath", args.streamPath).eq("attachmentId", args.attachment.id), ) .unique(); if (existing !== null) { return sameAttachment(existing, args) ? "existing" : "conflict"; } await ctx.db.insert("flueAttachments", { streamPath: args.streamPath, conversationId: args.conversationId, attachmentId: args.attachment.id, mimeType: args.attachment.mimeType, size: args.attachment.size, digest: args.attachment.digest, filename: args.attachment.filename, bytes: args.bytes, createdAt: Date.now(), }); return "inserted"; }, }); export const getAttachment = query({ args: { ...tokenArgs, streamPath: v.string(), conversationId: v.string(), attachmentId: v.string(), }, handler: async (ctx, args) => { assertToken(args.token); const attachment = await ctx.db .query("flueAttachments") .withIndex("by_streamPath_and_conversationId_and_attachmentId", (q) => q .eq("streamPath", args.streamPath) .eq("conversationId", args.conversationId) .eq("attachmentId", args.attachmentId), ) .unique(); return attachment === null ? null : toAttachmentWire(attachment); }, }); export const deleteAttachmentsForInstance = mutation({ args: { ...tokenArgs, streamPath: v.string() }, handler: async (ctx, args) => { assertToken(args.token); const attachments = await ctx.db .query("flueAttachments") .withIndex("by_streamPath", (q) => q.eq("streamPath", args.streamPath)) .collect(); for (const attachment of attachments) { await ctx.db.delete(attachment._id); } }, });