diff --git a/services/api/src/db/cursor-helpers.ts b/services/api/src/db/cursor-helpers.ts index b981c847..6510279e 100644 --- a/services/api/src/db/cursor-helpers.ts +++ b/services/api/src/db/cursor-helpers.ts @@ -1,3 +1,5 @@ +import type postgres from "postgres" +import type { Queryable } from "./client.js" export const CURSOR_UUID_RE = /^[0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12}$/i @@ -10,6 +12,10 @@ export interface TimeCursor { id: string } +export interface DecodedTimeCursor extends TimeCursor { + instant: string +} + export const MIN_UUID = "00000000-0000-0000-0000-000000000000" export const MAX_UUID = "ffffffff-ffff-ffff-ffff-ffffffffffff" @@ -19,37 +25,48 @@ export const CURSOR_ISO_RE = export const CURSOR_MIN_MS = Date.UTC(1970, 0, 1) export const CURSOR_MAX_MS = Date.UTC(2100, 0, 1) -function parseCursorInstant(iso: string): Date | null { +const CURSOR_FRACTION_RE = /\.(\d{1,6})/ + +function parseCursorInstant(iso: string): { at: Date; instant: string } | null { if (!CURSOR_ISO_RE.test(iso)) return null const at = new Date(iso) const ms = at.getTime() if (!Number.isFinite(ms) || ms < CURSOR_MIN_MS || ms > CURSOR_MAX_MS) return null - return at + const micros = (CURSOR_FRACTION_RE.exec(iso)?.[1] ?? "").padEnd(6, "0").slice(3) + return { at, instant: at.toISOString().replace(/Z$/, `${micros}Z`) } } export function parseTimeCursor( cursor: string | null | undefined, opts?: { requireUuid?: boolean; direction?: "asc" | "desc" }, -): TimeCursor | null { +): DecodedTimeCursor | null { if (cursor === null || cursor === undefined || cursor === "") return null const requireUuid = opts?.requireUuid ?? true const idx = cursor.indexOf("|") if (idx === 0) return null if (idx < 0) { - const at = parseCursorInstant(cursor) - if (at === null) return null - return { at, id: opts?.direction === "asc" ? MIN_UUID : MAX_UUID } + const parsed = parseCursorInstant(cursor) + if (parsed === null) return null + return { ...parsed, id: opts?.direction === "asc" ? MIN_UUID : MAX_UUID } } const id = cursor.slice(idx + 1) - const at = parseCursorInstant(cursor.slice(0, idx)) - if (at === null) return null + const parsed = parseCursorInstant(cursor.slice(0, idx)) + if (parsed === null) return null if (requireUuid && !CURSOR_UUID_RE.test(id)) return null if (id.length === 0) return null - return { at, id } + return { ...parsed, id } +} + +export function encodeTimeCursor(c: { at: Date | string; id: string }): string { + return `${typeof c.at === "string" ? c.at : c.at.toISOString()}|${c.id}` +} + +export function cursorInstantSql(sql: Queryable, column: postgres.Fragment): postgres.Fragment { + return sql`to_char(${column} AT TIME ZONE 'UTC', 'YYYY-MM-DD"T"HH24:MI:SS.US"Z"')` } -export function encodeTimeCursor(c: TimeCursor): string { - return `${c.at.toISOString()}|${c.id}` +export function cursorAtSql(sql: Queryable, cursor: { instant: string }): postgres.Fragment { + return sql`${cursor.instant}::text::timestamptz` } export interface NameCursor { diff --git a/services/api/src/services/admin/inbound-repository.drizzle.ts b/services/api/src/services/admin/inbound-repository.drizzle.ts index 671f6165..502f6cdf 100644 --- a/services/api/src/services/admin/inbound-repository.drizzle.ts +++ b/services/api/src/services/admin/inbound-repository.drizzle.ts @@ -1,5 +1,6 @@ import type { Sql } from "../../db/client.js" +import { cursorAtSql, cursorInstantSql } from "../../db/cursor-helpers.js" import { clampLimit, decodeCursor, encodeCursor } from "./pagination.js" import { likeContains } from "./like.js" import { writeAudit } from "./audit.js" @@ -142,7 +143,7 @@ export function makeDrizzleInboundRepository(sql: Sql): InboundRepository { const anchor = decodeCursor(query.cursor, true) const cursorFilter = anchor !== null - ? sql`AND (received_at, id) < (${anchor.createdAt}, ${anchor.id}::uuid)` + ? sql`AND (received_at, id) < (${cursorAtSql(sql, anchor)}, ${anchor.id}::uuid)` : sql`` const statusFilter = query.status === "unread" @@ -161,11 +162,12 @@ export function makeDrizzleInboundRepository(sql: Sql): InboundRepository { return sql`AND (from_addr ILIKE ${like} ESCAPE '\\' OR subject ILIKE ${like} ESCAPE '\\' OR recipient ILIKE ${like} ESCAPE '\\')` })() : sql`` - const rows = await sql` + const rows = await sql<(InboundListRowSelect & { cursor_at: string })[]>` SELECT id, from_addr, recipient, subject, left(body_text, ${PREVIEW_SOURCE_CHARS}) AS preview_text, left(body_html, ${HTML_PREVIEW_SOURCE_CHARS}) AS preview_html, has_attachments, status, received_at, + ${cursorInstantSql(sql, sql`received_at`)} AS cursor_at, headers->>${INBOUND_AUTH_VERDICT_HEADER}::text AS auth_verdict FROM inbound_emails WHERE true @@ -181,7 +183,7 @@ export function makeDrizzleInboundRepository(sql: Sql): InboundRepository { const items = page.map(toListItem) const last = page[page.length - 1] const nextCursor = - hasMore && last ? encodeCursor({ createdAt: last.received_at, id: last.id }) : null + hasMore && last ? encodeCursor({ createdAt: last.cursor_at, id: last.id }) : null return { items, nextCursor } }, diff --git a/services/api/src/services/admin/inbox-feed-repository.drizzle.ts b/services/api/src/services/admin/inbox-feed-repository.drizzle.ts index b4c3ba60..ff7cc7e8 100644 --- a/services/api/src/services/admin/inbox-feed-repository.drizzle.ts +++ b/services/api/src/services/admin/inbox-feed-repository.drizzle.ts @@ -1,4 +1,5 @@ import type { Sql } from "../../db/client.js" +import { cursorAtSql, cursorInstantSql } from "../../db/cursor-helpers.js" import { clampLimit, decodeCursor, encodeCursor } from "./pagination.js" import { HTML_PREVIEW_SOURCE_CHARS, PREVIEW_SOURCE_CHARS, toPreview } from "./mail-preview.js" import { normalizeAuthVerdict, replyPublication } from "./mail-mappers.js" @@ -37,12 +38,14 @@ export const INBOX_FEED_REPLY_FILTERS: ReadonlySet = new Set([ export interface InboxFeedEmailRow extends Omit { source: "email" ts: Date + cursor_at: string } export interface InboxFeedReplyRow { source: "reply" id: string ts: Date + cursor_at: string from_addr: string | null subject: string | null preview_text: string | null @@ -92,7 +95,7 @@ export function toInboxFeedPage(rows: readonly InboxFeedRow[], limit: number): I const hasMore = rows.length > limit && last !== undefined return { items: page.map(toInboxFeedItem), - nextCursor: hasMore ? encodeCursor({ createdAt: last.ts, id: last.id }) : null, + nextCursor: hasMore ? encodeCursor({ createdAt: last.cursor_at, id: last.id }) : null, } } @@ -108,7 +111,7 @@ export function makeDrizzleInboxFeedRepository(sql: Sql): InboxFeedRepository { const before = (ts: SqlFragment, id: SqlFragment): SqlFragment => anchor === null ? sql`` - : sql`AND (${ts}, ${id}) < (${anchor.createdAt}, ${anchor.id}::uuid)` + : sql`AND (${ts}, ${id}) < (${cursorAtSql(sql, anchor)}, ${anchor.id}::uuid)` const branches: SqlFragment[] = [] if (INBOX_FEED_EMAIL_FILTERS.has(filter)) { @@ -162,7 +165,8 @@ export function makeDrizzleInboxFeedRepository(sql: Sql): InboxFeedRepository { const union = branches.reduce((acc, branch) => sql`${acc} UNION ALL ${branch}`) const rows = await sql` - SELECT * FROM (${union}) feed + SELECT feed.*, ${cursorInstantSql(sql, sql`feed.ts`)} AS cursor_at + FROM (${union}) feed ORDER BY ts DESC, id DESC LIMIT ${limit + 1} ` diff --git a/services/api/src/services/admin/inbox-feed-repository.memory.ts b/services/api/src/services/admin/inbox-feed-repository.memory.ts index 3bf70bc9..89fbc036 100644 --- a/services/api/src/services/admin/inbox-feed-repository.memory.ts +++ b/services/api/src/services/admin/inbox-feed-repository.memory.ts @@ -46,6 +46,7 @@ export class InMemoryInboxFeedRepository implements InboxFeedRepository { source: "email", id: r.id, ts: r.receivedAt, + cursor_at: r.receivedAt.toISOString(), from_addr: r.fromAddr, recipient: r.recipient, subject: r.subject, @@ -70,6 +71,7 @@ export class InMemoryInboxFeedRepository implements InboxFeedRepository { source: "reply", id: m.id, ts: m.createdAt, + cursor_at: m.createdAt.toISOString(), from_addr: m.fromAddr, subject: m.subject ?? t.subject, preview_text: m.body, diff --git a/services/api/src/services/admin/pagination.ts b/services/api/src/services/admin/pagination.ts index ba89d4dc..196d5fb4 100644 --- a/services/api/src/services/admin/pagination.ts +++ b/services/api/src/services/admin/pagination.ts @@ -22,8 +22,12 @@ export interface CursorAnchor { id: string } +export interface DecodedCursorAnchor extends CursorAnchor { + instant: string +} + /** Encode a keyset anchor into the opaque "|" cursor string. */ -export function encodeCursor(anchor: CursorAnchor): string { +export function encodeCursor(anchor: { createdAt: Date | string; id: string }): string { return encodeTimeCursor({ at: anchor.createdAt, id: anchor.id }) } @@ -38,10 +42,10 @@ export function encodeCursor(anchor: CursorAnchor): string { export function decodeCursor( cursor: string | null | undefined, requireUuidId = false, -): CursorAnchor | null { +): DecodedCursorAnchor | null { const parsed = parseTimeCursor(cursor, { requireUuid: requireUuidId }) if (parsed === null) return null - return { createdAt: parsed.at, id: parsed.id } + return { createdAt: parsed.at, id: parsed.id, instant: parsed.instant } } /** diff --git a/services/api/test/integration/inbound-repository.test.ts b/services/api/test/integration/inbound-repository.test.ts index d74c693a..c47a9b09 100644 --- a/services/api/test/integration/inbound-repository.test.ts +++ b/services/api/test/integration/inbound-repository.test.ts @@ -80,6 +80,25 @@ describe.skipIf(!pg)("inbound repository (integration: real schema)", () => { expect(onlySupport.items).toHaveLength(2) }) + it("pages rows that share a millisecond by their microseconds, skipping none", async () => { + for (const micros of ["123100", "123500", "123900"]) { + const { id } = await repo.insertIdempotent(insert({ messageId: `` })) + await h.sql` + UPDATE inbound_emails SET received_at = ${`2027-01-01T00:00:00.${micros}Z`}::text::timestamptz + WHERE id = ${id} + ` + } + const walked: string[] = [] + let cursor: string | undefined + do { + const page = await repo.list({ limit: 1, ...(cursor === undefined ? {} : { cursor }) }) + walked.push(...page.items.map((i) => i.id)) + cursor = page.nextCursor ?? undefined + } while (cursor !== undefined) + expect(walked).toEqual((await repo.list({})).items.map((i) => i.id)) + expect(walked).toHaveLength(3) + }) + it("reads the stored auth verdict onto list rows and the detail", async () => { const headers = { "x-civfix-auth-verdict": "fail" } const r = await repo.insertIdempotent(insert({ messageId: "", headers })) diff --git a/services/api/test/integration/inbox-feed-repository.test.ts b/services/api/test/integration/inbox-feed-repository.test.ts index 3313a26e..9135de7c 100644 --- a/services/api/test/integration/inbox-feed-repository.test.ts +++ b/services/api/test/integration/inbox-feed-repository.test.ts @@ -153,4 +153,24 @@ describe.skipIf(!pg)("inbox feed repository (integration: real schema)", () => { expect(all).toHaveLength(6) expect(all.slice(0, 2).sort().reverse()).toEqual(all.slice(0, 2)) }) + + it("walks rows that share a millisecond by their microseconds, skipping none", async () => { + const s = await mailbox() + const microsecond = (micros: string) => `2027-01-01T00:00:00.${micros}Z` + await h.sql`UPDATE inbound_emails SET received_at = ${microsecond("123100")}::text::timestamptz` + await h.sql` + UPDATE mail_messages SET created_at = ${microsecond("123500")}::text::timestamptz + WHERE id IN (${s.withheld!}, ${s.published!}) + ` + const all = await ids({}) + const walked: string[] = [] + let cursor: string | null = null + do { + const page = await feed.list({ limit: 1, ...(cursor === null ? {} : { cursor }) }) + walked.push(...page.items.map((i) => i.id)) + cursor = page.nextCursor + } while (cursor !== null) + expect(walked).toEqual(all) + expect(all).toHaveLength(5) + }) }) diff --git a/services/api/test/unit/admin-inbox-keyset.test.ts b/services/api/test/unit/admin-inbox-keyset.test.ts new file mode 100644 index 00000000..8e754724 --- /dev/null +++ b/services/api/test/unit/admin-inbox-keyset.test.ts @@ -0,0 +1,47 @@ +import { describe, expect, it } from "vitest" +import type { Sql } from "../../src/db/client.js" +import { makeFakeSql } from "../helpers/fake-sql.js" +import { makeDrizzleInboundRepository } from "../../src/services/admin/inbound-repository.drizzle.js" +import { makeDrizzleInboxFeedRepository } from "../../src/services/admin/inbox-feed-repository.drizzle.js" + +const A = "aaaaaaaa-0000-4000-8000-000000000001" +const B = "aaaaaaaa-0000-4000-8000-000000000002" +const TS = new Date("2027-01-01T00:00:00.123Z") +const INSTANT = "2027-01-01T00:00:00.123456Z" + +const emailRow = (id: string) => ({ + source: "email", + id, + ts: TS, + received_at: TS, + cursor_at: INSTANT, + from_addr: "resident@example.com", + recipient: "support@civfix.org", + subject: "Probe", + preview_text: "Hello", + preview_html: null, + has_attachments: false, + status: "unread", + auth_verdict: null, +}) + +const repos = { + "the inbox feed": (sql: Sql) => makeDrizzleInboxFeedRepository(sql), + "the unmatched inbox": (sql: Sql) => makeDrizzleInboundRepository(sql), +} + +describe.each(Object.entries(repos))("%s keyset", (_name, make) => { + it("pages from the row's microsecond instant, not its millisecond Date", async () => { + const db = makeFakeSql([{ match: /inbound_emails/, rows: [emailRow(A), emailRow(B)] }]) + const repo = make(db.sql as unknown as Sql) + + const first = await repo.list({ limit: 1 }) + expect(first.nextCursor).toBe(`${INSTANT}|${A}`) + + await repo.list({ limit: 1, cursor: first.nextCursor! }) + const next = db.statements.at(-1)! + expect(next.values).toContain(INSTANT) + expect(next.values).not.toContainEqual(TS) + expect(next.sql).toContain("?::text::timestamptz") + }) +}) diff --git a/services/api/test/unit/cursor-helpers.test.ts b/services/api/test/unit/cursor-helpers.test.ts index 4e6c5df9..daae2d20 100644 --- a/services/api/test/unit/cursor-helpers.test.ts +++ b/services/api/test/unit/cursor-helpers.test.ts @@ -18,6 +18,7 @@ import { const UUID = "3f2504e0-4f89-11d3-9a0c-0305e82c3301" const UUID_2 = "9c858901-8a57-4791-81fe-4c455b099bc9" const ISO = "2026-07-24T12:34:56.789Z" +const INSTANT = "2026-07-24T12:34:56.789000Z" describe("isUuid / CURSOR_UUID_RE", () => { it("accepts a canonical uuid in either case and rejects near-misses", () => { @@ -53,10 +54,26 @@ describe("parseTimeCursor", () => { expect(parsed!.id).toBe(original.id) }) + it("keeps the cursor's microseconds as the instant the database compares against", () => { + const parsed = parseTimeCursor(`2026-07-24T12:34:56.789123Z|${UUID}`) + expect(parsed?.at.toISOString()).toBe(ISO) + expect(parsed?.instant).toBe("2026-07-24T12:34:56.789123Z") + expect(parseTimeCursor(`2026-07-24T14:34:56.7+02:00|${UUID}`)?.instant).toBe( + "2026-07-24T12:34:56.700000Z", + ) + }) + + it("encodes a database instant verbatim, so its microseconds survive the round trip", () => { + const encoded = encodeTimeCursor({ at: "2026-07-24T12:34:56.789123Z", id: UUID }) + expect(encoded).toBe(`2026-07-24T12:34:56.789123Z|${UUID}`) + expect(parseTimeCursor(encoded)?.instant).toBe("2026-07-24T12:34:56.789123Z") + }) + it("anchors a legacy timestamp-only cursor with a DIRECTION-AWARE sentinel so no boundary row is skipped", () => { - expect(parseTimeCursor(ISO)).toEqual({ at: new Date(ISO), id: MAX_UUID }) - expect(parseTimeCursor(ISO, { direction: "desc" })).toEqual({ at: new Date(ISO), id: MAX_UUID }) - expect(parseTimeCursor(ISO, { direction: "asc" })).toEqual({ at: new Date(ISO), id: MIN_UUID }) + const anchor = { at: new Date(ISO), instant: INSTANT } + expect(parseTimeCursor(ISO)).toEqual({ ...anchor, id: MAX_UUID }) + expect(parseTimeCursor(ISO, { direction: "desc" })).toEqual({ ...anchor, id: MAX_UUID }) + expect(parseTimeCursor(ISO, { direction: "asc" })).toEqual({ ...anchor, id: MIN_UUID }) expect(isUuid(parseTimeCursor(ISO)!.id)).toBe(true) expect(isUuid(parseTimeCursor(ISO, { direction: "asc" })!.id)).toBe(true) }) @@ -86,7 +103,7 @@ describe("parseTimeCursor", () => { it("with requireUuid:false accepts a non-uuid id, but still rejects an EMPTY one", () => { const parsed = parseTimeCursor(`${ISO}|room-42`, { requireUuid: false }) - expect(parsed).toEqual({ at: new Date(ISO), id: "room-42" }) + expect(parsed).toEqual({ at: new Date(ISO), instant: INSTANT, id: "room-42" }) expect(parseTimeCursor(`${ISO}|`, { requireUuid: false })).toBeNull() expect(parseTimeCursor("nope|room-42", { requireUuid: false })).toBeNull() }) @@ -232,7 +249,11 @@ describe("paginate", () => { const page = paginate(rows, 1, (r) => ({ at: r.at, id: r.id })) expect(page.items).toEqual([rows[0]]) expect(page.nextCursor).toBe(`${ISO}|${UUID}`) - expect(parseTimeCursor(page.nextCursor)).toEqual({ at: new Date(ISO), id: UUID }) + expect(parseTimeCursor(page.nextCursor)).toEqual({ + at: new Date(ISO), + instant: INSTANT, + id: UUID, + }) }) it("falls back to `createdAt` when the anchor has no `at`", () => {