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
39 changes: 28 additions & 11 deletions services/api/src/db/cursor-helpers.ts
Original file line number Diff line number Diff line change
@@ -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

Expand All @@ -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"

Expand All @@ -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`) }
Comment on lines +35 to +36

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P1 Preserve legacy cursor pages

A composite cursor issued before this change contains only milliseconds. This code pads a legacy value such as .123Z to .123000Z, and the new descending keyset queries then use that reconstructed value as an exact microsecond anchor. Continuing from a cursor for a row at .123900Z omits an older row at .123100Z, so operators continuing an existing inbox page after deployment can permanently miss messages. Keep millisecond-compatible continuation behavior for legacy composite cursors while using precise instants for newly emitted cursors.

Artifacts

Legacy cursor integration test

  • The exact focused test source seeds microsecond-separated rows and compares exact versus legacy cursor continuation; it exercises both affected repositories.

Exact-microsecond cursor control

  • Executed PostgreSQL integration test output for cursor anchors that preserve `.123900Z`; both inbound and inbox-feed continuations return all expected rows.

Legacy millisecond cursor failure

  • Executed PostgreSQL integration test output for legacy `.123Z|UUID` cursors; both repository paths omit the row at `.123100Z`, confirming the defect.

Combined cursor integration run

  • Executed all four focused database cases together; exact-microsecond controls pass and legacy-millisecond cases fail, confirming the incompatible behavior.

View artifacts

T-Rex Ran code and verified through T-Rex

Prompt To Fix With AI
This is a comment left during a code review.
Path: services/api/src/db/cursor-helpers.ts
Line: 35-36

Comment:
**Preserve legacy cursor pages**

A composite cursor issued before this change contains only milliseconds. This code pads a legacy value such as `.123Z` to `.123000Z`, and the new descending keyset queries then use that reconstructed value as an exact microsecond anchor. Continuing from a cursor for a row at `.123900Z` omits an older row at `.123100Z`, so operators continuing an existing inbox page after deployment can permanently miss messages. Keep millisecond-compatible continuation behavior for legacy composite cursors while using precise instants for newly emitted cursors.

---

For each issue above, determine whether it is valid and should be fixed. If so, fix it directly.

Fix in Claude Code

}

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 {
Expand Down
8 changes: 5 additions & 3 deletions services/api/src/services/admin/inbound-repository.drizzle.ts
Original file line number Diff line number Diff line change
@@ -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"
Expand Down Expand Up @@ -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"
Expand All @@ -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<InboundListRowSelect[]>`
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
Expand All @@ -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 }
},

Expand Down
10 changes: 7 additions & 3 deletions services/api/src/services/admin/inbox-feed-repository.drizzle.ts
Original file line number Diff line number Diff line change
@@ -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"
Expand Down Expand Up @@ -37,12 +38,14 @@ export const INBOX_FEED_REPLY_FILTERS: ReadonlySet<InboxFeedFilter> = new Set([
export interface InboxFeedEmailRow extends Omit<InboundListRowSelect, "received_at"> {
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
Expand Down Expand Up @@ -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,
}
}

Expand All @@ -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)) {
Expand Down Expand Up @@ -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<InboxFeedRow[]>`
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}
`
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand All @@ -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,
Expand Down
10 changes: 7 additions & 3 deletions services/api/src/services/admin/pagination.ts
Original file line number Diff line number Diff line change
Expand Up @@ -22,8 +22,12 @@ export interface CursorAnchor {
id: string
}

export interface DecodedCursorAnchor extends CursorAnchor {
instant: string
}

/** Encode a keyset anchor into the opaque "<iso>|<id>" 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 })
}

Expand All @@ -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 }
}

/**
Expand Down
19 changes: 19 additions & 0 deletions services/api/test/integration/inbound-repository.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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: `<us-${micros}@x>` }))
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: "<verdict@x>", headers }))
Expand Down
20 changes: 20 additions & 0 deletions services/api/test/integration/inbox-feed-repository.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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)
})
})
47 changes: 47 additions & 0 deletions services/api/test/unit/admin-inbox-keyset.test.ts
Original file line number Diff line number Diff line change
@@ -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")
})
})
31 changes: 26 additions & 5 deletions services/api/test/unit/cursor-helpers.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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", () => {
Expand Down Expand Up @@ -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)
})
Expand Down Expand Up @@ -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()
})
Expand Down Expand Up @@ -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`", () => {
Expand Down
Loading