diff --git a/packages/chat/src/store.ts b/packages/chat/src/store.ts index 32279e628..71c8abc68 100644 --- a/packages/chat/src/store.ts +++ b/packages/chat/src/store.ts @@ -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"; @@ -311,8 +311,8 @@ export function createDrizzleChatStore>( 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(); @@ -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; }, diff --git a/packages/chat/test/read-state.drizzle.test.ts b/packages/chat/test/read-state.drizzle.test.ts new file mode 100644 index 000000000..a8c844b93 --- /dev/null +++ b/packages/chat/test/read-state.drizzle.test.ts @@ -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(); + } + }); +}); diff --git a/packages/chat/test/store.test.ts b/packages/chat/test/store.test.ts index 06bc0a0d8..f7cbab837 100644 --- a/packages/chat/test/store.test.ts +++ b/packages/chat/test/store.test.ts @@ -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"); +});