Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
23 changes: 16 additions & 7 deletions packages/chat/src/store.ts
Original file line number Diff line number Diff line change
Expand Up @@ -10,7 +10,7 @@
// Routing against the interface (rather than a raw drizzle handle) keeps the
// route layer testable with a plain in-memory fake, with no database and no
// drizzle SQL-condition internals involved.
import { and, eq, inArray } from "drizzle-orm";
import { and, eq, inArray, sql } from "drizzle-orm";
import type { PostgresJsDatabase } from "drizzle-orm/postgres-js";

import { participantsOf } from "./workbench-settings";
Expand Down Expand Up @@ -311,8 +311,8 @@ export function createDrizzleChatStore<TSchema extends Record<string, unknown>>(
workbenchReadState.principalId,
],
set: {
lastSeenCreatedAt: input.lastSeenCreatedAt,
lastSeenId: input.lastSeenId,
lastSeenCreatedAt: sql`CASE WHEN excluded.last_seen_created_at >= ${workbenchReadState.lastSeenCreatedAt} THEN excluded.last_seen_created_at ELSE ${workbenchReadState.lastSeenCreatedAt} END`,
lastSeenId: sql`CASE WHEN excluded.last_seen_created_at >= ${workbenchReadState.lastSeenCreatedAt} THEN excluded.last_seen_id ELSE ${workbenchReadState.lastSeenId} END`,
},
})
.returning();
Expand Down Expand Up @@ -459,11 +459,20 @@ export function createInMemoryChatStore(): ChatStore {
},

async putReadState(input) {
const row: ReadStateRow = { ...input };
readStateByKey.set(
readStateKey(input.tenantId, input.workbenchId, input.principalId),
row,
const key = readStateKey(
input.tenantId,
input.workbenchId,
input.principalId,
);
const existing = readStateByKey.get(key);
if (
existing !== undefined &&
existing.lastSeenCreatedAt > input.lastSeenCreatedAt
) {
return existing;
}
const row: ReadStateRow = { ...input };
readStateByKey.set(key, row);
return row;
},

Expand Down
162 changes: 162 additions & 0 deletions packages/chat/test/read-state.drizzle.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,162 @@
// DB-gated: skipped when no DATABASE_URL is reachable (a fresh
// checkout still runs the unit gates), mirroring `migrations.test.ts`.
// Runs against its own scratch database.
//
// `store.test.ts` proves `putReadState`'s monotonicity guard against
// the in-memory store. This exercises the real `createDrizzleChatStore`
// path, where the guard is a conditional `ON CONFLICT DO UPDATE ... SET`
// rather than an in-process comparison — proving the SQL itself never
// regresses a reader's cursor.
import { afterAll, beforeAll, describe, expect, test } from "bun:test";
import { drizzle } from "drizzle-orm/postgres-js";
import postgres from "postgres";

import { e2eDatabaseUrl } from "../../../scripts/e2e/harness";
import { applyChatMigrations } from "../src/migrations";
import { createDrizzleChatStore } from "../src/store";

function scratchUrlFor(e2eUrl: string): string {
const url = new URL(e2eUrl);
const database = url.pathname.replace(/^\//, "");
url.pathname = `/${database}_chat_read_state_drizzle_test`;
return url.toString();
}

const databaseUrl = e2eDatabaseUrl();
const describeIfDb = databaseUrl === undefined ? describe.skip : describe;

const TENANT = "tnt_1";
const WORKBENCH = "run_workbench1";
const PRINCIPAL = "prn_alice";

describeIfDb("createDrizzleChatStore: putReadState monotonicity", () => {
const scratchUrl = scratchUrlFor(
databaseUrl ?? "postgres://localhost:5432/unused",
);
const scratchTarget = new URL(scratchUrl);
const scratchDatabase = scratchTarget.pathname.replace(/^\//, "");

beforeAll(async () => {
const maintenanceUrl = new URL(scratchUrl);
maintenanceUrl.pathname = "/postgres";
const maintenance = postgres(maintenanceUrl.toString(), {
max: 1,
onnotice: () => undefined,
});
try {
await maintenance.unsafe(`DROP DATABASE IF EXISTS "${scratchDatabase}"`);
await maintenance.unsafe(`CREATE DATABASE "${scratchDatabase}"`);
} finally {
await maintenance.end();
}
await applyChatMigrations(scratchUrl);
});

afterAll(async () => {
const maintenanceUrl = new URL(scratchUrl);
maintenanceUrl.pathname = "/postgres";
const maintenance = postgres(maintenanceUrl.toString(), {
max: 1,
onnotice: () => undefined,
});
try {
await maintenance.unsafe(`DROP DATABASE IF EXISTS "${scratchDatabase}"`);
} finally {
await maintenance.end();
}
});

test("a stale write landing after a newer one never moves the cursor backward", async () => {
const sql = postgres(scratchUrl, { max: 5, onnotice: () => undefined });
try {
const store = createDrizzleChatStore(drizzle(sql));

await store.putReadState({
tenantId: TENANT,
workbenchId: WORKBENCH,
principalId: PRINCIPAL,
lastSeenCreatedAt: new Date("2026-01-02T00:00:00.000Z"),
lastSeenId: "mail_2",
});

const result = await store.putReadState({
tenantId: TENANT,
workbenchId: WORKBENCH,
principalId: PRINCIPAL,
lastSeenCreatedAt: new Date("2026-01-01T00:00:00.000Z"),
lastSeenId: "mail_1",
});

expect(result.lastSeenId).toBe("mail_2");
expect(result.lastSeenCreatedAt).toEqual(
new Date("2026-01-02T00:00:00.000Z"),
);

const stored = await store.getReadState(TENANT, WORKBENCH, PRINCIPAL);
expect(stored?.lastSeenId).toBe("mail_2");
} finally {
await sql.end();
}
});

test("a newer write still moves the cursor forward", async () => {
const sql = postgres(scratchUrl, { max: 5, onnotice: () => undefined });
try {
const store = createDrizzleChatStore(drizzle(sql));

await store.putReadState({
tenantId: TENANT,
workbenchId: WORKBENCH,
principalId: "prn_bob",
lastSeenCreatedAt: new Date("2026-01-01T00:00:00.000Z"),
lastSeenId: "mail_1",
});

const result = await store.putReadState({
tenantId: TENANT,
workbenchId: WORKBENCH,
principalId: "prn_bob",
lastSeenCreatedAt: new Date("2026-01-03T00:00:00.000Z"),
lastSeenId: "mail_3",
});

expect(result.lastSeenId).toBe("mail_3");

const stored = await store.getReadState(TENANT, WORKBENCH, "prn_bob");
expect(stored?.lastSeenId).toBe("mail_3");
} finally {
await sql.end();
}
});

test("a same-millisecond forward move to a different message still lands", async () => {
const sql = postgres(scratchUrl, { max: 5, onnotice: () => undefined });
try {
const store = createDrizzleChatStore(drizzle(sql));
const sameCreatedAt = new Date("2026-01-04T00:00:00.001Z");

await store.putReadState({
tenantId: TENANT,
workbenchId: WORKBENCH,
principalId: "prn_carol",
lastSeenCreatedAt: sameCreatedAt,
lastSeenId: "mail_4",
});

const result = await store.putReadState({
tenantId: TENANT,
workbenchId: WORKBENCH,
principalId: "prn_carol",
lastSeenCreatedAt: sameCreatedAt,
lastSeenId: "mail_5",
});

expect(result.lastSeenId).toBe("mail_5");

const stored = await store.getReadState(TENANT, WORKBENCH, "prn_carol");
expect(stored?.lastSeenId).toBe("mail_5");
} finally {
await sql.end();
}
});
});
52 changes: 52 additions & 0 deletions packages/chat/test/store.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -122,3 +122,55 @@ test("putReadState upserts a per-principal cursor without disturbing other princ
const bob = await store.getReadState("tnt_1", "chn_1", "prn_bob");
expect(bob).toBeUndefined();
});

test("putReadState never moves the cursor backward when a stale write lands after a newer one", async () => {
const store = createInMemoryChatStore();
await store.putReadState({
tenantId: "tnt_1",
workbenchId: "chn_1",
principalId: "prn_alice",
lastSeenCreatedAt: new Date("2026-01-02T00:00:00.000Z"),
lastSeenId: "mail_2",
});

const result = await store.putReadState({
tenantId: "tnt_1",
workbenchId: "chn_1",
principalId: "prn_alice",
lastSeenCreatedAt: new Date("2026-01-01T00:00:00.000Z"),
lastSeenId: "mail_1",
});

expect(result.lastSeenId).toBe("mail_2");
expect(result.lastSeenCreatedAt).toEqual(
new Date("2026-01-02T00:00:00.000Z"),
);

const alice = await store.getReadState("tnt_1", "chn_1", "prn_alice");
expect(alice?.lastSeenId).toBe("mail_2");
});

test("putReadState still lands a same-millisecond forward move to a different message", async () => {
const store = createInMemoryChatStore();
const sameCreatedAt = new Date("2026-01-02T00:00:00.001Z");
await store.putReadState({
tenantId: "tnt_1",
workbenchId: "chn_1",
principalId: "prn_alice",
lastSeenCreatedAt: sameCreatedAt,
lastSeenId: "mail_2",
});

const result = await store.putReadState({
tenantId: "tnt_1",
workbenchId: "chn_1",
principalId: "prn_alice",
lastSeenCreatedAt: sameCreatedAt,
lastSeenId: "mail_3",
});

expect(result.lastSeenId).toBe("mail_3");

const alice = await store.getReadState("tnt_1", "chn_1", "prn_alice");
expect(alice?.lastSeenId).toBe("mail_3");
});
Loading