From 1624c86e348f5eb5fd7a71d99bb7e5cae605ed3e Mon Sep 17 00:00:00 2001 From: Theo Date: Wed, 23 Sep 2026 17:05:28 -0700 Subject: [PATCH] Audit the operator who publishes a verified city reply --- docs/inbound-mail-effects.md | 6 ++ services/api/src/routes/admin/mail.routes.ts | 12 ++-- .../src/services/admin/activity-service.ts | 2 +- .../api/src/services/admin/inbound-sweep.ts | 2 +- .../admin/inbound-thread-correlation.ts | 11 +++- .../admin/mail-reply-publish-service.ts | 34 ++++++++-- .../services/admin/mail-repository.drizzle.ts | 30 +++++++-- .../services/admin/mail-repository.memory.ts | 13 +++- .../api/src/services/admin/mail-repository.ts | 2 +- .../integration/admin-mail-repository.test.ts | 45 +++++++++++++ services/api/test/unit/admin-activity.test.ts | 2 +- .../unit/mail-reply-publish-service.test.ts | 65 ++++++++++++++++++- 12 files changed, 196 insertions(+), 28 deletions(-) diff --git a/docs/inbound-mail-effects.md b/docs/inbound-mail-effects.md index 21e8e6d0..bf68fe92 100644 --- a/docs/inbound-mail-effects.md +++ b/docs/inbound-mail-effects.md @@ -80,6 +80,12 @@ message of that thread qualifies (404 otherwise), and a thread with no report or or another runner holds the lease, the answer is `pending`: the cleared flag makes the reply eligible for the sweep, which finishes it. +The same endpoint finishes a verified reply whose effects are still owed (`pending`, for instance after +a failed chat insert). When the operator's call is the run that completes the effects, the transaction +that marks them applied also writes a `mail.reply_published` row naming that operator, unless the reply +already has one from its approval. A reply therefore carries at most one such row however many +operators or retries publish it, and a reply the sweep finishes carries none. + ## Stage 2 publishes the city's reply text (product decision) The `report_timeline` row (stage 1) and the reporter's push (stage 3) still carry **fixed copy only** — diff --git a/services/api/src/routes/admin/mail.routes.ts b/services/api/src/routes/admin/mail.routes.ts index 047c77df..98c40a2f 100644 --- a/services/api/src/routes/admin/mail.routes.ts +++ b/services/api/src/routes/admin/mail.routes.ts @@ -140,9 +140,11 @@ export async function registerAdminMailRoutes( (overrides): MailReplyPublishService => makeMailReplyPublishService({ repo: overrides.repo, - applyEffects: (thread, message) => { + applyEffects: (thread, message, publishedBy) => { const effects = overrides.inboundEffects ?? {} - return applyInboundEffects(container, effects, overrides.repo, thread, message) + return applyInboundEffects(container, effects, overrides.repo, thread, message, { + publishedBy, + }) }, logger: app.log, }), @@ -150,8 +152,10 @@ export async function registerAdminMailRoutes( const repo = makeDrizzleMailRepository(container.getDb().sql) return makeMailReplyPublishService({ repo, - applyEffects: (thread, message) => - applyInboundEffects(container, { logger: app.log }, repo, thread, message), + applyEffects: (thread, message, publishedBy) => + applyInboundEffects(container, { logger: app.log }, repo, thread, message, { + publishedBy, + }), logger: app.log, }) }, diff --git a/services/api/src/services/admin/activity-service.ts b/services/api/src/services/admin/activity-service.ts index 40408700..f28cbefb 100644 --- a/services/api/src/services/admin/activity-service.ts +++ b/services/api/src/services/admin/activity-service.ts @@ -207,7 +207,7 @@ export function describeAuditAction(action: string): string { "mail.status_changed": "Updated a mail thread", "mail.forward_template_set": "Updated the default forwarding template", "mail.reply_published_without_text": "A city reply was posted without its text", - "mail.reply_published": "Published a withheld city reply", + "mail.reply_published": "Published a city reply", "outreach.digest_sent": "Sent an outreach digest", "inbox.status_changed": "Updated an inbox message", // The L4 read audits are filtered out of the feed at the repo (AUDIT_READ_ACTIONS), so these labels diff --git a/services/api/src/services/admin/inbound-sweep.ts b/services/api/src/services/admin/inbound-sweep.ts index 7cd36725..0ca5f379 100644 --- a/services/api/src/services/admin/inbound-sweep.ts +++ b/services/api/src/services/admin/inbound-sweep.ts @@ -129,7 +129,7 @@ async function redriveEffects( mailRepo, thread, message, - () => now, + { now: () => now }, ) effectsRedriven += 1 } catch (err) { diff --git a/services/api/src/services/admin/inbound-thread-correlation.ts b/services/api/src/services/admin/inbound-thread-correlation.ts index 1b6a17fa..b3670498 100644 --- a/services/api/src/services/admin/inbound-thread-correlation.ts +++ b/services/api/src/services/admin/inbound-thread-correlation.ts @@ -2,6 +2,7 @@ import { createHash } from "node:crypto" import type { Container } from "../../di.js" import type { ParsedMail } from "@civfix/shared/interfaces" import type { + MailAuditInput, MailMessageRecord, MailRepository, MailThreadRecord, @@ -205,17 +206,23 @@ export async function isJurisdictionSender( return false } +export interface ApplyEffectsOptions { + now?: () => Date + publishedBy?: MailAuditInput +} + export async function applyInboundEffects( container: Container, injected: InboundEffectDeps, mailRepo: MailRepository, thread: MailThreadRecord, message: MailMessageRecord, - now: () => Date = () => new Date(), + opts: ApplyEffectsOptions = {}, ): Promise { if (message.direction !== "in") return if (message.unaffiliated) return if (thread.reportId === null && thread.cleanupId === null) return + const now = opts.now ?? (() => new Date()) const leaseBefore = new Date(now().getTime() - EFFECTS_LEASE_MS) const stage = await mailRepo.claimMessageEffects(message.id, { leaseBefore }) if (stage === null) return @@ -225,7 +232,7 @@ export async function applyInboundEffects( } else { await onEventReply(container, injected.cleanupRepo, mailRepo, thread, message, stage) } - await mailRepo.markMessageEffectsApplied(message.id) + await mailRepo.markMessageEffectsApplied(message.id, opts.publishedBy) } catch (err) { await mailRepo.releaseMessageEffects(message.id).catch((releaseErr: unknown) => { injected.logger?.warn( diff --git a/services/api/src/services/admin/mail-reply-publish-service.ts b/services/api/src/services/admin/mail-reply-publish-service.ts index 27c4042b..8efb226d 100644 --- a/services/api/src/services/admin/mail-reply-publish-service.ts +++ b/services/api/src/services/admin/mail-reply-publish-service.ts @@ -1,6 +1,11 @@ import { AppError, type PublishMailReplyResponse } from "@civfix/shared" import { domainOf } from "../../adapters/inbound-mail.cf.js" -import type { MailMessageRecord, MailRepository, MailThreadRecord } from "./mail-repository.js" +import type { + MailAuditInput, + MailMessageRecord, + MailRepository, + MailThreadRecord, +} from "./mail-repository.js" export const MAIL_REPLY_NOT_FOUND = "Reply not found on this thread." export const MAIL_REPLY_NOT_PUBLISHABLE = @@ -12,7 +17,11 @@ export interface MailReplyPublishLogger { export interface MailReplyPublishServiceDeps { repo: MailRepository - applyEffects: (thread: MailThreadRecord, message: MailMessageRecord) => Promise + applyEffects: ( + thread: MailThreadRecord, + message: MailMessageRecord, + publishedBy: MailAuditInput, + ) => Promise logger: MailReplyPublishLogger } @@ -37,12 +46,12 @@ export function makeMailReplyPublishService( return message } - async function approve( + function publishAudit( thread: MailThreadRecord, message: MailMessageRecord, actorId: string, - ): Promise { - const approved = await repo.approveWithheldReply(message.id, { + ): MailAuditInput { + return { actorId, action: "mail.reply_published", target: `mail:${thread.id}`, @@ -53,7 +62,18 @@ export function makeMailReplyPublishService( authVerdict: message.authVerdict, fromDomain: domainOf(message.fromAddr), }, - }) + } + } + + async function approve( + thread: MailThreadRecord, + message: MailMessageRecord, + actorId: string, + ): Promise { + const approved = await repo.approveWithheldReply( + message.id, + publishAudit(thread, message, actorId), + ) return approved ?? inboundMessage(thread.id, message.id) } @@ -69,7 +89,7 @@ export function makeMailReplyPublishService( const approved = message.unaffiliated ? await approve(thread, message, actorId) : message try { - await deps.applyEffects(thread, approved) + await deps.applyEffects(thread, approved, publishAudit(thread, approved, actorId)) } catch (err) { deps.logger.warn( { err, threadId, messageId }, diff --git a/services/api/src/services/admin/mail-repository.drizzle.ts b/services/api/src/services/admin/mail-repository.drizzle.ts index 36964b89..34cdaf01 100644 --- a/services/api/src/services/admin/mail-repository.drizzle.ts +++ b/services/api/src/services/admin/mail-repository.drizzle.ts @@ -616,12 +616,30 @@ export function makeDrizzleMailRepository(sql: Sql): MailRepository { ` }, - async markMessageEffectsApplied(id: string): Promise { - await sql` - UPDATE mail_messages - SET effects_applied_at = now() - WHERE id = ${id} AND effects_applied_at IS NULL - ` + async markMessageEffectsApplied(id: string, publishedBy?: MailAuditInput): Promise { + await sql.begin(async (tx) => { + const applied = await tx<{ id: string }[]>` + UPDATE mail_messages + SET effects_applied_at = now() + WHERE id = ${id} AND effects_applied_at IS NULL + RETURNING id + ` + if (applied.length === 0 || publishedBy === undefined) return + const audited = await tx<{ id: string }[]>` + SELECT id FROM audit_log + WHERE action = ${publishedBy.action} + AND target = ${publishedBy.target} + AND meta->>'messageId' = ${id} + LIMIT 1 + ` + if (audited.length > 0) return + await writeAudit(tx, { + actorId: publishedBy.actorId, + action: publishedBy.action, + target: publishedBy.target, + meta: publishedBy.meta ?? null, + }) + }) }, async releaseMessageEffects(id: string): Promise { diff --git a/services/api/src/services/admin/mail-repository.memory.ts b/services/api/src/services/admin/mail-repository.memory.ts index 6f70e929..e4033c25 100644 --- a/services/api/src/services/admin/mail-repository.memory.ts +++ b/services/api/src/services/admin/mail-repository.memory.ts @@ -441,9 +441,18 @@ export class InMemoryMailRepository implements MailRepository { return Promise.resolve() } - markMessageEffectsApplied(id: string): Promise { + markMessageEffectsApplied(id: string, publishedBy?: MailAuditInput): Promise { const message = this.messages.find((m) => m.id === id) - if (message && message.effectsAppliedAt === null) message.effectsAppliedAt = this.nextDate() + if (!message || message.effectsAppliedAt !== null) return Promise.resolve() + message.effectsAppliedAt = this.nextDate() + if (publishedBy === undefined) return Promise.resolve() + const audited = this.audits.some( + (a) => + a.action === publishedBy.action && + a.target === publishedBy.target && + a.meta?.["messageId"] === id, + ) + if (!audited) this.recordAudit(publishedBy) return Promise.resolve() } diff --git a/services/api/src/services/admin/mail-repository.ts b/services/api/src/services/admin/mail-repository.ts index bee19d44..0c78394d 100644 --- a/services/api/src/services/admin/mail-repository.ts +++ b/services/api/src/services/admin/mail-repository.ts @@ -188,7 +188,7 @@ export interface MailRepository { hasSendInFlight(threadId: string): Promise claimMessageEffects(id: string, input: ClaimEffectsInput): Promise setMessageEffectsStage(id: string, stage: number): Promise - markMessageEffectsApplied(id: string): Promise + markMessageEffectsApplied(id: string, publishedBy?: MailAuditInput): Promise releaseMessageEffects(id: string): Promise settleRepliedThread(input: SettleRepliedThreadInput): Promise settleThreadStatus(input: SettleThreadStatusInput): Promise diff --git a/services/api/test/integration/admin-mail-repository.test.ts b/services/api/test/integration/admin-mail-repository.test.ts index 1fce217c..555da0e4 100644 --- a/services/api/test/integration/admin-mail-repository.test.ts +++ b/services/api/test/integration/admin-mail-repository.test.ts @@ -561,4 +561,49 @@ describe.skipIf(!pg)("admin mail repository (integration: real schema)", () => { await h.sql`DELETE FROM reports WHERE id = ${report!.id}` }) + + it("writes an operator's publish audit with the applied mark, once per reply", async () => { + const [report] = await h.sql<{ id: string }[]>` + INSERT INTO reports (idempotency_key, geom, geom_source, category, status, h3_cell, jurisdiction_geoid) + VALUES (gen_random_uuid(), ST_SetSRID(ST_MakePoint(-118.25, 34.05), 4326), 'gps', 'graffiti', + 'published', '8a2a1072b59ffff', ${GEOID}) + RETURNING id + ` + const t = await repo.findOrCreateReportThread(report!.id, { subject: "Publish audit" }) + const inbound = (messageId: string, unaffiliated: boolean) => + repo.insertMessage({ threadId: t.id, direction: "in", messageId, unaffiliated }) + const verified = await inbound("", false) + const withheld = await inbound("", true) + const swept = await inbound("", false) + const audit = (messageId: string) => + ({ + actorId: null, + action: "mail.reply_published", + target: `mail:${t.id}`, + meta: { messageId }, + }) as const + const audited = async (messageId: string) => { + const rows = await h.sql` + SELECT 1 FROM audit_log + WHERE action = 'mail.reply_published' AND meta->>'messageId' = ${messageId} + ` + return rows.length + } + + await repo.markMessageEffectsApplied(verified!.id, audit(verified!.id)) + await repo.markMessageEffectsApplied(verified!.id, audit(verified!.id)) + await repo.approveWithheldReply(withheld!.id, audit(withheld!.id)) + await repo.markMessageEffectsApplied(withheld!.id, audit(withheld!.id)) + await repo.markMessageEffectsApplied(swept!.id) + await repo.markMessageEffectsApplied(swept!.id, audit(swept!.id)) + + expect([ + await audited(verified!.id), + await audited(withheld!.id), + await audited(swept!.id), + ]).toEqual([1, 1, 0]) + expect((await repo.findInboundMessage(t.id, swept!.id))?.effectsAppliedAt).not.toBeNull() + + await h.sql`DELETE FROM reports WHERE id = ${report!.id}` + }) }) diff --git a/services/api/test/unit/admin-activity.test.ts b/services/api/test/unit/admin-activity.test.ts index 507d7f3a..19552c9e 100644 --- a/services/api/test/unit/admin-activity.test.ts +++ b/services/api/test/unit/admin-activity.test.ts @@ -53,7 +53,7 @@ describe("describeAuditAction", () => { expect(describeAuditAction("mail.reply_published_without_text")).toBe( "A city reply was posted without its text", ) - expect(describeAuditAction("mail.reply_published")).toBe("Published a withheld city reply") + expect(describeAuditAction("mail.reply_published")).toBe("Published a city reply") }) }) diff --git a/services/api/test/unit/mail-reply-publish-service.test.ts b/services/api/test/unit/mail-reply-publish-service.test.ts index 252ecf44..9a2060da 100644 --- a/services/api/test/unit/mail-reply-publish-service.test.ts +++ b/services/api/test/unit/mail-reply-publish-service.test.ts @@ -51,8 +51,8 @@ function harness() { const container = { env: {} } as unknown as Container const service = makeMailReplyPublishService({ repo, - applyEffects: (thread, message) => - applyInboundEffects(container, effects, repo, thread, message), + applyEffects: (thread, message, publishedBy) => + applyInboundEffects(container, effects, repo, thread, message, { publishedBy }), logger: { warn: (_obj, msg) => warnings.push(msg) }, }) const withheld = (link: { reportId?: string; cleanupId?: string }) => { @@ -67,7 +67,19 @@ function harness() { }) return { threadId: thread.id, messageId: message.id, message, actorId: OPERATOR_ID } } - return { repo, reports, cleanups, notifier, emitted, chat, warnings, service, withheld } + return { + repo, + reports, + cleanups, + notifier, + emitted, + chat, + warnings, + service, + withheld, + container, + effects, + } } describe("makeMailReplyPublishService", () => { @@ -186,4 +198,51 @@ describe("makeMailReplyPublishService", () => { expect(h.repo.audits).toHaveLength(0) expect([input.message.unaffiliated, loose.message.unaffiliated]).toEqual([true, true]) }) + + it("audits the operator who publishes a verified reply whose effects are still owed", async () => { + const h = harness() + const input = h.withheld({ reportId: REPORT_ID }) + input.message.unaffiliated = false + input.message.authVerdict = "pass" + + expect(await h.service.publish(input)).toEqual({ publication: "published" }) + expect(await h.service.publish({ ...input, actorId: "operator-2" })).toEqual({ + publication: "published", + }) + expect(h.repo.audits).toEqual([ + expect.objectContaining({ + actorId: OPERATOR_ID, + action: "mail.reply_published", + meta: expect.objectContaining({ messageId: input.messageId, authVerdict: "pass" }), + }), + ]) + }) + + it("writes one publish audit row when operators publish the same reply at once", async () => { + const h = harness() + const verified = h.withheld({ reportId: REPORT_ID }) + verified.message.unaffiliated = false + const withheld = h.withheld({ cleanupId: "cleanup-1" }) + + for (const input of [verified, withheld]) { + const racers = ["operator-2", "operator-3"].map((actorId) => + h.service.publish({ ...input, actorId }), + ) + await Promise.all(racers) + expect(input.message.effectsAppliedAt).not.toBeNull() + } + const published = h.repo.audits.map((a) => a.meta?.["messageId"]) + expect(published.sort()).toEqual([verified.messageId, withheld.messageId].sort()) + }) + + it("leaves no operator audit on a reply the sweep published", async () => { + const h = harness() + const input = h.withheld({ reportId: REPORT_ID }) + input.message.unaffiliated = false + const thread = (await h.repo.getThreadRecord(input.threadId))! + await applyInboundEffects(h.container, h.effects, h.repo, thread, input.message) + + expect(await h.service.publish(input)).toEqual({ publication: "published" }) + expect(h.repo.audits).toHaveLength(0) + }) })