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
17 changes: 14 additions & 3 deletions docs/inbound-mail-effects.md
Original file line number Diff line number Diff line change
@@ -1,7 +1,7 @@
# Mail effects and the outbound send triad

**Audience:** internal (engineering). Not served publicly.
**Last updated:** 2026-09-23 (withheld replies flag their thread for review and an operator can publish one; a reply stripped to nothing is flagged too).
**Last updated:** 2026-09-23 (a thread in review stays there while it holds a withheld reply, whatever reply settles it, an operator's included; Mark replied keeps the reply dismissed; the stripped-reply audit row is written once).

An inbound message that correlates to a mail thread can drive **public** effects: a report status
transition, a public `report_timeline` row, a report-chat system message, and a push to the reporter.
Expand Down Expand Up @@ -60,6 +60,16 @@ thread still holds a withheld reply (`hasWithheldReply`): a resend, a follow-up
forward would otherwise drop the reply out of review. A thread with no report or event has nothing
public to publish to, so an unaffiliated reply there is stored without the flag.

A reply whose effects run, verified or published, settles the thread last (stage 4,
`settleRepliedThread`), and an operator's reply settles it the same way once it is delivered
(`settleThreadStatus`). The thread stays `needs_action` when it is already `needs_action` and still
holds a withheld reply; otherwise it becomes `replied`. Both are read under a lock on the thread row, in
one transaction (for an inbound reply, the one that records the stage), so a later verified reply, an
operator's reply, or publishing one of two withheld replies leaves the thread in review. The rule only
keeps a thread in review, it never puts one back: **Mark replied** is how an operator turns a withheld
reply down without publishing it, so a thread already set to `replied` stays `replied` and the reply
stays dismissed, whatever reply settles it next.

An operator publishes a withheld reply with `publishMailReply`
(`POST /admin/mail/:id/messages/:messageId/publish`: operator only, CSRF, 20 a minute per operator).
It is an explicit, audited override of the affiliation gate. One transaction clears `unaffiliated` and
Expand Down Expand Up @@ -135,7 +145,8 @@ If the result is empty (a reply that was nothing but quoted history), stage 2 fa
behavior: the note only, with `body: null`. When the stored body had text and the strip removed all of
it, for instance a one-line reply that mentions our reply address, the thread ends at `needs_action`
instead of `replied` and a `mail.reply_published_without_text` audit row with no actor records it, so an
operator can see that the city's words did not reach the chat and relay them. Like a send failure's, this
operator can see that the city's words did not reach the chat and relay them. The row is written in the
stage 4 transaction, so a re-drive after a crash never writes a second one. Like a send failure's, this
flag clears on the thread's next delivered outbound message.

The chat emitter is injectable (`InboundEffectDeps.chatEmitter`, plumbed through
Expand All @@ -153,7 +164,7 @@ The columns on `mail_messages` (migration `0100`) are a **lease**, not a flag:
|---|---|
| `effects_claimed_at` | A runner holds the message. **Reclaimable**: the sweep re-drives any claim older than `EFFECTS_LEASE_MS` (10 min). A process death between claim and completion — a deploy restart, OOM, the drain watchdog — must not strand the row, which is the exact failure the re-drive exists to prevent. |
| `effects_applied_at` | Set **only** on completion. `IS NULL` is the "still owed" set the partial index serves; the lease comparison stays in the query because `now()` is not `IMMUTABLE`. |
| `effects_stage` | How far the ordered pipeline got: `0` none, `1` timeline, `2` chat, `3` reporter notified. A re-drive resumes from here. |
| `effects_stage` | How far the ordered pipeline got: `0` none, `1` timeline, `2` chat, `3` reporter notified, `4` thread settled. An event reply goes from `1` to `4`. A re-drive resumes from here. |

`applyInboundEffects` (`services/api/src/services/admin/inbound-thread-correlation.ts`) claims, runs the
stages it still owes, then calls `markMessageEffectsApplied`. On a throw it calls
Expand Down
33 changes: 24 additions & 9 deletions services/api/src/services/admin/inbound-thread-correlation.ts
Original file line number Diff line number Diff line change
Expand Up @@ -36,6 +36,7 @@ export const EFFECTS_LEASE_MS = 10 * 60 * 1000
export const EFFECTS_STAGE_TIMELINE = 1
export const EFFECTS_STAGE_CHAT = 2
export const EFFECTS_STAGE_NOTIFIED = 3
export const EFFECTS_STAGE_SETTLED = 4

export const JURISDICTION_REPLY_NOTIFICATION_BODY =
"The city responded. See their reply in the report chat."
Expand Down Expand Up @@ -300,16 +301,24 @@ export async function onJurisdictionReply(
await mailRepo.setMessageEffectsStage(messageId, EFFECTS_STAGE_NOTIFIED)
}

if (chatBody === null && (message.body ?? "").trim() !== "") {
await mailRepo.setThreadStatus(thread.id, "needs_action", {
actorId: null,
action: "mail.reply_published_without_text",
target: `mail:${thread.id}`,
meta: { messageId, reportId },
if (stage < EFFECTS_STAGE_SETTLED) {
const emptied = chatBody === null && (message.body ?? "").trim() !== ""
await mailRepo.settleRepliedThread({
threadId: thread.id,
messageId,
stage: EFFECTS_STAGE_SETTLED,
...(emptied
? {
flag: {
actorId: null,
action: "mail.reply_published_without_text",
target: `mail:${thread.id}`,
meta: { messageId, reportId },
},
}
: {}),
})
return
}
await mailRepo.setThreadStatus(thread.id, "replied")
}

export async function onEventReply(
Expand All @@ -335,7 +344,13 @@ export async function onEventReply(
})
await mailRepo.setMessageEffectsStage(message.id, EFFECTS_STAGE_TIMELINE)
}
await mailRepo.setThreadStatus(thread.id, "replied")
if (stage < EFFECTS_STAGE_SETTLED) {
await mailRepo.settleRepliedThread({
threadId: thread.id,
messageId: message.id,
stage: EFFECTS_STAGE_SETTLED,
})
}
}

export function replyPreview(body: string): string {
Expand Down
67 changes: 55 additions & 12 deletions services/api/src/services/admin/mail-repository.drizzle.ts
Original file line number Diff line number Diff line change
Expand Up @@ -37,6 +37,8 @@ import {
type PendingEffectsQuery,
type RecordEventInput,
type RecordSendFailureInput,
type SettleRepliedThreadInput,
type SettleThreadStatusInput,
type ThreadInit,
} from "./mail-repository.js"
import type {
Expand Down Expand Up @@ -91,6 +93,41 @@ function messageColumns(sql: Queryable): SqlFragment {
`
}

function withheldReplyExpr(sql: Queryable, threadId: string): SqlFragment {
return sql`EXISTS (
SELECT 1
FROM mail_messages m
JOIN mail_threads t ON t.id = m.thread_id
WHERE m.thread_id = ${threadId}
AND m.direction = 'in'
AND m.unaffiliated = true
AND m.effects_applied_at IS NULL
AND (t.report_id IS NOT NULL OR t.cleanup_id IS NOT NULL)
)`
}

async function settleThreadStatusIn(tx: Queryable, input: SettleThreadStatusInput): Promise<void> {
await tx`SELECT id FROM mail_threads WHERE id = ${input.threadId} FOR UPDATE`
const settled = await tx<{ id: string }[]>`
UPDATE mail_threads
SET status = CASE
WHEN ${input.flag !== undefined}
OR (status = 'needs_action' AND ${withheldReplyExpr(tx, input.threadId)})
THEN 'needs_action'
ELSE 'replied'
END
WHERE id = ${input.threadId}
RETURNING id
`
if (settled.length === 0 || input.flag === undefined) return
await writeAudit(tx, {
actorId: input.flag.actorId,
action: input.flag.action,
target: input.flag.target,
meta: input.flag.meta ?? null,
})
}

function threadColumns(sql: Queryable, alias?: string): SqlFragment {
const p = alias === undefined ? sql`` : sql`${sql(alias)}.`
return sql`
Expand Down Expand Up @@ -595,6 +632,23 @@ export function makeDrizzleMailRepository(sql: Sql): MailRepository {
`
},

async settleRepliedThread(input: SettleRepliedThreadInput): Promise<void> {
await sql.begin(async (tx) => {
const advanced = await tx<{ id: string }[]>`
UPDATE mail_messages
SET effects_stage = ${input.stage}
WHERE id = ${input.messageId} AND effects_stage < ${input.stage}
RETURNING id
`
if (advanced.length === 0) return
await settleThreadStatusIn(tx, input)
})
},

async settleThreadStatus(input: SettleThreadStatusInput): Promise<void> {
await sql.begin((tx) => settleThreadStatusIn(tx, input))
},

async findMessagesPendingEffects(input: PendingEffectsQuery): Promise<PendingEffects[]> {
const rows = await sql<PendingEffectsRowSelect[]>`
SELECT m.id, m.thread_id, m.direction, m.from_addr, m.to_addr, m.subject,
Expand Down Expand Up @@ -637,18 +691,7 @@ export function makeDrizzleMailRepository(sql: Sql): MailRepository {
},

async hasWithheldReply(threadId: string): Promise<boolean> {
const rows = await sql<{ ok: boolean }[]>`
SELECT EXISTS (
SELECT 1
FROM mail_messages m
JOIN mail_threads t ON t.id = m.thread_id
WHERE m.thread_id = ${threadId}
AND m.direction = 'in'
AND m.unaffiliated = true
AND m.effects_applied_at IS NULL
AND (t.report_id IS NOT NULL OR t.cleanup_id IS NOT NULL)
) AS ok
`
const rows = await sql<{ ok: boolean }[]>`SELECT ${withheldReplyExpr(sql, threadId)} AS ok`
return rows[0]?.ok ?? false
},

Expand Down
42 changes: 31 additions & 11 deletions services/api/src/services/admin/mail-repository.memory.ts
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,8 @@ import {
type PendingEffectsQuery,
type RecordEventInput,
type RecordSendFailureInput,
type SettleRepliedThreadInput,
type SettleThreadStatusInput,
type ThreadInit,
} from "./mail-repository.js"
import {
Expand Down Expand Up @@ -451,6 +453,24 @@ export class InMemoryMailRepository implements MailRepository {
return Promise.resolve()
}

settleRepliedThread(input: SettleRepliedThreadInput): Promise<void> {
const message = this.messages.find((m) => m.id === input.messageId)
if (!message || message.effectsStage >= input.stage) return Promise.resolve()
message.effectsStage = input.stage
return this.settleThreadStatus(input)
}

settleThreadStatus(input: SettleThreadStatusInput): Promise<void> {
const thread = this.threads.get(input.threadId)
if (!thread) return Promise.resolve()
const needsAction =
input.flag !== undefined ||
(thread.status === "needs_action" && this.holdsWithheldReply(thread.id))
thread.status = needsAction ? "needs_action" : "replied"
this.recordAudit(input.flag)
return Promise.resolve()
}

findMessagesPendingEffects(input: PendingEffectsQuery): Promise<PendingEffects[]> {
const out: PendingEffects[] = []
for (const m of [...this.messages].sort((a, b) => cmpCreated(a, b))) {
Expand All @@ -471,18 +491,18 @@ export class InMemoryMailRepository implements MailRepository {
}

hasWithheldReply(threadId: string): Promise<boolean> {
return Promise.resolve(this.holdsWithheldReply(threadId))
}

private holdsWithheldReply(threadId: string): boolean {
const thread = this.threads.get(threadId)
if (!thread || (thread.reportId === null && thread.cleanupId === null)) {
return Promise.resolve(false)
}
return Promise.resolve(
this.messages.some(
(m) =>
m.threadId === threadId &&
m.direction === "in" &&
m.unaffiliated &&
m.effectsAppliedAt === null,
),
if (!thread || (thread.reportId === null && thread.cleanupId === null)) return false
return this.messages.some(
(m) =>
m.threadId === threadId &&
m.direction === "in" &&
m.unaffiliated &&
m.effectsAppliedAt === null,
)
}

Expand Down
12 changes: 12 additions & 0 deletions services/api/src/services/admin/mail-repository.ts
Original file line number Diff line number Diff line change
Expand Up @@ -190,6 +190,8 @@ export interface MailRepository {
setMessageEffectsStage(id: string, stage: number): Promise<void>
markMessageEffectsApplied(id: string): Promise<void>
releaseMessageEffects(id: string): Promise<void>
settleRepliedThread(input: SettleRepliedThreadInput): Promise<void>
settleThreadStatus(input: SettleThreadStatusInput): Promise<void>
findMessagesPendingEffects(input: PendingEffectsQuery): Promise<PendingEffects[]>
hasWithheldReply(threadId: string): Promise<boolean>
findInboundMessage(threadId: string, messageId: string): Promise<MailMessageRecord | null>
Expand All @@ -200,6 +202,16 @@ export interface ClaimEffectsInput {
leaseBefore: Date
}

export interface SettleThreadStatusInput {
threadId: string
flag?: MailAuditInput
}

export interface SettleRepliedThreadInput extends SettleThreadStatusInput {
messageId: string
stage: number
}

export interface PendingEffectsQuery {
before: Date
leaseBefore: Date
Expand Down
2 changes: 1 addition & 1 deletion services/api/src/services/admin/mail-service.ts
Original file line number Diff line number Diff line change
Expand Up @@ -137,7 +137,7 @@ export function makeMailService(deps: MailServiceDeps): MailService {
toAddr,
audit: { actorId, action: "mail.replied", meta: { to: toAddr } },
})
await repo.setThreadStatus(id, "replied")
await repo.settleThreadStatus({ threadId: id })
await repo.markThreadRead(id)
return requireThreadDTO(id)
},
Expand Down
64 changes: 64 additions & 0 deletions services/api/test/integration/admin-mail-repository.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -334,6 +334,70 @@ describe.skipIf(!pg)("admin mail repository (integration: real schema)", () => {
await h.sql`DELETE FROM reports WHERE id = ${report!.id}`
})

it("settles a replied thread from its withheld replies, once per message", 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: "Settle" })
const inbound = (messageId: string, unaffiliated: boolean) =>
repo.insertMessage({
threadId: t.id,
direction: "in",
messageId,
unaffiliated,
...(unaffiliated ? { threadStatus: "needs_action" as const } : {}),
})
const withheld = await inbound("<settle-w@vendor.example>", true)
const verified = await inbound("<settle-v@lacity.gov>", false)
const flag = {
actorId: null,
action: "mail.reply_published_without_text",
target: `mail:${t.id}`,
} as const
const settle = (messageId: string, withFlag = false) =>
repo.settleRepliedThread({
threadId: t.id,
messageId,
stage: 4,
...(withFlag ? { flag } : {}),
})
const status = async () => (await repo.getThreadRecord(t.id))?.status
const flags = async () => {
const rows = await h.sql`SELECT 1 FROM audit_log WHERE action = ${flag.action}`
return rows.length
}

await settle(verified!.id)
expect(await status()).toBe("needs_action")

await repo.setThreadStatus(t.id, "sent")
await settle(verified!.id)
expect(await status()).toBe("sent")

await repo.approveWithheldReply(withheld!.id, { ...flag, action: "mail.reply_published" })
await settle(withheld!.id, true)
await settle(withheld!.id, true)
expect([await status(), await flags()]).toEqual(["needs_action", 1])

const later = await inbound("<settle-later@lacity.gov>", false)
await settle(later!.id)
expect(await status()).toBe("replied")
expect((await repo.findMessageByMessageId("<settle-later@lacity.gov>"))?.effectsStage).toBe(4)

await inbound("<settle-dismissed@vendor.example>", true)
await repo.settleThreadStatus({ threadId: t.id })
expect(await status()).toBe("needs_action")
await repo.setThreadStatus(t.id, "replied")
const after = await inbound("<settle-after@lacity.gov>", false)
await settle(after!.id)
expect([await status(), await repo.hasWithheldReply(t.id)]).toEqual(["replied", true])

await h.sql`DELETE FROM reports WHERE id = ${report!.id}`
})

it("B4: an EXPIRED claim is reclaimable and keeps the stage it reached", async () => {
const t = await repo.createThread({ subject: "Lease" })
const msg = await repo.insertMessage({
Expand Down
21 changes: 21 additions & 0 deletions services/api/test/unit/admin-mail-repository.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -255,6 +255,27 @@ describe("InMemoryMailRepository: withheld replies", () => {
expect(await repo.hasWithheldReply(thread.id)).toBe(false)
})

it("settles a thread to replied unless it is still in review with a withheld reply", async () => {
const repo = new InMemoryMailRepository()
const thread = repo.seedThread({ reportId: "r1", status: "needs_action" })
repo.seedMessage({ threadId: thread.id, direction: "in", unaffiliated: true })
const status = async () => (await repo.getThreadRecord(thread.id))?.status
const flag = {
actorId: null,
action: "mail.reply_published_without_text",
target: `mail:${thread.id}`,
} as const

await repo.settleThreadStatus({ threadId: thread.id })
expect(await status()).toBe("needs_action")
await repo.setThreadStatus(thread.id, "replied")
await repo.settleThreadStatus({ threadId: thread.id })
expect(await status()).toBe("replied")
await repo.settleThreadStatus({ threadId: thread.id, flag })
expect(await status()).toBe("needs_action")
expect(repo.audits.map((a) => a.action)).toEqual([flag.action])
})

it("withholds nothing on a thread with no report or event", async () => {
const repo = new InMemoryMailRepository()
const thread = repo.seedThread()
Expand Down
Loading
Loading