/* eslint-disable unicorn/no-array-sort, no-await-in-loop, unicorn/no-await-expression-member, unicorn/filename-case, unicorn/prefer-at, unicorn/no-array-reduce, @typescript-eslint/no-non-null-assertion */ import { env } from "@code/env/convex"; import { v } from "convex/values"; import type { Doc } from "./_generated/dataModel"; import { mutation, query } from "./_generated/server"; import type { MutationCtx, QueryCtx } from "./_generated/server"; import { projectConversationRecords } from "./conversationProjections"; 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">; interface 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; } interface SettlementObligationRow { readonly submissionId: string; readonly sessionKey: string; readonly attemptId: string; readonly recordId: string; readonly recordJson: string; } interface AttemptMarkerRow { readonly submissionId: string; readonly attemptId: string; readonly createdAt: number; } interface ConversationStreamIdentity { readonly agentName: string; readonly instanceId: string; } interface ConversationProducerClaim { readonly producerId: string; readonly producerEpoch: number; readonly incarnation: string; readonly nextProducerSequence: number; readonly offset: string; } interface ConversationStreamRow { readonly identity: ConversationStreamIdentity; readonly incarnation: string; readonly nextOffset: number; readonly producerId: string | null; readonly producerEpoch: number; readonly nextProducerSequence: number; } interface ConversationBatchRow { readonly offset: number; readonly recordsJson: string; } interface EventEntryRow { readonly offset: number; readonly dataJson: string; } interface 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; } interface 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; } interface AttachmentRefWire { readonly id: string; readonly mimeType: string; readonly size: number; readonly digest: string; readonly filename?: string; } interface AttachmentWire { readonly attachment: AttachmentRefWire; readonly bytes: ArrayBuffer; } interface AdmitSubmissionResponse { readonly kind: "submission" | "retained_receipt" | "conflict"; readonly submission?: SubmissionRow; readonly receipt?: { readonly submissionId: string; readonly acceptedAt: number; }; } interface AppendResult { readonly offset: number; readonly appended: boolean; } interface ListRunsCursor { readonly startedAt: string; readonly runId: string; } interface 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({ digest: v.string(), filename: v.optional(v.string()), id: v.string(), mimeType: v.string(), size: v.number(), }); 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 => ({ abortRequestedAt: doc.abortRequestedAt ?? null, acceptedAt: doc.acceptedAt, attemptCount: doc.attemptCount, attemptId: doc.attemptId ?? null, canonicalReadyAt: doc.canonicalReadyAt ?? null, chunksJson: doc.chunksJson, error: doc.error ?? null, inputAppliedAt: doc.inputAppliedAt ?? null, inputJson: doc.inputJson, kind: doc.kind, leaseExpiresAt: doc.leaseExpiresAt, maxRetry: doc.maxRetry, ownerId: doc.ownerId ?? null, recoveryRequestedAt: doc.recoveryRequestedAt ?? null, sequence: doc.sequence, sessionKey: doc.sessionKey, startedAt: doc.startedAt ?? null, status: doc.status, submissionId: doc.submissionId, timeoutAt: doc.timeoutAt, traceCarrierJson: doc.traceCarrierJson ?? null, }); const toSettlementObligationRow = ( doc: SubmissionDoc ): SettlementObligationRow | null => { if ( doc.attemptId === undefined || doc.settlementRecordId === undefined || doc.settlementRecordJson === undefined ) { return null; } return { attemptId: doc.attemptId, recordId: doc.settlementRecordId, recordJson: doc.settlementRecordJson, sessionKey: doc.sessionKey, submissionId: doc.submissionId, }; }; const toAttemptMarkerRow = (doc: AttemptMarkerDoc): AttemptMarkerRow => ({ attemptId: doc.attemptId, createdAt: doc.createdAt, submissionId: doc.submissionId, }); const toConversationStreamRow = ( doc: ConversationStreamDoc ): ConversationStreamRow => ({ identity: safeJsonParse(doc.identityJson), incarnation: doc.incarnation, nextOffset: doc.nextOffset - 1, nextProducerSequence: doc.nextProducerSequence, producerEpoch: doc.producerEpoch, producerId: doc.producerId ?? null, }); const toConversationBatchRow = ( doc: ConversationBatchDoc ): ConversationBatchRow => ({ offset: doc.seq, recordsJson: doc.recordsJson, }); const toEventEntryRow = (doc: EventEntryDoc): EventEntryRow => ({ dataJson: doc.dataJson, offset: doc.seq, }); const toRunRow = (doc: RunDoc): RunRow => ({ durationMs: doc.durationMs ?? null, endedAt: doc.endedAt ?? null, errorJson: doc.errorJson ?? null, inputJson: doc.inputJson ?? null, isError: doc.isError ?? null, resultJson: doc.resultJson ?? null, runId: doc.runId, startedAt: doc.startedAt, status: doc.status, traceCarrierJson: doc.traceCarrierJson ?? null, workflowName: doc.workflowName, }); const toRunPointerRow = (doc: RunDoc): RunPointerRow => ({ durationMs: doc.durationMs ?? null, endedAt: doc.endedAt ?? null, isError: doc.isError ?? null, runId: doc.runId, startedAt: doc.startedAt, status: doc.status, workflowName: doc.workflowName, }); const toAttachmentWire = (doc: AttachmentDoc): AttachmentWire => ({ attachment: { digest: doc.digest, id: doc.attachmentId, mimeType: doc.mimeType, size: doc.size, ...(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 { // JSON parse failure means the stored value is not an array } 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 = ( ctx: ReadCtx, submissionId: string ): Promise => ctx.db .query("flueSubmissions") .withIndex("by_submissionId", (q) => q.eq("submissionId", submissionId)) .unique(); const getConversationStreamDoc = ( 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, attemptId: v.string(), leaseExpiresAt: v.optional(v.number()), nextAttemptId: v.string(), ownerId: v.optional(v.string()), 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 !== "running" || submission.attemptId !== args.attemptId ) { return null; } const now = Date.now(); await ctx.db.patch(submission._id, { attemptCount: submission.attemptCount + 1, attemptId: args.nextAttemptId, recoveryRequestedAt: undefined, startedAt: now, ...(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, clientRequestId: v.optional(v.string()), input: v.object({ acceptedAt: v.number(), chunksJson: v.optional(v.string()), inputJson: v.string(), kind: submissionKind, sessionKey: v.string(), submissionId: v.string(), traceCarrierJson: v.optional(v.string()), }), turnId: v.optional(v.id("conversationTurns")), }, handler: async (ctx, args): Promise => { assertToken(args.token); const correlatedTurn = args.clientRequestId === undefined || args.turnId === undefined ? null : await ctx.db.get(args.turnId); if ( correlatedTurn !== null && correlatedTurn.clientRequestId !== args.clientRequestId ) { throw new Error("[flue] Turn admission correlation does not match."); } 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: { acceptedAt: existing.acceptedAt, submissionId: existing.submissionId, }, }; } if ( correlatedTurn !== null && correlatedTurn.submissionId !== existing.submissionId ) { await ctx.db.patch(correlatedTurn._id, { error: undefined, leaseExpiresAt: undefined, leaseOwner: undefined, status: "running", submissionId: existing.submissionId, }); } 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", { acceptedAt: args.input.acceptedAt, attemptCount: 0, chunksJson: args.input.chunksJson ?? "[]", inputJson: args.input.inputJson, kind: args.input.kind, leaseExpiresAt: 0, maxRetry: DURABILITY_DEFAULT_MAX_ATTEMPTS, sequence: nextSequence, sessionKey: args.input.sessionKey, status: "queued", submissionId: args.input.submissionId, timeoutAt: 0, traceCarrierJson: args.input.traceCarrierJson, updatedAt: now, }); const created = await ctx.db.get(rowId); if (created === null) { throw new Error("[flue] Failed to create submission row."); } if (correlatedTurn !== null) { await ctx.db.patch(correlatedTurn._id, { error: undefined, leaseExpiresAt: undefined, leaseOwner: undefined, status: "running", submissionId: args.input.submissionId, }); } 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, attemptId: v.string(), leaseExpiresAt: v.number(), ownerId: v.string(), submissionId: v.string(), }, 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, { attemptCount: submission.attemptCount + 1, attemptId: args.attemptId, leaseExpiresAt: args.leaseExpiresAt, maxRetry: submission.maxRetry > 0 ? submission.maxRetry : DURABILITY_DEFAULT_MAX_ATTEMPTS, ownerId: args.ownerId, startedAt: now, status: "running", timeoutAt: submission.timeoutAt > 0 ? submission.timeoutAt : now + DURABILITY_DEFAULT_TIMEOUT_MS, updatedAt: now, }); const updated = await ctx.db.get(submission._id); return updated === null ? null : toSubmissionRow(updated); }, }); export const markSubmissionInputApplied = mutation({ args: { ...tokenArgs, attemptId: v.string(), maxRetry: v.optional(v.number()), submissionId: v.string(), 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, attemptId: v.string(), 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 !== "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, attemptId: v.string(), 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 !== "running" || submission.attemptId !== args.attemptId || submission.inputAppliedAt !== undefined ) { return false; } await ctx.db.patch(submission._id, { attemptId: undefined, leaseExpiresAt: 0, ownerId: undefined, recoveryRequestedAt: undefined, startedAt: undefined, status: "queued", updatedAt: Date.now(), }); return true; }, }); export const reserveSubmissionSettlement = mutation({ args: { ...tokenArgs, attemptId: v.string(), recordId: v.string(), recordJson: v.string(), 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) { 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, { settlementRecordId: args.recordId, settlementRecordJson: args.recordJson, status: "terminalizing", updatedAt: Date.now(), }); const updated = await ctx.db.get(submission._id); return updated === null ? null : toSettlementObligationRow(updated); }, }); export const finalizeSubmissionSettlement = mutation({ args: { ...tokenArgs, attemptId: v.string(), recordId: v.string(), 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 !== "terminalizing" || submission.attemptId !== args.attemptId || submission.settlementRecordId !== args.recordId ) { return false; } await ctx.db.patch(submission._id, { leaseExpiresAt: 0, ownerId: undefined, settledOutcome: parseSettledOutcome(submission.settlementRecordJson) ?? "completed", status: "settled", updatedAt: Date.now(), }); return true; }, }); export const settleSubmission = mutation({ args: { ...tokenArgs, attemptId: v.string(), errorJson: v.optional(v.string()), outcome: v.union(v.literal("completed"), v.literal("failed")), 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 !== "running" || submission.attemptId !== args.attemptId ) { return false; } await ctx.db.patch(submission._id, { error: args.errorJson, leaseExpiresAt: 0, ownerId: undefined, settledOutcome: args.outcome, status: "settled", updatedAt: Date.now(), }); return true; }, }); export const insertAttemptMarker = mutation({ args: { ...tokenArgs, attemptId: v.string(), submissionId: 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", { attemptId: args.attemptId, createdAt: Date.now(), submissionId: args.submissionId, }); 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, attemptId: v.string(), submissionId: 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, identity: conversationIdentity, path: v.string() }, 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", { closed: false, createdAt: Date.now(), identityJson, incarnation: crypto.randomUUID(), nextOffset: 0, nextProducerSequence: 0, path: args.path, producerEpoch: 0, }); }, }); 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, { nextProducerSequence: 0, producerEpoch, producerId: args.producerId, }); return { incarnation: stream.incarnation, nextProducerSequence: 0, offset: `${stream.nextOffset - 1}`, producerEpoch, producerId: args.producerId, }; }, }); export const appendConversationBatch = mutation({ args: { ...tokenArgs, incarnation: v.string(), path: v.string(), producerEpoch: v.number(), producerId: v.string(), producerSequence: v.number(), recordsJson: v.string(), submission: v.optional( v.object({ attemptId: v.string(), submissionId: 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 { appended: false, offset: existing.seq }; } 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", { appendedAt: Date.now(), attemptId: args.submission?.attemptId, path: args.path, producerEpoch: args.producerEpoch, producerId: args.producerId, producerSequence: args.producerSequence, recordsJson: args.recordsJson, seq, submissionId: args.submission?.submissionId, }); await ctx.db.patch(stream._id, { nextOffset: seq + 1, nextProducerSequence: stream.nextProducerSequence + 1, }); // Projection is a disposable product view. Canonical persistence must win // even if a future projector record shape is malformed. try { await projectConversationRecords( ctx, args.recordsJson, args.submission?.submissionId ); } catch { // The raw canonical batch above remains durable and replayable. } return { appended: true, offset: seq }; }, }); export const readConversationBatches = query({ args: { ...tokenArgs, afterOffset: v.number(), limit: v.number(), 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(); 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", { closed: false, createdAt: Date.now(), nextSeq: 0, path: args.path, }); }, }); export const appendEvent = mutation({ args: { ...tokenArgs, dataJson: v.string(), path: 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", { appendedAt: Date.now(), dataJson: args.dataJson, path: args.path, seq, }); await ctx.db.patch(stream._id, { nextSeq: seq + 1 }); return { appended: true, offset: seq }; }, }); export const appendEventOnce = mutation({ args: { ...tokenArgs, dataJson: v.string(), key: v.string(), path: 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 { appended: false, offset: existing.seq }; } const seq = stream.nextSeq; await ctx.db.insert("flueEventEntries", { appendedAt: Date.now(), dataJson: args.dataJson, onceKey: args.key, path: args.path, seq, }); await ctx.db.patch(stream._id, { nextSeq: seq + 1 }); return { appended: true, offset: seq }; }, }); export const readEvents = query({ args: { ...tokenArgs, afterOffset: v.number(), limit: v.number(), 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) { return { closed: false, events: [] as EventEntryRow[], nextOffset: -1, upToDate: true, }; } 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 { closed: stream.closed, events: page.map(toEventEntryRow), nextOffset: page.length > 0 ? page[page.length - 1]!.seq : args.afterOffset, upToDate: rows.length <= args.limit, }; }, }); 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 : { closed: stream.closed, nextOffset: stream.nextSeq - 1 }; }, }); export const createRun = mutation({ args: { ...tokenArgs, inputJson: v.string(), runId: v.string(), startedAt: v.string(), traceCarrierJson: v.optional(v.string()), workflowName: 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", { inputJson: args.inputJson, runId: args.runId, startedAt: args.startedAt, status: "active", traceCarrierJson: args.traceCarrierJson, workflowName: args.workflowName, }); }, }); export const endRun = mutation({ args: { ...tokenArgs, durationMs: v.number(), endedAt: v.string(), errorJson: v.optional(v.string()), isError: v.boolean(), resultJson: v.optional(v.string()), 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(); if (run === null) { return; } await ctx.db.patch(run._id, { durationMs: args.durationMs, endedAt: args.endedAt, errorJson: args.errorJson, isError: args.isError, resultJson: args.resultJson, status: args.isError ? "errored" : "completed", }); }, }); 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, cursor: v.optional(v.object({ runId: v.string(), startedAt: v.string() })), limit: v.number(), status: v.optional(runStatus), workflowName: v.optional(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 { hasMore: page.length > args.limit, rows: page.slice(0, args.limit), }; }, }); export const putAttachment = mutation({ args: { ...tokenArgs, attachment: attachmentRef, bytes: v.bytes(), conversationId: v.string(), streamPath: v.string(), }, 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", { attachmentId: args.attachment.id, bytes: args.bytes, conversationId: args.conversationId, createdAt: Date.now(), digest: args.attachment.digest, filename: args.attachment.filename, mimeType: args.attachment.mimeType, size: args.attachment.size, streamPath: args.streamPath, }); return "inserted"; }, }); export const getAttachment = query({ args: { ...tokenArgs, attachmentId: v.string(), conversationId: v.string(), streamPath: 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); } }, });