Convex control plane (Slices 1-4 prerequisites for Slice 5): - Versioned Definition/Design persistence: revisions preserve history instead of deleting prior slices/approvals - Fenced conversation turn queue: attempt-numbered leases prevent stale worker overwrites; expired turns reconciled by cron - Fenced Run/Attempt claiming: only queued attempts from running Runs can be claimed; lease expiry checks on every finish/checkpoint - Multi-slice progression: successful slice marks next slice ready instead of completing the whole Work; completed Runs blocked from retry - Typed AttemptClassification in schema and resolver decisions table - Durable resolver decisions for normal outcomes, cancellation, and lease-expiry recovery - Persistent artifacts and delivery metadata with provenance, source revision, verification status, and controlled delivery transitions - High-impact open questions block definition approval - Question submission gated to defining/awaiting-approval states only - Delivery status updates validate JSON metadata before persisting - Planner failures recorded as durable blocked events - Approval identity uses canonical tokenIdentifier Primitives: - decodeDefinition separates decode from validate - Design validation enforces 1-4 slices with unique IDs and valid deps - Artifact and delivery draft schemas with verified-revision binding - Delivery status transition table Package management: - Consolidated shared deps into single catalog (hono, valibot, streamdown, @types/react, @types/node, @tailwindcss/*, react-native) - Removed cross-version Hono boundary: flue() exported as Fetchable - Deleted bun.lock and regenerated from clean catalog resolution Lint/format: - Removed .eslintignore; oxc uses native ignorePatterns in oxlint.config.ts - Deleted redundant packages/backend/.oxlintrc.json - Added repos/** and scripts/** to ignore patterns - Convex-specific rule overrides for ES2022 target constraints - Fixed all oxc errors in fluePersistence.ts and agents/src/db.ts - Fixed all formatting across changed files
143 lines
4.2 KiB
TypeScript
143 lines
4.2 KiB
TypeScript
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",
|
|
]);
|
|
});
|
|
});
|