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
6 changes: 6 additions & 0 deletions docs/inbound-mail-effects.md
Original file line number Diff line number Diff line change
Expand Up @@ -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** —
Expand Down
12 changes: 8 additions & 4 deletions services/api/src/routes/admin/mail.routes.ts
Original file line number Diff line number Diff line change
Expand Up @@ -140,18 +140,22 @@ 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,
}),
() => {
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,
})
},
Expand Down
2 changes: 1 addition & 1 deletion services/api/src/services/admin/activity-service.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
2 changes: 1 addition & 1 deletion services/api/src/services/admin/inbound-sweep.ts
Original file line number Diff line number Diff line change
Expand Up @@ -129,7 +129,7 @@ async function redriveEffects(
mailRepo,
thread,
message,
() => now,
{ now: () => now },
)
effectsRedriven += 1
} catch (err) {
Expand Down
11 changes: 9 additions & 2 deletions services/api/src/services/admin/inbound-thread-correlation.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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<void> {
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
Expand All @@ -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(
Expand Down
34 changes: 27 additions & 7 deletions services/api/src/services/admin/mail-reply-publish-service.ts
Original file line number Diff line number Diff line change
@@ -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 =
Expand All @@ -12,7 +17,11 @@ export interface MailReplyPublishLogger {

export interface MailReplyPublishServiceDeps {
repo: MailRepository
applyEffects: (thread: MailThreadRecord, message: MailMessageRecord) => Promise<void>
applyEffects: (
thread: MailThreadRecord,
message: MailMessageRecord,
publishedBy: MailAuditInput,
) => Promise<void>
logger: MailReplyPublishLogger
}

Expand All @@ -37,12 +46,12 @@ export function makeMailReplyPublishService(
return message
}

async function approve(
function publishAudit(
thread: MailThreadRecord,
message: MailMessageRecord,
actorId: string,
): Promise<MailMessageRecord> {
const approved = await repo.approveWithheldReply(message.id, {
): MailAuditInput {
return {
actorId,
action: "mail.reply_published",
target: `mail:${thread.id}`,
Expand All @@ -53,7 +62,18 @@ export function makeMailReplyPublishService(
authVerdict: message.authVerdict,
fromDomain: domainOf(message.fromAddr),
},
})
}
}

async function approve(
thread: MailThreadRecord,
message: MailMessageRecord,
actorId: string,
): Promise<MailMessageRecord> {
const approved = await repo.approveWithheldReply(
message.id,
publishAudit(thread, message, actorId),
)
return approved ?? inboundMessage(thread.id, message.id)
}

Expand All @@ -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 },
Expand Down
30 changes: 24 additions & 6 deletions services/api/src/services/admin/mail-repository.drizzle.ts
Original file line number Diff line number Diff line change
Expand Up @@ -616,12 +616,30 @@ export function makeDrizzleMailRepository(sql: Sql): MailRepository {
`
},

async markMessageEffectsApplied(id: string): Promise<void> {
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<void> {
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<void> {
Expand Down
13 changes: 11 additions & 2 deletions services/api/src/services/admin/mail-repository.memory.ts
Original file line number Diff line number Diff line change
Expand Up @@ -441,9 +441,18 @@ export class InMemoryMailRepository implements MailRepository {
return Promise.resolve()
}

markMessageEffectsApplied(id: string): Promise<void> {
markMessageEffectsApplied(id: string, publishedBy?: MailAuditInput): Promise<void> {
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()
}

Expand Down
2 changes: 1 addition & 1 deletion services/api/src/services/admin/mail-repository.ts
Original file line number Diff line number Diff line change
Expand Up @@ -188,7 +188,7 @@ export interface MailRepository {
hasSendInFlight(threadId: string): Promise<boolean>
claimMessageEffects(id: string, input: ClaimEffectsInput): Promise<number | null>
setMessageEffectsStage(id: string, stage: number): Promise<void>
markMessageEffectsApplied(id: string): Promise<void>
markMessageEffectsApplied(id: string, publishedBy?: MailAuditInput): Promise<void>
releaseMessageEffects(id: string): Promise<void>
settleRepliedThread(input: SettleRepliedThreadInput): Promise<void>
settleThreadStatus(input: SettleThreadStatusInput): Promise<void>
Expand Down
45 changes: 45 additions & 0 deletions services/api/test/integration/admin-mail-repository.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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("<pub-v@lacity.gov>", false)
const withheld = await inbound("<pub-w@vendor.example>", true)
const swept = await inbound("<pub-s@lacity.gov>", 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}`
})
})
2 changes: 1 addition & 1 deletion services/api/test/unit/admin-activity.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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")
})
})

Expand Down
65 changes: 62 additions & 3 deletions services/api/test/unit/mail-reply-publish-service.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 }) => {
Expand All @@ -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", () => {
Expand Down Expand Up @@ -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)
})
})
Loading