Files
zopu-code/packages/backend/convex/fluePersistence.test.ts
2026-07-29 08:19:17 +05:30

190 lines
5.6 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("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",
]);
});
});