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 Promise>; } } 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("binds a direct admission to the product turn atomically", async () => { const t = convexTest({ modules, schema }); const organizationId = await t.run(async (ctx) => ctx.db.insert("organizations", { createdAt: 1, createdBy: "user-1", kind: "personal", name: "Test", }) ); const conversationId = await t.run(async (ctx) => ctx.db.insert("conversations", { createdAt: 1, organizationId }) ); const turnId = await t.run(async (ctx) => ctx.db.insert("conversationTurns", { attemptNumber: 1, clientRequestId: "request-1", conversationId, createdAt: 1, leaseExpiresAt: 100, leaseOwner: "worker", status: "dispatching", }) ); await t.mutation(api.fluePersistence.admitSubmission, { clientRequestId: "request-1", input: { acceptedAt: 1, chunksJson: "[]", inputJson: '{"kind":"direct"}', kind: "direct", sessionKey: "agent/instance/default", submissionId: "submission-1", }, token, turnId, }); const turn = await t.run((ctx) => ctx.db.get(turnId)); expect(turn).toMatchObject({ status: "running", submissionId: "submission-1", }); expect(turn).not.toHaveProperty("leaseExpiresAt"); expect(turn).not.toHaveProperty("leaseOwner"); }); 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", ]); }); });