3 Commits

Author SHA1 Message Date
-Puter
a5414d83bd Convex-only Slice 1: clients talk to Convex, Flue is a private worker
Normalize the topology so web/desktop/mobile clients communicate only with
Convex. Convex admits durable commands, dispatches private Flue turns through
FLUE_URL, and stores the product-facing result before clients observe it.

- Normalized relational schema: organizations, projects, conversations/turns/
  messages/attachments, signals/sources/constraints, works/events/attachments.
- Conversation turns queued in Convex; the agent runAction dispatches to Flue
  and persists assistant response, signals, and proposed work back into Convex.
- Browser chat moved off the Flue transport: use-chat-agent now reads reactive
  Convex rows and sends through conversationMessages.send.
- Fixed signed-in blank screen: gate all product queries on useConvexAuth
  (isAuthenticated && !isRefreshing) plus personal-org bootstrap.
- Pinned react/react-dom to 19.2.8 via catalog to fix SSR dual-React crash.
- Removed ~16.6k net lines: daemon, AgentOS/Orb, project issues/runs/artifacts,
  todos, browser Flue transport, mobile execution views, smoke scripts.

Verified end-to-end on cheaptricks: connect repo, send actionable message,
Flue turn completes, Signal + proposed Work created with exact provenance,
persists across refresh.
2026-07-27 21:32:17 +05:30
-Puter
4c741fbe06 Fix type errors, lint, and complexity from VPS architecture cleanup
- Fix brand-to-Convex-Id casts in use-project-workspace (as unknown as)
- Extract ProjectsLoading and ConnectProject components to drop SliceOnePage
  complexity below the lint threshold of 20
- Fix formatting in project-requests.ts
- Archive orphaned standalone chat files (route removed, referenced dead env)
- All 8 packages typecheck clean; 31 web tests pass; root check clean
2026-07-27 18:57:16 +05:30
Miniputer
1855735245 Auth flow refactor, project requests, and routing updates
- Add server-side auth token loader (auth.server.ts) with SSR cookie forwarding
- Add project-requests client for explicit/signal/existing issue dispatch
- Refactor auth layout and app layout for SSR auth gating
- Update auth-client and auth-provider for convex token flow
- Extend project-issue primitives with request result schema
- Wire projectIssues backend mutation for request handling
- Update slice-one page and project-workspace hook
- Adjust routes (remove unused), Caddy config, and env templates
2026-07-27 18:47:43 +05:30
50 changed files with 1051 additions and 5830 deletions

View File

@@ -26,14 +26,14 @@
"react-dom": "catalog:",
"react-router": "^8.1.0",
"sonner": "catalog:",
"streamdown": "catalog:"
"streamdown": "2.5.0"
},
"devDependencies": {
"@code/config": "workspace:*",
"@react-router/dev": "^8.1.0",
"@tailwindcss/vite": "catalog:",
"@types/node": "catalog:",
"@types/react": "catalog:",
"@tailwindcss/vite": "^4.3.2",
"@types/node": "^22.13.14",
"@types/react": "^19.2.17",
"@types/react-dom": "catalog:",
"react-router-devtools": "^6.2.1",
"tailwindcss": "catalog:",

View File

@@ -6,15 +6,11 @@ import {
import { Button } from "@code/ui/components/button";
import {
ChevronRight,
Check,
FolderGit2,
Hammer,
LoaderCircle,
Menu,
MessageSquareText,
ImagePlus,
Play,
RotateCcw,
Send,
Sparkles,
X,
@@ -38,67 +34,11 @@ const EMPTY_WORKS: readonly SliceWork[] = [];
interface WorkCardProps {
readonly onSourceSelect: (rawText: string) => void;
readonly work: SliceWork;
readonly slice: SliceOneState;
}
const starterDefinition = (work: SliceWork) => ({
acceptanceCriteria: ["The requested outcome is observable and documented"],
affectedUsers: ["Project users"],
assumptions: [],
constraints: [],
desiredOutcome: work.objective,
inScope: [work.objective],
outOfScope: ["Unrelated product changes"],
problem: work.objective,
questions: [],
requiredArtifacts: ["Simulation activity and terminal outcome"],
risk: "medium",
});
const starterDesign = (work: SliceWork) => ({
architectureSummary:
"Validate the approved Definition, then exercise one deterministic fake slice.",
callFlowDelta: [
"Work -> Run -> Attempt -> normalized events -> terminal outcome",
],
concerns: [],
evidenceRequirements: ["Terminal Run classification"],
fileTreeDelta: [],
impactMap: {
files: [],
modules: [],
risks: [],
summary: "Compact vertical-slice simulation",
},
invariants: [
"Simulation never claims implementation",
"Every Attempt reaches a terminal classification",
],
keyInterfaces: ["HarnessRuntime", "AttemptOutcome"],
slices: [
{
codeBoundaries: ["workExecution"],
dependsOn: [],
evidenceRequirements: ["Normalized activity events"],
id: "slice-1",
objective: work.objective,
observableBehavior: "A terminal fake Run is visible",
reviewRequired: false,
title: "Deterministic simulation",
verification: ["Run completes with a terminal classification"],
},
],
tradeoffs: ["Fake runtime proves contract before sandbox integration"],
});
// oxlint-disable-next-line complexity -- the expanded card intentionally keeps the three review sections together.
const WorkCard = ({ onSourceSelect, work, slice }: WorkCardProps) => {
const WorkCard = ({ onSourceSelect, work }: WorkCardProps) => {
const [sourcesOpen, setSourcesOpen] = useState(false);
const [expanded, setExpanded] = useState(false);
const sources = work.signals.flatMap((signal) => signal.sources);
const { definition } = work;
const { design } = work;
const [latestRun] = work.runs;
return (
<article className="border border-[#d7d3c7] bg-[#fffefa] p-4 text-[#20201d] shadow-[0_10px_30px_rgba(30,30,20,0.06)]">
<div className="flex items-start gap-3">
@@ -130,16 +70,6 @@ const WorkCard = ({ onSourceSelect, work, slice }: WorkCardProps) => {
className={`size-4 transition-transform ${sourcesOpen ? "rotate-90" : ""}`}
/>
</button>
<button
className="mt-2 flex w-full items-center justify-between border-t border-[#e7e3d9] pt-3 text-left text-xs font-medium text-[#20201d]"
onClick={() => setExpanded((open) => !open)}
type="button"
>
<span>{expanded ? "Hide Work details" : "Open Work details"}</span>
<ChevronRight
className={`size-4 transition-transform ${expanded ? "rotate-90" : ""}`}
/>
</button>
{sourcesOpen ? (
<div className="mt-3 space-y-2">
{sources.map((source) => (
@@ -155,156 +85,6 @@ const WorkCard = ({ onSourceSelect, work, slice }: WorkCardProps) => {
))}
</div>
) : null}
{expanded ? (
<div className="mt-4 space-y-4 border-t border-[#e7e3d9] pt-4 text-xs">
<section>
<p className="font-semibold uppercase tracking-wide text-[#65713a]">
Outcome
</p>
<p className="mt-1 leading-5 text-[#626057]">{work.objective}</p>
<p className="mt-2 text-[#747168]">
Risk: {definition?.risk ?? "not defined"}
</p>
{definition?.questions?.length ? (
<p className="mt-1 text-amber-800">
{
definition.questions.filter(
(question) => question.status === "open"
).length
}{" "}
open question(s)
</p>
) : null}
<div className="mt-2 flex flex-wrap gap-2">
{work.status === "proposed" ? (
<Button
size="sm"
onClick={() => void slice.requestDefinition(work._id)}
>
<Sparkles className="size-3.5" /> Define
</Button>
) : null}
{work.status === "defining" ? (
<Button
size="sm"
variant="outline"
onClick={() =>
void slice.saveDefinition(work._id, starterDefinition(work))
}
>
<Hammer className="size-3.5" /> Save definition
</Button>
) : null}
{work.status === "awaiting-definition-approval" &&
work.definitionVersion ? (
<Button
size="sm"
onClick={() =>
void slice.approveDefinition(
work._id,
work.definitionVersion as number
)
}
>
<Check className="size-3.5" /> Approve
</Button>
) : null}
</div>
</section>
<section>
<p className="font-semibold uppercase tracking-wide text-[#65713a]">
Design
</p>
<p className="mt-1 leading-5 text-[#626057]">
{design?.architectureSummary ?? "No Design Packet yet."}
</p>
{design?.slices?.map((item) => (
<p className="mt-1 text-[#747168]" key={item.id}>
{item.title}: {item.observableBehavior}
</p>
))}
<div className="mt-2 flex flex-wrap gap-2">
{work.status === "designing" ? (
<Button
size="sm"
variant="outline"
onClick={() =>
void slice.saveDesign(work._id, starterDesign(work))
}
>
<Hammer className="size-3.5" /> Save design
</Button>
) : null}
{work.status === "awaiting-design-approval" &&
work.definitionApprovalVersion &&
work.designVersion ? (
<Button
size="sm"
onClick={() =>
void slice.approveDesign(
work._id,
work.definitionApprovalVersion as number,
work.designVersion as number
)
}
>
<Check className="size-3.5" /> Approve design
</Button>
) : null}
</div>
</section>
<section>
<p className="font-semibold uppercase tracking-wide text-[#65713a]">
Build
</p>
{latestRun ? (
<p className="mt-1 leading-5 text-[#626057]">
Run {latestRun.status}:{" "}
{latestRun.terminalSummary ??
latestRun.terminalClassification ??
"activity is still arriving"}
</p>
) : (
<p className="mt-1 text-[#747168]">No simulation Run yet.</p>
)}
<div className="mt-2 flex flex-wrap gap-2">
{work.status === "ready" ? (
<Button
size="sm"
onClick={() =>
void slice.startSimulation(
work._id,
"success",
design?.slices?.[0]?.id
)
}
>
<Play className="size-3.5" /> Simulate
</Button>
) : null}
{latestRun?.status === "running" ? (
<Button
size="sm"
variant="outline"
onClick={() => void slice.cancelSimulation(latestRun._id)}
>
<X className="size-3.5" /> Cancel
</Button>
) : null}
{latestRun?.status === "terminal" &&
latestRun.terminalClassification === "RetryableFailure" ? (
<Button
size="sm"
variant="outline"
onClick={() => void slice.retrySimulation(latestRun._id)}
>
<RotateCcw className="size-3.5" /> Retry
</Button>
) : null}
</div>
</section>
</div>
) : null}
</article>
);
};
@@ -498,7 +278,6 @@ export const SliceOnePage = () => {
</p>
<WorkCard
onSourceSelect={revealSourceMessage}
slice={slice}
work={work}
/>
</div>
@@ -608,7 +387,6 @@ export const SliceOnePage = () => {
<WorkCard
key={work._id}
onSourceSelect={revealSourceMessage}
slice={slice}
work={work}
/>
))}
@@ -644,7 +422,6 @@ export const SliceOnePage = () => {
<WorkCard
key={work._id}
onSourceSelect={revealSourceMessage}
slice={slice}
work={work}
/>
))}

View File

@@ -1,7 +1,6 @@
import { api } from "@code/backend/convex/_generated/api";
import type { Id } from "@code/backend/convex/_generated/dataModel";
import { useAction, useMutation, useQuery } from "convex/react";
import { makeFunctionReference } from "convex/server";
import { useAction, useQuery } from "convex/react";
import { useMemo, useState } from "react";
import { useOrganizationChatAgent } from "@/hooks/chat/use-chat-agent";
@@ -10,93 +9,6 @@ import { usePersonalOrganization } from "@/hooks/use-personal-organization";
const toError = (error: unknown) =>
error instanceof Error ? error : new Error(String(error));
interface WorkRecord {
readonly _id: Id<"works">;
readonly title: string;
readonly objective: string;
readonly status: string;
readonly definitionVersion?: number;
readonly definitionApprovalVersion?: number;
readonly designVersion?: number;
readonly designApprovalVersion?: number;
readonly signals: readonly {
sources: readonly { messageId: string; rawText: string }[];
}[];
readonly events: readonly { _id: string; createdAt: number; kind: string }[];
readonly definitions: readonly unknown[];
readonly designs: readonly unknown[];
readonly slices: readonly unknown[];
readonly runs: readonly {
_id: Id<"workRuns">;
status: string;
terminalClassification?: string;
terminalSummary?: string;
}[];
readonly definition: {
risk?: string;
questions?: { status: string }[];
} | null;
readonly design: {
architectureSummary?: string;
slices?: { id: string; title: string; observableBehavior: string }[];
} | null;
}
const workListRef = makeFunctionReference<
"query",
{ projectId: Id<"projects"> },
WorkRecord[]
>("workPlanning:listForProject");
const requestDefinitionRef = makeFunctionReference<
"mutation",
{ workId: Id<"works"> },
unknown
>("workPlanning:requestDefinition");
const saveDefinitionRef = makeFunctionReference<
"mutation",
{ workId: Id<"works">; payloadJson: string },
unknown
>("workPlanning:saveDefinitionProposal");
const approveDefinitionRef = makeFunctionReference<
"mutation",
{ workId: Id<"works">; version: number },
unknown
>("workPlanning:approveDefinition");
const saveDesignRef = makeFunctionReference<
"mutation",
{ workId: Id<"works">; payloadJson: string },
unknown
>("workPlanning:saveDesignProposal");
const approveDesignRef = makeFunctionReference<
"mutation",
{ workId: Id<"works">; definitionVersion: number; designVersion: number },
unknown
>("workPlanning:approveDesign");
const startSimulationRef = makeFunctionReference<
"mutation",
{
workId: Id<"works">;
scenario:
| "success"
| "transient-failure-then-success"
| "needs-input"
| "permanent-failure"
| "cancelled";
sliceId?: string;
},
unknown
>("workExecution:startSimulatedExecution");
const cancelSimulationRef = makeFunctionReference<
"mutation",
{ runId: Id<"workRuns"> },
unknown
>("workExecution:cancelSimulatedExecution");
const retrySimulationRef = makeFunctionReference<
"mutation",
{ runId: Id<"workRuns"> },
unknown
>("workExecution:retrySimulatedExecution");
export const useSliceOne = () => {
const organization = usePersonalOrganization();
const projects = useQuery(
@@ -117,17 +29,9 @@ export const useSliceOne = () => {
? selectedProjectId
: ((projects?.[0]?.id as unknown as Id<"projects"> | undefined) ?? null);
const works = useQuery(
workListRef,
api.works.listForProject,
activeProjectId ? { projectId: activeProjectId } : "skip"
);
const requestDefinitionMutation = useMutation(requestDefinitionRef);
const saveDefinitionMutation = useMutation(saveDefinitionRef);
const approveDefinitionMutation = useMutation(approveDefinitionRef);
const saveDesignMutation = useMutation(saveDesignRef);
const approveDesignMutation = useMutation(approveDesignRef);
const startSimulationMutation = useMutation(startSimulationRef);
const cancelSimulationMutation = useMutation(cancelSimulationRef);
const retrySimulationMutation = useMutation(retrySimulationRef);
const selectedProject = useMemo(
() =>
projects?.find(
@@ -158,52 +62,16 @@ export const useSliceOne = () => {
setSelectedProjectId(projectId as unknown as Id<"projects">);
};
const requestDefinition = (workId: Id<"works">) =>
requestDefinitionMutation({ workId });
const saveDefinition = (workId: Id<"works">, payload: unknown) =>
saveDefinitionMutation({ payloadJson: JSON.stringify(payload), workId });
const approveDefinition = (workId: Id<"works">, version: number) =>
approveDefinitionMutation({ version, workId });
const saveDesign = (workId: Id<"works">, payload: unknown) =>
saveDesignMutation({ payloadJson: JSON.stringify(payload), workId });
const approveDesign = (
workId: Id<"works">,
definitionVersion: number,
designVersion: number
) => approveDesignMutation({ definitionVersion, designVersion, workId });
const startSimulation = (
workId: Id<"works">,
scenario:
| "success"
| "transient-failure-then-success"
| "needs-input"
| "permanent-failure"
| "cancelled",
sliceId?: string
) => startSimulationMutation({ scenario, sliceId, workId });
const cancelSimulation = (runId: Id<"workRuns">) =>
cancelSimulationMutation({ runId });
const retrySimulation = (runId: Id<"workRuns">) =>
retrySimulationMutation({ runId });
return {
agent,
approveDefinition,
approveDesign,
cancelSimulation,
connectRepository,
error,
pending,
projects,
repository,
requestDefinition,
retrySimulation,
saveDefinition,
saveDesign,
selectProject,
selectedProject,
setRepository,
startSimulation,
works,
} as const;
};

680
bun.lock

File diff suppressed because it is too large Load Diff

View File

@@ -5,28 +5,5 @@ import remix from "ultracite/oxlint/remix";
export default defineConfig({
extends: [core, react, remix],
ignorePatterns: [...core.ignorePatterns, "repos/**", "scripts/**"],
overrides: [
{
files: ["**/convex/**/*.ts", "**/convex/**/*.test.ts"],
rules: {
"unicorn/filename-case": "off",
// Convex targets ES2022: no toSorted/at; sequential awaits ensure atomicity
"unicorn/no-array-sort": "off",
"no-await-in-loop": "off",
// Convex query builders use `any` in index callbacks
"@typescript-eslint/no-explicit-any": "off",
"@typescript-eslint/no-non-null-assertion": "off",
"eslint/complexity": "off",
},
},
{
files: ["**/convex/**/*.test.ts"],
rules: {
"unicorn/no-await-expression-member": "off",
"eslint/require-await": "off",
"sort-keys": "off",
},
},
],
ignorePatterns: core.ignorePatterns,
});

View File

@@ -37,15 +37,7 @@
"heroui-native": "^1.0.5",
"vite": "^7.3.6",
"vitest": "^4.1.10",
"convex-test": "^0.0.54",
"react-native": "0.86.0",
"@types/react": "^19.2.17",
"@types/node": "^22.13.14",
"hono": "^4.8.3",
"valibot": "^1.4.2",
"streamdown": "2.5.0",
"@tailwindcss/postcss": "^4.3.2",
"@tailwindcss/vite": "^4.3.2"
"convex-test": "^0.0.54"
}
},
"type": "module",
@@ -53,8 +45,8 @@
"dev": "vp run -r dev",
"build": "vp run -r build",
"check-types": "vp run -r check-types",
"check": "ultracite check package.json vite.config.ts apps/web/package.json apps/web/src/root.tsx apps/web/src/index.css apps/web/src/components/chat apps/web/src/components/slice-one apps/web/src/hooks/chat apps/web/src/hooks/slice-one apps/web/src/lib/chat apps/web/src/lib/slice-one apps/web/src/routes.ts apps/web/src/routes/app/dashboard/page.tsx apps/web/src/routes/app/mobile/page.tsx packages/agents/package.json packages/agents/src/agents/zopu.ts packages/agents/src/app.ts packages/agents/src/auth.ts packages/agents/src/tools/slice-one.ts packages/backend/package.json packages/backend/convex/conversationMessages.ts packages/backend/convex/conversationMessages.test.ts packages/backend/convex/fluePersistence.test.ts packages/backend/convex/projects.ts packages/backend/convex/schema.ts packages/backend/convex/signalRouting.ts packages/backend/convex/signalRouting.test.ts packages/backend/convex/works.ts packages/backend/convex/works.test.ts packages/backend/convex/workArtifacts.ts packages/backend/convex/workArtifacts.test.ts packages/backend/convex/workExecution.ts packages/backend/convex/workExecution.test.ts packages/backend/convex/workPlanning.ts packages/backend/convex/workPlanning.test.ts packages/backend/convex/crons.ts packages/primitives/src/work.ts packages/primitives/src/work.test.ts packages/primitives/src/work-artifact.ts packages/primitives/src/work-artifact.test.ts packages/primitives/src/resolver.ts packages/primitives/src/work-lifecycle.ts packages/primitives/src/work-resolution.test.ts",
"lint": "oxlint --disable-nested-config",
"check": "ultracite check package.json vite.config.ts apps/web/package.json apps/web/src/root.tsx apps/web/src/index.css apps/web/src/components/chat apps/web/src/components/slice-one apps/web/src/hooks/chat apps/web/src/hooks/slice-one apps/web/src/lib/chat apps/web/src/lib/slice-one apps/web/src/routes.ts apps/web/src/routes/app/dashboard/page.tsx apps/web/src/routes/app/mobile/page.tsx packages/agents/package.json packages/agents/src/agents/zopu.ts packages/agents/src/app.ts packages/agents/src/auth.ts packages/agents/src/tools/slice-one.ts packages/backend/package.json packages/backend/convex/conversationMessages.ts packages/backend/convex/conversationMessages.test.ts packages/backend/convex/projects.ts packages/backend/convex/schema.ts packages/backend/convex/signalRouting.ts packages/backend/convex/signalRouting.test.ts packages/backend/convex/works.ts packages/backend/convex/works.test.ts packages/primitives/src/work.ts packages/primitives/src/work.test.ts",
"lint": "vp lint",
"format": "vp fmt",
"staged": "vp staged",
"hooks:setup": "vp config",
@@ -80,7 +72,7 @@
"@code/primitives": "workspace:*",
"@flue/sdk": "1.0.0-beta.9",
"@types/bun": "catalog:",
"@types/node": "catalog:",
"@types/node": "^22.13.14",
"convex": "catalog:",
"convex-test": "catalog:",
"effect": "catalog:",

View File

@@ -9,16 +9,15 @@
"dev": "bun --env-file=../../.env flue dev",
"dev:tailscale": "bun --env-file=../../.env flue dev",
"run": "bun --env-file=../../.env flue run",
"run:zopu": "bun --env-file=../../.env flue run zopu",
"run:work-planner": "bun --env-file=../../.env flue run work-planner"
"run:zopu": "bun --env-file=../../.env flue run zopu"
},
"dependencies": {
"@code/backend": "workspace:*",
"@code/env": "workspace:*",
"@flue/runtime": "latest",
"convex": "catalog:",
"hono": "catalog:",
"valibot": "catalog:"
"hono": "4.12.31",
"valibot": "^1.4.2"
},
"devDependencies": {
"@code/config": "workspace:*",

View File

@@ -1,30 +0,0 @@
import { parseAgentEnv } from "@code/env/agent";
import { defineAgent } from "@flue/runtime";
import { createWorkPlannerTools } from "../tools/work-planner";
const INSTRUCTIONS = `You are Zopu's private Work Planner.
You may only submit typed proposals through tools:
- a Work Definition proposal;
- a Design Packet proposal containing 1-4 observable, independently verifiable slices;
- structured questions.
You must never approve a Definition or Design, change Work status, start execution, claim implementation, create Git changes, or claim that simulation implemented the Work. Convex validates and stores every proposal. Ask a focused structured question when a high-impact decision is unresolved.`;
export {
convexAgentRoute as attachments,
convexAgentRoute as route,
} from "../auth";
export default defineAgent(({ env }) => {
const { AGENT_MODEL_NAME, AGENT_MODEL_PROVIDER } = parseAgentEnv(env);
return {
description:
"Proposes Work Definitions, Design Packets, and structured questions.",
instructions: INSTRUCTIONS,
model: `${AGENT_MODEL_PROVIDER}/${AGENT_MODEL_NAME}`,
thinkingLevel: "medium",
tools: createWorkPlannerTools(env),
};
});

View File

@@ -3,7 +3,7 @@ import { defineAgent } from "@flue/runtime";
import { createSliceOneTools } from "../tools/slice-one";
const INSTRUCTIONS = `You are Zopu for product Slice 1 and the Work planning handoff.
const INSTRUCTIONS = `You are Zopu for product Slice 1: Conversation to Signal to proposed Work.
## Your role
@@ -46,8 +46,7 @@ When a user sends a message, follow this decision flow:
- Repeated delivery of the same message must not create duplicate Signals or attachments. The backend is idempotent.
- Never claim you created a Signal until the tool call returns successfully.
- Do not start implementation, planning, sandboxes, Git, verification, or delivery. Those are explicitly outside Slice 1.
- After creating proposed Work, the system may invoke the private work-planner. Never claim that a Definition, Design, approval, or implementation exists unless Convex reports it.
- Proposed Work is the only Work status you directly create.`;
- Proposed Work is the only Work status in this slice.`;
export {
convexAgentRoute as attachments,

View File

@@ -1,7 +1,7 @@
import { parseAgentEnv } from "@code/env/agent";
import { registerProvider } from "@flue/runtime";
import type { Fetchable } from "@flue/runtime/routing";
import { flue } from "@flue/runtime/routing";
import { Hono } from "hono";
const agentEnv = parseAgentEnv(process.env);
@@ -24,5 +24,7 @@ registerProvider(agentEnv.AGENT_MODEL_PROVIDER, {
},
});
// flue() returns a complete Hono app; the Fetchable contract accepts it.
export default flue() satisfies Fetchable;
const app = new Hono();
app.route("/", flue() as unknown as Hono);
export default app;

View File

@@ -1,4 +1,3 @@
/* eslint-disable max-classes-per-file, @typescript-eslint/no-invalid-void-type, class-methods-use-this, @typescript-eslint/parameter-properties */
/**
* Convex-backed Flue {@link PersistenceAdapter} for the agents package.
*
@@ -26,11 +25,70 @@
*/
import { env } from "@code/env/server";
import { AttachmentConflictError, DEFAULT_LIST_LIMIT, DEFAULT_READ_LIMIT, MAX_LIST_LIMIT, MAX_READ_LIMIT, StreamListenerRegistry, assertSupportedFlueSchemaVersion, clampLimit, copyAttachmentBytes, createSessionStorageKey, decodeRunCursor, encodeRunCursor, formatOffset, hydratePersistedDirectSubmission, parseOffset, prepareDirectSubmission, verifyAttachmentBytes } from '@flue/runtime/adapter';
import type { AgentAttemptMarker, AgentDispatchAdmission, AgentDispatchReceipt, AgentExecutionStore, AgentSubmission, AgentSubmissionStore, AttachmentRef, AttachmentStore, ConversationProducerClaim, ConversationRecord, ConversationStreamBatch, ConversationStreamIdentity, ConversationStreamMeta, ConversationStreamReadResult, ConversationStreamStore, CreateRunInput, DirectAgentSubmissionInput, DispatchAgentSubmissionInput, DispatchInput, EndRunInput, EventStreamMeta, EventStreamReadResult, EventStreamStore, GetAttachmentInput, PersistedChunkRow, PersistenceAdapter, PersistenceStores, PutAttachmentInput, RunPointer, RunRecord, RunStore, RunStatus, StoredAttachment, SubmissionAttemptRef, SubmissionClaimRef, SubmissionDurability, SubmissionSettlementObligation, SubmissionSettledRecord } from '@flue/runtime/adapter';
import {
type AgentAttemptMarker,
type AgentDispatchAdmission,
type AgentDispatchReceipt,
type AgentExecutionStore,
type AgentSubmission,
type AgentSubmissionStore,
AttachmentConflictError,
type AttachmentRef,
type AttachmentStore,
type ConversationProducerClaim,
type ConversationRecord,
type ConversationStreamBatch,
type ConversationStreamIdentity,
type ConversationStreamMeta,
type ConversationStreamReadResult,
type ConversationStreamStore,
type CreateRunInput,
DEFAULT_LIST_LIMIT,
DEFAULT_READ_LIMIT,
type DirectAgentSubmissionInput,
type DispatchAgentSubmissionInput,
type DispatchInput,
type EndRunInput,
type EventStreamMeta,
type EventStreamReadResult,
type EventStreamStore,
type GetAttachmentInput,
MAX_LIST_LIMIT,
MAX_READ_LIMIT,
type PersistedChunkRow,
type PersistenceAdapter,
type PersistenceStores,
type PutAttachmentInput,
type RunPointer,
type RunRecord,
type RunStore,
type RunStatus,
StreamListenerRegistry,
type StoredAttachment,
type SubmissionAttemptRef,
type SubmissionClaimRef,
type SubmissionDurability,
type SubmissionSettlementObligation,
type SubmissionSettledRecord,
assertSupportedFlueSchemaVersion,
clampLimit,
copyAttachmentBytes,
createSessionStorageKey,
decodeRunCursor,
encodeRunCursor,
formatOffset,
hydratePersistedDirectSubmission,
parseOffset,
prepareDirectSubmission,
verifyAttachmentBytes,
} from "@flue/runtime/adapter";
import { ConvexHttpClient } from "convex/browser";
import { makeFunctionReference } from 'convex/server';
import type { FunctionArgs, FunctionReference, FunctionReturnType } from 'convex/server';
import {
type FunctionArgs,
type FunctionReference,
type FunctionReturnType,
makeFunctionReference,
} from "convex/server";
// ---------------------------------------------------------------------------
// Token + arg types
@@ -42,7 +100,7 @@ import type { FunctionArgs, FunctionReference, FunctionReturnType } from 'convex
// ---------------------------------------------------------------------------
/** Auth token present on every Convex function call. */
interface TokenArgs { readonly token: string; readonly [key: string]: unknown }
type TokenArgs = { readonly token: string; readonly [key: string]: unknown };
const TOKEN = (): string => env.FLUE_DB_TOKEN;
@@ -54,7 +112,7 @@ type TraceCarrier = NonNullable<DirectAgentSubmissionInput["traceCarrier"]>;
// ---------------------------------------------------------------------------
/** Common fields stored for every submission row. */
interface SubmissionRow {
type SubmissionRow = {
readonly sequence: number;
readonly submissionId: string;
readonly sessionKey: string;
@@ -76,42 +134,42 @@ interface SubmissionRow {
readonly ownerId: string | null;
readonly leaseExpiresAt: number;
readonly traceCarrierJson: string | null;
}
};
interface SettlementObligationRow {
type SettlementObligationRow = {
readonly submissionId: string;
readonly sessionKey: string;
readonly attemptId: string;
readonly recordId: string;
readonly recordJson: string;
}
};
interface AttemptMarkerRow {
type AttemptMarkerRow = {
readonly submissionId: string;
readonly attemptId: string;
readonly createdAt: number;
}
};
interface ConversationStreamRow {
type ConversationStreamRow = {
readonly identity: ConversationStreamIdentity;
readonly incarnation: string;
readonly nextOffset: number;
readonly producerId: string | null;
readonly producerEpoch: number;
readonly nextProducerSequence: number;
}
};
interface ConversationBatchRow {
type ConversationBatchRow = {
readonly offset: number;
readonly recordsJson: string;
}
};
interface EventEntryRow {
type EventEntryRow = {
readonly offset: number;
readonly dataJson: string;
}
};
interface RunRow {
type RunRow = {
readonly runId: string;
readonly workflowName: string;
readonly status: RunStatus;
@@ -123,9 +181,9 @@ interface RunRow {
readonly resultJson: string | null;
readonly errorJson: string | null;
readonly traceCarrierJson: string | null;
}
};
interface RunPointerRow {
type RunPointerRow = {
readonly runId: string;
readonly workflowName: string;
readonly status: RunStatus;
@@ -133,15 +191,15 @@ interface RunPointerRow {
readonly endedAt: string | null;
readonly durationMs: number | null;
readonly isError: boolean | null;
}
};
interface AttachmentRefWire {
type AttachmentRefWire = {
readonly id: string;
readonly mimeType: string;
readonly size: number;
readonly digest: string;
readonly filename?: string;
}
};
// ---------------------------------------------------------------------------
// Typed function references
@@ -205,7 +263,7 @@ const replaceSubmissionAttempt = makeFunctionReference<
>("fluePersistence:replaceSubmissionAttempt");
/** Unified admission envelope — `kind` discriminates dispatch vs direct. */
interface AdmitSubmissionEnvelope {
type AdmitSubmissionEnvelope = {
readonly kind: "dispatch" | "direct";
readonly submissionId: string;
readonly sessionKey: string;
@@ -213,12 +271,12 @@ interface AdmitSubmissionEnvelope {
readonly inputJson: string;
readonly chunksJson?: string;
readonly traceCarrierJson?: string;
}
interface AdmitSubmissionResponse {
};
type AdmitSubmissionResponse = {
readonly kind: "submission" | "retained_receipt" | "conflict";
readonly submission?: SubmissionRow;
readonly receipt?: AgentDispatchReceipt;
}
};
const admitSubmission = makeFunctionReference<
"mutation",
TokenArgs & { readonly input: AdmitSubmissionEnvelope },
@@ -352,7 +410,7 @@ type AppendBatchArgs = TokenArgs & {
readonly submission?: { readonly submissionId: string; readonly attemptId: string };
readonly recordsJson: string;
};
interface AppendResult { readonly offset: number; readonly appended: boolean }
type AppendResult = { readonly offset: number; readonly appended: boolean };
const appendConversationBatch = makeFunctionReference<
"mutation",
AppendBatchArgs,
@@ -364,11 +422,11 @@ type ReadBatchesArgs = TokenArgs & {
readonly afterOffset: number;
readonly limit: number;
};
interface ReadBatchesResult {
type ReadBatchesResult = {
readonly batches: ConversationBatchRow[];
readonly nextOffset: number;
readonly upToDate: boolean;
}
};
const readConversationBatches = makeFunctionReference<
"query",
ReadBatchesArgs,
@@ -420,12 +478,12 @@ type ReadEventsArgs = TokenArgs & {
readonly afterOffset: number;
readonly limit: number;
};
interface ReadEventsResult {
type ReadEventsResult = {
readonly events: EventEntryRow[];
readonly nextOffset: number;
readonly upToDate: boolean;
readonly closed: boolean;
}
};
const readEventsFn = makeFunctionReference<
"query",
ReadEventsArgs,
@@ -438,7 +496,7 @@ const closeEventStream = makeFunctionReference<
void
>("fluePersistence:closeEventStream");
interface EventStreamMetaRow { readonly nextOffset: number; readonly closed: boolean }
type EventStreamMetaRow = { readonly nextOffset: number; readonly closed: boolean };
const getEventStreamMeta = makeFunctionReference<
"query",
TokenArgs & { readonly path: string },
@@ -487,10 +545,10 @@ type ListRunsArgs = TokenArgs & {
readonly limit: number;
readonly cursor?: { readonly startedAt: string; readonly runId: string };
};
interface ListRunsResult {
type ListRunsResult = {
readonly rows: RunPointerRow[];
readonly hasMore: boolean;
}
};
const listRuns = makeFunctionReference<"query", ListRunsArgs, ListRunsResult>(
"fluePersistence:listRuns",
);
@@ -513,10 +571,10 @@ type GetAttachmentArgs = TokenArgs & {
readonly conversationId: string;
readonly attachmentId: string;
};
interface AttachmentWire {
type AttachmentWire = {
readonly attachment: AttachmentRefWire;
readonly bytes: ArrayBuffer;
}
};
const getAttachment = makeFunctionReference<
"query",
GetAttachmentArgs,
@@ -534,7 +592,7 @@ const deleteAttachmentsForInstance = makeFunctionReference<
// ---------------------------------------------------------------------------
const safeJsonParse = <T>(text: string | null | undefined): T | undefined => {
if (text === null || text === undefined || text === "") {return undefined;}
if (text === null || text === undefined || text === "") return undefined;
return JSON.parse(text) as T;
};
@@ -543,7 +601,7 @@ const jsonStringify = (value: unknown): string => JSON.stringify(value);
const parseAcceptedAt = (value: string, label: string): number => {
const ms = Date.parse(value);
if (!Number.isFinite(ms)) {
throw new TypeError(
throw new Error(
`[flue] ${label} produced non-finite acceptedAt: ${value}`,
);
}
@@ -555,18 +613,18 @@ const afterOffset = (offset: string | undefined): number =>
offset === undefined ? -1 : parseOffset(offset);
const wireRefToAttachmentRef = (wire: AttachmentRefWire): AttachmentRef => ({
digest: wire.digest,
id: wire.id,
mimeType: wire.mimeType,
size: wire.size,
digest: wire.digest,
...(wire.filename === undefined ? {} : { filename: wire.filename }),
});
const attachmentRefToWire = (ref: AttachmentRef): AttachmentRefWire => ({
digest: ref.digest,
id: ref.id,
mimeType: ref.mimeType,
size: ref.size,
digest: ref.digest,
...(ref.filename === undefined ? {} : { filename: ref.filename }),
});
@@ -578,17 +636,17 @@ const attachmentRefToWire = (ref: AttachmentRef): AttachmentRefWire => ({
*/
const hydrateSubmission = (row: SubmissionRow): AgentSubmission => {
const baseSubmission: Omit<AgentSubmission, "input"> = {
acceptedAt: row.acceptedAt,
attemptCount: row.attemptCount,
canonicalReadyAt: row.canonicalReadyAt,
kind: row.kind,
leaseExpiresAt: row.leaseExpiresAt,
maxRetry: row.maxRetry,
sequence: row.sequence,
sessionKey: row.sessionKey,
status: row.status,
submissionId: row.submissionId,
sessionKey: row.sessionKey,
kind: row.kind,
status: row.status,
acceptedAt: row.acceptedAt,
canonicalReadyAt: row.canonicalReadyAt,
attemptCount: row.attemptCount,
maxRetry: row.maxRetry,
timeoutAt: row.timeoutAt,
leaseExpiresAt: row.leaseExpiresAt,
...(row.attemptId === null ? {} : { attemptId: row.attemptId }),
...(row.inputAppliedAt === null
? {}
@@ -641,19 +699,19 @@ const hydrateObligation = (
);
}
return {
attemptId: row.attemptId,
record,
recordId: row.recordId,
sessionKey: row.sessionKey,
submissionId: row.submissionId,
sessionKey: row.sessionKey,
attemptId: row.attemptId,
recordId: row.recordId,
record,
};
};
const toRunRecord = (row: RunRow): RunRecord => ({
runId: row.runId,
startedAt: row.startedAt,
status: row.status,
workflowName: row.workflowName,
status: row.status,
startedAt: row.startedAt,
...(row.endedAt === null ? {} : { endedAt: row.endedAt }),
...(row.isError === null ? {} : { isError: row.isError }),
...(row.durationMs === null ? {} : { durationMs: row.durationMs }),
@@ -667,9 +725,9 @@ const toRunRecord = (row: RunRow): RunRecord => ({
const toRunPointer = (row: RunPointerRow): RunPointer => ({
runId: row.runId,
startedAt: row.startedAt,
status: row.status,
workflowName: row.workflowName,
status: row.status,
startedAt: row.startedAt,
...(row.endedAt === null ? {} : { endedAt: row.endedAt }),
...(row.durationMs === null ? {} : { durationMs: row.durationMs }),
...(row.isError === null ? {} : { isError: row.isError }),
@@ -714,13 +772,13 @@ class ConvexAgentSubmissionStore implements AgentSubmissionStore {
async getSubmission(submissionId: string): Promise<AgentSubmission | null> {
const row = await this.convex.query(getSubmission, {
submissionId,
token: TOKEN(),
submissionId,
});
return row === null ? null : hydrateSubmission(row);
}
hasUnsettledSubmissions(): Promise<boolean> {
async hasUnsettledSubmissions(): Promise<boolean> {
return this.convex.query(hasUnsettledSubmissions, { token: TOKEN() });
}
@@ -760,10 +818,10 @@ class ConvexAgentSubmissionStore implements AgentSubmissionStore {
lease?: { ownerId: string; leaseExpiresAt: number },
): Promise<AgentSubmission | null> {
const row = await this.convex.mutation(replaceSubmissionAttempt, {
token: TOKEN(),
submissionId: attempt.submissionId,
attemptId: attempt.attemptId,
nextAttemptId,
submissionId: attempt.submissionId,
token: TOKEN(),
...(lease === undefined ? {} : lease),
});
return row === null ? null : hydrateSubmission(row);
@@ -775,25 +833,25 @@ class ConvexAgentSubmissionStore implements AgentSubmissionStore {
DispatchAgentSubmissionInput,
"kind" | "submissionId"
> = {
acceptedAt: input.acceptedAt,
agent: input.agent,
dispatchId: input.dispatchId,
agent: input.agent,
id: input.id,
input: input.input,
acceptedAt: input.acceptedAt,
};
const res = await this.convex.mutation(admitSubmission, {
token: TOKEN(),
input: {
kind: "dispatch",
submissionId: input.dispatchId,
sessionKey,
acceptedAt: parseAcceptedAt(input.acceptedAt, "admitDispatch"),
inputJson: jsonStringify({
kind: "dispatch",
submissionId: input.dispatchId,
...submissionInput,
}),
kind: "dispatch",
sessionKey,
submissionId: input.dispatchId,
},
token: TOKEN(),
});
if (res.kind === "submission" && res.submission !== undefined) {
return {
@@ -813,18 +871,18 @@ class ConvexAgentSubmissionStore implements AgentSubmissionStore {
const sessionKey = createSessionStorageKey(input.id, "default", "default");
const extracted = prepareDirectSubmission(input);
const res = await this.convex.mutation(admitSubmission, {
token: TOKEN(),
input: {
acceptedAt: parseAcceptedAt(input.acceptedAt, "admitDirect"),
chunksJson: jsonStringify(extracted.chunks),
inputJson: jsonStringify(extracted.value),
kind: "direct",
sessionKey,
submissionId: input.submissionId,
sessionKey,
acceptedAt: parseAcceptedAt(input.acceptedAt, "admitDirect"),
inputJson: jsonStringify(extracted.value),
chunksJson: jsonStringify(extracted.chunks),
...(input.traceCarrier === undefined
? {}
: { traceCarrierJson: jsonStringify(input.traceCarrier) }),
},
token: TOKEN(),
});
if (res.kind === "submission" && res.submission !== undefined) {
return hydrateSubmission(res.submission);
@@ -841,8 +899,8 @@ class ConvexAgentSubmissionStore implements AgentSubmissionStore {
submissionId: string,
): Promise<AgentSubmission | null> {
const row = await this.convex.mutation(markSubmissionCanonicalReady, {
submissionId,
token: TOKEN(),
submissionId,
});
return row === null ? null : hydrateSubmission(row);
}
@@ -851,51 +909,51 @@ class ConvexAgentSubmissionStore implements AgentSubmissionStore {
claim: SubmissionClaimRef,
): Promise<AgentSubmission | null> {
const row = await this.convex.mutation(claimSubmission, {
attemptId: claim.attemptId,
leaseExpiresAt: claim.leaseExpiresAt,
ownerId: claim.ownerId,
submissionId: claim.submissionId,
token: TOKEN(),
submissionId: claim.submissionId,
attemptId: claim.attemptId,
ownerId: claim.ownerId,
leaseExpiresAt: claim.leaseExpiresAt,
});
return row === null ? null : hydrateSubmission(row);
}
markSubmissionInputApplied(
async markSubmissionInputApplied(
attempt: SubmissionAttemptRef,
durability?: SubmissionDurability,
): Promise<boolean> {
return this.convex.mutation(markSubmissionInputApplied, {
attemptId: attempt.attemptId,
submissionId: attempt.submissionId,
token: TOKEN(),
submissionId: attempt.submissionId,
attemptId: attempt.attemptId,
...(durability === undefined ? {} : durability),
});
}
requestSubmissionRecovery(
async requestSubmissionRecovery(
attempt: SubmissionAttemptRef,
): Promise<boolean> {
return this.convex.mutation(requestSubmissionRecovery, {
attemptId: attempt.attemptId,
token: TOKEN(),
submissionId: attempt.submissionId,
token: TOKEN(),
attemptId: attempt.attemptId,
});
}
requestSessionAbort(sessionKey: string): Promise<string[]> {
async requestSessionAbort(sessionKey: string): Promise<string[]> {
return this.convex.mutation(requestSessionAbort, {
sessionKey,
token: TOKEN(),
sessionKey,
});
}
requeueSubmissionBeforeInputApplied(
async requeueSubmissionBeforeInputApplied(
attempt: SubmissionAttemptRef,
): Promise<boolean> {
return this.convex.mutation(requeueSubmissionBeforeInputApplied, {
attemptId: attempt.attemptId,
submissionId: attempt.submissionId,
token: TOKEN(),
submissionId: attempt.submissionId,
attemptId: attempt.attemptId,
});
}
@@ -904,62 +962,62 @@ class ConvexAgentSubmissionStore implements AgentSubmissionStore {
settlement: { recordId: string; record: SubmissionSettledRecord },
): Promise<SubmissionSettlementObligation | null> {
const row = await this.convex.mutation(reserveSubmissionSettlement, {
token: TOKEN(),
submissionId: attempt.submissionId,
attemptId: attempt.attemptId,
recordId: settlement.recordId,
recordJson: jsonStringify(settlement.record),
submissionId: attempt.submissionId,
token: TOKEN(),
});
return row === null ? null : hydrateObligation(row);
}
finalizeSubmissionSettlement(
async finalizeSubmissionSettlement(
attempt: SubmissionAttemptRef,
recordId: string,
): Promise<boolean> {
return this.convex.mutation(finalizeSubmissionSettlement, {
token: TOKEN(),
submissionId: attempt.submissionId,
attemptId: attempt.attemptId,
recordId,
submissionId: attempt.submissionId,
token: TOKEN(),
});
}
completeSubmission(attempt: SubmissionAttemptRef): Promise<boolean> {
async completeSubmission(attempt: SubmissionAttemptRef): Promise<boolean> {
return this.convex.mutation(settleSubmission, {
token: TOKEN(),
submissionId: attempt.submissionId,
attemptId: attempt.attemptId,
outcome: "completed",
submissionId: attempt.submissionId,
token: TOKEN(),
});
}
failSubmission(
async failSubmission(
attempt: SubmissionAttemptRef,
error: unknown,
): Promise<boolean> {
return this.convex.mutation(settleSubmission, {
attemptId: attempt.attemptId,
errorJson: jsonStringify(error),
outcome: "failed",
submissionId: attempt.submissionId,
token: TOKEN(),
submissionId: attempt.submissionId,
attemptId: attempt.attemptId,
outcome: "failed",
errorJson: jsonStringify(error),
});
}
async insertAttemptMarker(attempt: SubmissionAttemptRef): Promise<void> {
await this.convex.mutation(insertAttemptMarker, {
attemptId: attempt.attemptId,
submissionId: attempt.submissionId,
token: TOKEN(),
submissionId: attempt.submissionId,
attemptId: attempt.attemptId,
});
}
async deleteAttemptMarker(attempt: SubmissionAttemptRef): Promise<void> {
await this.convex.mutation(deleteAttemptMarker, {
attemptId: attempt.attemptId,
submissionId: attempt.submissionId,
token: TOKEN(),
submissionId: attempt.submissionId,
attemptId: attempt.attemptId,
});
}
@@ -968,17 +1026,17 @@ class ConvexAgentSubmissionStore implements AgentSubmissionStore {
token: TOKEN(),
});
return rows.map((row) => ({
submissionId: row.submissionId,
attemptId: row.attemptId,
createdAt: row.createdAt,
submissionId: row.submissionId,
}));
}
async renewLeases(ownerId: string, submissionIds: string[]): Promise<void> {
await this.convex.mutation(renewLeases, {
token: TOKEN(),
ownerId,
submissionIds,
token: TOKEN(),
});
}
@@ -1004,20 +1062,20 @@ class ConvexConversationStreamStore implements ConversationStreamStore {
identity: ConversationStreamIdentity,
): Promise<void> {
await this.convex.mutation(createConversationStream, {
identity,
path,
token: TOKEN(),
path,
identity,
});
}
acquireProducer(
async acquireProducer(
path: string,
producerId: string,
): Promise<ConversationProducerClaim> {
return this.convex.mutation(acquireConversationProducer, {
token: TOKEN(),
path,
producerId,
token: TOKEN(),
});
}
@@ -1031,14 +1089,14 @@ class ConvexConversationStreamStore implements ConversationStreamStore {
records: readonly ConversationRecord[];
}): Promise<{ offset: string }> {
const result = await this.convex.mutation(appendConversationBatch, {
incarnation: input.incarnation,
path: input.path,
producerEpoch: input.producerEpoch,
producerId: input.producerId,
producerSequence: input.producerSequence,
recordsJson: jsonStringify(input.records),
token: TOKEN(),
path: input.path,
producerId: input.producerId,
producerEpoch: input.producerEpoch,
incarnation: input.incarnation,
producerSequence: input.producerSequence,
...(input.submission === undefined ? {} : { submission: input.submission }),
recordsJson: jsonStringify(input.records),
});
if (result.appended) {
this.listeners.notify(input.path);
@@ -1056,10 +1114,10 @@ class ConvexConversationStreamStore implements ConversationStreamStore {
MAX_READ_LIMIT,
);
const result = await this.convex.query(readConversationBatches, {
token: TOKEN(),
path,
afterOffset: afterOffset(options?.offset),
limit,
path,
token: TOKEN(),
});
const batches: ConversationStreamBatch[] = result.batches.map((row) => ({
offset: formatOffset(row.offset),
@@ -1074,24 +1132,24 @@ class ConvexConversationStreamStore implements ConversationStreamStore {
async getMeta(path: string): Promise<ConversationStreamMeta | null> {
const row = await this.convex.query(getConversationStreamMeta, {
path,
token: TOKEN(),
path,
});
if (row === null) {return null;}
if (row === null) return null;
return {
identity: row.identity,
incarnation: row.incarnation,
nextOffset: formatOffset(row.nextOffset),
nextProducerSequence: row.nextProducerSequence,
producerEpoch: row.producerEpoch,
producerId: row.producerId,
producerEpoch: row.producerEpoch,
nextProducerSequence: row.nextProducerSequence,
};
}
async delete(path: string): Promise<void> {
await this.convex.mutation(deleteConversationStream, {
path,
token: TOKEN(),
path,
});
}
@@ -1111,16 +1169,16 @@ class ConvexEventStreamStore implements EventStreamStore {
async createStream(path: string): Promise<void> {
await this.convex.mutation(createEventStream, {
path,
token: TOKEN(),
path,
});
}
async appendEvent(path: string, event: unknown): Promise<string> {
const result = await this.convex.mutation(appendEventFn, {
dataJson: jsonStringify(event),
path,
token: TOKEN(),
path,
dataJson: jsonStringify(event),
});
if (result.appended) {
this.listeners.notify(path);
@@ -1134,10 +1192,10 @@ class ConvexEventStreamStore implements EventStreamStore {
event: unknown,
): Promise<string> {
const result = await this.convex.mutation(appendEventOnce, {
dataJson: jsonStringify(event),
key,
path,
token: TOKEN(),
path,
key,
dataJson: jsonStringify(event),
});
if (result.appended) {
this.listeners.notify(path);
@@ -1151,38 +1209,38 @@ class ConvexEventStreamStore implements EventStreamStore {
): Promise<EventStreamReadResult> {
const limit = clampLimit(opts?.limit, DEFAULT_READ_LIMIT, MAX_READ_LIMIT);
const result = await this.convex.query(readEventsFn, {
token: TOKEN(),
path,
afterOffset: afterOffset(opts?.offset),
limit,
path,
token: TOKEN(),
});
return {
closed: result.closed,
events: result.events.map((row) => ({
data: safeJsonParse(row.dataJson),
offset: formatOffset(row.offset),
data: safeJsonParse(row.dataJson),
})),
nextOffset: formatOffset(result.nextOffset),
upToDate: result.upToDate,
closed: result.closed,
};
}
async closeStream(path: string): Promise<void> {
await this.convex.mutation(closeEventStream, {
path,
token: TOKEN(),
path,
});
}
async getStreamMeta(path: string): Promise<EventStreamMeta | null> {
const row = await this.convex.query(getEventStreamMeta, {
path,
token: TOKEN(),
path,
});
if (row === null) {return null;}
if (row === null) return null;
return {
closed: row.closed,
nextOffset: formatOffset(row.nextOffset),
closed: row.closed,
};
}
@@ -1200,11 +1258,11 @@ class ConvexRunStore implements RunStore {
async createRun(input: CreateRunInput): Promise<void> {
await this.convex.mutation(createRun, {
inputJson: jsonStringify(input.input),
runId: input.runId,
startedAt: input.startedAt,
token: TOKEN(),
runId: input.runId,
workflowName: input.workflowName,
startedAt: input.startedAt,
inputJson: jsonStringify(input.input),
...(input.traceCarrier === undefined
? {}
: { traceCarrierJson: jsonStringify(input.traceCarrier) }),
@@ -1213,11 +1271,11 @@ class ConvexRunStore implements RunStore {
async endRun(input: EndRunInput): Promise<void> {
await this.convex.mutation(endRun, {
durationMs: input.durationMs,
token: TOKEN(),
runId: input.runId,
endedAt: input.endedAt,
isError: input.isError,
runId: input.runId,
token: TOKEN(),
durationMs: input.durationMs,
...(input.result === undefined
? {}
: { resultJson: jsonStringify(input.result) }),
@@ -1228,14 +1286,14 @@ class ConvexRunStore implements RunStore {
}
async getRun(runId: string): Promise<RunRecord | null> {
const row = await this.convex.query(getRun, { runId, token: TOKEN() });
const row = await this.convex.query(getRun, { token: TOKEN(), runId });
return row === null ? null : toRunRecord(row);
}
lookupRun(
async lookupRun(
runId: string,
): Promise<{ runId: string; workflowName: string } | null> {
return this.convex.query(lookupRun, { runId, token: TOKEN() });
return this.convex.query(lookupRun, { token: TOKEN(), runId });
}
async listRuns(opts?: {
@@ -1256,10 +1314,10 @@ class ConvexRunStore implements RunStore {
...(decoded === undefined ? {} : { cursor: decoded }),
});
const runs = result.rows.map(toRunPointer);
const last = runs.at(-1);
const last = runs[runs.length - 1];
const nextCursor =
result.hasMore && last !== undefined
? encodeRunCursor({ runId: last.runId, startedAt: last.startedAt })
? encodeRunCursor({ startedAt: last.startedAt, runId: last.runId })
: undefined;
return { runs, ...(nextCursor === undefined ? {} : { nextCursor }) };
}
@@ -1278,31 +1336,31 @@ class ConvexAttachmentStore implements AttachmentStore {
// payload cannot be mutated after the call returns.
const bytes = copyAttachmentBytes(input.bytes);
const status = await this.convex.mutation(putAttachment, {
token: TOKEN(),
streamPath: input.streamPath,
conversationId: input.conversationId,
attachment: attachmentRefToWire(input.attachment),
bytes: bytes.buffer.slice(
bytes.byteOffset,
bytes.byteOffset + bytes.byteLength,
) as ArrayBuffer,
conversationId: input.conversationId,
streamPath: input.streamPath,
token: TOKEN(),
});
if (status === "conflict") {
throw new AttachmentConflictError({
attachmentId: input.attachment.id,
path: input.streamPath,
attachmentId: input.attachment.id,
});
}
}
async get(input: GetAttachmentInput): Promise<StoredAttachment | null> {
const row = await this.convex.query(getAttachment, {
attachmentId: input.attachmentId,
conversationId: input.conversationId,
streamPath: input.streamPath,
token: TOKEN(),
streamPath: input.streamPath,
conversationId: input.conversationId,
attachmentId: input.attachmentId,
});
if (row === null) {return null;}
if (row === null) return null;
const attachment = wireRefToAttachmentRef(row.attachment);
const bytes = new Uint8Array(row.bytes);
// Verify integrity before handing bytes back to the runtime; throws
@@ -1316,8 +1374,8 @@ class ConvexAttachmentStore implements AttachmentStore {
async deleteForInstance(streamPath: string): Promise<void> {
await this.convex.mutation(deleteAttachmentsForInstance, {
streamPath,
token: TOKEN(),
streamPath,
});
}
}
@@ -1351,11 +1409,11 @@ class ConvexPersistenceAdapterImpl implements PersistenceAdapter {
const submissions = new ConvexAgentSubmissionStore(this.convex);
const executionStore: AgentExecutionStore = { submissions };
return {
attachmentStore: new ConvexAttachmentStore(this.convex),
conversationStreamStore: new ConvexConversationStreamStore(this.convex),
eventStreamStore: new ConvexEventStreamStore(this.convex),
executionStore,
runStore: new ConvexRunStore(this.convex),
eventStreamStore: new ConvexEventStreamStore(this.convex),
conversationStreamStore: new ConvexConversationStreamStore(this.convex),
attachmentStore: new ConvexAttachmentStore(this.convex),
};
}

View File

@@ -1,69 +0,0 @@
import { parseAgentEnv } from "@code/env/agent";
import { defineTool } from "@flue/runtime";
import { ConvexHttpClient } from "convex/browser";
import { makeFunctionReference } from "convex/server";
import * as v from "valibot";
const submitDefinitionRef = makeFunctionReference<
"mutation",
{ workId: string; payloadJson: string; token: string },
{ version: number }
>("workPlanning:submitDefinitionProposal");
const submitDesignRef = makeFunctionReference<
"mutation",
{ workId: string; payloadJson: string; token: string },
{ version: number }
>("workPlanning:submitDesignProposal");
const submitQuestionRef = makeFunctionReference<
"mutation",
{ workId: string; questionJson: string; token: string },
{ accepted: boolean }
>("workPlanning:submitQuestion");
export const createWorkPlannerTools = (
runtimeEnv: Record<string, string | undefined>
) => {
const env = parseAgentEnv(runtimeEnv);
const client = new ConvexHttpClient(env.CONVEX_URL);
const token = env.FLUE_DB_TOKEN;
return [
defineTool({
description:
"Submit a versioned Work Definition proposal. This never approves the Definition.",
input: v.object({ definitionJson: v.string(), workId: v.string() }),
name: "submit_definition_proposal",
async run({ input }) {
return await client.mutation(submitDefinitionRef, {
payloadJson: input.definitionJson,
token,
workId: input.workId,
});
},
}),
defineTool({
description:
"Submit a Design Packet with observable independently verifiable slices. This never approves the Design.",
input: v.object({ designJson: v.string(), workId: v.string() }),
name: "submit_design_proposal",
async run({ input }) {
return await client.mutation(submitDesignRef, {
payloadJson: input.designJson,
token,
workId: input.workId,
});
},
}),
defineTool({
description: "Submit one structured question for human attention.",
input: v.object({ questionJson: v.string(), workId: v.string() }),
name: "submit_question",
async run({ input }) {
return await client.mutation(submitQuestionRef, {
questionJson: input.questionJson,
token,
workId: input.workId,
});
},
}),
];
};

View File

@@ -23,13 +23,13 @@
"expo-secure-store": "~57.0.0",
"heroui-native": "catalog:",
"react": "catalog:",
"react-native": "catalog:",
"react-native": "0.86.0",
"sonner": "catalog:",
"zod": "catalog:"
},
"devDependencies": {
"@code/config": "workspace:*",
"@types/react": "catalog:",
"@types/react": "~19.2.17",
"typescript": "catalog:",
"vite": "catalog:",
"vitest": "catalog:"

View File

@@ -0,0 +1,10 @@
{
"overrides": [
{
"files": ["convex/**/*.ts"],
"rules": {
"unicorn/filename-case": "off"
}
}
]
}

View File

@@ -1,6 +1,8 @@
# Convex Backend
Convex is the only application API used by web, desktop, and mobile clients. The active backend covers the durable control plane through Slice 4: tenancy, projects, conversation turns, Signals, Work planning, approved Design slices, simulated Runs and Attempts, artifacts/delivery metadata, and Flue persistence.
Convex is the only application API used by web, desktop, and mobile clients.
The active Slice 1 backend is intentionally limited to tenancy, projects,
conversation turns, Signals, Work, and Flue's required persistence adapter.
## Domain relations
@@ -9,9 +11,6 @@ Convex is the only application API used by web, desktop, and mobile clients. The
- `conversations`, `conversationTurns`, `conversationMessages`, and `conversationAttachments` own chat admission and reactive history.
- `signals`, `signalConstraints`, and `signalSources` preserve structured intent and exact evidence.
- `works`, `signalWorkAttachments`, and `workEvents` own proposed Work and provenance.
- `workDefinitions`, `workQuestions`, and `workApprovals` own versioned outcome contracts and approvals.
- `designPackets` and `workSlices` preserve versioned implementation intent and executable slice order.
- `workRuns`, `workAttempts`, `workAttemptEvents`, and `resolverDecisions` own bounded simulated execution and recovery.
- `workArtifacts` and `workDeliveries` preserve evidence and delivery metadata without performing Git or sandbox work.
The `flue*` tables are infrastructure tables required by Flue's persistence contract. They are deliberately separate from the product relations.
The `flue*` tables are infrastructure tables required by Flue's persistence
contract. They are deliberately separate from the product relations.

View File

@@ -1,5 +1,5 @@
import { convexTest } from "convex-test";
import { anyApi, makeFunctionReference } from "convex/server";
import { anyApi } from "convex/server";
import { describe, expect, test } from "vitest";
import schema from "./schema";
@@ -15,33 +15,6 @@ const api = anyApi;
const identityA = { tokenIdentifier: "https://convex.test|user-a" };
const identityB = { tokenIdentifier: "https://convex.test|user-b" };
const newTest = () => convexTest({ modules, schema });
const markProcessingRef = makeFunctionReference<
"mutation",
{ turnId: string; attempt: number; leaseOwner: string },
boolean
>("conversationMessages:markProcessing");
const completeTurnRef = makeFunctionReference<
"mutation",
{
turnId: string;
attempt: number;
leaseOwner: string;
submissionId: string;
text: string;
},
boolean
>("conversationMessages:completeTurn");
const failTurnRef = makeFunctionReference<
"mutation",
{
turnId: string;
attempt: number;
leaseOwner: string;
error: string;
retry: boolean;
},
boolean
>("conversationMessages:failTurn");
const ensureOrg = async (
t: ReturnType<typeof newTest>,
@@ -136,66 +109,4 @@ describe("conversationMessages", () => {
})
).rejects.toThrow(/Organization membership required/u);
});
test("fences stale turn attempts from overwriting a retry", async () => {
const t = newTest();
const organization = await ensureOrg(t, identityA);
const sent = await t
.withIdentity(identityA)
.mutation(api.conversationMessages.send, {
clientRequestId: "request-fenced",
images: [],
organizationId: organization._id,
rawText: "Keep only the current response",
});
expect(
await t.mutation(markProcessingRef, {
attempt: 1,
leaseOwner: "worker-1",
turnId: sent.turnId,
})
).toBe(true);
expect(
await t.mutation(failTurnRef, {
attempt: 1,
error: "retry",
leaseOwner: "worker-1",
retry: true,
turnId: sent.turnId,
})
).toBe(true);
expect(
await t.mutation(completeTurnRef, {
attempt: 1,
leaseOwner: "worker-1",
submissionId: "stale",
text: "stale response",
turnId: sent.turnId,
})
).toBe(false);
expect(
await t.mutation(markProcessingRef, {
attempt: 2,
leaseOwner: "worker-2",
turnId: sent.turnId,
})
).toBe(true);
expect(
await t.mutation(completeTurnRef, {
attempt: 2,
leaseOwner: "worker-2",
submissionId: "current",
text: "current response",
turnId: sent.turnId,
})
).toBe(true);
const messages = await t
.withIdentity(identityA)
.query(api.conversationMessages.listForCurrentOrganization, {
organizationId: organization._id,
});
expect(messages[1]?.rawText).toBe("current response");
});
});

View File

@@ -31,34 +31,18 @@ const getTurnRef = makeFunctionReference<
>("conversationMessages:getTurn");
const markProcessingRef = makeFunctionReference<
"mutation",
{
turnId: Id<"conversationTurns">;
attempt: number;
leaseOwner: string;
},
{ turnId: Id<"conversationTurns"> },
boolean
>("conversationMessages:markProcessing");
const completeTurnRef = makeFunctionReference<
"mutation",
{
turnId: Id<"conversationTurns">;
attempt: number;
leaseOwner: string;
submissionId: string;
text: string;
},
boolean
{ turnId: Id<"conversationTurns">; submissionId: string; text: string },
null
>("conversationMessages:completeTurn");
const failTurnRef = makeFunctionReference<
"mutation",
{
turnId: Id<"conversationTurns">;
attempt: number;
leaseOwner: string;
error: string;
retry: boolean;
},
boolean
{ turnId: Id<"conversationTurns">; error: string; retry: boolean },
null
>("conversationMessages:failTurn");
export const generateUploadUrl = mutation({
@@ -73,16 +57,16 @@ export const generateUploadUrl = mutation({
export const send = mutation({
args: {
organizationId: v.id("organizations"),
clientRequestId: v.string(),
rawText: v.string(),
images: v.array(
v.object({
storageId: v.id("_storage"),
filename: v.optional(v.string()),
mimeType: v.string(),
storageId: v.id("_storage"),
})
),
organizationId: v.id("organizations"),
rawText: v.string(),
},
handler: async (ctx, args) => {
await requireOrganizationMember(ctx, args.organizationId);
@@ -131,7 +115,6 @@ export const send = mutation({
.first();
const ordinal = (lastMessage?.ordinal ?? -1) + 1;
const turnId = await ctx.db.insert("conversationTurns", {
attemptNumber: 1,
clientRequestId: args.clientRequestId,
conversationId: conversation._id,
createdAt,
@@ -250,49 +233,24 @@ export const getTurn = internalQuery({
});
export const markProcessing = internalMutation({
args: {
attempt: v.number(),
leaseOwner: v.string(),
turnId: v.id("conversationTurns"),
},
args: { turnId: v.id("conversationTurns") },
handler: async (ctx, args): Promise<boolean> => {
const turn = await ctx.db.get(args.turnId);
if (
!turn ||
turn.status !== "queued" ||
(turn.attemptNumber ?? 1) !== args.attempt
) {
if (!turn || turn.status === "completed") {
return false;
}
await ctx.db.patch(turn._id, {
error: undefined,
leaseExpiresAt: Date.now() + 60_000,
leaseOwner: args.leaseOwner,
status: "processing",
});
await ctx.db.patch(turn._id, { error: undefined, status: "processing" });
return true;
},
});
export const completeTurn = internalMutation({
args: {
attempt: v.number(),
leaseOwner: v.string(),
turnId: v.id("conversationTurns"),
submissionId: v.string(),
text: v.string(),
turnId: v.id("conversationTurns"),
},
handler: async (ctx, args): Promise<boolean> => {
const turn = await ctx.db.get(args.turnId);
if (
!turn ||
turn.status !== "processing" ||
(turn.attemptNumber ?? 1) !== args.attempt ||
turn.leaseOwner !== args.leaseOwner ||
(turn.leaseExpiresAt !== undefined && turn.leaseExpiresAt < Date.now())
) {
return false;
}
handler: async (ctx, args): Promise<null> => {
const assistant = await ctx.db
.query("conversationMessages")
.withIndex("by_turnId_and_role", (q) =>
@@ -305,43 +263,29 @@ export const completeTurn = internalMutation({
await ctx.db.patch(args.turnId, {
completedAt: Date.now(),
error: undefined,
leaseExpiresAt: undefined,
leaseOwner: undefined,
status: "completed",
submissionId: args.submissionId,
});
return true;
return null;
},
});
export const failTurn = internalMutation({
args: {
attempt: v.number(),
error: v.string(),
leaseOwner: v.string(),
retry: v.boolean(),
turnId: v.id("conversationTurns"),
error: v.string(),
retry: v.boolean(),
},
handler: async (ctx, args): Promise<boolean> => {
handler: async (ctx, args): Promise<null> => {
const turn = await ctx.db.get(args.turnId);
if (
!turn ||
turn.status !== "processing" ||
(turn.attemptNumber ?? 1) !== args.attempt ||
turn.leaseOwner !== args.leaseOwner ||
(turn.leaseExpiresAt !== undefined && turn.leaseExpiresAt < Date.now())
) {
return false;
if (turn && turn.status !== "completed") {
await ctx.db.patch(turn._id, {
completedAt: args.retry ? undefined : Date.now(),
error: args.error,
status: args.retry ? "queued" : "failed",
});
}
await ctx.db.patch(turn._id, {
attemptNumber: args.retry ? args.attempt + 1 : args.attempt,
completedAt: args.retry ? undefined : Date.now(),
error: args.error,
leaseExpiresAt: undefined,
leaseOwner: undefined,
status: args.retry ? "queued" : "failed",
});
return true;
return null;
},
});
@@ -354,17 +298,12 @@ const toBase64 = (buffer: ArrayBuffer): string => {
};
export const runTurn = internalAction({
args: { attempt: v.number(), turnId: v.id("conversationTurns") },
args: { turnId: v.id("conversationTurns"), attempt: v.number() },
handler: async (ctx, args): Promise<null> => {
const leaseOwner = `conversation:${args.turnId}:${args.attempt}`;
const turn = await ctx.runQuery(getTurnRef, { turnId: args.turnId });
if (
!turn ||
!(await ctx.runMutation(markProcessingRef, {
attempt: args.attempt,
leaseOwner,
turnId: args.turnId,
}))
!(await ctx.runMutation(markProcessingRef, { turnId: args.turnId }))
) {
return null;
}
@@ -419,8 +358,6 @@ export const runTurn = internalAction({
throw new Error(`Flue turn failed (${response.status})`);
}
await ctx.runMutation(completeTurnRef, {
attempt: args.attempt,
leaseOwner,
submissionId: payload.submissionId,
text: payload.result.text,
turnId: args.turnId,
@@ -428,9 +365,7 @@ export const runTurn = internalAction({
} catch (error) {
const retry = args.attempt < MAX_ATTEMPTS;
await ctx.runMutation(failTurnRef, {
attempt: args.attempt,
error: error instanceof Error ? error.message : String(error),
leaseOwner,
retry,
turnId: args.turnId,
});
@@ -444,34 +379,3 @@ export const runTurn = internalAction({
return null;
},
});
export const reconcileExpiredTurns = internalMutation({
args: {},
handler: async (ctx): Promise<{ reconciled: number }> => {
const expired = await ctx.db
.query("conversationTurns")
.withIndex("by_status_and_leaseExpiresAt", (q) =>
q.eq("status", "processing").lt("leaseExpiresAt", Date.now())
)
.collect();
for (const turn of expired) {
const attempt = turn.attemptNumber ?? 1;
const retry = attempt < MAX_ATTEMPTS;
await ctx.db.patch(turn._id, {
attemptNumber: retry ? attempt + 1 : attempt,
completedAt: retry ? undefined : Date.now(),
error: "Conversation worker lease expired",
leaseExpiresAt: undefined,
leaseOwner: undefined,
status: retry ? "queued" : "failed",
});
if (retry) {
await ctx.scheduler.runAfter(0, runTurnRef, {
attempt: attempt + 1,
turnId: turn._id,
});
}
}
return { reconciled: expired.length };
},
});

View File

@@ -4,12 +4,12 @@ import { v } from "convex/values";
const app = defineApp({
env: {
FLUE_DB_TOKEN: v.string(),
FLUE_URL: v.optional(v.string()),
GITEA_TOKEN: v.optional(v.string()),
GITEA_URL: v.optional(v.string()),
NATIVE_APP_URL: v.optional(v.string()),
SITE_URL: v.string(),
FLUE_DB_TOKEN: v.string(),
FLUE_URL: v.optional(v.string()),
GITEA_URL: v.optional(v.string()),
GITEA_TOKEN: v.optional(v.string()),
},
});
app.use(betterAuth);

View File

@@ -1,29 +0,0 @@
import { cronJobs, makeFunctionReference } from "convex/server";
// Reconciles attempts whose worker lease expired (crash/restart). The target
// is referenced by string so this file does not depend on regenerated codegen.
const reconcileRef = makeFunctionReference<
"mutation",
Record<string, never>,
{ reconciled: number } | null
>("workExecution:reconcileExpiredAttempts");
const reconcileConversationTurnsRef = makeFunctionReference<
"mutation",
Record<string, never>,
{ reconciled: number }
>("conversationMessages:reconcileExpiredTurns");
const crons = cronJobs();
crons.interval(
"reconcile expired work attempts",
{ seconds: 30 },
reconcileRef
);
crons.interval(
"reconcile expired conversation turns",
{ seconds: 30 },
reconcileConversationTurnsRef
);
export default crons;

View File

@@ -1,142 +0,0 @@
import { env } from "@code/env/convex";
import { convexTest } from "convex-test";
import { anyApi } from "convex/server";
import { describe, expect, test } from "vitest";
import schema from "./schema";
declare global {
interface ImportMeta {
readonly glob: (pattern: string) => Record<string, () => Promise<unknown>>;
}
}
const modules = import.meta.glob("./**/*.ts");
const api = anyApi;
const token = env.FLUE_DB_TOKEN;
describe("Flue Convex persistence", () => {
test("admits idempotently and claims only the session head", async () => {
const t = convexTest({ modules, schema });
const firstInput = {
acceptedAt: 1,
inputJson: '{"message":"first"}',
kind: "dispatch" as const,
sessionKey: "agent/instance/default",
submissionId: "dispatch-1",
};
const first = await t.mutation(api.fluePersistence.admitSubmission, {
input: firstInput,
token,
});
const replay = await t.mutation(api.fluePersistence.admitSubmission, {
input: firstInput,
token,
});
expect(first.kind).toBe("submission");
expect(replay.kind).toBe("retained_receipt");
await t.mutation(api.fluePersistence.admitSubmission, {
input: {
...firstInput,
acceptedAt: 2,
inputJson: '{"message":"second"}',
submissionId: "dispatch-2",
},
token,
});
await t.mutation(api.fluePersistence.markSubmissionCanonicalReady, {
submissionId: "dispatch-1",
token,
});
await t.mutation(api.fluePersistence.markSubmissionCanonicalReady, {
submissionId: "dispatch-2",
token,
});
const secondClaim = await t.mutation(api.fluePersistence.claimSubmission, {
attemptId: "attempt-2",
leaseExpiresAt: 100,
ownerId: "owner",
submissionId: "dispatch-2",
token,
});
expect(secondClaim).toBeNull();
const firstClaim = await t.mutation(api.fluePersistence.claimSubmission, {
attemptId: "attempt-1",
leaseExpiresAt: 100,
ownerId: "owner",
submissionId: "dispatch-1",
token,
});
expect(firstClaim).toMatchObject({
attemptId: "attempt-1",
status: "running",
});
});
test("fences stale conversation producers and conflicting attachments", async () => {
const t = convexTest({ modules, schema });
await t.mutation(api.fluePersistence.createConversationStream, {
identity: { agentName: "work-planner", instanceId: "org-1" },
path: "agents/work-planner/org-1",
token,
});
const first = await t.mutation(
api.fluePersistence.acquireConversationProducer,
{
path: "agents/work-planner/org-1",
producerId: "producer-1",
token,
}
);
await t.mutation(api.fluePersistence.acquireConversationProducer, {
path: "agents/work-planner/org-1",
producerId: "producer-2",
token,
});
await expect(
t.mutation(api.fluePersistence.appendConversationBatch, {
incarnation: first.incarnation,
path: "agents/work-planner/org-1",
producerEpoch: first.producerEpoch,
producerId: "producer-1",
producerSequence: 0,
recordsJson: '[{"id":"record-1","type":"message"}]',
token,
})
).rejects.toThrow(/producer ownership is stale/u);
const attachment = {
digest: "digest",
id: "attachment-1",
mimeType: "text/plain",
size: 3,
};
const inserted = await t.mutation(api.fluePersistence.putAttachment, {
attachment,
bytes: new TextEncoder().encode("one").buffer,
conversationId: "conversation-1",
streamPath: "agents/work-planner/org-1",
token,
});
const replay = await t.mutation(api.fluePersistence.putAttachment, {
attachment,
bytes: new TextEncoder().encode("one").buffer,
conversationId: "conversation-1",
streamPath: "agents/work-planner/org-1",
token,
});
const conflict = await t.mutation(api.fluePersistence.putAttachment, {
attachment,
bytes: new TextEncoder().encode("two").buffer,
conversationId: "conversation-1",
streamPath: "agents/work-planner/org-1",
token,
});
expect([inserted, replay, conflict]).toEqual([
"inserted",
"existing",
"conflict",
]);
});
});

File diff suppressed because it is too large Load Diff

View File

@@ -1,7 +1,7 @@
import { query } from "./_generated/server";
export const get = query({
handler: () =>
"OK"
,
handler: async () => {
return "OK";
},
});

View File

@@ -23,7 +23,7 @@ const ID_B = "https://convex.test|user-b";
const identityA = { tokenIdentifier: ID_A };
const identityB = { tokenIdentifier: ID_B };
const newTest = () => convexTest({ modules, schema });
const newTest = () => convexTest({ schema, modules });
describe("organizations", () => {
test("first ensure creates one organization and owner membership", async () => {

View File

@@ -1,10 +1,8 @@
import {
decodePublicGitImportResult,
preparePublicGitSource,
} from "@code/primitives/project";
import type {
ProjectImportOutcome,
ProjectView,
type ProjectImportOutcome,
type ProjectView,
} from "@code/primitives/project";
import { ConvexError, v } from "convex/values";
import { Effect } from "effect";
@@ -56,11 +54,18 @@ const toProjectView = async (
export const persistPublicGitImport = internalMutation({
args: {
userId: v.string(),
source: v.object({
host: v.string(),
projectName: v.string(),
repositoryPath: v.string(),
normalizedUrl: v.string(),
url: v.string(),
}),
remote: v.object({
defaultBranch: v.optional(v.string()),
documents: v.array(
v.object({
content: v.string(),
kind: v.union(
v.literal("readme"),
v.literal("agents"),
@@ -70,18 +75,11 @@ export const persistPublicGitImport = internalMutation({
v.literal("tech")
),
path: v.string(),
content: v.string(),
})
),
warnings: v.array(v.object({ message: v.string(), path: v.string() })),
warnings: v.array(v.object({ path: v.string(), message: v.string() })),
}),
source: v.object({
host: v.string(),
normalizedUrl: v.string(),
projectName: v.string(),
repositoryPath: v.string(),
url: v.string(),
}),
userId: v.string(),
},
handler: async (ctx, args): Promise<ProjectImportOutcome> => {
const organization = await ctx.db
@@ -108,9 +106,9 @@ export const persistPublicGitImport = internalMutation({
await ctx.db.patch(projectId, {
defaultBranch: args.remote.defaultBranch,
name: args.source.projectName,
repositoryPath: args.source.repositoryPath,
sourceHost: args.source.host,
sourceUrl: args.source.url,
repositoryPath: args.source.repositoryPath,
updatedAt: timestamp,
});
const oldDocuments = await ctx.db

View File

@@ -1,5 +1,11 @@
import { CONTEXT_KINDS, contextPathForKind, MAX_CONTEXT_DOCUMENT_CHARACTERS } from '@code/primitives/project';
import type { PreparedPublicGitSource, PublicGitImportResult, RepositoryContextDocument } from '@code/primitives/project';
import {
CONTEXT_KINDS,
contextPathForKind,
MAX_CONTEXT_DOCUMENT_CHARACTERS,
type PreparedPublicGitSource,
type PublicGitImportResult,
type RepositoryContextDocument,
} from "@code/primitives/project";
const GITHUB_HOST = "github.com";
const GITEA_HOST = "git.openputer.com";

View File

@@ -1,59 +1,32 @@
import { defineSchema, defineTable } from "convex/server";
import { v } from "convex/values";
/* eslint-disable sort-keys */
const attemptClassification = v.union(
v.literal("Succeeded"),
v.literal("RetryableFailure"),
v.literal("NeedsInput"),
v.literal("Blocked"),
v.literal("VerificationFailed"),
v.literal("BudgetExhausted"),
v.literal("Cancelled"),
v.literal("PermanentFailure")
);
const workStatus = 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")
);
export default defineSchema({
organizations: defineTable({
createdAt: v.number(),
createdBy: v.string(),
kind: v.union(v.literal("personal"), v.literal("team")),
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({
createdAt: v.number(),
organizationId: v.id("organizations"),
role: v.union(v.literal("owner"), v.literal("member")),
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({
createdAt: v.number(),
defaultBranch: v.optional(v.string()),
description: v.optional(v.string()),
name: v.string(),
normalizedSourceUrl: v.string(),
organizationId: v.id("organizations"),
repositoryPath: v.string(),
sourceHost: v.string(),
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"])
@@ -62,8 +35,7 @@ export default defineSchema({
"normalizedSourceUrl",
]),
projectContextDocuments: defineTable({
content: v.string(),
createdAt: v.number(),
projectId: v.id("projects"),
kind: v.union(
v.literal("readme"),
v.literal("agents"),
@@ -72,25 +44,20 @@ export default defineSchema({
v.literal("design"),
v.literal("tech")
),
origin: v.literal("repository"),
path: v.string(),
projectId: v.id("projects"),
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({
createdAt: v.number(),
organizationId: v.id("organizations"),
createdAt: v.number(),
}).index("by_organizationId", ["organizationId"]),
conversationTurns: defineTable({
attemptNumber: v.optional(v.number()),
clientRequestId: v.string(),
completedAt: v.optional(v.number()),
conversationId: v.id("conversations"),
createdAt: v.number(),
error: v.optional(v.string()),
leaseExpiresAt: v.optional(v.number()),
leaseOwner: v.optional(v.string()),
clientRequestId: v.string(),
status: v.union(
v.literal("queued"),
v.literal("processing"),
@@ -98,368 +65,90 @@ export default defineSchema({
v.literal("failed")
),
submissionId: v.optional(v.string()),
})
.index("by_conversationId_and_clientRequestId", [
"conversationId",
"clientRequestId",
])
.index("by_status_and_leaseExpiresAt", ["status", "leaseExpiresAt"]),
conversationMessages: defineTable({
content: v.string(),
conversationId: v.id("conversations"),
error: v.optional(v.string()),
createdAt: v.number(),
ordinal: v.number(),
role: v.union(v.literal("user"), v.literal("assistant")),
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({
createdAt: v.number(),
filename: v.optional(v.string()),
messageId: v.id("conversationMessages"),
mimeType: v.string(),
storageId: v.id("_storage"),
filename: v.optional(v.string()),
mimeType: v.string(),
createdAt: v.number(),
}).index("by_messageId", ["messageId"]),
signals: defineTable({
conversationId: v.id("conversations"),
createdAt: v.number(),
desiredOutcome: v.string(),
organizationId: v.id("organizations"),
processedByAgentInstanceId: v.string(),
processedByAgentName: v.string(),
projectId: v.id("projects"),
conversationId: v.id("conversations"),
sourceKey: v.string(),
summary: 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({
ordinal: v.number(),
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(),
signalId: v.id("signals"),
sourceCreatedAt: v.number(),
})
.index("by_signalId_and_ordinal", ["signalId", "ordinal"])
.index("by_messageId", ["messageId"]),
works: defineTable({
createdAt: v.number(),
definitionApprovalVersion: v.optional(v.number()),
definitionVersion: v.optional(v.number()),
designApprovalVersion: v.optional(v.number()),
designVersion: v.optional(v.number()),
objective: v.string(),
organizationId: v.id("organizations"),
projectId: v.id("projects"),
status: workStatus,
title: v.string(),
objective: v.string(),
status: v.literal("proposed"),
createdAt: v.number(),
updatedAt: v.number(),
})
.index("by_project_and_createdAt", ["projectId", "createdAt"])
.index("by_organization_and_createdAt", ["organizationId", "createdAt"]),
signalWorkAttachments: defineTable({
createdAt: v.number(),
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({
createdAt: v.number(),
idempotencyKey: v.string(),
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.created"),
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("planner.failed"),
v.literal("slice.started"),
v.literal("slice.completed"),
v.literal("slice.ready"),
v.literal("run.started"),
v.literal("run.completed"),
v.literal("run.cancelled"),
v.literal("attempt.claimed"),
v.literal("attempt.event"),
v.literal("attempt.completed"),
v.literal("attempt.reconciled"),
v.literal("resolver.decided"),
v.literal("artifact.recorded"),
v.literal("delivery.recorded"),
v.literal("delivery.updated")
),
payloadJson: v.optional(v.string()),
referenceId: v.optional(v.string()),
signalId: v.optional(v.id("signals")),
workId: v.id("works"),
signalId: v.id("signals"),
kind: v.union(v.literal("work.proposed"), v.literal("signal.attached")),
idempotencyKey: v.string(),
createdAt: v.number(),
})
.index("by_work_and_createdAt", ["workId", "createdAt"])
.index("by_work_and_idempotencyKey", ["workId", "idempotencyKey"]),
workDefinitions: defineTable({
createdAt: v.number(),
createdBy: v.string(),
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")
),
version: v.number(),
workId: v.id("works"),
})
.index("by_work_and_version", ["workId", "version"])
.index("by_work_and_status", ["workId", "status"]),
workQuestions: defineTable({
alternativesJson: v.string(),
answer: v.optional(v.string()),
createdAt: v.number(),
definitionVersion: v.number(),
impact: v.union(v.literal("low"), v.literal("medium"), v.literal("high")),
prompt: v.string(),
questionId: v.string(),
recommendation: v.optional(v.string()),
status: v.union(
v.literal("open"),
v.literal("answered"),
v.literal("withdrawn")
),
workId: v.id("works"),
})
.index("by_workId_and_definitionVersion", ["workId", "definitionVersion"])
.index("by_workId_and_definitionVersion_and_questionId", [
"workId",
"definitionVersion",
"questionId",
]),
workApprovals: defineTable({
approvedAt: v.number(),
approvedBy: v.string(),
definitionVersion: v.number(),
designVersion: v.optional(v.number()),
kind: v.union(v.literal("definition"), v.literal("design")),
status: v.union(v.literal("active"), v.literal("invalidated")),
workId: v.id("works"),
})
.index("by_workId_and_kind", ["workId", "kind"])
.index("by_workId_and_kind_and_definitionVersion_and_designVersion", [
"workId",
"kind",
"definitionVersion",
"designVersion",
]),
designPackets: defineTable({
createdAt: v.number(),
createdBy: v.string(),
definitionVersion: v.number(),
payloadJson: v.string(),
status: v.union(
v.literal("proposed"),
v.literal("current"),
v.literal("superseded")
),
version: v.number(),
workId: v.id("works"),
}).index("by_work_and_version", ["workId", "version"]),
workSlices: defineTable({
createdAt: v.optional(v.number()),
designVersion: v.number(),
objective: v.string(),
observableBehavior: v.string(),
ordinal: v.number(),
payloadJson: v.string(),
sliceId: v.string(),
status: v.union(
v.literal("planned"),
v.literal("ready"),
v.literal("running"),
v.literal("completed"),
v.literal("blocked")
),
title: v.string(),
workId: v.id("works"),
})
.index("by_workId_and_designVersion", ["workId", "designVersion"])
.index("by_workId_and_designVersion_and_sliceId", [
"workId",
"designVersion",
"sliceId",
]),
workRuns: defineTable({
createdAt: v.number(),
designVersion: v.optional(v.number()),
endedAt: v.optional(v.number()),
kitId: v.string(),
kitVersion: v.string(),
scenario: v.union(
v.literal("success"),
v.literal("transient-failure-then-success"),
v.literal("needs-input"),
v.literal("permanent-failure"),
v.literal("cancelled")
),
sliceId: v.optional(v.string()),
sliceRowId: v.optional(v.id("workSlices")),
startedAt: v.optional(v.number()),
status: v.union(
v.literal("ready"),
v.literal("running"),
v.literal("terminal"),
v.literal("cancelled")
),
terminalClassification: v.optional(attemptClassification),
terminalSummary: v.optional(v.string()),
workId: v.id("works"),
}).index("by_work_and_createdAt", ["workId", "createdAt"]),
workAttempts: defineTable({
classification: v.optional(attemptClassification),
endedAt: v.optional(v.number()),
leaseExpiresAt: v.optional(v.number()),
leaseOwner: v.optional(v.string()),
number: v.number(),
runId: v.id("workRuns"),
startedAt: v.optional(v.number()),
status: v.union(
v.literal("queued"),
v.literal("claimed"),
v.literal("running"),
v.literal("terminal")
),
summary: v.optional(v.string()),
workId: v.id("works"),
})
.index("by_runId_and_number", ["runId", "number"])
.index("by_status_and_leaseExpiresAt", ["status", "leaseExpiresAt"]),
workAttemptEvents: defineTable({
attemptId: v.id("workAttempts"),
kind: v.string(),
message: v.string(),
metadataJson: v.string(),
occurredAt: v.number(),
sequence: v.number(),
}).index("by_attempt_and_sequence", ["attemptId", "sequence"]),
resolverDecisions: defineTable({
attemptId: v.id("workAttempts"),
attemptNumber: v.number(),
classification: attemptClassification,
createdAt: v.number(),
decision: v.union(v.literal("retry"), v.literal("terminal")),
resultingWorkStatus: v.optional(workStatus),
runId: v.id("workRuns"),
summary: v.string(),
workId: v.id("works"),
})
.index("by_runId_and_attemptNumber", ["runId", "attemptNumber"])
.index("by_attemptId", ["attemptId"]),
workArtifacts: defineTable({
attemptId: v.optional(v.id("workAttempts")),
contentHash: v.optional(v.string()),
createdAt: v.number(),
designVersion: v.optional(v.number()),
environmentId: v.optional(v.string()),
idempotencyKey: v.string(),
kind: v.union(
v.literal("definition"),
v.literal("design"),
v.literal("diff"),
v.literal("test-report"),
v.literal("verification-report"),
v.literal("runtime-log"),
v.literal("screenshot"),
v.literal("video"),
v.literal("commit"),
v.literal("branch"),
v.literal("pull-request"),
v.literal("preview"),
v.literal("deployment"),
v.literal("other")
),
metadataJson: v.string(),
organizationId: v.id("organizations"),
producer: v.string(),
projectId: v.id("projects"),
provenanceJson: v.string(),
runId: v.optional(v.id("workRuns")),
sliceId: v.optional(v.string()),
sourceRevision: v.optional(v.string()),
title: v.string(),
uri: v.optional(v.string()),
verificationStatus: v.union(
v.literal("unverified"),
v.literal("verified"),
v.literal("rejected")
),
workId: v.id("works"),
})
.index("by_workId_and_createdAt", ["workId", "createdAt"])
.index("by_workId_and_idempotencyKey", ["workId", "idempotencyKey"])
.index("by_runId_and_createdAt", ["runId", "createdAt"]),
workDeliveries: defineTable({
artifactId: v.optional(v.id("workArtifacts")),
createdAt: v.number(),
externalId: v.string(),
idempotencyKey: v.string(),
kind: v.union(
v.literal("branch"),
v.literal("commit"),
v.literal("pull-request"),
v.literal("preview"),
v.literal("deployment")
),
metadataJson: v.string(),
organizationId: v.id("organizations"),
projectId: v.id("projects"),
provider: v.string(),
sourceRevision: v.string(),
status: v.union(
v.literal("recorded"),
v.literal("ready"),
v.literal("approved"),
v.literal("delivered"),
v.literal("failed"),
v.literal("cancelled")
),
target: v.string(),
updatedAt: v.number(),
url: v.optional(v.string()),
workId: v.id("works"),
})
.index("by_workId_and_createdAt", ["workId", "createdAt"])
.index("by_workId_and_idempotencyKey", ["workId", "idempotencyKey"]),
// -----------------------------------------------------------------
// Flue persistence stores (schema/format version 4).
// -----------------------------------------------------------------
@@ -526,38 +215,38 @@ export default defineSchema({
// 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(),
submissionId: v.string(),
})
.index("by_submissionId_and_attemptId", ["submissionId", "attemptId"])
.index("by_submissionId", ["submissionId"]),
// Conversation-stream metadata: one row per stream path.
flueConversationStreams: defineTable({
closed: v.boolean(),
createdAt: v.number(),
path: v.string(),
identityJson: v.string(),
incarnation: v.string(),
nextOffset: v.number(),
nextProducerSequence: v.number(),
path: v.string(),
producerEpoch: v.number(),
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({
appendedAt: v.number(),
attemptId: v.optional(v.string()),
path: v.string(),
producerEpoch: v.number(),
seq: v.number(),
producerId: v.string(),
producerEpoch: v.number(),
producerSequence: v.number(),
recordsJson: v.string(),
seq: v.number(),
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", [
@@ -569,41 +258,41 @@ export default defineSchema({
// Event-stream metadata: one row per stream path.
flueEventStreams: defineTable({
path: v.string(),
nextSeq: v.number(),
closed: v.boolean(),
createdAt: v.number(),
nextSeq: v.number(),
path: v.string(),
}).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({
appendedAt: v.number(),
dataJson: v.string(),
onceKey: v.optional(v.string()),
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({
durationMs: v.optional(v.number()),
endedAt: v.optional(v.string()),
errorJson: v.optional(v.string()),
inputJson: v.optional(v.string()),
isError: v.optional(v.boolean()),
resultJson: v.optional(v.string()),
runId: v.string(),
startedAt: 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()),
workflowName: 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"])
@@ -623,15 +312,15 @@ export default defineSchema({
// Immutable attachment bytes. Identity is (streamPath, attachmentId); reads
// are additionally scoped by conversationId.
flueAttachments: defineTable({
attachmentId: v.string(),
bytes: v.bytes(),
streamPath: v.string(),
conversationId: v.string(),
createdAt: v.number(),
digest: v.string(),
filename: v.optional(v.string()),
attachmentId: v.string(),
mimeType: v.string(),
size: v.number(),
streamPath: v.string(),
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"])

View File

@@ -82,16 +82,16 @@ export const listEvidence = query({
export const createSignal = mutation({
args: {
messageIds: v.array(v.string()),
organizationId: v.id("organizations"),
projectId: v.id("projects"),
messageIds: v.array(v.string()),
problemStatement: v.object({
constraints: v.array(v.string()),
desiredOutcome: v.string(),
summary: v.string(),
title: v.string(),
summary: v.string(),
desiredOutcome: v.string(),
constraints: v.array(v.string()),
}),
processedByAgentInstanceId: v.string(),
projectId: v.id("projects"),
token: v.string(),
},
handler: async (ctx, args) => {

View File

@@ -1,108 +0,0 @@
import { env } from "@code/env/convex";
import { convexTest } from "convex-test";
import { anyApi } from "convex/server";
import { describe, expect, test } from "vitest";
import schema from "./schema";
declare global {
interface ImportMeta {
readonly glob: (pattern: string) => Record<string, () => Promise<unknown>>;
}
}
const modules = import.meta.glob("./**/*.ts");
const api = anyApi;
describe("Work artifacts and delivery", () => {
test("records exact idempotent artifact and delivery metadata", async () => {
const t = convexTest({ modules, schema });
const workId = await t.run(async (ctx) => {
const organizationId = await ctx.db.insert("organizations", {
createdAt: 1,
createdBy: "user",
kind: "personal",
name: "Personal",
});
const projectId = await ctx.db.insert("projects", {
createdAt: 1,
name: "Zopu",
normalizedSourceUrl: "https://example.com/zopu",
organizationId,
repositoryPath: "puter/zopu",
sourceHost: "example.com",
sourceUrl: "https://example.com/zopu",
updatedAt: 1,
});
return await ctx.db.insert("works", {
createdAt: 1,
objective: "Persist evidence",
organizationId,
projectId,
status: "executing",
title: "Artifact proof",
updatedAt: 1,
});
});
const draft = {
idempotencyKey: "attempt-1:test-report",
kind: "test-report" as const,
metadataJson: "{}",
producer: "fake-harness",
provenanceJson: "{}",
sourceRevision: "abc123",
title: "Focused tests",
verificationStatus: "verified" as const,
};
const first = await t.mutation(api.workArtifacts.recordArtifact, {
draft,
token: env.FLUE_DB_TOKEN,
workId,
});
const replay = await t.mutation(api.workArtifacts.recordArtifact, {
draft,
token: env.FLUE_DB_TOKEN,
workId,
});
expect(first.created).toBe(true);
expect(replay).toEqual({ artifactId: first.artifactId, created: false });
const delivery = await t.mutation(api.workArtifacts.recordDelivery, {
artifactId: first.artifactId,
draft: {
externalId: "42",
idempotencyKey: "pr:42",
kind: "pull-request",
metadataJson: "{}",
provider: "gitea",
sourceRevision: "abc123",
status: "ready",
target: "main",
url: "https://git.example/pulls/42",
},
token: env.FLUE_DB_TOKEN,
workId,
});
expect(delivery.created).toBe(true);
const updated = await t.mutation(api.workArtifacts.updateDeliveryStatus, {
deliveryId: delivery.deliveryId,
metadataJson: '{"approvedBy":"user"}',
status: "approved",
token: env.FLUE_DB_TOKEN,
});
expect(updated).toEqual({ changed: true, status: "approved" });
const stored = await t.run(async (ctx) => ({
artifacts: await ctx.db.query("workArtifacts").collect(),
deliveries: await ctx.db.query("workDeliveries").collect(),
events: await ctx.db.query("workEvents").collect(),
}));
expect(stored.artifacts).toHaveLength(1);
expect(stored.deliveries).toHaveLength(1);
expect(stored.events.map((event) => event.kind)).toEqual([
"artifact.recorded",
"delivery.recorded",
"delivery.updated",
]);
});
});

View File

@@ -1,370 +0,0 @@
import { env } from "@code/env/convex";
import {
canTransitionWorkDelivery,
decodeWorkArtifactDraft,
decodeWorkDeliveryDraft,
} from "@code/primitives/work-artifact";
import { ConvexError, v } from "convex/values";
import { Effect } from "effect";
import type { Id } from "./_generated/dataModel";
import { mutation, query } from "./_generated/server";
import type { MutationCtx } from "./_generated/server";
import { requireProjectMember } from "./authz";
const artifactKind = v.union(
v.literal("definition"),
v.literal("design"),
v.literal("diff"),
v.literal("test-report"),
v.literal("verification-report"),
v.literal("runtime-log"),
v.literal("screenshot"),
v.literal("video"),
v.literal("commit"),
v.literal("branch"),
v.literal("pull-request"),
v.literal("preview"),
v.literal("deployment"),
v.literal("other")
);
const verificationStatus = v.union(
v.literal("unverified"),
v.literal("verified"),
v.literal("rejected")
);
const deliveryKind = v.union(
v.literal("branch"),
v.literal("commit"),
v.literal("pull-request"),
v.literal("preview"),
v.literal("deployment")
);
const deliveryStatus = v.union(
v.literal("recorded"),
v.literal("ready"),
v.literal("approved"),
v.literal("delivered"),
v.literal("failed"),
v.literal("cancelled")
);
const requireAgent = (token: string): void => {
if (token !== env.FLUE_DB_TOKEN) {
throw new ConvexError("Invalid agent control token");
}
};
const appendEvent = async (
ctx: MutationCtx,
workId: Id<"works">,
kind: "artifact.recorded" | "delivery.recorded" | "delivery.updated",
idempotencyKey: string,
referenceId: string,
payloadJson: string
): Promise<void> => {
const existing = await ctx.db
.query("workEvents")
.withIndex("by_work_and_idempotencyKey", (q) =>
q.eq("workId", workId).eq("idempotencyKey", idempotencyKey)
)
.unique();
if (existing) {
return;
}
await ctx.db.insert("workEvents", {
createdAt: Date.now(),
idempotencyKey,
kind,
payloadJson,
referenceId,
workId,
});
};
export const recordArtifact = mutation({
args: {
attemptId: v.optional(v.id("workAttempts")),
designVersion: v.optional(v.number()),
draft: v.object({
contentHash: v.optional(v.string()),
environmentId: v.optional(v.string()),
idempotencyKey: v.string(),
kind: artifactKind,
metadataJson: v.string(),
producer: v.string(),
provenanceJson: v.string(),
sourceRevision: v.optional(v.string()),
title: v.string(),
uri: v.optional(v.string()),
verificationStatus,
}),
runId: v.optional(v.id("workRuns")),
sliceId: v.optional(v.string()),
token: v.string(),
workId: v.id("works"),
},
handler: async (ctx, args) => {
requireAgent(args.token);
const work = await ctx.db.get(args.workId);
if (!work) {
throw new ConvexError("Work not found");
}
const draft = await Effect.runPromise(
decodeWorkArtifactDraft(args.draft)
).catch((error: unknown) => {
throw new ConvexError(
error instanceof Error ? error.message : "Invalid Work artifact"
);
});
const attempt = args.attemptId ? await ctx.db.get(args.attemptId) : null;
if (args.attemptId && !attempt) {
throw new ConvexError("Attempt not found");
}
const referencedRunId = args.runId ?? attempt?.runId;
const run = referencedRunId ? await ctx.db.get(referencedRunId) : null;
if (referencedRunId && !run) {
throw new ConvexError("Run not found");
}
if (run && run.workId !== work._id) {
throw new ConvexError("Run does not belong to Work");
}
if (
attempt &&
(attempt.workId !== work._id ||
(run !== null && attempt.runId !== run._id))
) {
throw new ConvexError("Attempt does not belong to Work Run");
}
if (
run &&
((args.designVersion !== undefined &&
args.designVersion !== run.designVersion) ||
(args.sliceId !== undefined && args.sliceId !== run.sliceId))
) {
throw new ConvexError("Artifact provenance conflicts with Work Run");
}
const designVersion = args.designVersion ?? run?.designVersion;
const sliceId = args.sliceId ?? run?.sliceId;
if (designVersion !== undefined && sliceId !== undefined) {
const slice = await ctx.db
.query("workSlices")
.withIndex("by_workId_and_designVersion_and_sliceId", (q) =>
q
.eq("workId", work._id)
.eq("designVersion", designVersion)
.eq("sliceId", sliceId)
)
.unique();
if (!slice) {
throw new ConvexError("Artifact slice not found");
}
} else if (designVersion !== undefined) {
const design = await ctx.db
.query("designPackets")
.withIndex("by_work_and_version", (q) =>
q.eq("workId", work._id).eq("version", designVersion)
)
.unique();
if (!design) {
throw new ConvexError("Artifact Design not found");
}
}
const existing = await ctx.db
.query("workArtifacts")
.withIndex("by_workId_and_idempotencyKey", (q) =>
q.eq("workId", work._id).eq("idempotencyKey", draft.idempotencyKey)
)
.unique();
if (existing) {
const exactReplay =
existing.kind === draft.kind &&
existing.title === draft.title &&
existing.uri === draft.uri &&
existing.contentHash === draft.contentHash &&
existing.sourceRevision === draft.sourceRevision &&
existing.environmentId === draft.environmentId &&
existing.producer === draft.producer &&
existing.provenanceJson === draft.provenanceJson &&
existing.metadataJson === draft.metadataJson &&
existing.verificationStatus === draft.verificationStatus &&
existing.designVersion === designVersion &&
existing.sliceId === sliceId &&
existing.runId === run?._id &&
existing.attemptId === attempt?._id;
if (!exactReplay) {
throw new ConvexError("Artifact idempotency key has conflicting data");
}
return { artifactId: existing._id, created: false };
}
const createdAt = Date.now();
const artifactId = await ctx.db.insert("workArtifacts", {
...draft,
attemptId: attempt?._id,
createdAt,
designVersion,
organizationId: work.organizationId,
projectId: work.projectId,
runId: run?._id,
sliceId,
workId: work._id,
});
await appendEvent(
ctx,
work._id,
"artifact.recorded",
`artifact:${draft.idempotencyKey}`,
String(artifactId),
JSON.stringify({ kind: draft.kind, title: draft.title })
);
return { artifactId, created: true };
},
});
export const recordDelivery = mutation({
args: {
artifactId: v.optional(v.id("workArtifacts")),
draft: v.object({
externalId: v.string(),
idempotencyKey: v.string(),
kind: deliveryKind,
metadataJson: v.string(),
provider: v.string(),
sourceRevision: v.string(),
status: deliveryStatus,
target: v.string(),
url: v.optional(v.string()),
}),
token: v.string(),
workId: v.id("works"),
},
handler: async (ctx, args) => {
requireAgent(args.token);
const work = await ctx.db.get(args.workId);
if (!work) {
throw new ConvexError("Work not found");
}
const draft = await Effect.runPromise(
decodeWorkDeliveryDraft(args.draft)
).catch((error: unknown) => {
throw new ConvexError(
error instanceof Error ? error.message : "Invalid Work delivery"
);
});
const artifact = args.artifactId ? await ctx.db.get(args.artifactId) : null;
if (args.artifactId && !artifact) {
throw new ConvexError("Artifact not found");
}
if (artifact && artifact.workId !== work._id) {
throw new ConvexError("Artifact does not belong to Work");
}
const existing = await ctx.db
.query("workDeliveries")
.withIndex("by_workId_and_idempotencyKey", (q) =>
q.eq("workId", work._id).eq("idempotencyKey", draft.idempotencyKey)
)
.unique();
if (existing) {
const exactReplay =
existing.kind === draft.kind &&
existing.provider === draft.provider &&
existing.externalId === draft.externalId &&
existing.url === draft.url &&
existing.sourceRevision === draft.sourceRevision &&
existing.target === draft.target &&
existing.status === draft.status &&
existing.metadataJson === draft.metadataJson &&
existing.artifactId === artifact?._id;
if (!exactReplay) {
throw new ConvexError("Delivery idempotency key has conflicting data");
}
return { created: false, deliveryId: existing._id };
}
const timestamp = Date.now();
const deliveryId = await ctx.db.insert("workDeliveries", {
...draft,
artifactId: artifact?._id,
createdAt: timestamp,
organizationId: work.organizationId,
projectId: work.projectId,
updatedAt: timestamp,
workId: work._id,
});
await appendEvent(
ctx,
work._id,
"delivery.recorded",
`delivery:${draft.idempotencyKey}`,
String(deliveryId),
JSON.stringify({ kind: draft.kind, sourceRevision: draft.sourceRevision })
);
return { created: true, deliveryId };
},
});
export const updateDeliveryStatus = mutation({
args: {
deliveryId: v.id("workDeliveries"),
metadataJson: v.string(),
status: deliveryStatus,
token: v.string(),
},
handler: async (ctx, args) => {
requireAgent(args.token);
// Validate JSON before persisting, matching recordDelivery's decode path.
const metadata = JSON.parse(args.metadataJson) as unknown;
if (typeof metadata !== "object" || metadata === null) {
throw new ConvexError("Delivery metadata must be a JSON object");
}
const delivery = await ctx.db.get(args.deliveryId);
if (!delivery) {
throw new ConvexError("Delivery not found");
}
if (delivery.status === args.status) {
return { changed: false, status: delivery.status };
}
if (!canTransitionWorkDelivery(delivery.status, args.status)) {
throw new ConvexError(
`Invalid delivery transition: ${delivery.status} -> ${args.status}`
);
}
await ctx.db.patch(delivery._id, {
metadataJson: args.metadataJson,
status: args.status,
updatedAt: Date.now(),
});
await appendEvent(
ctx,
delivery.workId,
"delivery.updated",
`delivery-updated:${delivery._id}:${args.status}`,
String(delivery._id),
JSON.stringify({ from: delivery.status, to: args.status })
);
return { changed: true, status: args.status };
},
});
export const listForWork = query({
args: { workId: v.id("works") },
handler: async (ctx, args) => {
const work = await ctx.db.get(args.workId);
if (!work) {
return null;
}
await requireProjectMember(ctx, work.projectId);
const [artifacts, deliveries] = await Promise.all([
ctx.db
.query("workArtifacts")
.withIndex("by_workId_and_createdAt", (q) => q.eq("workId", work._id))
.order("desc")
.take(100),
ctx.db
.query("workDeliveries")
.withIndex("by_workId_and_createdAt", (q) => q.eq("workId", work._id))
.order("desc")
.take(50),
]);
return { artifacts, deliveries };
},
});

View File

@@ -1,373 +0,0 @@
import { convexTest } from "convex-test";
import { anyApi, makeFunctionReference } from "convex/server";
import { describe, expect, test } from "vitest";
import type { Id } from "./_generated/dataModel";
import schema from "./schema";
// Execution-only module set: avoids loading workPlanning (and its env read) so
// these tests run without Convex env vars.
const modules = {
"./_generated/api.ts": () => import("./_generated/api"),
"./_generated/server.ts": () => import("./_generated/server"),
"./authz.ts": () => import("./authz"),
"./workExecution.ts": () => import("./workExecution"),
};
const identity = { tokenIdentifier: "https://convex.test|exec-user" };
const api = anyApi;
const claimAttemptRef = makeFunctionReference<
"mutation",
{ attemptId: Id<"workAttempts">; leaseMs: number; owner: string },
{ number: number } | null
>("workExecution:claimAttempt");
const makeTest = () => convexTest({ modules, schema });
type TestT = ReturnType<typeof makeTest>;
type Scenario =
| "success"
| "transient-failure-then-success"
| "needs-input"
| "permanent-failure"
| "cancelled";
const seedReadyWork = async (options: { secondSlice?: boolean } = {}) => {
const t = convexTest({ modules, schema });
const { workId } = await t.withIdentity(identity).run(async (ctx) => {
const organizationId = await ctx.db.insert("organizations", {
createdAt: 1,
createdBy: identity.tokenIdentifier,
kind: "personal",
name: "Personal",
});
await ctx.db.insert("organizationMembers", {
createdAt: 1,
organizationId,
role: "owner",
userId: identity.tokenIdentifier,
});
const projectId = await ctx.db.insert("projects", {
createdAt: 1,
name: "Zopu",
normalizedSourceUrl: "https://example.com/zopu",
organizationId,
repositoryPath: "puter/zopu",
sourceHost: "example.com",
sourceUrl: "https://example.com/zopu",
updatedAt: 1,
});
const id = await ctx.db.insert("works", {
createdAt: 1,
definitionApprovalVersion: 1,
definitionVersion: 1,
designApprovalVersion: 1,
designVersion: 1,
objective: "Prove execution",
organizationId,
projectId,
status: "ready",
title: "Executable work",
updatedAt: 1,
});
await ctx.db.insert("designPackets", {
createdAt: 1,
createdBy: "test",
definitionVersion: 1,
payloadJson: "{}",
status: "current",
version: 1,
workId: id,
});
await ctx.db.insert("workSlices", {
createdAt: 1,
designVersion: 1,
objective: "execute the slice",
observableBehavior: "terminal run",
ordinal: 0,
payloadJson: "{}",
sliceId: "slice-1",
status: "ready",
title: "first",
workId: id,
});
if (options.secondSlice) {
await ctx.db.insert("workSlices", {
createdAt: 1,
designVersion: 1,
objective: "execute the second slice",
observableBehavior: "second terminal run",
ordinal: 1,
payloadJson: "{}",
sliceId: "slice-2",
status: "planned",
title: "second",
workId: id,
});
}
return { workId: id };
});
return { t, workId };
};
const runAttempt = async (
t: TestT,
attemptId: Id<"workAttempts">,
scenario: Scenario
) => t.action(api.workExecution.executeFakeAttempt, { attemptId, scenario });
const attemptNumber = async (
t: TestT,
runId: Id<"workRuns">,
number: number
): Promise<Id<"workAttempts">> => {
const attempt = await t.run(async (ctx) =>
ctx.db
.query("workAttempts")
.withIndex("by_runId_and_number", (q) =>
q.eq("runId", runId).eq("number", number)
)
.unique()
);
if (!attempt) {
throw new Error(`attempt ${number} not found`);
}
return attempt._id;
};
const snapshot = async (t: TestT, workId: Id<"works">, runId: Id<"workRuns">) =>
t.run(async (ctx) => ({
attempts: await ctx.db
.query("workAttempts")
.withIndex("by_runId_and_number", (q) => q.eq("runId", runId))
.collect(),
run: await ctx.db.get(runId),
slice: await ctx.db
.query("workSlices")
.withIndex("by_workId_and_designVersion", (q) =>
q.eq("workId", workId).eq("designVersion", 1)
)
.first(),
slices: await ctx.db
.query("workSlices")
.withIndex("by_workId_and_designVersion", (q) =>
q.eq("workId", workId).eq("designVersion", 1)
)
.collect(),
work: await ctx.db.get(workId),
}));
describe("simulated execution resolver", () => {
test("success settles Work as completed and marks the slice completed", async () => {
const { t, workId } = await seedReadyWork();
const started = await t
.withIdentity(identity)
.mutation(api.workExecution.startSimulatedExecution, {
scenario: "success",
sliceId: "slice-1",
workId,
});
await runAttempt(t, started.attemptId, "success");
const state = await snapshot(t, workId, started.runId);
expect(state.work?.status).toBe("completed");
expect(state.run?.terminalClassification).toBe("Succeeded");
expect(state.attempts).toHaveLength(1);
expect(state.slice?.status).toBe("completed");
const decisions = await t.run(async (ctx) =>
ctx.db.query("resolverDecisions").collect()
);
expect(decisions).toHaveLength(1);
expect(decisions[0]).toMatchObject({
classification: "Succeeded",
decision: "terminal",
resultingWorkStatus: "completed",
});
});
test("completes Work only after every approved slice succeeds", async () => {
const { t, workId } = await seedReadyWork({ secondSlice: true });
const first = await t
.withIdentity(identity)
.mutation(api.workExecution.startSimulatedExecution, {
scenario: "success",
workId,
});
await runAttempt(t, first.attemptId, "success");
const afterFirst = await snapshot(t, workId, first.runId);
expect(afterFirst.work?.status).toBe("ready");
expect(
afterFirst.slices.find((slice) => slice.sliceId === "slice-1")?.status
).toBe("completed");
expect(
afterFirst.slices.find((slice) => slice.sliceId === "slice-2")?.status
).toBe("ready");
await expect(
t
.withIdentity(identity)
.mutation(api.workExecution.retrySimulatedExecution, {
runId: first.runId,
})
).rejects.toThrow(/retryable state/u);
const second = await t
.withIdentity(identity)
.mutation(api.workExecution.startSimulatedExecution, {
scenario: "success",
sliceId: "slice-2",
workId,
});
await runAttempt(t, second.attemptId, "success");
const completed = await snapshot(t, workId, second.runId);
expect(completed.work?.status).toBe("completed");
expect(
completed.slices.every((slice) => slice.status === "completed")
).toBe(true);
});
test("allows only one owner to claim a queued attempt", async () => {
const { t, workId } = await seedReadyWork();
const started = await t
.withIdentity(identity)
.mutation(api.workExecution.startSimulatedExecution, {
scenario: "success",
workId,
});
const first = await t.mutation(claimAttemptRef, {
attemptId: started.attemptId,
leaseMs: 60_000,
owner: "worker-a",
});
const second = await t.mutation(claimAttemptRef, {
attemptId: started.attemptId,
leaseMs: 60_000,
owner: "worker-b",
});
expect(first?.number).toBe(1);
expect(second).toBeNull();
});
test("transient failure retries within the kit budget and then succeeds", async () => {
const { t, workId } = await seedReadyWork();
const started = await t
.withIdentity(identity)
.mutation(api.workExecution.startSimulatedExecution, {
scenario: "transient-failure-then-success",
workId,
});
await runAttempt(t, started.attemptId, "transient-failure-then-success");
// resolver auto-scheduled attempt 2; drive it manually in the test runtime
const second = await attemptNumber(t, started.runId, 2);
await runAttempt(t, second, "transient-failure-then-success");
const state = await snapshot(t, workId, started.runId);
expect(state.work?.status).toBe("completed");
expect(state.run?.terminalClassification).toBe("Succeeded");
expect(state.attempts).toHaveLength(2);
expect(state.attempts[0]?.classification).toBe("RetryableFailure");
});
test("permanent failure settles Work as failed, not silently runnable", async () => {
const { t, workId } = await seedReadyWork();
const started = await t
.withIdentity(identity)
.mutation(api.workExecution.startSimulatedExecution, {
scenario: "permanent-failure",
workId,
});
await runAttempt(t, started.attemptId, "permanent-failure");
const state = await snapshot(t, workId, started.runId);
expect(state.work?.status).toBe("failed");
expect(state.run?.terminalClassification).toBe("PermanentFailure");
});
test("needs-input surfaces a blocker instead of looping", async () => {
const { t, workId } = await seedReadyWork();
const started = await t
.withIdentity(identity)
.mutation(api.workExecution.startSimulatedExecution, {
scenario: "needs-input",
workId,
});
await runAttempt(t, started.attemptId, "needs-input");
const state = await snapshot(t, workId, started.runId);
expect(state.work?.status).toBe("needs-input");
expect(state.run?.terminalClassification).toBe("NeedsInput");
});
test("cancellation terminates the attempt and returns Work to ready", async () => {
const { t, workId } = await seedReadyWork();
const started = await t
.withIdentity(identity)
.mutation(api.workExecution.startSimulatedExecution, {
scenario: "success",
workId,
});
await t
.withIdentity(identity)
.mutation(api.workExecution.cancelSimulatedExecution, {
runId: started.runId,
});
const state = await snapshot(t, workId, started.runId);
expect(state.work?.status).toBe("ready");
expect(state.run?.status).toBe("cancelled");
expect(state.slice?.status).toBe("ready");
});
test("a slice that is not part of the current approved Design is rejected", async () => {
const { t, workId } = await seedReadyWork();
await expect(
t
.withIdentity(identity)
.mutation(api.workExecution.startSimulatedExecution, {
scenario: "success",
sliceId: "not-a-real-slice",
workId,
})
).rejects.toThrow(/current approved Design/u);
});
test("manual retry restarts a failed Run", async () => {
const { t, workId } = await seedReadyWork();
const started = await t
.withIdentity(identity)
.mutation(api.workExecution.startSimulatedExecution, {
scenario: "permanent-failure",
workId,
});
await runAttempt(t, started.attemptId, "permanent-failure");
expect((await snapshot(t, workId, started.runId)).work?.status).toBe(
"failed"
);
const retried = await t
.withIdentity(identity)
.mutation(api.workExecution.retrySimulatedExecution, {
runId: started.runId,
});
await runAttempt(t, retried.attemptId as Id<"workAttempts">, "success");
const state = await snapshot(t, workId, started.runId);
expect(state.work?.status).toBe("completed");
expect(retried.number).toBe(2);
});
test("reconciling an expired lease resumes within budget", async () => {
const { t, workId } = await seedReadyWork();
const started = await t
.withIdentity(identity)
.mutation(api.workExecution.startSimulatedExecution, {
scenario: "success",
workId,
});
// simulate a worker that claimed but died before finishing
await t.run(async (ctx) => {
await ctx.db.patch(started.attemptId, {
leaseExpiresAt: Date.now() - 1000,
leaseOwner: "dead-worker",
status: "claimed",
});
});
await t.mutation(api.workExecution.reconcileExpiredAttempts, {});
const after = await snapshot(t, workId, started.runId);
expect(after.attempts.find((a) => a.number === 1)?.status).toBe("terminal");
expect(
after.attempts.find((a) => a.number === 2 && a.status === "queued")
).toBeTruthy();
});
});

View File

@@ -1,772 +0,0 @@
import { FakeHarnessLive } from "@code/primitives/harness-runtime";
import type { FakeScenario } from "@code/primitives/harness-runtime";
import { defaultCodingKitV0, resolveOutcome } from "@code/primitives/resolver";
import type { WorkEventKind } from "@code/primitives/work";
import { makeFunctionReference } from "convex/server";
import { ConvexError, v } from "convex/values";
import { Effect } from "effect";
import type { Doc, Id } from "./_generated/dataModel";
import {
internalAction,
internalMutation,
mutation,
query,
} from "./_generated/server";
import type { MutationCtx } from "./_generated/server";
import { requireProjectMember } from "./authz";
const executeRef = makeFunctionReference<
"action",
{ attemptId: Id<"workAttempts">; scenario: FakeScenario }
>("workExecution:executeFakeAttempt");
const claimRef = makeFunctionReference<
"mutation",
{ attemptId: Id<"workAttempts">; leaseMs: number; owner: string },
{ number: number } | null
>("workExecution:claimAttempt");
const checkpointRef = makeFunctionReference<
"mutation",
{
attemptId: Id<"workAttempts">;
kind: string;
message: string;
metadataJson: string;
owner: string;
sequence: number;
},
unknown
>("workExecution:checkpointAttempt");
const finishRef = makeFunctionReference<
"mutation",
{
attemptId: Id<"workAttempts">;
classification: string;
owner: string;
retryable: boolean;
summary: string;
},
unknown
>("workExecution:finishAttempt");
const now = () => Date.now();
const attemptClassification = v.union(
v.literal("Succeeded"),
v.literal("RetryableFailure"),
v.literal("NeedsInput"),
v.literal("Blocked"),
v.literal("VerificationFailed"),
v.literal("BudgetExhausted"),
v.literal("Cancelled"),
v.literal("PermanentFailure")
);
const requireWorkForMember = async (
ctx: MutationCtx,
workId: Id<"works">
): Promise<Doc<"works">> => {
const work = await ctx.db.get(workId);
if (!work) {
throw new ConvexError("Work not found");
}
await requireProjectMember(ctx, work.projectId);
return work;
};
const appendWorkEvent = async (
ctx: MutationCtx,
workId: Id<"works">,
kind: WorkEventKind,
idempotencyKey: string,
referenceId?: string,
payloadJson?: string
) => {
const existing = await ctx.db
.query("workEvents")
.withIndex("by_work_and_idempotencyKey", (q) =>
q.eq("workId", workId).eq("idempotencyKey", idempotencyKey)
)
.unique();
if (existing) {
return;
}
await ctx.db.insert("workEvents", {
createdAt: now(),
idempotencyKey,
kind,
workId,
...(referenceId ? { referenceId } : {}),
...(payloadJson ? { payloadJson } : {}),
});
};
// Slice the Run targets for the currently approved Design. `sliceId` is
// optional only to let the resolver default to the first ready slice.
const resolveRunSlice = async (
ctx: MutationCtx,
work: Doc<"works">,
sliceId?: string
): Promise<Doc<"workSlices">> => {
const currentDesignVersion = work.designVersion;
if (currentDesignVersion === undefined) {
throw new ConvexError("Work has no approved Design to execute");
}
const slices = await ctx.db
.query("workSlices")
.withIndex("by_workId_and_designVersion", (q) =>
q.eq("workId", work._id).eq("designVersion", currentDesignVersion)
)
.collect();
if (slices.length === 0) {
throw new ConvexError("No slices found for the current approved Design");
}
const slice = sliceId
? slices.find((row) => row.sliceId === sliceId)
: slices
.sort((a, b) => a.ordinal - b.ordinal)
.find((row) => row.status === "ready");
if (!slice) {
throw new ConvexError(
"Requested slice does not belong to the current approved Design"
);
}
if (slice.status !== "ready") {
throw new ConvexError("Only the next ready slice can be executed");
}
return slice;
};
const getRunSlice = async (
ctx: MutationCtx,
run: Doc<"workRuns">
): Promise<Doc<"workSlices"> | null> => {
if (run.sliceRowId) {
return await ctx.db.get(run.sliceRowId);
}
if (run.designVersion === undefined || run.sliceId === undefined) {
return null;
}
return await ctx.db
.query("workSlices")
.withIndex("by_workId_and_designVersion_and_sliceId", (q) =>
q
.eq("workId", run.workId)
.eq("designVersion", run.designVersion!)
.eq("sliceId", run.sliceId!)
)
.unique();
};
const RETRYABLE_WORK_STATUSES = new Set(["failed", "needs-input", "blocked"]);
export const startSimulatedExecution = mutation({
args: {
scenario: v.union(
v.literal("success"),
v.literal("transient-failure-then-success"),
v.literal("needs-input"),
v.literal("permanent-failure"),
v.literal("cancelled")
),
sliceId: v.optional(v.string()),
workId: v.id("works"),
},
handler: async (ctx, args) => {
const work = await requireWorkForMember(ctx, args.workId);
if (work.status !== "ready") {
throw new ConvexError("Work must be Ready before simulated execution");
}
if (
work.definitionApprovalVersion !== work.definitionVersion ||
work.designApprovalVersion !== work.designVersion
) {
throw new ConvexError(
"Execution requires exact approved Definition and Design versions"
);
}
const slice = await resolveRunSlice(ctx, work, args.sliceId);
const createdAt = now();
const runId = await ctx.db.insert("workRuns", {
createdAt,
designVersion: slice.designVersion,
kitId: defaultCodingKitV0.id,
kitVersion: defaultCodingKitV0.version,
scenario: args.scenario,
sliceId: slice.sliceId,
sliceRowId: slice._id,
status: "ready",
workId: work._id,
});
const attemptId = await ctx.db.insert("workAttempts", {
number: 1,
runId,
status: "queued",
workId: work._id,
});
await ctx.db.patch(work._id, {
status: "executing",
updatedAt: createdAt,
});
await ctx.db.patch(slice._id, { status: "running" });
await ctx.db.patch(runId, { startedAt: createdAt, status: "running" });
await appendWorkEvent(
ctx,
work._id,
"run.started",
`run-started:${runId}`,
String(runId)
);
await appendWorkEvent(
ctx,
work._id,
"slice.started",
`slice-started:${runId}:${slice.sliceId}`,
String(slice._id)
);
await ctx.scheduler.runAfter(0, executeRef, {
attemptId,
scenario: args.scenario,
});
return { attemptId, runId };
},
});
export const cancelSimulatedExecution = mutation({
args: { runId: v.id("workRuns") },
handler: async (ctx, args) => {
const run = await ctx.db.get(args.runId);
if (!run) {
throw new ConvexError("Run not found");
}
const work = await requireWorkForMember(ctx, run.workId);
if (run.status === "terminal" || run.status === "cancelled") {
return { cancelled: false };
}
await ctx.db.patch(run._id, {
endedAt: now(),
status: "cancelled",
terminalClassification: "Cancelled",
terminalSummary: "Simulation cancelled",
});
const attempts = await ctx.db
.query("workAttempts")
.withIndex("by_runId_and_number", (q) => q.eq("runId", run._id))
.collect();
for (const attempt of attempts) {
if (attempt.status !== "terminal") {
await ctx.db.patch(attempt._id, {
classification: "Cancelled",
endedAt: now(),
status: "terminal",
summary: "Simulation cancelled",
});
await ctx.db.insert("resolverDecisions", {
attemptId: attempt._id,
attemptNumber: attempt.number,
classification: "Cancelled",
createdAt: now(),
decision: "terminal",
resultingWorkStatus: "ready",
runId: run._id,
summary: "Simulation cancelled",
workId: work._id,
});
await appendWorkEvent(
ctx,
work._id,
"resolver.decided",
`resolver-decided:${attempt._id}`,
String(attempt._id),
JSON.stringify({ classification: "Cancelled", decision: "terminal" })
);
}
}
const slice = await getRunSlice(ctx, run);
if (slice) {
await ctx.db.patch(slice._id, { status: "ready" });
}
await ctx.db.patch(work._id, { status: "ready", updatedAt: now() });
await appendWorkEvent(
ctx,
work._id,
"run.cancelled",
`run-cancelled:${run._id}`,
String(run._id)
);
return { cancelled: true };
},
});
// User-initiated restart of a terminal Run. Unlike the automatic resolver
// retry, this is an explicit decision and may exceed the kit budget.
export const retrySimulatedExecution = mutation({
args: { runId: v.id("workRuns") },
handler: async (ctx, args) => {
const run = await ctx.db.get(args.runId);
if (!run) {
throw new ConvexError("Run not found");
}
const work = await requireWorkForMember(ctx, run.workId);
if (run.status !== "terminal") {
throw new ConvexError("Only terminal Runs can be retried");
}
if (!RETRYABLE_WORK_STATUSES.has(work.status)) {
throw new ConvexError("Work is not in a retryable state");
}
const attempts = await ctx.db
.query("workAttempts")
.withIndex("by_runId_and_number", (q) => q.eq("runId", run._id))
.collect();
const number = attempts.length + 1;
const attemptId = await ctx.db.insert("workAttempts", {
number,
runId: run._id,
status: "queued",
workId: work._id,
});
await ctx.db.patch(run._id, {
endedAt: undefined,
startedAt: now(),
status: "running",
terminalClassification: undefined,
terminalSummary: undefined,
});
await ctx.db.patch(work._id, { status: "executing", updatedAt: now() });
const slice = await getRunSlice(ctx, run);
if (!slice) {
throw new ConvexError("Run slice not found");
}
await ctx.db.patch(slice._id, { status: "running" });
await ctx.scheduler.runAfter(0, executeRef, {
attemptId,
scenario: run.scenario,
});
return { attemptId, number };
},
});
export const claimAttempt = internalMutation({
args: {
attemptId: v.id("workAttempts"),
leaseMs: v.number(),
owner: v.string(),
},
handler: async (ctx, args) => {
const attempt = await ctx.db.get(args.attemptId);
if (!attempt || attempt.status !== "queued") {
return null;
}
const run = await ctx.db.get(attempt.runId);
if (!run || run.status !== "running") {
return null;
}
const claimedAt = now();
await ctx.db.patch(attempt._id, {
leaseExpiresAt: claimedAt + args.leaseMs,
leaseOwner: args.owner,
startedAt: attempt.startedAt ?? claimedAt,
status: "claimed",
});
await ctx.db.patch(attempt.runId, {
startedAt: claimedAt,
status: "running",
});
return { ...attempt, leaseOwner: args.owner, status: "claimed" as const };
},
});
export const checkpointAttempt = internalMutation({
args: {
attemptId: v.id("workAttempts"),
kind: v.string(),
message: v.string(),
metadataJson: v.string(),
owner: v.string(),
sequence: v.number(),
},
handler: async (ctx, args) => {
const attempt = await ctx.db.get(args.attemptId);
if (
!attempt ||
attempt.leaseOwner !== args.owner ||
attempt.status === "terminal" ||
(attempt.leaseExpiresAt !== undefined && attempt.leaseExpiresAt < now())
) {
return false;
}
const existing = await ctx.db
.query("workAttemptEvents")
.withIndex("by_attempt_and_sequence", (q: any) =>
q.eq("attemptId", attempt._id).eq("sequence", args.sequence)
)
.unique();
if (!existing) {
await ctx.db.insert("workAttemptEvents", {
attemptId: attempt._id,
kind: args.kind,
message: args.message,
metadataJson: args.metadataJson,
occurredAt: now(),
sequence: args.sequence,
});
}
await ctx.db.patch(attempt._id, {
leaseExpiresAt: now() + 60_000,
status: "running",
});
return true;
},
});
// Single resolver mutation: record the attempt outcome, then either schedule
// the next attempt within the kit retry policy or settle the Run/Work at the
// terminal classification. Replaces the old completeAttempt that always reset
// Work to "ready" and ignored the retry policy.
export const finishAttempt = internalMutation({
args: {
attemptId: v.id("workAttempts"),
classification: attemptClassification,
owner: v.string(),
retryable: v.boolean(),
summary: v.string(),
},
handler: async (ctx, args) => {
const attempt = await ctx.db.get(args.attemptId);
if (
!attempt ||
attempt.leaseOwner !== args.owner ||
attempt.status === "terminal" ||
(attempt.leaseExpiresAt !== undefined && attempt.leaseExpiresAt < now())
) {
return null;
}
const endedAt = now();
await ctx.db.patch(attempt._id, {
classification: args.classification,
endedAt,
leaseExpiresAt: undefined,
status: "terminal",
summary: args.summary,
});
const outcome = {
classification: args.classification,
retryable: args.retryable,
summary: args.summary,
};
const resolution = resolveOutcome(
outcome,
attempt.number,
defaultCodingKitV0.retryPolicy
);
const run = await ctx.db.get(attempt.runId);
if (!run) {
return null;
}
await appendWorkEvent(
ctx,
attempt.workId,
"attempt.completed",
`attempt-completed:${attempt._id}`,
String(attempt._id),
JSON.stringify({
classification: args.classification,
retried: resolution.kind === "retry",
summary: args.summary,
})
);
await appendWorkEvent(
ctx,
attempt.workId,
"resolver.decided",
`resolver-decided:${attempt._id}`,
String(attempt._id),
JSON.stringify({
classification: args.classification,
decision: resolution.kind,
})
);
if (resolution.kind === "retry") {
await ctx.db.insert("resolverDecisions", {
attemptId: attempt._id,
attemptNumber: attempt.number,
classification: args.classification,
createdAt: endedAt,
decision: "retry",
runId: run._id,
summary: args.summary,
workId: attempt.workId,
});
// Keep Run running and Work executing; spin the next attempt.
const nextAttemptId = await ctx.db.insert("workAttempts", {
number: attempt.number + 1,
runId: run._id,
status: "queued",
workId: attempt.workId,
});
await ctx.scheduler.runAfter(0, executeRef, {
attemptId: nextAttemptId,
scenario: run.scenario,
});
return { nextAttemptId, retried: true };
}
await ctx.db.patch(run._id, {
endedAt,
status: "terminal",
terminalClassification: args.classification,
terminalSummary: args.summary,
});
const work = await ctx.db.get(run.workId);
const slice = await getRunSlice(ctx, run);
let resultingWorkStatus = resolution.workStatus;
if (args.classification === "Succeeded" && slice) {
await ctx.db.patch(slice._id, { status: "completed" });
await appendWorkEvent(
ctx,
attempt.workId,
"slice.completed",
`slice-completed:${run._id}:${slice.sliceId}`,
String(slice._id)
);
const slices = await ctx.db
.query("workSlices")
.withIndex("by_workId_and_designVersion", (q) =>
q
.eq("workId", attempt.workId)
.eq("designVersion", slice.designVersion)
)
.collect();
const nextSlice = slices
.sort((left, right) => left.ordinal - right.ordinal)
.find((candidate) => candidate.status === "planned");
if (nextSlice) {
await ctx.db.patch(nextSlice._id, { status: "ready" });
resultingWorkStatus = "ready";
await appendWorkEvent(
ctx,
attempt.workId,
"slice.ready",
`slice-ready:${run._id}:${nextSlice.sliceId}`,
String(nextSlice._id)
);
}
} else if (slice) {
const blocked = ["NeedsInput", "Blocked", "VerificationFailed"].includes(
args.classification
);
await ctx.db.patch(slice._id, { status: blocked ? "blocked" : "ready" });
}
if (work) {
await ctx.db.patch(work._id, {
status: resultingWorkStatus,
updatedAt: endedAt,
});
}
await ctx.db.insert("resolverDecisions", {
attemptId: attempt._id,
attemptNumber: attempt.number,
classification: args.classification,
createdAt: endedAt,
decision: "terminal",
resultingWorkStatus,
runId: run._id,
summary: args.summary,
workId: attempt.workId,
});
await appendWorkEvent(
ctx,
attempt.workId,
"run.completed",
`run-completed:${run._id}`,
String(run._id),
JSON.stringify({
classification: args.classification,
workStatus: resultingWorkStatus,
})
);
return { retried: false };
},
});
export const executeFakeAttempt = internalAction({
args: {
attemptId: v.id("workAttempts"),
scenario: v.union(
v.literal("success"),
v.literal("transient-failure-then-success"),
v.literal("needs-input"),
v.literal("permanent-failure"),
v.literal("cancelled")
),
},
handler: async (ctx, args) => {
const owner = `fake-worker:${args.attemptId}`;
const claimed = await ctx.runMutation(claimRef, {
attemptId: args.attemptId,
leaseMs: 60_000,
owner,
});
if (!claimed) {
return null;
}
const [events, outcome] = await Effect.runPromise(
FakeHarnessLive().run({
attemptNumber: claimed.number,
scenario: args.scenario,
})
);
for (const item of events) {
await ctx.runMutation(checkpointRef, {
attemptId: args.attemptId,
kind: item.kind,
message: item.message,
metadataJson: JSON.stringify(item.metadata),
owner,
sequence: item.sequence,
});
}
await ctx.runMutation(finishRef, {
attemptId: args.attemptId,
classification: outcome.classification,
owner,
retryable: outcome.retryable,
summary: outcome.summary,
});
return null;
},
});
// Bounded recovery: a claimed/running attempt whose lease expired means the
// worker died mid-flight. Terminate that attempt, then either resume within
// the kit budget or settle the Run/Work so nothing stays "running" forever.
export const reconcileExpiredAttempts = internalMutation({
args: {},
handler: async (ctx) => {
const timestamp = now();
const [claimed, running] = await Promise.all([
ctx.db
.query("workAttempts")
.withIndex("by_status_and_leaseExpiresAt", (q) =>
q.eq("status", "claimed").lt("leaseExpiresAt", timestamp)
)
.collect(),
ctx.db
.query("workAttempts")
.withIndex("by_status_and_leaseExpiresAt", (q) =>
q.eq("status", "running").lt("leaseExpiresAt", timestamp)
)
.collect(),
]);
const active = [...claimed, ...running];
let reconciled = 0;
for (const attempt of active) {
if (
(attempt.status === "claimed" || attempt.status === "running") &&
attempt.leaseExpiresAt !== undefined &&
attempt.leaseExpiresAt < now()
) {
await ctx.db.patch(attempt._id, {
classification: "Blocked",
endedAt: now(),
leaseExpiresAt: undefined,
leaseOwner: undefined,
status: "terminal",
summary: "Attempt lease expired and was reconciled",
});
const run = await ctx.db.get(attempt.runId);
await appendWorkEvent(
ctx,
attempt.workId,
"attempt.reconciled",
`attempt-reconciled:${attempt._id}`,
String(attempt._id),
JSON.stringify({ attemptNumber: attempt.number })
);
reconciled += 1;
if (!run) {
continue;
}
const attempts = await ctx.db
.query("workAttempts")
.withIndex("by_runId_and_number", (q) => q.eq("runId", run._id))
.collect();
const withinBudget =
attempts.length < defaultCodingKitV0.retryPolicy.maxAttempts;
const retry = run.status === "running" && withinBudget;
await ctx.db.insert("resolverDecisions", {
attemptId: attempt._id,
attemptNumber: attempt.number,
classification: "Blocked",
createdAt: now(),
decision: retry ? "retry" : "terminal",
resultingWorkStatus: retry ? undefined : "blocked",
runId: run._id,
summary: "Attempt lease expired and was reconciled",
workId: attempt.workId,
});
await appendWorkEvent(
ctx,
attempt.workId,
"resolver.decided",
`resolver-decided:${attempt._id}`,
String(attempt._id),
JSON.stringify({
classification: "Blocked",
decision: retry ? "retry" : "terminal",
})
);
if (retry) {
const nextAttemptId = await ctx.db.insert("workAttempts", {
number: attempts.length + 1,
runId: run._id,
status: "queued",
workId: attempt.workId,
});
await ctx.scheduler.runAfter(0, executeRef, {
attemptId: nextAttemptId,
scenario: run.scenario,
});
} else {
await ctx.db.patch(run._id, {
endedAt: now(),
status: "terminal",
terminalClassification: "Blocked",
terminalSummary: "Run reconciled after expired lease",
});
const work = await ctx.db.get(run.workId);
if (work) {
await ctx.db.patch(work._id, {
status: "blocked",
updatedAt: now(),
});
}
const slice = await getRunSlice(ctx, run);
if (slice) {
await ctx.db.patch(slice._id, { status: "blocked" });
}
}
}
}
return { reconciled };
},
});
export const listRunEvents = query({
args: { attemptId: v.id("workAttempts") },
handler: async (ctx, args) => {
const attempt = await ctx.db.get(args.attemptId);
if (!attempt) {
return [];
}
const work = await ctx.db.get(attempt.workId);
if (!work) {
return [];
}
await requireProjectMember(ctx, work.projectId);
return await ctx.db
.query("workAttemptEvents")
.withIndex("by_attempt_and_sequence", (q) =>
q.eq("attemptId", args.attemptId)
)
.order("asc")
.collect();
},
});

View File

@@ -1,414 +0,0 @@
import { convexTest } from "convex-test";
import { anyApi, makeFunctionReference } from "convex/server";
import { describe, expect, test } from "vitest";
import type { Id } from "./_generated/dataModel";
import schema from "./schema";
declare global {
interface ImportMeta {
readonly glob: (pattern: string) => Record<string, () => Promise<unknown>>;
}
}
const modules = import.meta.glob("./**/*.ts");
const api = anyApi;
const executeFakeAttemptRef = makeFunctionReference<
"action",
{ attemptId: string; scenario: "success" }
>("workExecution:executeFakeAttempt");
const identity = { tokenIdentifier: "https://convex.test|planner-user" };
describe("Work planning and execution commands", () => {
test("binds approvals to exact versions and admits a fake Run only from Ready", async () => {
const t = convexTest({ modules, schema });
const fixture = await t.withIdentity(identity).run(async (ctx) => {
const organizationId = await ctx.db.insert("organizations", {
createdAt: 1,
createdBy: identity.tokenIdentifier,
kind: "personal",
name: "Personal",
});
await ctx.db.insert("organizationMembers", {
createdAt: 1,
organizationId,
role: "owner",
userId: identity.tokenIdentifier,
});
const projectId = await ctx.db.insert("projects", {
createdAt: 1,
name: "Zopu",
normalizedSourceUrl: "https://example.com/zopu",
organizationId,
repositoryPath: "puter/zopu",
sourceHost: "example.com",
sourceUrl: "https://example.com/zopu",
updatedAt: 1,
});
const workId = await ctx.db.insert("works", {
createdAt: 1,
objective: "Prove the flow",
organizationId,
projectId,
status: "proposed",
title: "Executable work",
updatedAt: 1,
});
return { projectId, workId };
});
await t
.withIdentity(identity)
.mutation(api.workPlanning.requestDefinition, { workId: fixture.workId });
const definition = await t
.withIdentity(identity)
.mutation(api.workPlanning.saveDefinitionProposal, {
payloadJson: JSON.stringify({
acceptanceCriteria: ["terminal"],
affectedUsers: ["user"],
assumptions: [],
constraints: [],
desiredOutcome: "A terminal simulation",
inScope: ["fake"],
outOfScope: [],
problem: "Need a proof",
questions: [],
requiredArtifacts: ["events"],
risk: "low",
}),
workId: fixture.workId,
});
await t
.withIdentity(identity)
.mutation(api.workPlanning.approveDefinition, {
version: definition.version,
workId: fixture.workId,
});
const design = await t
.withIdentity(identity)
.mutation(api.workPlanning.saveDesignProposal, {
payloadJson: JSON.stringify({
architectureSummary: "fake",
callFlowDelta: [],
concerns: [],
evidenceRequirements: ["event"],
fileTreeDelta: [],
impactMap: { files: [], modules: [], risks: [], summary: "small" },
invariants: ["terminal"],
keyInterfaces: [],
slices: [
{
codeBoundaries: ["runtime"],
dependsOn: [],
evidenceRequirements: ["event"],
id: "slice-1",
objective: "fake",
observableBehavior: "event",
reviewRequired: false,
title: "fake",
verification: ["assert"],
},
],
tradeoffs: [],
}),
workId: fixture.workId,
});
await t.withIdentity(identity).mutation(api.workPlanning.approveDesign, {
definitionVersion: definition.version,
designVersion: design.version,
workId: fixture.workId,
});
const started = await t
.withIdentity(identity)
.mutation(api.workExecution.startSimulatedExecution, {
scenario: "success",
sliceId: "slice-1",
workId: fixture.workId,
});
expect(started.runId).toBeTruthy();
const state = await t.run(async (ctx) => ctx.db.get(fixture.workId));
expect(state?.status).toBe("executing");
expect(state?.definitionApprovalVersion).toBe(definition.version);
expect(state?.designApprovalVersion).toBe(design.version);
await t.action(executeFakeAttemptRef, {
attemptId: String(started.attemptId),
scenario: "success",
});
const completed = await t.run(async (ctx) => ({
attempt: await ctx.db
.query("workAttempts")
.withIndex("by_runId_and_number", (q) => q.eq("runId", started.runId))
.unique(),
events: await ctx.db.query("workAttemptEvents").collect(),
run: await ctx.db.get(started.runId as Id<"workRuns">),
work: await ctx.db.get(fixture.workId),
}));
expect(completed.work?.status).toBe("completed");
expect(completed.run?.terminalClassification).toBe("Succeeded");
expect(completed.attempt?.status).toBe("terminal");
expect(completed.events.length).toBeGreaterThan(0);
});
});
const validDefinitionPayload = () => ({
acceptanceCriteria: ["terminal outcome"],
affectedUsers: ["user"],
assumptions: [],
constraints: [],
desiredOutcome: "A terminal simulation",
inScope: ["fake"],
outOfScope: [],
problem: "Need a proof",
questions: [],
requiredArtifacts: ["events"],
risk: "low",
});
const validDesignPayload = (sliceId = "slice-1") => ({
architectureSummary: "fake",
callFlowDelta: [],
concerns: [],
evidenceRequirements: ["event"],
fileTreeDelta: [],
impactMap: { files: [], modules: [], risks: [], summary: "small" },
invariants: ["terminal"],
keyInterfaces: [],
slices: [
{
codeBoundaries: ["runtime"],
dependsOn: [],
evidenceRequirements: ["event"],
id: sliceId,
objective: "fake",
observableBehavior: "event",
reviewRequired: false,
title: "fake",
verification: ["assert"],
},
],
tradeoffs: [],
});
describe("Work planning guards", () => {
test("persists high-impact questions but blocks approval", async () => {
const t = convexTest({ modules, schema });
const { workId } = await t.withIdentity(identity).run(async (ctx) => {
const organizationId = await ctx.db.insert("organizations", {
createdAt: 1,
createdBy: identity.tokenIdentifier,
kind: "personal",
name: "Personal",
});
await ctx.db.insert("organizationMembers", {
createdAt: 1,
organizationId,
role: "owner",
userId: identity.tokenIdentifier,
});
const projectId = await ctx.db.insert("projects", {
createdAt: 1,
name: "Zopu",
normalizedSourceUrl: "https://example.com/zopu",
organizationId,
repositoryPath: "puter/zopu",
sourceHost: "example.com",
sourceUrl: "https://example.com/zopu",
updatedAt: 1,
});
const id = await ctx.db.insert("works", {
createdAt: 1,
objective: "Resolve a critical question",
organizationId,
projectId,
status: "defining",
title: "Question gate",
updatedAt: 1,
});
return { workId: id };
});
const saved = await t
.withIdentity(identity)
.mutation(api.workPlanning.saveDefinitionProposal, {
payloadJson: JSON.stringify({
...validDefinitionPayload(),
questions: [
{
alternatives: ["public", "private"],
id: "visibility",
impact: "high",
prompt: "Should the endpoint be public?",
status: "open",
},
],
}),
workId,
});
await expect(
t.withIdentity(identity).mutation(api.workPlanning.approveDefinition, {
version: saved.version,
workId,
})
).rejects.toThrow(/High-impact open questions/u);
const stored = await t.run(async (ctx) =>
ctx.db.query("workQuestions").collect()
);
expect(stored).toHaveLength(1);
expect(stored[0]?.status).toBe("open");
});
test("keeps prior Design slices as versioned history", async () => {
const t = convexTest({ modules, schema });
const { workId } = await t.withIdentity(identity).run(async (ctx) => {
const organizationId = await ctx.db.insert("organizations", {
createdAt: 1,
createdBy: identity.tokenIdentifier,
kind: "personal",
name: "Personal",
});
await ctx.db.insert("organizationMembers", {
createdAt: 1,
organizationId,
role: "owner",
userId: identity.tokenIdentifier,
});
const projectId = await ctx.db.insert("projects", {
createdAt: 1,
name: "Zopu",
normalizedSourceUrl: "https://example.com/zopu",
organizationId,
repositoryPath: "puter/zopu",
sourceHost: "example.com",
sourceUrl: "https://example.com/zopu",
updatedAt: 1,
});
const id = await ctx.db.insert("works", {
createdAt: 1,
definitionApprovalVersion: 1,
definitionVersion: 1,
objective: "Preserve Design history",
organizationId,
projectId,
status: "designing",
title: "Versioned Design",
updatedAt: 1,
});
return { workId: id };
});
await t
.withIdentity(identity)
.mutation(api.workPlanning.saveDesignProposal, {
payloadJson: JSON.stringify(validDesignPayload("slice-v1")),
workId,
});
await t
.withIdentity(identity)
.mutation(api.workPlanning.saveDesignProposal, {
payloadJson: JSON.stringify(validDesignPayload("slice-v2")),
workId,
});
const slices = await t.run(async (ctx) =>
ctx.db.query("workSlices").collect()
);
expect(slices.map((slice) => slice.designVersion).sort()).toEqual([1, 2]);
});
test("a Definition cannot be revised while a Run is executing", async () => {
const t = convexTest({ modules, schema });
const { workId } = await t.withIdentity(identity).run(async (ctx) => {
const organizationId = await ctx.db.insert("organizations", {
createdAt: 1,
createdBy: identity.tokenIdentifier,
kind: "personal",
name: "Personal",
});
await ctx.db.insert("organizationMembers", {
createdAt: 1,
organizationId,
role: "owner",
userId: identity.tokenIdentifier,
});
const projectId = await ctx.db.insert("projects", {
createdAt: 1,
name: "Zopu",
normalizedSourceUrl: "https://example.com/zopu",
organizationId,
repositoryPath: "puter/zopu",
sourceHost: "example.com",
sourceUrl: "https://example.com/zopu",
updatedAt: 1,
});
const id = await ctx.db.insert("works", {
createdAt: 1,
definitionApprovalVersion: 1,
definitionVersion: 1,
designApprovalVersion: 1,
designVersion: 1,
objective: "Prove the guard",
organizationId,
projectId,
status: "executing",
title: "Guarded work",
updatedAt: 1,
});
return { workId: id };
});
await expect(
t
.withIdentity(identity)
.mutation(api.workPlanning.saveDefinitionProposal, {
payloadJson: JSON.stringify(validDefinitionPayload()),
workId,
})
).rejects.toThrow(/while a Run is executing/u);
});
test("planner questions are persisted as durable Work questions", async () => {
const t = convexTest({ modules, schema });
const { workId } = await t.withIdentity(identity).run(async (ctx) => {
const organizationId = await ctx.db.insert("organizations", {
createdAt: 1,
createdBy: identity.tokenIdentifier,
kind: "personal",
name: "Personal",
});
const projectId = await ctx.db.insert("projects", {
createdAt: 1,
name: "Zopu",
normalizedSourceUrl: "https://example.com/zopu",
organizationId,
repositoryPath: "puter/zopu",
sourceHost: "example.com",
sourceUrl: "https://example.com/zopu",
updatedAt: 1,
});
const id = await ctx.db.insert("works", {
createdAt: 1,
definitionVersion: 1,
objective: "Prove questions",
organizationId,
projectId,
status: "awaiting-definition-approval",
title: "Questionable work",
updatedAt: 1,
});
return { workId: id };
});
await t.mutation(api.workPlanning.submitQuestion, {
questionJson: JSON.stringify({
alternatives: [],
id: "q1",
impact: "medium",
prompt: "Which runtime?",
status: "open",
}),
token: "test-token",
workId,
});
const questions = await t.run(async (ctx) =>
ctx.db.query("workQuestions").collect()
);
expect(questions).toHaveLength(1);
expect(questions[0]?.questionId).toBe("q1");
expect(questions[0]?.definitionVersion).toBe(1);
});
});

View File

@@ -1,842 +0,0 @@
import type { WorkEventKind } from "@code/primitives/work";
import {
decodeDefinition,
validateDefinition,
WorkQuestion,
} from "@code/primitives/work-definition";
import type { WorkDefinition } from "@code/primitives/work-definition";
import { validateDesignPacket } from "@code/primitives/work-design";
import { makeFunctionReference } from "convex/server";
import { ConvexError, v } from "convex/values";
import { Effect, Schema } from "effect";
import type { Doc, Id } from "./_generated/dataModel";
import {
env,
internalAction,
internalMutation,
internalQuery,
mutation,
query,
} from "./_generated/server";
import type { MutationCtx, QueryCtx } from "./_generated/server";
import { requireProjectMember } from "./authz";
const plannerRef = makeFunctionReference<
"action",
{ organizationId: Id<"organizations">; workId: Id<"works"> }
>("workPlanning:runPlanner");
const plannerContextRef = makeFunctionReference<
"query",
{ workId: Id<"works"> },
{ objective: string; status: string; title: string } | null
>("workPlanning:getPlannerContext");
const plannerFailureRef = makeFunctionReference<
"mutation",
{ workId: Id<"works">; error: string },
null
>("workPlanning:recordPlannerFailure");
const parsePayload = (payloadJson: string): unknown => {
try {
return JSON.parse(payloadJson) as unknown;
} catch {
throw new ConvexError("Proposal payload must be valid JSON");
}
};
const objectPayload = (payloadJson: string): Record<string, unknown> => {
const value = parsePayload(payloadJson);
if (typeof value !== "object" || value === null || Array.isArray(value)) {
throw new ConvexError("Proposal payload must be an object");
}
return value as Record<string, unknown>;
};
const requireWork = async (
ctx: QueryCtx | MutationCtx,
workId: Id<"works">
): Promise<Doc<"works">> => {
const work = await ctx.db.get(workId);
if (!work) {
throw new ConvexError("Work not found");
}
return work;
};
const appendEvent = async (
ctx: MutationCtx,
workId: Id<"works">,
kind: WorkEventKind,
idempotencyKey: string,
referenceId?: string,
payloadJson?: string
) => {
const existing = await ctx.db
.query("workEvents")
.withIndex("by_work_and_idempotencyKey", (q) =>
q.eq("workId", workId).eq("idempotencyKey", idempotencyKey)
)
.unique();
if (existing) {
return existing._id;
}
return await ctx.db.insert("workEvents", {
createdAt: Date.now(),
idempotencyKey,
kind,
workId,
...(referenceId ? { referenceId } : {}),
...(payloadJson ? { payloadJson } : {}),
});
};
const invalidateApprovalsAndDesign = async (
ctx: MutationCtx,
workId: Id<"works">
) => {
const approvals = await ctx.db
.query("workApprovals")
.withIndex("by_workId_and_kind", (q) => q.eq("workId", workId))
.collect();
for (const approval of approvals) {
if (approval.status === "active") {
await ctx.db.patch(approval._id, { status: "invalidated" });
}
}
const designs = await ctx.db
.query("designPackets")
.withIndex("by_work_and_version", (q) => q.eq("workId", workId))
.collect();
for (const design of designs) {
if (design.status === "current") {
await ctx.db.patch(design._id, { status: "superseded" });
}
}
};
export const requestDefinition = mutation({
args: { workId: v.id("works") },
handler: async (ctx, args) => {
const work = await requireWork(ctx, args.workId);
await requireProjectMember(ctx, work.projectId);
if (
work.status !== "proposed" &&
work.status !== "defining" &&
!(
work.status === "blocked" &&
work.definitionApprovalVersion !== work.definitionVersion
)
) {
throw new ConvexError("Work is not available for definition");
}
const now = Date.now();
await ctx.db.patch(work._id, { status: "defining", updatedAt: now });
await appendEvent(
ctx,
work._id,
"definition.requested",
`definition-requested:${work._id}`
);
await ctx.scheduler.runAfter(0, plannerRef, {
organizationId: work.organizationId,
workId: work._id,
});
return { status: "defining" as const };
},
});
const saveDefinition = async (
ctx: MutationCtx,
work: Doc<"works">,
payloadJson: string,
createdBy: string
) => {
if (work.status === "executing") {
throw new ConvexError("Cannot revise Work while a Run is executing");
}
const decoded = await Effect.runPromise(
decodeDefinition({
...objectPayload(payloadJson),
version: (work.definitionVersion ?? 0) + 1,
})
).catch((error: unknown) => {
throw new ConvexError(
error instanceof Error ? error.message : "Invalid Definition"
);
});
const { version } = decoded;
await invalidateApprovalsAndDesign(ctx, work._id);
const previous = await ctx.db
.query("workDefinitions")
.withIndex("by_work_and_version", (q) => q.eq("workId", work._id))
.collect();
for (const row of previous) {
if (row.status === "current") {
await ctx.db.patch(row._id, { status: "superseded" });
}
}
const definitionId = await ctx.db.insert("workDefinitions", {
createdAt: Date.now(),
createdBy,
payloadJson: JSON.stringify(decoded),
risk: decoded.risk,
status: "current",
version,
workId: work._id,
});
for (const question of decoded.questions) {
await ctx.db.insert("workQuestions", {
alternativesJson: JSON.stringify(question.alternatives),
answer: question.answer,
createdAt: Date.now(),
definitionVersion: version,
impact: question.impact,
prompt: question.prompt,
questionId: question.id,
recommendation: question.recommendation,
status: question.status,
workId: work._id,
});
}
await ctx.db.patch(work._id, {
definitionApprovalVersion: undefined,
definitionVersion: version,
designApprovalVersion: undefined,
designVersion: undefined,
status: "awaiting-definition-approval",
updatedAt: Date.now(),
});
await appendEvent(
ctx,
work._id,
work.definitionVersion ? "definition.revised" : "definition.saved",
`definition:${version}`,
String(definitionId),
JSON.stringify(decoded)
);
return { definitionId, version };
};
export const saveDefinitionProposal = mutation({
args: { payloadJson: v.string(), workId: v.id("works") },
handler: async (ctx, args) => {
const work = await requireWork(ctx, args.workId);
await requireProjectMember(ctx, work.projectId);
return await saveDefinition(ctx, work, args.payloadJson, "user");
},
});
export const reviseDefinition = saveDefinitionProposal;
export const approveDefinition = mutation({
args: { version: v.number(), workId: v.id("works") },
handler: async (ctx, args) => {
const work = await requireWork(ctx, args.workId);
await requireProjectMember(ctx, work.projectId);
if (
work.definitionVersion !== args.version ||
work.status !== "awaiting-definition-approval"
) {
throw new ConvexError("Definition version is not current or approvable");
}
const row = await ctx.db
.query("workDefinitions")
.withIndex("by_work_and_version", (q) =>
q.eq("workId", work._id).eq("version", args.version)
)
.unique();
if (!row) {
throw new ConvexError("Definition not found");
}
const definition = JSON.parse(row.payloadJson) as WorkDefinition;
const valid = await Effect.runPromise(validateDefinition(definition)).catch(
(error: unknown) => {
throw new ConvexError(
error instanceof Error
? error.message
: "Definition cannot be approved"
);
}
);
const questions = await ctx.db
.query("workQuestions")
.withIndex("by_workId_and_definitionVersion", (q) =>
q.eq("workId", work._id).eq("definitionVersion", args.version)
)
.collect();
if (
questions.some(
(question) => question.status === "open" && question.impact === "high"
)
) {
throw new ConvexError(
"High-impact open questions must be resolved before approval"
);
}
const identity = await ctx.auth.getUserIdentity();
if (!identity) {
throw new ConvexError("Authentication required");
}
await ctx.db.insert("workApprovals", {
approvedAt: Date.now(),
approvedBy: identity.tokenIdentifier,
definitionVersion: args.version,
kind: "definition",
status: "active",
workId: work._id,
});
await ctx.db.patch(work._id, {
definitionApprovalVersion: valid.version,
status: "designing",
updatedAt: Date.now(),
});
await appendEvent(
ctx,
work._id,
"definition.approved",
`definition-approved:${args.version}`
);
await ctx.scheduler.runAfter(0, plannerRef, {
organizationId: work.organizationId,
workId: work._id,
});
return { status: "designing" as const, version: args.version };
},
});
const reviseQuestion = async (
ctx: MutationCtx,
work: Doc<"works">,
questionId: string,
update: { status: "answered" | "withdrawn"; answer?: string }
) => {
const current = await ctx.db
.query("workDefinitions")
.withIndex("by_work_and_version", (q) =>
q.eq("workId", work._id).eq("version", work.definitionVersion ?? 0)
)
.unique();
if (!current) {
throw new ConvexError("Current Definition not found");
}
const definition = JSON.parse(current.payloadJson) as WorkDefinition;
const persistedQuestions = await ctx.db
.query("workQuestions")
.withIndex("by_workId_and_definitionVersion", (q) =>
q.eq("workId", work._id).eq("definitionVersion", current.version)
)
.collect();
if (
!persistedQuestions.some((question) => question.questionId === questionId)
) {
throw new ConvexError("Question not found");
}
const next = {
...definition,
questions: persistedQuestions.map((question) => ({
alternatives: JSON.parse(question.alternativesJson) as string[],
answer:
question.questionId === questionId ? update.answer : question.answer,
id: question.questionId,
impact: question.impact,
prompt: question.prompt,
recommendation: question.recommendation,
status:
question.questionId === questionId ? update.status : question.status,
})),
};
return await saveDefinition(ctx, work, JSON.stringify(next), "user");
};
export const answerQuestion = mutation({
args: { answer: v.string(), questionId: v.string(), workId: v.id("works") },
handler: async (ctx, args) => {
const work = await requireWork(ctx, args.workId);
await requireProjectMember(ctx, work.projectId);
const result = await reviseQuestion(ctx, work, args.questionId, {
answer: args.answer,
status: "answered",
});
await appendEvent(
ctx,
work._id,
"question.answered",
`question-answered:${args.questionId}:${result.version}`
);
return result;
},
});
export const withdrawQuestion = mutation({
args: { questionId: v.string(), workId: v.id("works") },
handler: async (ctx, args) => {
const work = await requireWork(ctx, args.workId);
await requireProjectMember(ctx, work.projectId);
const result = await reviseQuestion(ctx, work, args.questionId, {
status: "withdrawn",
});
await appendEvent(
ctx,
work._id,
"question.withdrawn",
`question-withdrawn:${args.questionId}:${result.version}`
);
return result;
},
});
export const requestDesign = mutation({
args: { workId: v.id("works") },
handler: async (ctx, args) => {
const work = await requireWork(ctx, args.workId);
await requireProjectMember(ctx, work.projectId);
if (
(work.status !== "designing" && work.status !== "blocked") ||
work.definitionApprovalVersion !== work.definitionVersion
) {
throw new ConvexError("Definition must be approved before Design");
}
if (work.status === "blocked") {
await ctx.db.patch(work._id, {
status: "designing",
updatedAt: Date.now(),
});
}
await appendEvent(
ctx,
work._id,
"design.requested",
`design-requested:${work.definitionVersion}`
);
await ctx.scheduler.runAfter(0, plannerRef, {
organizationId: work.organizationId,
workId: work._id,
});
return { status: "designing" as const };
},
});
const saveDesign = async (
ctx: MutationCtx,
work: Doc<"works">,
payloadJson: string,
createdBy: string
) => {
if (work.status === "executing") {
throw new ConvexError("Cannot revise Work while a Run is executing");
}
if (work.definitionApprovalVersion !== work.definitionVersion) {
throw new ConvexError("Design must bind an approved Definition version");
}
const decoded = await Effect.runPromise(
validateDesignPacket({
...objectPayload(payloadJson),
definitionVersion: work.definitionVersion,
version: (work.designVersion ?? 0) + 1,
})
).catch((error: unknown) => {
throw new ConvexError(
error instanceof Error ? error.message : "Invalid Design Packet"
);
});
const prior = await ctx.db
.query("designPackets")
.withIndex("by_work_and_version", (q) => q.eq("workId", work._id))
.collect();
for (const row of prior) {
if (row.status === "current") {
await ctx.db.patch(row._id, { status: "superseded" });
}
}
const designApprovals = await ctx.db
.query("workApprovals")
.withIndex("by_workId_and_kind", (q) =>
q.eq("workId", work._id).eq("kind", "design")
)
.collect();
for (const approval of designApprovals) {
if (approval.status === "active") {
await ctx.db.patch(approval._id, { status: "invalidated" });
}
}
const designId = await ctx.db.insert("designPackets", {
createdAt: Date.now(),
createdBy,
definitionVersion: decoded.definitionVersion,
payloadJson: JSON.stringify(decoded),
status: "current",
version: decoded.version,
workId: work._id,
});
for (const [ordinal, slice] of decoded.slices.entries()) {
await ctx.db.insert("workSlices", {
createdAt: Date.now(),
designVersion: decoded.version,
objective: slice.objective,
observableBehavior: slice.observableBehavior,
ordinal,
payloadJson: JSON.stringify(slice),
sliceId: slice.id,
status: ordinal === 0 ? "ready" : "planned",
title: slice.title,
workId: work._id,
});
}
await ctx.db.patch(work._id, {
designApprovalVersion: undefined,
designVersion: decoded.version,
status: "awaiting-design-approval",
updatedAt: Date.now(),
});
await appendEvent(
ctx,
work._id,
work.designVersion ? "design.revised" : "design.saved",
`design:${decoded.version}`,
String(designId),
JSON.stringify(decoded)
);
return { designId, version: decoded.version };
};
export const saveDesignProposal = mutation({
args: { payloadJson: v.string(), workId: v.id("works") },
handler: async (ctx, args) => {
const work = await requireWork(ctx, args.workId);
await requireProjectMember(ctx, work.projectId);
return await saveDesign(ctx, work, args.payloadJson, "user");
},
});
export const reviseDesign = saveDesignProposal;
export const approveDesign = mutation({
args: {
definitionVersion: v.number(),
designVersion: v.number(),
workId: v.id("works"),
},
handler: async (ctx, args) => {
const work = await requireWork(ctx, args.workId);
await requireProjectMember(ctx, work.projectId);
if (
work.definitionApprovalVersion !== args.definitionVersion ||
work.designVersion !== args.designVersion ||
work.status !== "awaiting-design-approval"
) {
throw new ConvexError(
"Design approval must bind current Definition and Design versions"
);
}
const identity = await ctx.auth.getUserIdentity();
if (!identity) {
throw new ConvexError("Authentication required");
}
await ctx.db.insert("workApprovals", {
approvedAt: Date.now(),
approvedBy: identity.tokenIdentifier,
definitionVersion: args.definitionVersion,
designVersion: args.designVersion,
kind: "design",
status: "active",
workId: work._id,
});
await ctx.db.patch(work._id, {
designApprovalVersion: args.designVersion,
status: "ready",
updatedAt: Date.now(),
});
await appendEvent(
ctx,
work._id,
"design.approved",
`design-approved:${args.definitionVersion}:${args.designVersion}`
);
return { status: "ready" as const };
},
});
export const submitDefinitionProposal = mutation({
args: { payloadJson: v.string(), token: v.string(), workId: v.id("works") },
handler: async (ctx, args) => {
if (args.token !== env.FLUE_DB_TOKEN) {
throw new ConvexError("Invalid agent control token");
}
const work = await requireWork(ctx, args.workId);
return await saveDefinition(ctx, work, args.payloadJson, "work-planner");
},
});
export const submitDesignProposal = mutation({
args: { payloadJson: v.string(), token: v.string(), workId: v.id("works") },
handler: async (ctx, args) => {
if (args.token !== env.FLUE_DB_TOKEN) {
throw new ConvexError("Invalid agent control token");
}
const work = await requireWork(ctx, args.workId);
return await saveDesign(ctx, work, args.payloadJson, "work-planner");
},
});
export const submitQuestion = mutation({
args: { questionJson: v.string(), token: v.string(), workId: v.id("works") },
handler: async (ctx, args) => {
if (args.token !== env.FLUE_DB_TOKEN) {
throw new ConvexError("Invalid agent control token");
}
const work = await requireWork(ctx, args.workId);
if (work.definitionVersion === undefined) {
throw new ConvexError("Work has no Definition to attach a question to");
}
if (
work.status !== "defining" &&
work.status !== "awaiting-definition-approval"
) {
throw new ConvexError(
"Questions can only be submitted while the Work is being defined"
);
}
const question = await Effect.runPromise(
Schema.decodeUnknownEffect(WorkQuestion)(parsePayload(args.questionJson))
).catch((error: unknown) => {
throw new ConvexError(
error instanceof Error ? error.message : "Invalid question"
);
});
const existing = await ctx.db
.query("workQuestions")
.withIndex("by_workId_and_definitionVersion_and_questionId", (q) =>
q
.eq("workId", work._id)
.eq("definitionVersion", work.definitionVersion!)
.eq("questionId", question.id)
)
.unique();
const questionRow = {
alternativesJson: JSON.stringify(question.alternatives),
answer: question.answer,
createdAt: Date.now(),
definitionVersion: work.definitionVersion,
impact: question.impact,
prompt: question.prompt,
questionId: question.id,
recommendation: question.recommendation,
status: question.status,
workId: work._id,
};
if (existing) {
const same =
existing.prompt === questionRow.prompt &&
existing.impact === questionRow.impact &&
existing.recommendation === questionRow.recommendation &&
existing.alternativesJson === questionRow.alternativesJson &&
existing.status === questionRow.status &&
existing.answer === questionRow.answer;
if (!same) {
throw new ConvexError("Question ID already has different content");
}
return { accepted: false };
}
await ctx.db.insert("workQuestions", questionRow);
await appendEvent(
ctx,
work._id,
"question.created",
`question-created:${work.definitionVersion}:${question.id}`,
question.id,
JSON.stringify(question)
);
return { accepted: true };
},
});
export const recordPlannerFailure = internalMutation({
args: { error: v.string(), workId: v.id("works") },
handler: async (ctx, args): Promise<null> => {
const work = await ctx.db.get(args.workId);
if (!work || (work.status !== "defining" && work.status !== "designing")) {
return null;
}
await ctx.db.patch(work._id, { status: "blocked", updatedAt: Date.now() });
await appendEvent(
ctx,
work._id,
"planner.failed",
`planner-failed:${work._id}:${work.status}`,
undefined,
JSON.stringify({ error: args.error, phase: work.status })
);
return null;
},
});
export const runPlanner = internalAction({
args: { organizationId: v.id("organizations"), workId: v.id("works") },
handler: async (ctx, args) => {
const flueUrl = (env as typeof env & { readonly FLUE_URL?: string })
.FLUE_URL;
if (!flueUrl) {
await ctx.runMutation(plannerFailureRef, {
error: "FLUE_URL is not configured in Convex",
workId: args.workId,
});
return null;
}
const context = await ctx.runQuery(plannerContextRef, {
workId: args.workId,
});
if (!context) {
return null;
}
const endpoint = new URL(
`agents/work-planner/${encodeURIComponent(String(args.organizationId))}`,
`${flueUrl.replace(/\/+$/u, "")}/`
);
endpoint.searchParams.set("wait", "result");
try {
const response = await fetch(endpoint, {
body: JSON.stringify({
message: `Plan Work ${String(args.workId)} titled "${context.title}" with objective "${context.objective}". Current state is ${context.status}. If the Work is defining, submit only a Definition proposal and questions. If the Work is designing, submit only a Design Packet proposal bound to the approved Definition. Never approve or execute.`,
}),
headers: {
authorization: `Bearer ${env.FLUE_DB_TOKEN}`,
"content-type": "application/json",
"x-zopu-organization-id": String(args.organizationId),
"x-zopu-request-id": `work-planner:${args.workId}`,
},
method: "POST",
});
if (response.ok) {
return null;
}
await ctx.runMutation(plannerFailureRef, {
error: `Work Planner request failed (${response.status})`,
workId: args.workId,
});
} catch (error) {
await ctx.runMutation(plannerFailureRef, {
error: error instanceof Error ? error.message : String(error),
workId: args.workId,
});
}
// The private worker submits typed proposals through Convex mutations; it
// never approves or advances Work itself.
return null;
},
});
export const getPlannerContext = internalQuery({
args: { workId: v.id("works") },
handler: async (ctx, args) => {
const work = await ctx.db.get(args.workId);
return work
? { objective: work.objective, status: work.status, title: work.title }
: null;
},
});
export const listForProject = query({
args: { projectId: v.id("projects") },
handler: async (ctx, args) => {
await requireProjectMember(ctx, args.projectId);
const works = await ctx.db
.query("works")
.withIndex("by_project_and_createdAt", (q: any) =>
q.eq("projectId", args.projectId)
)
.order("desc")
.take(100);
return await Promise.all(
works.map(async (work) => {
const definitions = await ctx.db
.query("workDefinitions")
.withIndex("by_work_and_version", (q: any) =>
q.eq("workId", work._id)
)
.order("desc")
.take(10);
const designs = await ctx.db
.query("designPackets")
.withIndex("by_work_and_version", (q: any) =>
q.eq("workId", work._id)
)
.order("desc")
.take(10);
const slices = await ctx.db
.query("workSlices")
.withIndex("by_workId_and_designVersion", (q) =>
q
.eq("workId", work._id)
.eq("designVersion", work.designVersion ?? 0)
)
.collect();
const runs = await ctx.db
.query("workRuns")
.withIndex("by_work_and_createdAt", (q: any) =>
q.eq("workId", work._id)
)
.order("desc")
.take(10);
const events = await ctx.db
.query("workEvents")
.withIndex("by_work_and_createdAt", (q: any) =>
q.eq("workId", work._id)
)
.order("desc")
.take(100);
const attachments = await ctx.db
.query("signalWorkAttachments")
.withIndex("by_work", (q: any) => q.eq("workId", work._id))
.collect();
const signals = await Promise.all(
attachments.map(async (attachment) => {
const signal = await ctx.db.get(attachment.signalId);
if (!signal) {
return null;
}
const sources = await ctx.db
.query("signalSources")
.withIndex("by_signalId_and_ordinal", (q) =>
q.eq("signalId", signal._id)
)
.collect();
return {
createdAt: signal.createdAt,
signalId: signal._id,
sources: sources
.sort((a, b) => a.ordinal - b.ordinal)
.map((source) => ({
createdAt: source.sourceCreatedAt,
messageId: String(source.messageId),
rawText: source.rawTextSnapshot,
submissionId: null,
})),
summary: signal.summary,
title: signal.title,
};
})
);
const currentDefinition = definitions.find(
(definition) => definition.status === "current"
);
const currentDesign = designs.find(
(design) => design.status === "current"
);
return {
...work,
definition: currentDefinition
? JSON.parse(currentDefinition.payloadJson)
: null,
definitions,
design: currentDesign ? JSON.parse(currentDesign.payloadJson) : null,
designs,
events,
runs,
signals: signals.filter((signal) => signal !== null),
slices: slices.sort((a, b) => a.ordinal - b.ordinal),
};
})
);
},
});

View File

@@ -8,8 +8,7 @@ import { ConvexError, v } from "convex/values";
import { Effect } from "effect";
import type { Doc, Id } from "./_generated/dataModel";
import { mutation, query } from "./_generated/server";
import type { MutationCtx } from "./_generated/server";
import { mutation, type MutationCtx, query } from "./_generated/server";
import { requireProjectMember } from "./authz";
const requireAgent = (token: string) => {
@@ -242,6 +241,7 @@ export const listForProject = query({
return {
createdAt: signal.createdAt,
signalId: signal._id,
summary: signal.summary,
sources: sources
.sort((a, b) => a.ordinal - b.ordinal)
.map((source) => ({
@@ -250,7 +250,6 @@ export const listForProject = query({
rawText: source.rawTextSnapshot,
submissionId: null,
})),
summary: signal.summary,
title: signal.title,
};
})

View File

@@ -8,8 +8,8 @@
"dev": "convex dev --env-file ../../.env",
"dev:setup": "convex dev --env-file ../../.env --configure --until-success",
"check-types": "tsc --noEmit -p convex/tsconfig.json",
"test": "CONVEX_SITE_URL=https://convex.test FLUE_DB_TOKEN=test-token SITE_URL=https://site.test vitest run",
"test:watch": "CONVEX_SITE_URL=https://convex.test FLUE_DB_TOKEN=test-token SITE_URL=https://site.test vitest"
"test": "vitest run",
"test:watch": "vitest"
},
"dependencies": {
"@better-auth/expo": "catalog:",
@@ -23,7 +23,7 @@
},
"devDependencies": {
"@code/config": "workspace:*",
"@types/node": "catalog:",
"@types/node": "^24.3.0",
"convex-test": "catalog:",
"typescript": "catalog:",
"vitest": "catalog:"

View File

@@ -19,7 +19,7 @@
},
"devDependencies": {
"@code/config": "workspace:*",
"@types/node": "catalog:",
"@types/node": "^22.13.14",
"typescript": "catalog:"
}
}

View File

@@ -15,13 +15,7 @@
"./project-workspace": "./src/project-workspace.ts",
"./signal": "./src/signal.ts",
"./smoke": "./src/smoke.ts",
"./work": "./src/work.ts",
"./work-artifact": "./src/work-artifact.ts",
"./work-definition": "./src/work-definition.ts",
"./work-design": "./src/work-design.ts",
"./work-lifecycle": "./src/work-lifecycle.ts",
"./resolver": "./src/resolver.ts",
"./harness-runtime": "./src/harness-runtime.ts"
"./work": "./src/work.ts"
},
"scripts": {
"check-types": "tsc --noEmit",
@@ -36,7 +30,7 @@
},
"devDependencies": {
"@code/config": "workspace:*",
"@types/node": "catalog:",
"@types/node": "^22.13.14",
"typescript": "catalog:",
"vitest": "catalog:"
}

View File

@@ -0,0 +1,13 @@
{
"name": "@code/primitives",
"exports": {
".": "./src/index.ts",
"./agent-os": "./src/agent-os.ts",
"./git": "./src/git.ts",
"./project": "./src/project.ts",
"./project-issue": "./src/project-issue.ts",
"./project-workspace": "./src/project-workspace.ts",
"./project-work": "./src/project-work.ts",
"./signal": "./src/signal.ts",
"./smoke": "./src/smoke.ts"
},

View File

@@ -1,167 +0,0 @@
import { Effect, Schema } from "effect";
import type { AttemptOutcome } from "./resolver";
const Text = Schema.String.check(
Schema.makeFilter((value) => value.trim().length > 0, {
expected: "a non-empty string",
})
);
export const HarnessEventKind = Schema.Literals([
"started",
"progress",
"question",
"log",
"completed",
"failed",
"cancelled",
]);
export type HarnessEventKind = typeof HarnessEventKind.Type;
export const HarnessEvent = Schema.Struct({
kind: HarnessEventKind,
message: Text,
metadata: Schema.Record(Schema.String, Schema.String),
occurredAt: Schema.Number,
sequence: Schema.Int,
});
export type HarnessEvent = typeof HarnessEvent.Type;
export const FakeScenario = Schema.Literals([
"success",
"transient-failure-then-success",
"needs-input",
"permanent-failure",
"cancelled",
]);
export type FakeScenario = typeof FakeScenario.Type;
export interface HarnessRuntime {
readonly name: string;
readonly run: (input: {
readonly attemptNumber?: number;
readonly scenario: FakeScenario;
readonly signal?: AbortSignal;
}) => Effect.Effect<readonly [readonly HarnessEvent[], AttemptOutcome]>;
}
const event = (
sequence: number,
kind: HarnessEventKind,
message: string,
occurredAt: number
): HarnessEvent => ({ kind, message, metadata: {}, occurredAt, sequence });
export const FakeHarnessLive = (
now: () => number = Date.now
): HarnessRuntime => ({
name: "fake",
run: ({ attemptNumber = 1, scenario, signal }) =>
Effect.sync(() => {
const started = now();
const events: HarnessEvent[] = [
event(0, "started", `fake harness started: ${scenario}`, started),
];
if (signal?.aborted || scenario === "cancelled") {
events.push(event(1, "cancelled", "simulation cancelled", now()));
return [
events,
{
classification: "Cancelled",
retryable: false,
summary: "Simulation cancelled",
},
] as const;
}
events.push(
event(1, "progress", "validated approved Definition and Design", now())
);
switch (scenario) {
case "success": {
events.push(
event(
2,
"completed",
"fake implementation evidence recorded",
now()
)
);
return [
events,
{
classification: "Succeeded",
retryable: false,
summary: "Simulation succeeded; no implementation was claimed",
},
] as const;
}
case "transient-failure-then-success": {
if (attemptNumber > 1) {
events.push(
event(2, "completed", "transient failure recovered", now())
);
return [
events,
{
classification: "Succeeded",
retryable: false,
summary:
"Simulation recovered after a transient failure; no implementation was claimed",
},
] as const;
}
events.push(event(2, "failed", "transient provider failure", now()));
return [
events,
{
classification: "RetryableFailure",
retryable: true,
summary: "Transient simulation failure",
},
] as const;
}
case "needs-input": {
events.push(
event(2, "question", "simulation needs a human decision", now())
);
return [
events,
{
classification: "NeedsInput",
retryable: false,
summary: "Simulation needs input",
},
] as const;
}
case "permanent-failure": {
events.push(
event(
2,
"failed",
"deterministic permanent simulation failure",
now()
)
);
return [
events,
{
classification: "PermanentFailure",
retryable: false,
summary: "Permanent simulation failure",
},
] as const;
}
default: {
return [
events,
{
classification: "PermanentFailure",
retryable: false,
summary: "Unknown simulation scenario",
},
] as const;
}
}
}),
});

View File

@@ -10,9 +10,3 @@ export * from "./project-work";
export * from "./signal";
export * from "./smoke";
export * from "./work";
export * from "./work-artifact";
export * from "./work-definition";
export * from "./work-design";
export * from "./work-lifecycle";
export * from "./resolver";
export * from "./harness-runtime";

View File

@@ -1,108 +0,0 @@
import { Schema } from "effect";
import type { WorkStatus } from "./work-lifecycle";
const Text = Schema.String.check(
Schema.makeFilter((value) => value.trim().length > 0, {
expected: "a non-empty string",
})
);
export const RunStatus = Schema.Literals([
"ready",
"running",
"terminal",
"cancelled",
]);
export type RunStatus = typeof RunStatus.Type;
export const AttemptStatus = Schema.Literals([
"queued",
"claimed",
"running",
"terminal",
]);
export type AttemptStatus = typeof AttemptStatus.Type;
export const AttemptClassification = Schema.Literals([
"Succeeded",
"RetryableFailure",
"NeedsInput",
"Blocked",
"VerificationFailed",
"BudgetExhausted",
"Cancelled",
"PermanentFailure",
]);
export type AttemptClassification = typeof AttemptClassification.Type;
export const AttemptOutcome = Schema.Struct({
classification: AttemptClassification,
retryable: Schema.Boolean,
summary: Text,
});
export type AttemptOutcome = typeof AttemptOutcome.Type;
export const RetryPolicy = Schema.Struct({
maxAttempts: Schema.Int,
retryableClassifications: Schema.Array(AttemptClassification),
});
export type RetryPolicy = typeof RetryPolicy.Type;
export const CodingKitV0 = Schema.Struct({
harness: Schema.Literal("fake"),
id: Schema.Literal("coding-kit-v0"),
retryPolicy: RetryPolicy,
version: Schema.Literal("0.1.0"),
});
export type CodingKitV0 = typeof CodingKitV0.Type;
export const shouldRetry = (
outcome: AttemptOutcome,
attemptNumber: number,
policy: RetryPolicy
): boolean =>
outcome.retryable &&
attemptNumber < policy.maxAttempts &&
policy.retryableClassifications.includes(outcome.classification);
// Terminal Work status the resolver settles on once it will run no further
// attempt. RetryableFailure lands here only after the retry budget is gone.
export const WORK_STATUS_FOR_OUTCOME: Readonly<
Record<AttemptClassification, WorkStatus>
> = {
Blocked: "blocked",
BudgetExhausted: "failed",
Cancelled: "ready",
NeedsInput: "needs-input",
PermanentFailure: "failed",
RetryableFailure: "failed",
Succeeded: "completed",
VerificationFailed: "blocked",
};
export type Resolution =
| { readonly kind: "retry" }
| { readonly kind: "terminal"; readonly workStatus: WorkStatus };
// Resolver decision after an attempt finishes: retry within policy, or settle
// Work at the terminal status mapped from the attempt classification.
export const resolveOutcome = (
outcome: AttemptOutcome,
attemptNumber: number,
policy: RetryPolicy
): Resolution =>
shouldRetry(outcome, attemptNumber, policy)
? { kind: "retry" }
: {
kind: "terminal",
workStatus: WORK_STATUS_FOR_OUTCOME[outcome.classification],
};
export const defaultCodingKitV0: CodingKitV0 = {
harness: "fake",
id: "coding-kit-v0",
retryPolicy: {
maxAttempts: 3,
retryableClassifications: ["RetryableFailure"],
},
version: "0.1.0",
};

View File

@@ -1,58 +0,0 @@
import { Effect } from "effect";
import { describe, expect, test } from "vitest";
import {
decodeWorkArtifactDraft,
decodeWorkDeliveryDraft,
} from "./work-artifact";
describe("Work artifact contracts", () => {
test("accepts traceable artifact and delivery metadata", async () => {
await expect(
Effect.runPromise(
decodeWorkArtifactDraft({
idempotencyKey: "attempt:1:test-report",
kind: "test-report",
metadataJson: "{}",
producer: "fake-harness",
provenanceJson: "{}",
sourceRevision: "abc123",
title: "Focused test report",
verificationStatus: "verified",
})
)
).resolves.toMatchObject({ kind: "test-report" });
await expect(
Effect.runPromise(
decodeWorkDeliveryDraft({
externalId: "42",
idempotencyKey: "delivery:pr:42",
kind: "pull-request",
metadataJson: "{}",
provider: "gitea",
sourceRevision: "abc123",
status: "ready",
target: "main",
url: "https://git.example/pulls/42",
})
)
).resolves.toMatchObject({ kind: "pull-request" });
});
test("rejects unverifiable evidence metadata", async () => {
await expect(
Effect.runPromise(
decodeWorkArtifactDraft({
idempotencyKey: "attempt:1:test-report",
kind: "test-report",
metadataJson: "not-json",
producer: "fake-harness",
provenanceJson: "{}",
title: "Focused test report",
verificationStatus: "verified",
})
)
).rejects.toThrow();
});
});

View File

@@ -1,141 +0,0 @@
import { Effect, Schema } from "effect";
const Text = Schema.String.check(
Schema.makeFilter((value) => value.trim().length > 0, {
expected: "a non-empty string",
})
);
const JsonText = Schema.String.check(
Schema.makeFilter(
(value) => {
try {
JSON.parse(value);
return true;
} catch {
return false;
}
},
{ expected: "valid JSON text" }
)
);
export const WorkArtifactKind = Schema.Literals([
"definition",
"design",
"diff",
"test-report",
"verification-report",
"runtime-log",
"screenshot",
"video",
"commit",
"branch",
"pull-request",
"preview",
"deployment",
"other",
]);
export type WorkArtifactKind = typeof WorkArtifactKind.Type;
export const ArtifactVerificationStatus = Schema.Literals([
"unverified",
"verified",
"rejected",
]);
export type ArtifactVerificationStatus = typeof ArtifactVerificationStatus.Type;
export const WorkArtifactDraft = Schema.Struct({
contentHash: Schema.optional(Text),
environmentId: Schema.optional(Text),
idempotencyKey: Text,
kind: WorkArtifactKind,
metadataJson: JsonText,
producer: Text,
provenanceJson: JsonText,
sourceRevision: Schema.optional(Text),
title: Text,
uri: Schema.optional(Text),
verificationStatus: ArtifactVerificationStatus,
});
export type WorkArtifactDraft = typeof WorkArtifactDraft.Type;
export const WorkDeliveryKind = Schema.Literals([
"branch",
"commit",
"pull-request",
"preview",
"deployment",
]);
export type WorkDeliveryKind = typeof WorkDeliveryKind.Type;
export const WorkDeliveryStatus = Schema.Literals([
"recorded",
"ready",
"approved",
"delivered",
"failed",
"cancelled",
]);
export type WorkDeliveryStatus = typeof WorkDeliveryStatus.Type;
export const WorkDeliveryDraft = Schema.Struct({
externalId: Text,
idempotencyKey: Text,
kind: WorkDeliveryKind,
metadataJson: JsonText,
provider: Text,
sourceRevision: Text,
status: WorkDeliveryStatus,
target: Text,
url: Schema.optional(Text),
});
export type WorkDeliveryDraft = typeof WorkDeliveryDraft.Type;
export class WorkArtifactError extends Schema.TaggedErrorClass<WorkArtifactError>()(
"WorkArtifactError",
{ message: Schema.String }
) {}
const decode = <A, I, R>(schema: Schema.Codec<A, I, R>, input: unknown) =>
Schema.decodeUnknownEffect(schema)(input).pipe(
Effect.mapError(
(cause) => new WorkArtifactError({ message: cause.message })
)
);
export const decodeWorkArtifactDraft = Effect.fn("Work.decodeArtifactDraft")(
function* decodeArtifactDraft(input: unknown) {
const draft = yield* decode(WorkArtifactDraft, input);
if (
draft.verificationStatus === "verified" &&
draft.sourceRevision === undefined
) {
return yield* Effect.fail(
new WorkArtifactError({
message: "Verified artifacts must bind an exact source revision",
})
);
}
return draft;
}
);
export const decodeWorkDeliveryDraft = Effect.fn("Work.decodeDeliveryDraft")(
(input: unknown) => decode(WorkDeliveryDraft, input)
);
const DELIVERY_TRANSITIONS: Readonly<
Record<WorkDeliveryStatus, readonly WorkDeliveryStatus[]>
> = {
approved: ["delivered", "failed", "cancelled"],
cancelled: [],
delivered: [],
failed: [],
ready: ["approved", "delivered", "failed", "cancelled"],
recorded: ["ready", "failed", "cancelled"],
};
export const canTransitionWorkDelivery = (
from: WorkDeliveryStatus,
to: WorkDeliveryStatus
): boolean => DELIVERY_TRANSITIONS[from].includes(to);

View File

@@ -1,110 +0,0 @@
import { Effect, Schema } from "effect";
const Text = Schema.String.check(
Schema.makeFilter((value) => value.trim().length > 0, {
expected: "a non-empty string",
})
);
export const DefinitionRisk = Schema.Literals(["low", "medium", "high"]);
export type DefinitionRisk = typeof DefinitionRisk.Type;
export const QuestionImpact = Schema.Literals(["low", "medium", "high"]);
export type QuestionImpact = typeof QuestionImpact.Type;
export const WorkQuestion = Schema.Struct({
alternatives: Schema.Array(Text),
answer: Schema.optional(Text),
id: Text,
impact: QuestionImpact,
prompt: Text,
recommendation: Schema.optional(Text),
status: Schema.Literals(["open", "answered", "withdrawn"]),
});
export type WorkQuestion = typeof WorkQuestion.Type;
export const WorkDefinition = Schema.Struct({
acceptanceCriteria: Schema.NonEmptyArray(Text),
affectedUsers: Schema.Array(Text),
assumptions: Schema.Array(Text),
constraints: Schema.Array(Text),
desiredOutcome: Text,
inScope: Schema.Array(Text),
outOfScope: Schema.Array(Text),
problem: Text,
questions: Schema.Array(WorkQuestion),
requiredArtifacts: Schema.Array(Text),
risk: DefinitionRisk,
rollback: Schema.optional(Text),
rollout: Schema.optional(Text),
version: Schema.Int,
});
export type WorkDefinition = typeof WorkDefinition.Type;
export const DefinitionApproval = Schema.Struct({
approvedAt: Schema.Number,
approvedBy: Text,
definitionVersion: Schema.Int,
});
export type DefinitionApproval = typeof DefinitionApproval.Type;
export const DefinitionProposal = WorkDefinition;
export type DefinitionProposal = typeof DefinitionProposal.Type;
export class DefinitionError extends Schema.TaggedErrorClass<DefinitionError>()(
"DefinitionError",
{
message: Schema.String,
reason: Schema.Literals(["Invalid", "OpenHighImpactQuestion"]),
}
) {}
export const decodeDefinition = Effect.fn("Work.decodeDefinition")(
function* decodeDefinition(input: unknown) {
const definition = yield* Schema.decodeUnknownEffect(WorkDefinition)(
input
).pipe(
Effect.mapError(
(cause) =>
new DefinitionError({ message: cause.message, reason: "Invalid" })
)
);
const questionIds = new Set<string>();
for (const question of definition.questions) {
if (questionIds.has(question.id)) {
return yield* Effect.fail(
new DefinitionError({
message: `Question ID must be unique: ${question.id}`,
reason: "Invalid",
})
);
}
questionIds.add(question.id);
}
return definition;
}
);
export const validateDefinition = Effect.fn("Work.validateDefinition")(
function* validateDefinition(input: unknown) {
const definition = yield* decodeDefinition(input);
const blocked = definition.questions.some(
(question) => question.status === "open" && question.impact === "high"
);
if (blocked) {
return yield* Effect.fail(
new DefinitionError({
message:
"High-impact open questions must be resolved before approval",
reason: "OpenHighImpactQuestion",
})
);
}
return definition;
}
);
export const canApproveDefinition = (definition: WorkDefinition): boolean =>
definition.questions.every(
(question) => question.status !== "open" || question.impact !== "high"
);

View File

@@ -1,119 +0,0 @@
import { Effect, Schema } from "effect";
const Text = Schema.String.check(
Schema.makeFilter((value) => value.trim().length > 0, {
expected: "a non-empty string",
})
);
export const ImpactMap = Schema.Struct({
files: Schema.Array(Text),
modules: Schema.Array(Text),
risks: Schema.Array(Text),
summary: Text,
});
export type ImpactMap = typeof ImpactMap.Type;
export const VerticalSlice = Schema.Struct({
codeBoundaries: Schema.Array(Text),
dependsOn: Schema.Array(Text),
evidenceRequirements: Schema.NonEmptyArray(Text),
id: Text,
objective: Text,
observableBehavior: Text,
reviewRequired: Schema.Boolean,
title: Text,
verification: Schema.NonEmptyArray(Text),
});
export type VerticalSlice = typeof VerticalSlice.Type;
export const DesignPacket = Schema.Struct({
architectureSummary: Text,
callFlowDelta: Schema.Array(Text),
concerns: Schema.Array(Text),
definitionVersion: Schema.Int,
evidenceRequirements: Schema.NonEmptyArray(Text),
fileTreeDelta: Schema.Array(Text),
impactMap: ImpactMap,
invariants: Schema.NonEmptyArray(Text),
keyInterfaces: Schema.Array(Text),
slices: Schema.Array(VerticalSlice),
tradeoffs: Schema.Array(Text),
version: Schema.Int,
});
export type DesignPacket = typeof DesignPacket.Type;
export const DesignApproval = Schema.Struct({
approvedAt: Schema.Number,
approvedBy: Text,
definitionVersion: Schema.Int,
designVersion: Schema.Int,
});
export type DesignApproval = typeof DesignApproval.Type;
export class DesignError extends Schema.TaggedErrorClass<DesignError>()(
"DesignError",
{
message: Schema.String,
reason: Schema.Literals([
"Invalid",
"NoObservableSlices",
"DefinitionMismatch",
"InvalidSliceGraph",
]),
}
) {}
export const validateDesignPacket = Effect.fn("Work.validateDesignPacket")(
function* validateDesignPacket(input: unknown) {
const design = yield* Schema.decodeUnknownEffect(DesignPacket)(input).pipe(
Effect.mapError(
(cause) =>
new DesignError({ message: cause.message, reason: "Invalid" })
)
);
if (
design.slices.length === 0 ||
design.slices.length > 4 ||
design.slices.some(
(slice) =>
slice.codeBoundaries.length === 0 ||
slice.observableBehavior.trim().length === 0 ||
slice.evidenceRequirements.length === 0 ||
slice.verification.length === 0
)
) {
return yield* Effect.fail(
new DesignError({
message:
"Every Design must contain observable, independently verifiable slices",
reason: "NoObservableSlices",
})
);
}
const seen = new Set<string>();
for (const slice of design.slices) {
if (
seen.has(slice.id) ||
slice.dependsOn.some(
(dependency) => dependency === slice.id || !seen.has(dependency)
)
) {
return yield* Effect.fail(
new DesignError({
message:
"Slice IDs must be unique and dependencies must reference earlier slices",
reason: "InvalidSliceGraph",
})
);
}
seen.add(slice.id);
}
return design;
}
);
export const isObservableSlice = (slice: VerticalSlice): boolean =>
slice.observableBehavior.trim().length > 0 &&
slice.evidenceRequirements.length > 0 &&
slice.verification.length > 0;

View File

@@ -1,65 +0,0 @@
import { Schema } from "effect";
export const WorkStatus = Schema.Literals([
"proposed",
"defining",
"awaiting-definition-approval",
"designing",
"awaiting-design-approval",
"ready",
"executing",
"needs-input",
"blocked",
"completed",
"failed",
"cancelled",
]);
export type WorkStatus = typeof WorkStatus.Type;
export const WORK_TRANSITIONS: Readonly<
Record<WorkStatus, readonly WorkStatus[]>
> = {
"awaiting-definition-approval": ["designing", "defining"],
"awaiting-design-approval": ["ready", "designing"],
blocked: ["ready", "executing", "cancelled"],
cancelled: ["ready"],
completed: [],
defining: ["awaiting-definition-approval", "proposed"],
designing: ["awaiting-design-approval", "awaiting-definition-approval"],
executing: [
"ready",
"needs-input",
"blocked",
"completed",
"failed",
"cancelled",
],
failed: ["ready", "executing", "cancelled"],
"needs-input": ["ready", "executing", "cancelled"],
proposed: ["defining"],
ready: ["executing", "designing", "cancelled"],
};
export const canTransitionWork = (from: WorkStatus, to: WorkStatus): boolean =>
WORK_TRANSITIONS[from].includes(to);
export const assertWorkTransition = (
from: WorkStatus,
to: WorkStatus
): void => {
if (!canTransitionWork(from, to)) {
throw new Error(`Invalid Work transition: ${from} -> ${to}`);
}
};
export const definitionRevisionInvalidation = {
definitionApproval: "invalidate" as const,
dependentDesign: "invalidate" as const,
designApproval: "invalidate" as const,
workStatus: "defining" as const,
};
export const designRevisionInvalidation = {
designApproval: "invalidate" as const,
workStatus: "designing" as const,
};

View File

@@ -1,174 +0,0 @@
import { Effect } from "effect";
import { describe, expect, test } from "vitest";
import { FakeHarnessLive } from "./harness-runtime";
import { defaultCodingKitV0, resolveOutcome, shouldRetry } from "./resolver";
import { decodeDefinition, validateDefinition } from "./work-definition";
import { validateDesignPacket } from "./work-design";
import { canTransitionWork } from "./work-lifecycle";
const definition = {
acceptanceCriteria: ["terminal outcome"],
affectedUsers: ["founders"],
assumptions: [],
constraints: [],
desiredOutcome: "A visible terminal simulation exists",
inScope: ["fake execution"],
outOfScope: ["real sandboxing"],
problem: "A deterministic build path is needed",
questions: [
{
alternatives: [],
id: "q1",
impact: "high",
prompt: "Which runtime?",
status: "open",
},
],
requiredArtifacts: ["events"],
risk: "medium",
version: 1,
};
const probeOutcome = (
classification: Parameters<typeof resolveOutcome>[0]["classification"],
retryable = false
) => ({ classification, retryable, summary: "probe" });
describe("work resolution contracts", () => {
test("high-impact open questions block definition approval", async () => {
await expect(
Effect.runPromise(decodeDefinition(definition))
).resolves.toEqual(definition);
await expect(
Effect.runPromise(validateDefinition(definition))
).rejects.toThrow(/High-impact/u);
});
test("question identities are unique within a Definition version", async () => {
await expect(
Effect.runPromise(
decodeDefinition({
...definition,
questions: [definition.questions[0], definition.questions[0]],
})
)
).rejects.toThrow(/Question ID must be unique/u);
});
test("observable slices and terminal fake outcomes are enforced", async () => {
const approved = await Effect.runPromise(
validateDefinition({ ...definition, questions: [] })
);
const design = await Effect.runPromise(
validateDesignPacket({
architectureSummary: "small",
callFlowDelta: [],
concerns: [],
definitionVersion: approved.version,
evidenceRequirements: ["event"],
fileTreeDelta: [],
impactMap: { files: [], modules: [], risks: [], summary: "small" },
invariants: ["terminal"],
keyInterfaces: [],
slices: [
{
codeBoundaries: ["runtime"],
dependsOn: [],
evidenceRequirements: ["event"],
id: "s1",
objective: "one",
observableBehavior: "visible",
reviewRequired: false,
title: "one",
verification: ["assert"],
},
],
tradeoffs: [],
version: 1,
})
);
expect(design.slices).toHaveLength(1);
const [, outcome] = await Effect.runPromise(
FakeHarnessLive(() => 1).run({ scenario: "success" })
);
expect(outcome.classification).toBe("Succeeded");
expect(outcome.summary).toMatch(/no implementation/u);
const [, firstTransient] = await Effect.runPromise(
FakeHarnessLive(() => 1).run({
attemptNumber: 1,
scenario: "transient-failure-then-success",
})
);
const [, recoveredTransient] = await Effect.runPromise(
FakeHarnessLive(() => 1).run({
attemptNumber: 2,
scenario: "transient-failure-then-success",
})
);
expect(firstTransient.classification).toBe("RetryableFailure");
expect(recoveredTransient.classification).toBe("Succeeded");
expect(
shouldRetry(
{
classification: "RetryableFailure",
retryable: true,
summary: "retry",
},
1,
defaultCodingKitV0.retryPolicy
)
).toBe(true);
expect(canTransitionWork("ready", "executing")).toBe(true);
});
});
describe("resolver decisions", () => {
const policy = defaultCodingKitV0.retryPolicy;
test("retries a retryable failure while the budget remains", () => {
expect(
resolveOutcome(probeOutcome("RetryableFailure", true), 1, policy)
).toEqual({ kind: "retry" });
});
test("settles RetryableFailure as failed once the budget is exhausted", () => {
expect(
resolveOutcome(probeOutcome("RetryableFailure", true), 3, policy)
).toEqual({
kind: "terminal",
workStatus: "failed",
});
});
test("maps each terminal classification to its Work status", () => {
expect(resolveOutcome(probeOutcome("Succeeded"), 1, policy)).toEqual({
kind: "terminal",
workStatus: "completed",
});
expect(resolveOutcome(probeOutcome("PermanentFailure"), 1, policy)).toEqual(
{
kind: "terminal",
workStatus: "failed",
}
);
expect(resolveOutcome(probeOutcome("NeedsInput"), 1, policy)).toEqual({
kind: "terminal",
workStatus: "needs-input",
});
expect(resolveOutcome(probeOutcome("Blocked"), 1, policy)).toEqual({
kind: "terminal",
workStatus: "blocked",
});
expect(resolveOutcome(probeOutcome("Cancelled"), 1, policy)).toEqual({
kind: "terminal",
workStatus: "ready",
});
});
test("completed and cancelled are not silently reopened", () => {
expect(canTransitionWork("completed", "defining")).toBe(false);
expect(canTransitionWork("cancelled", "proposed")).toBe(false);
expect(canTransitionWork("ready", "executing")).toBe(true);
});
});

View File

@@ -1,14 +1,14 @@
import { Array as EffectArray, Effect, Order, Schema } from "effect";
export { WorkStatus } from "./work-lifecycle";
export type { WorkStatus as WorkLifecycleStatus } from "./work-lifecycle";
const MeaningfulString = Schema.String.check(
Schema.makeFilter((value) => value.trim().length > 0, {
expected: "a non-empty string",
})
);
export const WorkStatus = Schema.Literal("proposed");
export type WorkStatus = typeof WorkStatus.Type;
export const WorkDraft = Schema.Struct({
objective: MeaningfulString,
title: MeaningfulString,
@@ -91,34 +91,6 @@ const WorkAttachmentInput = Schema.Struct({
export const WorkEventKind = Schema.Literals([
"work.proposed",
"signal.attached",
"definition.requested",
"definition.saved",
"definition.revised",
"definition.approved",
"definition.invalidated",
"question.created",
"question.answered",
"question.withdrawn",
"design.requested",
"design.saved",
"design.revised",
"design.approved",
"design.invalidated",
"planner.failed",
"slice.started",
"slice.completed",
"slice.ready",
"run.started",
"run.completed",
"run.cancelled",
"attempt.claimed",
"attempt.event",
"attempt.completed",
"attempt.reconciled",
"resolver.decided",
"artifact.recorded",
"delivery.recorded",
"delivery.updated",
]);
export type WorkEventKind = typeof WorkEventKind.Type;

View File

@@ -30,15 +30,15 @@
"shadcn": "^4.12.0",
"shiki": "^4.3.1",
"sonner": "catalog:",
"streamdown": "catalog:",
"streamdown": "^2.5.0",
"tailwind-merge": "catalog:",
"tw-animate-css": "^1.4.0",
"use-stick-to-bottom": "^1.1.6"
},
"devDependencies": {
"@code/config": "workspace:*",
"@tailwindcss/postcss": "catalog:",
"@types/react": "catalog:",
"@tailwindcss/postcss": "^4.3.2",
"@types/react": "^19.2.17",
"@types/react-dom": "catalog:",
"tailwindcss": "catalog:",
"typescript": "catalog:"

View File

@@ -7,7 +7,6 @@ export default defineConfig({
"**/node_modules/**",
".agents/skills/**",
"repos/**",
"scripts/**",
"apps/daemon/**",
"apps/desktop/**",
"apps/native/**",
@@ -29,7 +28,6 @@ export default defineConfig({
"**/node_modules/**",
".agents/skills/**",
"repos/**",
"scripts/**",
"apps/daemon/**",
"apps/desktop/**",
"apps/native/**",