Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
21 commits
Select commit Hold shift + click to select a range
4b48fc9
one cursor module: admin pagination folds into the db cursor helpers
theobong Sep 23, 2026
de761e9
shared lib helpers: hash, bounded concurrency, sleep, time units, pag…
theobong Sep 23, 2026
16a669a
callers adopt the shared helpers for hashing, sleep, concurrency, tim…
theobong Sep 23, 2026
e40480f
drop the temporary mapWithLimit re-export
theobong Sep 23, 2026
8249375
chat notification copy helpers move from routes to services
theobong Sep 23, 2026
e03cc1a
one timeout module, one capped body reader, a shared chat message cor…
theobong Sep 23, 2026
5068964
one queue-name module, pg error codes in db, PR 11 delegates removed,…
theobong Sep 23, 2026
97d9311
db schema table files named after their tables
theobong Sep 23, 2026
bee9274
worker jobs seam module and factories named make*
theobong Sep 23, 2026
4c36b8f
api factories named make*, data builders keep build*
theobong Sep 23, 2026
010a91f
worker-facing persistence files and types named -repository
theobong Sep 23, 2026
52f0d3f
repository part files and binders named after their repositories
theobong Sep 23, 2026
cf5e725
repository contracts live in <name>-repository.ts
theobong Sep 23, 2026
233a4d1
repository contracts move out of services and types modules
theobong Sep 23, 2026
3bef075
repository contracts move out of the pg implementations
theobong Sep 23, 2026
e0e150d
worker-facing repository contracts get type-only subpaths
theobong Sep 23, 2026
0bcdc08
test-only in-memory repositories move to test helpers
theobong Sep 23, 2026
f951124
worker repository methods every implementation has are required
theobong Sep 23, 2026
0ec5bd6
comments follow the consolidated helpers; ADMIN_DEFAULT_LIMIT is modu…
theobong Sep 23, 2026
6ed89bd
org and team rules leave the contract files; comments and README foll…
theobong Sep 23, 2026
f2ac768
merge 11-sql-repositories
theobong Sep 23, 2026
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
The table of contents is too big for display.
Diff view
Diff view
  •  
  •  
  •  
16 changes: 12 additions & 4 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -125,10 +125,18 @@ postgres-js client, the R2 adapter, or the GlitchTip reporter. It depends on `@c
`services/api/package.json`:

- `@civfix/api/db` - the schema barrel + `makeDb` + `Db`/`Sql` types.
- `@civfix/api/media-repo` - the richer media-worker persistence seam (`MediaWorkerRepo`:
find / applyResult / insertAbuseFlag / findOrphans / deleteById), its Drizzle impl, the
`MEDIA_CHECKS_JOB` name, and `ensureNextMonthChatPartition` (so partition bounds/naming have one
source shared with the migrations).
- `@civfix/api/media-repo` - the Drizzle impl of the richer media-worker persistence seam and
`ensureNextMonthChatPartition` (so partition bounds/naming have one source shared with the
migrations).
- `@civfix/api/media-worker-repository`, `@civfix/api/anon-hold-release-repository`,
`@civfix/api/inbound-retention-repository` - the type-only repository contracts the worker's jobs
depend on (`MediaWorkerRepository`: find / applyResult / insertAbuseFlag / findOrphans /
deleteById; `AnonHoldReleaseRepository`; `InboundRetentionRepository`). Their Drizzle impls stay
on `media-repo`, `anon-hold-repo` and `inbound-retention-repo`.
- `@civfix/api/queue-names` - every pg-boss queue name and the shared queue policy.
- `@civfix/api/time`, `@civfix/api/concurrency`, `@civfix/api/env-parsers`, `@civfix/api/timeout`,
`@civfix/api/capped-body` - import-free helpers the worker shares with the api (time units,
`mapWithLimit`, `parseBool`, `settleWithin`, the size-capped stream reader).
- `@civfix/api/adapters/storage`, `@civfix/api/adapters/abuse-checks`, `@civfix/api/errors`,
`@civfix/api/migrate` - the R2 adapter, the real AbuseChecks adapter, the GlitchTip helper, and the
migration runner (the last reused only by the Docker-gated worker integration harness).
Expand Down
2 changes: 1 addition & 1 deletion docs/erasure-behavior.md
Original file line number Diff line number Diff line change
Expand Up @@ -215,7 +215,7 @@ tombstone (`softDeleteAndAnonymize`, `services/api/src/auth/pg-stores.ts`):
`total_hours`, `entry_count`, the period bounds and `document_sha256`.
3. **Best-effort deletes the R2 object** for each row, after the commit. Same
primitive the holder's own `revoke()` uses. Wired in production through
`buildAuthServicesFromContainer`; if the object store is not wired the row
`makeAuthServicesFromContainer`; if the object store is not wired the row
scrub still happens and the skipped object is logged.

Already-revoked rows are included in the scrub: their PII is no less PII, and
Expand Down
2 changes: 1 addition & 1 deletion docs/out-of-band-indexes.md
Original file line number Diff line number Diff line change
Expand Up @@ -172,7 +172,7 @@ through the exported `explainFeedCandidates` and requires
### `media_assets_orphan_sweep_idx` (migration 0098, audit H13)

Back the hourly orphan sweep's candidate scan (`findOrphans` /`deleteOrphan` in
`services/api/src/services/media-worker-repo.ts`), which looks for media rows
`services/api/src/services/media-worker-repository.drizzle.ts`), which looks for media rows
with no binding older than the orphan TTL. Without it every run sequentially
scans `media_assets`.

Expand Down
2 changes: 1 addition & 1 deletion scripts/check-dynamic-sql.mjs
Original file line number Diff line number Diff line change
Expand Up @@ -9,7 +9,7 @@ const ALLOW_DYNAMIC_SQL = new Set([
"services/api/src/db/migrate.ts:unsafe",
"services/api/src/db/backfill-keyset.ts:unsafe",
"services/api/src/db/sql/jurisdiction.ts:unsafe",
"services/api/src/services/media-worker-repo.ts:unsafe",
"services/api/src/services/media-worker-repository.drizzle.ts:unsafe",
])

const NOT_SQL = new Set(["services/media-worker/src/sandbox/phash.ts:raw"])
Expand Down
15 changes: 12 additions & 3 deletions services/api/package.json
Original file line number Diff line number Diff line change
Expand Up @@ -10,17 +10,26 @@
".": "./dist/server.js",
"./server": "./dist/server.js",
"./db": "./src/db/worker-db.ts",
"./media-repo": "./src/services/media-worker-repo.ts",
"./media-repo": "./src/services/media-worker-repository.drizzle.ts",
"./media-worker-repository": "./src/services/media-worker-repository.ts",
"./inbound-retention-repo": "./src/services/admin/inbound-retention-repository.drizzle.ts",
"./inbound-retention-repository": "./src/services/admin/inbound-retention-repository.ts",
"./retention-repo": "./src/services/retention-repository.drizzle.ts",
"./geocode-cache": "./src/services/geocode-cache.ts",
"./anon-hold-release": "./src/services/anon-hold-release.ts",
"./anon-hold-repo": "./src/services/anon-hold-release-repo.drizzle.ts",
"./anon-hold-repo": "./src/services/anon-hold-release-repository.drizzle.ts",
"./anon-hold-release-repository": "./src/services/anon-hold-release-repository.ts",
"./adapters/storage": "./src/adapters/storage.r2.ts",
"./adapters/storage-local": "./src/adapters/storage.local.ts",
"./adapters/abuse-checks": "./src/adapters/abuse-checks.ts",
"./errors": "./src/errors/glitchtip.ts",
"./migrate": "./src/db/migrate.ts"
"./migrate": "./src/db/migrate.ts",
"./env-parsers": "./src/env/parsers.ts",
"./concurrency": "./src/lib/concurrency.ts",
"./time": "./src/lib/time.ts",
"./capped-body": "./src/lib/capped-body.ts",
"./timeout": "./src/lib/timeout.ts",
"./queue-names": "./src/lib/queue-names.ts"
},
"scripts": {
"build": "tsup",
Expand Down
2 changes: 1 addition & 1 deletion services/api/scripts/render-email-gallery.ts
Original file line number Diff line number Diff line change
Expand Up @@ -14,7 +14,7 @@ import {
orgVerificationDecisionEmailVars,
} from "../src/services/host/organization-service.js"
import { teamInviteEmailVars } from "../src/services/host/host-team-service.js"
import type { AdminReportRecord } from "../src/services/admin/admin-report-types.js"
import type { AdminReportRecord } from "../src/services/admin/admin-report-repository.js"

interface GalleryEntry {
name: string
Expand Down
12 changes: 3 additions & 9 deletions services/api/src/adapters/chat-send-dedupe.redis.ts
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,8 @@ import {
type SendDedupeStore,
type SendReservation,
} from "../ws/send-resilience.js"
import { unrefSleep } from "../lib/sleep.js"
import { MS_PER_SECOND } from "../lib/time.js"

export interface RedisSendDedupeOptions {
ttlSeconds?: number
Expand All @@ -16,18 +18,10 @@ export interface RedisSendDedupeOptions {
sleep?: (ms: number) => Promise<void>
}

const MS_PER_SECOND = 1000

function wholeSecondsToMs(seconds: number): number {
return Math.max(1, Math.ceil(seconds)) * MS_PER_SECOND
}

const realSleep = (ms: number): Promise<void> =>
new Promise((resolve) => {
const timer = setTimeout(resolve, ms)
if (typeof timer.unref === "function") timer.unref()
})

// Every Redis failure fails open: dedupe only suppresses retried duplicates, so an outage must not
// block chat sends, and the shared client's error handler already logs the outage itself.
const OPEN: SendReservation = { state: "open" }
Expand All @@ -44,7 +38,7 @@ export class RedisSendDedupeStore implements SendDedupeStore {
this.ttlMs = wholeSecondsToMs(opts.ttlSeconds ?? SEND_DEDUPE_TTL_SECONDS)
this.pendingTtlMs = wholeSecondsToMs(opts.pendingTtlSeconds ?? SEND_DEDUPE_PENDING_TTL_SECONDS)
this.inFlightAttempts = opts.inFlightAttempts ?? SEND_DEDUPE_INFLIGHT_ATTEMPTS
this.sleep = opts.sleep ?? realSleep
this.sleep = opts.sleep ?? unrefSleep
}

async reserve(key: string): Promise<SendReservation> {
Expand Down
2 changes: 1 addition & 1 deletion services/api/src/adapters/chat-service.ws.ts
Original file line number Diff line number Diff line change
Expand Up @@ -8,7 +8,7 @@ import type { ChatMessageDTO, WsServerMessage } from "@civfix/shared"
import { randomUUID } from "node:crypto"
import { chatChannel, type ChatPubSub } from "./chat-pubsub.js"
import { RefCountedSubscriptions } from "./ref-counted-subscriptions.js"
import type { ChatRepository } from "../services/chat-repository.drizzle.js"
import type { ChatRepository } from "../services/chat-repository.js"

export interface WsChatServiceDeps {
repo: ChatRepository
Expand Down
6 changes: 2 additions & 4 deletions services/api/src/adapters/geocoder.tiger.ts
Original file line number Diff line number Diff line change
Expand Up @@ -14,10 +14,8 @@
import { AppError } from "@civfix/shared"
import type { Geocoder } from "@civfix/shared/interfaces"
import type { Sql } from "../db/client.js"
import {
makeDrizzleJurisdictionRepository,
type JurisdictionRepository,
} from "../services/jurisdiction-repository.drizzle.js"
import { makeDrizzleJurisdictionRepository } from "../services/jurisdiction-repository.drizzle.js"
import type { JurisdictionRepository } from "../services/jurisdiction-repository.js"

/** A fixed public Census code list, so it lives in-process rather than in a lookup table. */
const STATE_FIPS_TO_USPS: Readonly<Record<string, string>> = {
Expand Down
28 changes: 8 additions & 20 deletions services/api/src/adapters/http-fetch.ts
Original file line number Diff line number Diff line change
@@ -1,3 +1,5 @@
import { concatChunks, readCappedChunks, type CappedChunks } from "../lib/capped-body.js"

export type FetchJsonResult<T> =
| { ok: true; status: number; json: T }
| { ok: false; kind: "http"; status: number }
Expand Down Expand Up @@ -61,33 +63,19 @@ export async function fetchJsonWithTimeout<T>(
}
}
const reader = body.getReader()
const chunks: Uint8Array[] = []
let total = 0
let read: CappedChunks | null
try {
for (;;) {
const { done, value } = await reader.read()
if (done) break
if (value) {
total += value.byteLength
if (total > maxBytes) {
controller.abort()
return { ok: false, kind: "body", status, error: new JsonBodyTooLargeError(maxBytes) }
}
chunks.push(value)
}
}
read = await readCappedChunks(reader, maxBytes, () => controller.abort())
} catch (error) {
return { ok: false, kind: "body", status, error }
} finally {
reader.releaseLock?.()
}
if (read === null) {
return { ok: false, kind: "body", status, error: new JsonBodyTooLargeError(maxBytes) }
}
try {
const buf = new Uint8Array(total)
let offset = 0
for (const c of chunks) {
buf.set(c, offset)
offset += c.byteLength
}
const buf = concatChunks(read.chunks, read.total)
return { ok: true, status, json: JSON.parse(new TextDecoder().decode(buf)) as T }
} catch (error) {
return { ok: false, kind: "body", status, error }
Expand Down
6 changes: 2 additions & 4 deletions services/api/src/adapters/inbound-mail.cf.ts
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@ import type {
ParsedMailAddress,
ParsedMailAttachment,
} from "@civfix/shared/interfaces"
import { collapseWhitespace } from "@civfix/shared"
import type { AddressObject, Attachment, EmailAddress, HeaderLines } from "mailparser"
import { getDomain } from "tldts"
import { domainOfOrNull } from "./mail-text.js"
Expand Down Expand Up @@ -264,10 +265,7 @@ function singleFromMailbox(
}

function headerLineValue(header: HeaderLines[number]): string {
return header.line
.slice(header.line.indexOf(":") + 1)
.replace(/\s+/g, " ")
.trim()
return collapseWhitespace(header.line.slice(header.line.indexOf(":") + 1))
}

function toAddress(value: EmailAddress): ParsedMailAddress {
Expand Down
55 changes: 34 additions & 21 deletions services/api/src/adapters/jobs.pgboss.ts
Original file line number Diff line number Diff line change
@@ -1,20 +1,32 @@
import type PgBoss from "pg-boss"
import type { Jobs, EnqueueOptions, JobHandler } from "@civfix/shared/interfaces"
import { REGISTRATION_QUEUE_NAMES } from "../services/host/registration-queues.js"
import { COMMS_QUEUE_NAMES } from "../services/host/broadcast-queues.js"
import type { JobHandlerArgWithAttempt } from "../services/job-attempt.js"
import { MEDIA_CHECKS_JOB } from "../services/media-intake-service.js"
import { JURISDICTION_DISCOVERY_JOB } from "../services/jurisdiction-service.js"
import { OUTREACH_DIGEST_JOB } from "../services/admin/jurisdiction-contacts-types.js"
import { INBOUND_SWEEP_JOB } from "../services/admin/inbound-jobs.js"
import { REPORT_AUTOFORWARD_JOB } from "../services/report-service.types.js"
import { DATA_EXPORT_JOB } from "../services/data-export-jobs.js"
import { CLEANUP_CANCEL_FANOUT_JOB } from "../services/cleanup-notifications.js"
import {
ANON_HOLD_RELEASE_JOB,
BROADCAST_CHUNK_JOB,
BROADCAST_PLAN_JOB,
BROADCAST_SCHEDULE_SWEEP_JOB,
CHAT_ROOM_FANOUT_JOB,
CHECKIN_NOSHOW_SWEEP_JOB,
CLEANUP_CANCEL_FANOUT_JOB,
CLEANUP_GUEST_UPDATE_FANOUT_JOB,
DATA_EXPORT_JOB,
EVENT_METRICS_ROLLUP_JOB,
EVENT_REMINDERS_SWEEP_JOB,
GUEST_RETENTION_SWEEP_JOB,
} from "../services/guest-rsvp-service.js"
import { CHAT_ROOM_FANOUT_JOB } from "../services/chat-fanout-jobs.js"
HOST_EXPORT_JOB,
HOST_EXPORT_REAP_JOB,
HOST_RETENTION_SWEEP_JOB,
INBOUND_SWEEP_JOB,
JURISDICTION_DISCOVERY_JOB,
MEDIA_CHECKS_JOB,
OUTREACH_DIGEST_JOB,
REPORT_AUTOFORWARD_JOB,
SHARED_QUEUE_POLICY,
WAITLIST_EXPIRE_SWEEP_JOB,
WAITLIST_PROMOTE_JOB,
} from "../lib/queue-names.js"
import { SECONDS_PER_MINUTE } from "../lib/time.js"

export interface PgBossJobsLogger {
error(obj: unknown, msg?: string): void
Expand All @@ -32,13 +44,6 @@ export interface PgBossJobsConfig {
logger?: PgBossJobsLogger
}

// Enqueued by the claim route and worked by the media worker, which owns the exported constant.
const ANON_HOLD_RELEASE_JOB = "anon.hold.release"

// The media worker creates media.checks and anon.hold.release with this same policy. pg-boss keeps the
// last createQueue/updateQueue policy, so a mismatch silently breaks singletonKey dedup on those queues.
const SHARED_QUEUE_POLICY = "short"

export const API_QUEUE_NAMES = [
MEDIA_CHECKS_JOB,
JURISDICTION_DISCOVERY_JOB,
Expand All @@ -51,15 +56,23 @@ export const API_QUEUE_NAMES = [
CLEANUP_GUEST_UPDATE_FANOUT_JOB,
GUEST_RETENTION_SWEEP_JOB,
CHAT_ROOM_FANOUT_JOB,
...REGISTRATION_QUEUE_NAMES,
...COMMS_QUEUE_NAMES,
WAITLIST_PROMOTE_JOB,
WAITLIST_EXPIRE_SWEEP_JOB,
CHECKIN_NOSHOW_SWEEP_JOB,
BROADCAST_PLAN_JOB,
BROADCAST_CHUNK_JOB,
BROADCAST_SCHEDULE_SWEEP_JOB,
EVENT_REMINDERS_SWEEP_JOB,
EVENT_METRICS_ROLLUP_JOB,
HOST_EXPORT_JOB,
HOST_EXPORT_REAP_JOB,
HOST_RETENTION_SWEEP_JOB,
] as const

type ApiQueueName = (typeof API_QUEUE_NAMES)[number]

type QueueRetryPolicy = Required<Pick<PgBoss.Queue, "retryLimit" | "retryDelay" | "retryBackoff">>

const SECONDS_PER_MINUTE = 60
const DATA_EXPORT_RETRY_LIMIT = 10

// A data export that fails on a mail credential or approved-sender fault has to wait for an operator to
Expand Down
3 changes: 3 additions & 0 deletions services/api/src/adapters/mailer-defaults.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,3 @@
// A leaf module so env/core-env can share the default without importing mailer.oci, which pulls the
// native argon2 binding in through auth/otp.
export const OCI_MAILER_DEFAULT_TIMEOUT_MS = 15_000
6 changes: 2 additions & 4 deletions services/api/src/adapters/mailer.oci.ts
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,8 @@ import { renderEmailBody } from "./email-layout.js"
import { renderMessage } from "../i18n/renderMessage.js"
import { resolveLocale, type Locale } from "../i18n/locales.js"
import { OTP_TTL_SECONDS } from "../auth/otp.js"
import { SECONDS_PER_MINUTE } from "../lib/time.js"
import { OCI_MAILER_DEFAULT_TIMEOUT_MS } from "./mailer-defaults.js"

const CRLF_RE = /[\r\n\0]/
const CRLF_GLOBAL_RE = /[\r\n\0]/g
Expand Down Expand Up @@ -60,12 +62,8 @@ const consoleLogger: OciMailerLogger = {
warn: (obj, msg) => console.warn(msg ?? "", obj),
}

export const OCI_MAILER_DEFAULT_TIMEOUT_MS = 15_000

const SMTPS_IMPLICIT_TLS_PORT = 465

const SECONDS_PER_MINUTE = 60

const DEFAULT_EVENT_TITLE = "the event"

const GUEST_OTP_DEFAULT_MINUTES = "5"
Expand Down
5 changes: 1 addition & 4 deletions services/api/src/adapters/push-expo.ts
Original file line number Diff line number Diff line change
@@ -1,5 +1,6 @@
import type { PushLogger, PlatformDispatcher } from "./push-sender.js"
import { fetchJsonWithTimeout, type FetchJsonResult } from "./http-fetch.js"
import { sleep } from "../lib/sleep.js"

export interface ExpoPushConfig {
accessToken?: string
Expand All @@ -26,10 +27,6 @@ function isRetryableStatus(status: number): boolean {
return status === HTTP_TOO_MANY_REQUESTS || status >= HTTP_SERVER_ERROR_MIN
}

function sleep(ms: number): Promise<void> {
return new Promise((resolve) => setTimeout(resolve, ms))
}

export function isExpoPushToken(token: string): boolean {
return EXPO_TOKEN_PREFIXES.some((prefix) => token.startsWith(prefix))
}
Expand Down
6 changes: 3 additions & 3 deletions services/api/src/adapters/push-sender.ts
Original file line number Diff line number Diff line change
@@ -1,4 +1,3 @@
import { createHash } from "node:crypto"
import type { PushSender, PushPayload, PushPlatform } from "@civfix/shared/interfaces"
import type { Db } from "../db/client.js"
import { pushTokens } from "../db/schema/push_tokens.js"
Expand All @@ -7,7 +6,8 @@ import { makeApnsDispatcher } from "./push-apns.js"
import { makeFcmDispatcher } from "./push-fcm.js"
import { makeWebPushDispatcher } from "./push-webpush.js"
import { makeExpoDispatcher, isExpoPushToken, type ExpoPushConfig } from "./push-expo.js"
import { mapWithLimit } from "../services/media-presign.js"
import { mapWithLimit } from "../lib/concurrency.js"
import { sha256HexSync } from "../lib/hash.js"
import type { CounterStore } from "../abuse/counter-store.js"
import type { PushAddressResolver } from "../services/push-token-policy.js"

Expand Down Expand Up @@ -43,7 +43,7 @@ export interface PushLogger {
const LOG_HASH_HEX_CHARS = 12

export function hashForLog(value: string): string {
return createHash("sha256").update(value).digest("hex").slice(0, LOG_HASH_HEX_CHARS)
return sha256HexSync(value).slice(0, LOG_HASH_HEX_CHARS)
}

const ACTIVE_TOKEN_SCAN_CAP_PER_USER = 20
Expand Down
Loading
Loading