Files
zopu-code/packages/backend/convex/schema.ts
-Puter 005b26fa32 Slices 2-4: Work becomes executable
Domain primitives: versioned WorkDefinition/DesignPacket with risk, questions,
and approval gates; lifecycle transition table with invalidation rules;
Resolver (Run/Attempt/outcomes/CodingKitV0); provider-neutral HarnessRuntime
with normalized events and a deterministic FakeHarnessLive that always reaches
a terminal classification and never claims implementation.

Convex durable model: workDefinitions, workQuestions, workApprovals,
designPackets, workSlices, workRuns, workAttempts, workAttemptEvents tables;
broadened works.status union; generalized workEvents with optional signalId,
referenceId, and payloadJson. Authenticated commands for definition request/
save/revise/approve, question answer/withdraw, design save/revise/approve, and
simulated execution start/cancel/retry with leased attempts, checkpointed
events, and expired-lease reconciliation.

Private work-planner FLUE agent with proposal-only tools (definition, design,
question). Convex validates and stores every proposal; FLUE never approves or
advances Work.

Expanded Work card with Outcome, Design, and Build sections wired to the new
mutations. Slice 1 conversation and provenance preserved.
2026-07-27 22:52:14 +05:30

500 lines
16 KiB
TypeScript

import { defineSchema, defineTable } from "convex/server";
import { v } from "convex/values";
export default defineSchema({
organizations: defineTable({
name: v.string(),
kind: v.union(v.literal("personal"), v.literal("team")),
createdBy: v.string(),
createdAt: v.number(),
}).index("by_createdBy_and_kind", ["createdBy", "kind"]),
organizationMembers: defineTable({
organizationId: v.id("organizations"),
userId: v.string(),
role: v.union(v.literal("owner"), v.literal("member")),
createdAt: v.number(),
})
.index("by_userId", ["userId"])
.index("by_organizationId", ["organizationId"])
.index("by_organizationId_and_userId", ["organizationId", "userId"]),
projects: defineTable({
organizationId: v.id("organizations"),
name: v.string(),
description: v.optional(v.string()),
sourceUrl: v.string(),
normalizedSourceUrl: v.string(),
sourceHost: v.string(),
repositoryPath: v.string(),
defaultBranch: v.optional(v.string()),
createdAt: v.number(),
updatedAt: v.number(),
})
.index("by_organizationId_and_createdAt", ["organizationId", "createdAt"])
.index("by_organizationId_and_normalizedSourceUrl", [
"organizationId",
"normalizedSourceUrl",
]),
projectContextDocuments: defineTable({
projectId: v.id("projects"),
kind: v.union(
v.literal("readme"),
v.literal("agents"),
v.literal("product"),
v.literal("business"),
v.literal("design"),
v.literal("tech")
),
path: v.string(),
content: v.string(),
origin: v.literal("repository"),
sourceUrl: v.string(),
createdAt: v.number(),
updatedAt: v.number(),
}).index("by_projectId_and_path", ["projectId", "path"]),
conversations: defineTable({
organizationId: v.id("organizations"),
createdAt: v.number(),
}).index("by_organizationId", ["organizationId"]),
conversationTurns: defineTable({
conversationId: v.id("conversations"),
clientRequestId: v.string(),
status: v.union(
v.literal("queued"),
v.literal("processing"),
v.literal("completed"),
v.literal("failed")
),
submissionId: v.optional(v.string()),
error: v.optional(v.string()),
createdAt: v.number(),
completedAt: v.optional(v.number()),
}).index("by_conversationId_and_clientRequestId", [
"conversationId",
"clientRequestId",
]),
conversationMessages: defineTable({
conversationId: v.id("conversations"),
turnId: v.id("conversationTurns"),
role: v.union(v.literal("user"), v.literal("assistant")),
content: v.string(),
ordinal: v.number(),
createdAt: v.number(),
})
.index("by_conversationId_and_ordinal", ["conversationId", "ordinal"])
.index("by_turnId_and_role", ["turnId", "role"]),
conversationAttachments: defineTable({
messageId: v.id("conversationMessages"),
storageId: v.id("_storage"),
filename: v.optional(v.string()),
mimeType: v.string(),
createdAt: v.number(),
}).index("by_messageId", ["messageId"]),
signals: defineTable({
organizationId: v.id("organizations"),
projectId: v.id("projects"),
conversationId: v.id("conversations"),
sourceKey: v.string(),
title: v.string(),
summary: v.string(),
desiredOutcome: v.string(),
processedByAgentName: v.string(),
processedByAgentInstanceId: v.string(),
createdAt: v.number(),
})
.index("by_organization_and_createdAt", ["organizationId", "createdAt"])
.index("by_organization_and_sourceKey", ["organizationId", "sourceKey"])
.index("by_project_and_createdAt", ["projectId", "createdAt"]),
signalConstraints: defineTable({
signalId: v.id("signals"),
ordinal: v.number(),
value: v.string(),
}).index("by_signalId_and_ordinal", ["signalId", "ordinal"]),
signalSources: defineTable({
signalId: v.id("signals"),
messageId: v.id("conversationMessages"),
ordinal: v.number(),
rawTextSnapshot: v.string(),
sourceCreatedAt: v.number(),
})
.index("by_signalId_and_ordinal", ["signalId", "ordinal"])
.index("by_messageId", ["messageId"]),
works: defineTable({
organizationId: v.id("organizations"),
projectId: v.id("projects"),
title: v.string(),
objective: v.string(),
status: v.union(
v.literal("proposed"),
v.literal("defining"),
v.literal("awaiting-definition-approval"),
v.literal("designing"),
v.literal("awaiting-design-approval"),
v.literal("ready"),
v.literal("executing"),
v.literal("needs-input"),
v.literal("blocked"),
v.literal("completed"),
v.literal("failed"),
v.literal("cancelled")
),
definitionVersion: v.optional(v.number()),
definitionApprovalVersion: v.optional(v.number()),
designVersion: v.optional(v.number()),
designApprovalVersion: v.optional(v.number()),
createdAt: v.number(),
updatedAt: v.number(),
})
.index("by_project_and_createdAt", ["projectId", "createdAt"])
.index("by_organization_and_createdAt", ["organizationId", "createdAt"]),
signalWorkAttachments: defineTable({
signalId: v.id("signals"),
workId: v.id("works"),
createdAt: v.number(),
})
.index("by_signal", ["signalId"])
.index("by_work", ["workId"])
.index("by_signal_and_work", ["signalId", "workId"]),
workEvents: defineTable({
workId: v.id("works"),
signalId: v.optional(v.id("signals")),
kind: v.union(
v.literal("work.proposed"),
v.literal("signal.attached"),
v.literal("definition.requested"),
v.literal("definition.saved"),
v.literal("definition.revised"),
v.literal("definition.approved"),
v.literal("definition.invalidated"),
v.literal("question.answered"),
v.literal("question.withdrawn"),
v.literal("design.requested"),
v.literal("design.saved"),
v.literal("design.revised"),
v.literal("design.approved"),
v.literal("design.invalidated"),
v.literal("run.started"),
v.literal("run.cancelled"),
v.literal("attempt.claimed"),
v.literal("attempt.event"),
v.literal("attempt.completed"),
v.literal("attempt.reconciled")
),
referenceId: v.optional(v.string()),
payloadJson: v.optional(v.string()),
idempotencyKey: v.string(),
createdAt: v.number(),
})
.index("by_work_and_createdAt", ["workId", "createdAt"])
.index("by_work_and_idempotencyKey", ["workId", "idempotencyKey"]),
workDefinitions: defineTable({
workId: v.id("works"),
version: v.number(),
payloadJson: v.string(),
risk: v.union(v.literal("low"), v.literal("medium"), v.literal("high")),
status: v.union(
v.literal("proposed"),
v.literal("current"),
v.literal("superseded")
),
createdBy: v.string(),
createdAt: v.number(),
})
.index("by_work_and_version", ["workId", "version"])
.index("by_work_and_status", ["workId", "status"]),
workQuestions: defineTable({
workId: v.id("works"),
definitionVersion: v.number(),
questionId: v.string(),
prompt: v.string(),
impact: v.union(v.literal("low"), v.literal("medium"), v.literal("high")),
recommendation: v.optional(v.string()),
alternativesJson: v.string(),
status: v.union(
v.literal("open"),
v.literal("answered"),
v.literal("withdrawn")
),
answer: v.optional(v.string()),
createdAt: v.number(),
}).index("by_work_and_definitionVersion", ["workId", "definitionVersion"]),
workApprovals: defineTable({
workId: v.id("works"),
kind: v.union(v.literal("definition"), v.literal("design")),
definitionVersion: v.number(),
designVersion: v.optional(v.number()),
approvedBy: v.string(),
approvedAt: v.number(),
status: v.union(v.literal("active"), v.literal("invalidated")),
}).index("by_work_and_kind", ["workId", "kind"]),
designPackets: defineTable({
workId: v.id("works"),
version: v.number(),
definitionVersion: v.number(),
payloadJson: v.string(),
status: v.union(
v.literal("proposed"),
v.literal("current"),
v.literal("superseded")
),
createdBy: v.string(),
createdAt: v.number(),
}).index("by_work_and_version", ["workId", "version"]),
workSlices: defineTable({
workId: v.id("works"),
designVersion: v.number(),
sliceId: v.string(),
ordinal: v.number(),
title: v.string(),
objective: v.string(),
observableBehavior: v.string(),
payloadJson: v.string(),
status: v.union(
v.literal("planned"),
v.literal("ready"),
v.literal("running"),
v.literal("completed"),
v.literal("blocked")
),
}).index("by_work_and_designVersion", ["workId", "designVersion"]),
workRuns: defineTable({
workId: v.id("works"),
sliceId: v.optional(v.string()),
status: v.union(
v.literal("ready"),
v.literal("running"),
v.literal("terminal"),
v.literal("cancelled")
),
scenario: v.union(
v.literal("success"),
v.literal("transient-failure-then-success"),
v.literal("needs-input"),
v.literal("permanent-failure"),
v.literal("cancelled")
),
kitId: v.string(),
kitVersion: v.string(),
createdAt: v.number(),
startedAt: v.optional(v.number()),
endedAt: v.optional(v.number()),
terminalClassification: v.optional(v.string()),
terminalSummary: v.optional(v.string()),
}).index("by_work_and_createdAt", ["workId", "createdAt"]),
workAttempts: defineTable({
runId: v.id("workRuns"),
workId: v.id("works"),
number: v.number(),
status: v.union(
v.literal("queued"),
v.literal("claimed"),
v.literal("running"),
v.literal("terminal")
),
leaseOwner: v.optional(v.string()),
leaseExpiresAt: v.optional(v.number()),
startedAt: v.optional(v.number()),
endedAt: v.optional(v.number()),
classification: v.optional(v.string()),
summary: v.optional(v.string()),
}).index("by_run_and_number", ["runId", "number"]),
workAttemptEvents: defineTable({
attemptId: v.id("workAttempts"),
sequence: v.number(),
kind: v.string(),
message: v.string(),
metadataJson: v.string(),
occurredAt: v.number(),
}).index("by_attempt_and_sequence", ["attemptId", "sequence"]),
// -----------------------------------------------------------------
// Flue persistence stores (schema/format version 4).
// -----------------------------------------------------------------
// One-row key/value metadata table, keyed with `schema_version` for the
// persisted-store versioning obligation.
flueMeta: defineTable({
key: v.string(),
value: v.string(),
}).index("by_key", ["key"]),
// Durable agent-submission lifecycle: one row per submission. Settlement
// obligations live on the same row so reserve/finalize are atomic with the
// status transition they gate.
flueSubmissions: defineTable({
submissionId: v.string(),
sessionKey: v.string(),
sequence: v.number(),
kind: v.union(v.literal("dispatch"), v.literal("direct")),
// Client-supplied canonical input/chunks encoded by Flue helpers; exact
// string compared on idempotent replay.
inputJson: v.string(),
chunksJson: v.string(),
traceCarrierJson: v.optional(v.string()),
// Admit-time canonical record (Flue ConversationCreatedRecord) for the
// default session of the default harness.
recordJson: v.optional(v.string()),
status: v.union(
v.literal("queued"),
v.literal("running"),
v.literal("terminalizing"),
v.literal("settled")
),
acceptedAt: v.number(),
canonicalReadyAt: v.optional(v.number()),
attemptId: v.optional(v.string()),
inputAppliedAt: v.optional(v.number()),
recoveryRequestedAt: v.optional(v.number()),
abortRequestedAt: v.optional(v.number()),
startedAt: v.optional(v.number()),
error: v.optional(v.string()),
attemptCount: v.number(),
maxRetry: v.number(),
timeoutAt: v.number(),
ownerId: v.optional(v.string()),
leaseExpiresAt: v.number(),
// Settlement obligation (reserved/finalized on this row).
settlementRecordId: v.optional(v.string()),
settlementRecordJson: v.optional(v.string()),
settledOutcome: v.optional(
v.union(v.literal("completed"), v.literal("failed"), v.literal("aborted"))
),
updatedAt: v.number(),
})
.index("by_submissionId", ["submissionId"])
.index("by_sessionKey_and_sequence", ["sessionKey", "sequence"])
.index("by_status_and_sequence", ["status", "sequence"])
.index("by_status_and_leaseExpiresAt", ["status", "leaseExpiresAt"])
.index("by_status_and_settlementRecordId", [
"status",
"settlementRecordId",
]),
// Durable evidence that a submission attempt started and has not yet
// settled. Append-once by (submissionId, attemptId).
flueAttemptMarkers: defineTable({
submissionId: v.string(),
attemptId: v.string(),
createdAt: v.number(),
})
.index("by_submissionId_and_attemptId", ["submissionId", "attemptId"])
.index("by_submissionId", ["submissionId"]),
// Conversation-stream metadata: one row per stream path.
flueConversationStreams: defineTable({
path: v.string(),
identityJson: v.string(),
incarnation: v.string(),
producerId: v.optional(v.string()),
producerEpoch: v.number(),
nextProducerSequence: v.number(),
nextOffset: v.number(),
closed: v.boolean(),
createdAt: v.number(),
}).index("by_path", ["path"]),
// Conversation-stream batches: one row per appended batch. Offset is
// 0-based; `seq` is the row's position in the stream.
flueConversationBatches: defineTable({
path: v.string(),
seq: v.number(),
producerId: v.string(),
producerEpoch: v.number(),
producerSequence: v.number(),
recordsJson: v.string(),
submissionId: v.optional(v.string()),
attemptId: v.optional(v.string()),
appendedAt: v.number(),
})
.index("by_path_and_seq", ["path", "seq"])
.index("by_path_producer_epoch_producerSequence", [
"path",
"producerId",
"producerEpoch",
"producerSequence",
]),
// Event-stream metadata: one row per stream path.
flueEventStreams: defineTable({
path: v.string(),
nextSeq: v.number(),
closed: v.boolean(),
createdAt: v.number(),
}).index("by_path", ["path"]),
// Event-stream entries: one row per appended event. `onceKey` carries the
// idempotency key for `appendEventOnce` and is null for plain appends.
flueEventEntries: defineTable({
path: v.string(),
seq: v.number(),
onceKey: v.optional(v.string()),
dataJson: v.string(),
appendedAt: v.number(),
})
.index("by_path_and_seq", ["path", "seq"])
.index("by_path_and_onceKey", ["path", "onceKey"]),
// Workflow run records.
flueRuns: defineTable({
runId: v.string(),
workflowName: v.string(),
status: v.union(
v.literal("active"),
v.literal("completed"),
v.literal("errored")
),
startedAt: v.string(),
inputJson: v.optional(v.string()),
traceCarrierJson: v.optional(v.string()),
endedAt: v.optional(v.string()),
isError: v.optional(v.boolean()),
durationMs: v.optional(v.number()),
resultJson: v.optional(v.string()),
errorJson: v.optional(v.string()),
})
.index("by_runId", ["runId"])
.index("by_startedAt_and_runId", ["startedAt", "runId"])
.index("by_workflowName_and_startedAt_and_runId", [
"workflowName",
"startedAt",
"runId",
])
.index("by_status_startedAt_runId", ["status", "startedAt", "runId"])
.index("by_status_workflowName_startedAt_runId", [
"status",
"workflowName",
"startedAt",
"runId",
]),
// Immutable attachment bytes. Identity is (streamPath, attachmentId); reads
// are additionally scoped by conversationId.
flueAttachments: defineTable({
streamPath: v.string(),
conversationId: v.string(),
attachmentId: v.string(),
mimeType: v.string(),
size: v.number(),
digest: v.string(),
filename: v.optional(v.string()),
bytes: v.bytes(),
createdAt: v.number(),
})
.index("by_streamPath", ["streamPath"])
.index("by_streamPath_and_attachmentId", ["streamPath", "attachmentId"])
.index("by_streamPath_and_conversationId_and_attachmentId", [
"streamPath",
"conversationId",
"attachmentId",
]),
});