Files
accounted/lib/webhooks/dispatcher.ts
T
MattssonandClaude Fable 5 f24b26a139 fix: similar-sweep currency remediation, security hardening and v1 API fixes (#1215)
* fix(security): gate replace_sie_import behind owner/admin membership

The RPC was SECURITY DEFINER with EXECUTE granted to PUBLIC and anon, no
company_members lookup, no auth.uid() reference and no unauthorized raise,
while setting gnubok.allow_delete to disarm the BFL immutability and
retention triggers. Any caller holding a company_id and an import id could
hard delete another tenant's verifikationer. Confirmed live in production.

Applies the same fail closed owner/admin guard that undo_sie_import already
carries (migration 20260624120000), resolving the actor from
COALESCE(p_user_id, auth.uid()) so it denies when the role is NULL, then
revokes EXECUTE from PUBLIC and anon. search_path and the raised
statement_timeout are restated, since CREATE OR REPLACE drops settings that
are not repeated.

userId is a required parameter on replaceSIEImport: the service client has a
NULL auth.uid(), so a caller without an explicit actor now fails to compile
rather than hitting the closed gate at runtime.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>

* fix(security): validate arcim OAuth callback state server side

The callback route is skipAuth and decoded the state parameter as plain
base64url JSON, trusting consentId and provider from it. A one time code was
minted at flow start and never read. An unauthenticated attacker who learned
a consent id could run an OAuth flow on their own provider account and post
the callback with a forged state, landing their tokens on another tenant's
consent, so the victim's next migration imported the attacker's ledger.

State is now an opaque randomBytes(32) pointer to a provider_otc row,
consumed by a single atomic UPDATE guarded on used_at IS NULL and
expires_at, so a replay loses the row lock race and updates nothing.
provider is read from provider_consents rather than trusted from the client.
provider_otc already existed for exactly this purpose and was never wired up.

Also scopes getConsent to an owning company, closing a cross tenant status
oracle where the preview and migrate paths echoed a consent's status before
the scoped check ran.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>

* fix(security): scope documents storage to company_id (phase A)

The documents bucket policies matched on auth.uid(), and upload keys were
documents/{userId}/..., so company membership was never consulted. Removing a
member revoked nothing: their session still authenticated and they kept
direct Storage read access to every receipt, supplier invoice and bank
statement they had uploaded. The same bug was fixed for sie-files in
20260416120000; this bucket was left behind.

Phase A is additive. Company scoped policies are added alongside the
uploader scoped ones, uploads move to documents/{companyId}/{userId}/..., and
reads accept either layout so nothing breaks mid migration. Phase C, which
drops the old policies, is gated on the backfill reporting zero remaining
legacy prefix objects.

The policy compares the company segment as text rather than casting to uuid
the way sie-files does: this bucket holds keys whose second segment is not a
uuid (MCP audit packages), and Postgres does not guarantee the bucket prefix
qual runs before the cast, so a planner reordering would raise 22P02 and fail
the whole query instead of filtering the row out.

deleteDocument now removes both candidate keys. Removing only the stored
pointer would leave a readable orphan copy of a document the user asked to
erase.

The backfill script is included but has never been run. It defaults to dry
run, refuses .env.local by name, and verifies each copy is readable and
SHA-256 identical before repointing the row.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>

* fix(security): enforce events:read scope and membership on /api/events

This was the only one of the three validateApiKey call sites with no
downstream guard: v1 and the MCP server both check scope and re-verify
company membership, this route did neither. An events:read scope existed and
was documented as gating the endpoint but was never called, so a legacy key
falling back to DEFAULT_SCOPES read the full log. The bound company id went
straight from the api_keys row into a service role query, so a key whose user
had been removed from the company kept reading.

Adds the scope check before any database access, re-verifies company_members
with archived_at IS NULL, honours test mode by stamping X-Gnubok-Mode instead
of ignoring it, applies minimisePayload so the pull surface can never return
a wider payload than the push surface, and replaces the three flat error
strings with the canonical envelope.

Test key reads are served rather than blocked: TEST_KEY_WRITE_BLOCKED is
gated on mutations in with-api-v1, so a read gets the same treatment as every
other v1 read endpoint.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>

* perf(bookkeeping): sweep remaining journal_entries!inner embeds

A previous refactor removed this pattern from lib/reports and introduced
fetchEntryLines, but the class was never swept. Seventeen sites remained and
had become the top application consumer of production database time:
measured across the resulting query shapes, 32,694 calls and 25,848 seconds
of execution, mean 790ms, with shapes averaging 2.6s and 3.0s and maxing at
7,962ms against the 8s statement_timeout, which surfaced to users as 500s on
the booking path.

PostgREST compiles an embed with filters on the embedded side into a
correlated INNER JOIN LATERAL with a parameterized LIMIT, which stops
Postgres reordering the join, so each query walked the whole
journal_entry_lines table across all tenants. Driving from the entries side
instead turns that into two indexed round trips.

Converted sites keep their existing shape: the helper reattaches the parent
entry under the same key the embed produced. Several conversions also remove
a latent silent truncation where an unpaginated query was capped at
PostgREST's 1000 row ceiling.

Two deliberate exceptions. The free text ilike legs of the MCP display query
stay on the embed, because each is capped at legLimit and that cap drives the
truncation contract the tool reports, while the helper is unbounded. The
accounts route moves to the existing get_account_usage_counts RPC instead,
since its embed was a head count and the helper returns rows.

commitEntry's write path is untouched: the change there is confined to the
read query of the pre-commit dimension rule check.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>

* fix(api): anchor v1 list cursors on created_at

Page two returned page one, forever, while still advertising a fresh
next_cursor. The three routes sorted by and encoded a Postgres date column,
which serializes as YYYY-MM-DD, but decodeDefaultCursor validates the cursor
timestamp as full ISO-8601 and returned null, so the keyset filter was never
applied and has_more never went false. An integrator syncing verifikat looped
on the newest rows indefinitely.

The transactions route already solved this and its comment names the trap;
the fix was never ported. All three now order and encode on created_at with
an id tie break, matching the transactions keyset predicate exactly.
ISO_TIMESTAMP is deliberately left alone: relaxing it would silently change
sort semantics on the route that currently works.

Default ordering therefore moves from business date to insert order. Every
business date is still on the row, and the invoices list gains date_from and
date_to filters so a date range is still reachable; the other two already had
them.

The tests use an in-memory PostgREST that actually evaluates the filters,
because the repo's pass-through mock cannot catch this class of bug: the bug
is that the filter is never sent. They walk to exhaustion with a hard
iteration cap, so an unterminated walk fails instead of hanging.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>

* fix(api): separate dry run from commit in the idempotency hash

The request hash was built from url.pathname, which excludes the query
string, so a dry run and its commit hashed identically. Following the flow
documented in dry-run.ts, re-issuing the request with the same
Idempotency-Key returned the cached preview with Idempotent-Replayed set and
wrote nothing, while reporting 200. An agent or integrator saw success for a
write that never happened.

dry_run is folded into the hash only when true, not as an unconditional
boolean. Including it as false would change the hash of every ordinary write,
and with a 24h idempotency TTL any key in flight across the deploy would fail
the request_hash comparison and 409 on a legitimate retry. Both hash call
sites now go through one shared helper so they cannot drift into a permanent
cache miss, and dry run responses are no longer stored at all.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>

* ci: install the Bedrock SDK out of tree in the compliance review

The Swedish accounting compliance gate had failed ten consecutive runs and so
was posting nothing. With --no-package-lock npm discarded the lockfile and
re-resolved the whole tree from package.json, floating @hookform/resolvers to
5.4.3, whose valibot ^1 peer conflicts with the pinned valibot 0.39.0.

Installing into the parent of the checkout resolves only that one package, so
an unrelated peer conflict can never take the gate down again. Node still
finds it because ESM bare specifiers walk up parent node_modules; NODE_PATH
would not have worked, as it is CommonJS only. --legacy-peer-deps was
rejected because it masks future genuine peer conflicts and still reifies the
full tree.

The same step's SDK version is aligned from 0.31.0 back to the 0.29.1 that
package.json and check:guards enforce after the streaming outage. That drift
went unnoticed because the pin guard only inspects package.json and the
lockfile, never workflow files.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>

* build(docker): generate crontabs from vercel.json

vercel.json defines 16 cron jobs; both Docker crontabs carried 9, and were
byte identical to each other. Self hosted deployments therefore never sent
recurring invoices, never dispatched webhooks and never cleaned up
idempotency keys. tax-deadlines also ran once a year on 2 January instead of
daily, and documents/verify weekly instead of daily.

Extension crons are included rather than excluded. The Dockerfile copies the
whole tree before building, so every extension cron route is compiled into
the image regardless of the enabled preset, and each returns 200 when its
extension is unconfigured, so curl -sf logs no failure. Two such entries were
already present in the crontab for extensions absent from the preset, which
settles the intent.

documents/verify is treated as drift rather than a self hosted concession:
the weekly cadence was present in the hosted crontab too, and the run is
capped at 200 documents walking a nulls-first queue, so weekly drains the
integrity queue seven times slower on a check that exists for BFL retention.

webhooks/dispatch keeps its per minute cadence, adding 1,440 requests a day
on self hosted. A gentler tick would silently stretch the first retry, since
the retry ladder opens at 60 seconds. SCHEDULE_OVERRIDES is the one line
place to change that.

A parity test asserts the path sets match minus a documented exclusion list,
and ratchets three cron routes that are currently scheduled nowhere so they
are named rather than silently rotting.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>

* chore(observability): add a provider agnostic error sink

There is no error tracking in this codebase: logs go to console and Vercel
retention and nowhere else, nothing alerts on the 16 cron jobs, and seven
code comments across lib, app, components and extensions asserted that Sentry
captures errors when Sentry is not a dependency. The two most recent bug
fixes on this repo were both discovered by customer email.

This adds the sink, not a vendor. No dependency is taken: the interface has a
no-op default and a registration point, so behaviour is unchanged until an
adapter is registered. Releases are tagged from the build id already inlined
by next.config.ts.

Redaction moved out of lib/logger.ts into a leaf module that both the logger
and the sink import, so there is one denylist and no path from application
data to a third party can skip the personnummer regex, including direct sink
calls that bypass the logger. That matters here because these logs carry
personnummer and financial data.

verifyCronSecret now reports its own 401s, which covers all 16 jobs without
touching a route file and catches the case where CRON_SECRET is rotated
without updating the scheduler and every job silently 401s forever. The
threshold is one failure rather than the backup alert's three: suppressing
the first occurrence is precisely how an outage stays invisible.

The seven misleading comments are corrected to describe what the code
actually does, including the two cases that still are not covered: the client
side one, since the sink is server side, and a warn level call that is not
forwarded.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>

* fix: remediate the 2026-07-26 similar-sweep findings across all surfaces

Resolves the ~150-finding sweep (dev_docs/similar-sweep-2026-07-26.md) with
one agent per finding; every behavioural fix carries a regression test proven
to fail at HEAD. Full status, corrections to the sweep, refusals and open
decisions in dev_docs/similar-sweep-2026-07-26-remediation-status.md.

Structural roots closed:
- resolveSekAmountOrNull(): honest SEK resolution refuses instead of booking
  1:1; four duplicated toSek closures now refuse via INVOICE_FX_RATE_MISSING
- ledger-line-amount.ts: journal_entry_lines.currency labels the document,
  not the amount; SQL pre-filter decoy proven and fixed
- sparse-patch.ts: .partial() does not strip .default() in Zod 4.4.3; the
  exploitable salary payslip-line PATCH and KPI preferences sinks fixed
- tests/schema: migration-replay phantom-column guard (13k+ refs, closed
  CHECK sets, onConflict targets); found 28 real defects, all fixed, all
  four baselines now empty
- three new ratchet guards: sek-labelled-amount, cross-extension-import,
  ungated-extension-route

Highlights: lawful VAT-rate set on all seven invoice surfaces (ML 6 kap),
RC input VAT mismatch wired on web + both MCP callers, missing-underlag
resource delegates to the shared RPC predicate, push-notifications consent
polarity fail-closed, deadlines undo honours requested state, silent-failure
and read-side-fabrication classes fixed across settings/KPI/inbox/Stripe/
Arcim/kassaflodesanalys, error-envelope stringification fixed at 10+ sites
with isSwedishUserMessage extended.

Also includes the parallel session's MCP invoice tools (update_invoice,
recurring schedules, invoice deliveries) which share files with the sweep
work and are verified green together.

13 new migrations are NOT applied anywhere; they apply via branch merge.
20260726120000 backfills 1247 supplier-invoice rows. pg tests for new
DDL are written but unrun (no local Postgres).

Verified: 11088 tests / 881 files green, tsc 0 non-test errors, lint 0
errors, check:guards passing, MCP payload 57475/57500.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>

* fix(migrations): rename replace_sie_import migration off main's 20260726090000 version

origin/main shipped 20260726090000_agent_quota_rpc_caller_guard.sql; keeping
our replace_sie_import migration on the same version would abort the Supabase
apply with a schema_migrations_pkey duplicate at merge time.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>

* fix(review): remediate pre-publish deep-review findings across all slices

A 13-agent review of the full branch diff surfaced 1 critical, 5 high and
~45 further findings; this commit resolves them in one pass:

- replace_sie_import / undo_sie_import: p_user_id honored only for
  service_role callers; any other caller is pinned to auth.uid()
  (impersonation gate bypass), authz raise errcode 42501 mapped to a
  Swedish 403 in the route, new caller-guard migration for undo
- bulk_book_transactions refuses homogeneous non-SEK batches instead of
  writing foreign magnitudes into SEK ledger columns
- credit-note cap trigger: company-match on credited_invoice_id, no
  cross-tenant figures in exception text
- link_voucher RPCs resolve NULL invoice currency as SEK end to end
- personal-number ciphertext CHECK split into NOT VALID + VALIDATE
- same-currency foreign settlements clear 1510 at booking rate and book
  realized diff to 3960/7960; rate-less foreign write paths refuse
- receivables revaluation covers partially_paid and outstanding amounts
- period lock guard paginates candidates past the PostgREST 1000 cap
- documents: service-client storage removals after authz, dual-layout
  reads in integrity cron and archive export, backfill delete-source
  sweep actually deletes with hash verification and shared-key grouping
- invoice matching normalizes NULL/lowercase currencies (regression),
  duplicate candidates stop claiming amount matches they never ran
- match-invoice aborts on any booking failure (no paid-without-verifikat)
- refresh-exchange-rate reverts on concurrent booking (TOCTOU window)
- KPI preferences upsert arbiter aligned to the company-scoped constraint
- personnummer_last4 stripped from all salary responses incl. MCP tools
- worked-hours batch restores destroyed rows on conflict and error paths
- MCP: shared duplicate-claim builder (no more 'null kr'), short-circuit
  on tag_journal_lines overflow, auto_send schedules stage as high risk
- observability sink redacts emails/IBANs/API keys and keeps redacted
  stacks in prod; assorted small guards (safe-return-to /@, dry_run=True,
  cursor helper off-by-one, OAuth state TTL 10 min, arcim saveMappings
  call removed)

Full dispositions, deferred items and hand-verified accounting numbers
are documented in the PR body and DECISIONS.md.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>

* feat(personnummer): implement masking and encryption for personal numbers with tests

* fix(review): address CI and compliance-bot findings for PR #1215

pg-real: the CI image's auth shim reads the legacy request.jwt.claim.role
GUC, so both service-role simulations (runAsServiceRole and the
invoice-delivery test's local helper) never satisfied auth.role() =
'service_role' and every legitimate p_user_id path failed closed; the
shared helper now sets both GUC shapes plus SET LOCAL ROLE with a
fail-loud sanity check, and the delivery test reuses it. The link-voucher
migration had recreated both RPCs from pre-rewrite file text,
reintroducing the NULL-unsafe membership pattern the
null-safe-tenant-guards ratchet bans; both guards now use
public.caller_is_company_member() with all currency changes preserved.

Compliance bots: the customers export now emits the standard masked form
instead of raw AES-256-GCM ciphertext in the Org-/personnummer column,
and maskCustomerRow returns a non-round-trippable placeholder on decrypt
failure instead of 500ing the list. MCP parity: gnubok_lock_period's
staging pre-check now runs the exact countUnbookedInPeriod the commit
path enforces (exported from period-service; local mirror deleted), and
gnubok_agi_status resolves AGI state run-scoped so a correction run no
longer renders as already filed.

Declined with evidence: PR-Agent's opening-balances null-zeroing concern
(all mergeable columns are NOT NULL with defaults per 20260713101000).

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>

* fix(review): address codex review findings on PR #1215

- restore 20260726140000 to its preview-recorded content and restate the
  NULL-safe tenant guard under 20260727130000: a recorded migration version
  never re-runs, so the in-place edit could not reach the preview branch
- replace toFixed() with sv-SE two-decimal formatting in the ROT/RUT cap
  warning texts and update the pinned test expectations
- drop the em dash in the fiscal-periods route comment
- strip trailing whitespace in import-existing.test.ts

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>

* test(reports): raise timeout on real PDF render tests

renderToBuffer does real @react-pdf layout work and exceeds the 5s
default when the full suite saturates the CPU; tests pass in isolation.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>

---------

Co-authored-by: Claude Opus 5 (1M context) <noreply@anthropic.com>
2026-07-27 03:34:56 +02:00

653 lines
23 KiB
TypeScript
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
/**
* Webhook delivery dispatcher.
*
* Invoked from the per-minute cron at /api/webhooks/dispatch/cron. Picks up
* pending + retry-due deliveries (FOR UPDATE SKIP LOCKED so multiple cron
* invocations don't double-deliver), POSTs each one with HMAC signature,
* and updates the row to one of:
*
* - delivered (2xx response) : terminal
* - failed (5xx / network / 4xx : non-terminal until attempts
* other than 410) exhausted; bumps next_attempt_at
* by exponential backoff
* - dead (HTTP 410 OR : terminal
* attempts exhausted)
*
* The receiver is expected to respond within 10 seconds; we time out
* aggressively so a slow receiver doesn't block the per-minute cron.
*
* On HTTP 410 we additionally disable the webhook (sets disabled_at +
* disabled_reason='HTTP 410 from receiver') so future events don't even
* enqueue against it.
*/
import type { SupabaseClient } from '@supabase/supabase-js'
import { signPayload } from './signing'
import { pinnedHttpsFetch, type PinnedFetchResult } from './pinned-fetch'
import { createLogger } from '@/lib/logger'
const log = createLogger('webhooks/dispatcher')
/** 7 retries over ~87h (≈3.6 days). Index = attempts BEFORE this one. */
const RETRY_BACKOFF_SECONDS: ReadonlyArray<number> = [
60, // 1m: first retry
5 * 60, // 5m
30 * 60, // 30m
2 * 60 * 60, // 2h
12 * 60 * 60, // 12h
24 * 60 * 60, // 24h
48 * 60 * 60, // 48h: final retry
]
const MAX_ATTEMPTS = RETRY_BACKOFF_SECONDS.length + 1 // initial + 7 retries = 8 total
const REQUEST_TIMEOUT_MS = 10_000
const MAX_RESPONSE_BODY_BYTES = 4096
interface DueDelivery {
id: string
webhook_id: string
company_id: string
event_type: string
payload: Record<string, unknown>
previous_attributes: Record<string, unknown> | null
api_version: string
attempts: number
}
interface WebhookForDelivery {
id: string
company_id: string
webhook_url: string
secret: string
}
export interface DispatchSummary {
picked: number
delivered: number
failed: number
dead: number
}
/**
* Run one dispatch cycle. Picks up to `batchSize` due deliveries and
* processes them sequentially (the per-minute cadence + small batch size
* makes parallelism unnecessary; in-process serial is also gentler on the
* receiver if many events fan out to the same URL).
*/
export async function dispatchDueDeliveries(args: {
supabase: SupabaseClient
/** Max rows to claim per cron tick. Default 50. */
batchSize?: number
/** Override for tests. */
now?: Date
/** Override for tests; injected pinned-fetch implementation. */
pinnedFetchImpl?: typeof pinnedHttpsFetch
}): Promise<DispatchSummary> {
const batchSize = args.batchSize ?? 50
const now = args.now ?? new Date()
const pinnedFetchImpl = args.pinnedFetchImpl ?? pinnedHttpsFetch
const summary: DispatchSummary = { picked: 0, delivered: 0, failed: 0, dead: 0 }
// Recover stuck in_flight rows: a previous tick that was killed mid-flight
// (Vercel function timeout, hard crash, manual termination) leaves rows
// marked in_flight forever otherwise. Sweep them back to 'failed' so the
// retry loop picks them up at next_attempt_at.
//
// Threshold = 2× REQUEST_TIMEOUT_MS. A live attempt takes at most
// REQUEST_TIMEOUT_MS plus the body read; doubling that gives an
// unambiguous "this is stuck, not in-flight" boundary.
await recoverStuckInFlight(args.supabase, now)
const due = await claimDueDeliveries(args.supabase, batchSize, now)
summary.picked = due.length
if (due.length === 0) return summary
// Dedupe webhook lookups within a single cycle.
const webhookIds = Array.from(new Set(due.map((d) => d.webhook_id)))
const webhookMap = await loadWebhooksByIds(args.supabase, webhookIds)
for (const delivery of due) {
const webhook = webhookMap.get(delivery.webhook_id)
if (!webhook) {
// The webhook was deleted between enqueue and dispatch. Mark dead;
// there's no receiver to deliver to. The webhook_deliveries.webhook_id
// FK is ON DELETE SET NULL (migration 20260515170000), so the row
// stays in the audit trail under status='dead'.
await markDead(args.supabase, delivery.id, 'webhook_deleted')
summary.dead++
continue
}
// Defense-in-depth tenancy check: the webhook the delivery row points
// at MUST belong to the same company as the delivery row. Mismatch
// indicates a poisoned row: refuse to dispatch (which would sign with
// the wrong tenant's secret and POST to the wrong receiver).
if (webhook.company_id !== delivery.company_id) {
log.error('cross-tenant delivery refused', new Error('company_id mismatch'), {
deliveryId: delivery.id,
deliveryCompanyId: delivery.company_id,
webhookId: webhook.id,
webhookCompanyId: webhook.company_id,
})
await markDead(args.supabase, delivery.id, 'cross_tenant_mismatch')
summary.dead++
continue
}
const outcome = await attemptDelivery({
delivery,
webhook,
pinnedFetchImpl,
now,
})
// Structured per-delivery outcome log. Keeps companyId / webhookId /
// deliveryId available in log aggregation for per-tenant audit-trail
// reconstruction without grepping through individual mark*-helper
// writes (V16, security event correlation).
const logCtx = {
deliveryId: delivery.id,
webhookId: webhook.id,
companyId: delivery.company_id,
eventType: delivery.event_type,
attempt: delivery.attempts + 1,
}
switch (outcome.kind) {
case 'delivered':
await markDelivered(args.supabase, delivery.id, outcome)
log.info('delivery succeeded', { ...logCtx, responseStatus: outcome.responseStatus })
summary.delivered++
break
case 'dead':
await markDead(args.supabase, delivery.id, outcome.reason, outcome)
log.warn('delivery dead', { ...logCtx, reason: outcome.reason, responseStatus: outcome.responseStatus })
summary.dead++
if (outcome.disableWebhook) {
await disableWebhook(args.supabase, webhook.id, webhook.company_id, outcome.reason)
log.warn('webhook auto-disabled', { ...logCtx, reason: outcome.reason })
}
break
case 'failed':
if (delivery.attempts + 1 >= MAX_ATTEMPTS) {
await markDead(args.supabase, delivery.id, 'attempts_exhausted', outcome)
log.warn('delivery dead: attempts exhausted', { ...logCtx, lastError: outcome.error })
summary.dead++
} else {
await markFailedForRetry(args.supabase, delivery.id, delivery.attempts, outcome, now)
log.info('delivery failed: retry scheduled', { ...logCtx, error: outcome.error, responseStatus: outcome.responseStatus })
summary.failed++
}
break
}
}
return summary
}
// ──────────────────────────────────────────────────────────────────────
// DB ops
// ──────────────────────────────────────────────────────────────────────
/**
* Mark in_flight rows whose updated_at is older than the stuck-threshold
* back to 'failed' with next_attempt_at = now so they re-enter the
* dispatch queue. Best-effort: a write failure here is logged but
* doesn't block the rest of the cycle.
*/
async function recoverStuckInFlight(supabase: SupabaseClient, now: Date): Promise<void> {
const stuckBefore = new Date(now.getTime() - 2 * REQUEST_TIMEOUT_MS)
// Under READ COMMITTED (Postgres default), UPDATE re-evaluates the WHERE
// clause against each row's current value when it acquires the row lock.
// A row that raced from 'in_flight' to 'delivered'/'dead' between scan
// and lock will fail status='in_flight' on re-evaluation and be skipped
// entirely: the immutability trigger never fires, so a mid-flight
// terminal flip cannot abort the bulk update.
const { data, error } = await supabase
.from('webhook_deliveries')
.update({
status: 'failed',
next_attempt_at: now.toISOString(),
error: 'recovered_from_in_flight_timeout',
})
.eq('status', 'in_flight')
.lt('updated_at', stuckBefore.toISOString())
.select('id')
if (error) {
log.warn('stuck in_flight recovery failed', { code: error.code })
return
}
if (data && data.length > 0) {
log.warn('recovered stuck in_flight rows', { count: data.length })
}
}
async function claimDueDeliveries(
supabase: SupabaseClient,
batchSize: number,
now: Date,
): Promise<DueDelivery[]> {
// Atomic FOR UPDATE SKIP LOCKED claim via the SQL function shipped in
// migration 20260515220000. PostgREST can't express SKIP LOCKED through
// the JS client, so the function form is the documented entry point:
// see the migration comment for the full rationale (one round trip,
// no CAS contention, rows locked by a concurrent tick are simply
// invisible to the second caller).
//
// All filter semantics from the previous JS path are preserved inside
// the function: status IN ('pending','failed'), next_attempt_at <= now,
// webhook_id IS NOT NULL, ORDER BY next_attempt_at ASC, LIMIT batchSize.
const { data, error } = await supabase.rpc('claim_due_webhook_deliveries', {
p_batch_size: batchSize,
p_now: now.toISOString(),
})
if (error) {
log.error('claim_due_webhook_deliveries rpc failed', error as Error)
return []
}
return (data ?? []) as DueDelivery[]
}
async function loadWebhooksByIds(
supabase: SupabaseClient,
ids: string[],
): Promise<Map<string, WebhookForDelivery>> {
// Include company_id so the dispatch loop can assert that the delivery
// row's company_id matches the webhook's: defense in depth against a
// poisoned delivery row pointing at another tenant's webhook
// (compromised service-role path, faulty INSERT in a future code path,
// etc.). The DB trigger added in 20260515190000 enforces the same
// invariant at INSERT time; this is the application-layer mirror.
const { data, error } = await supabase
.from('webhooks')
.select('id, company_id, webhook_url, secret')
.in('id', ids)
if (error || !data) {
log.error('webhook lookup for dispatch failed', error as Error)
return new Map()
}
return new Map((data as WebhookForDelivery[]).map((w) => [w.id, w]))
}
async function markDelivered(
supabase: SupabaseClient,
id: string,
outcome: DeliveredOutcome,
): Promise<void> {
const { error } = await supabase
.from('webhook_deliveries')
.update({
status: 'delivered',
delivered_at: new Date().toISOString(),
attempts: outcome.attempts,
response_status: outcome.responseStatus,
response_body: outcome.responseBody,
response_headers: outcome.responseHeaders,
error: null,
})
.eq('id', id)
if (error) log.warn('mark delivered update failed', { id, code: error.code })
}
async function markFailedForRetry(
supabase: SupabaseClient,
id: string,
priorAttempts: number,
outcome: FailedOutcome,
now: Date,
): Promise<void> {
const nextAttemptIndex = priorAttempts // 0-indexed lookup into RETRY_BACKOFF_SECONDS
const backoffSeconds = RETRY_BACKOFF_SECONDS[Math.min(nextAttemptIndex, RETRY_BACKOFF_SECONDS.length - 1)]
const nextAttemptAt = new Date(now.getTime() + backoffSeconds * 1000)
const { error } = await supabase
.from('webhook_deliveries')
.update({
status: 'failed',
attempts: priorAttempts + 1,
next_attempt_at: nextAttemptAt.toISOString(),
response_status: outcome.responseStatus ?? null,
response_body: outcome.responseBody ?? null,
response_headers: outcome.responseHeaders ?? null,
error: outcome.error,
})
.eq('id', id)
if (error) log.warn('mark failed-for-retry update failed', { id, code: error.code })
}
async function markDead(
supabase: SupabaseClient,
id: string,
reason: string,
outcome?: AttemptOutcome,
): Promise<void> {
// delivered_at means "the receiver acknowledged the event". For dead
// rows (HTTP 410, attempts exhausted, webhook deleted, cross-tenant
// mismatch, unsafe URL) the receiver did NOT acknowledge: leaving
// delivered_at NULL keeps the audit semantics clean. An auditor
// querying `WHERE delivered_at IS NOT NULL` correctly sees only
// genuinely delivered rows. The terminal-state timestamp lives on
// `updated_at` (auto-stamped by the table's BEFORE UPDATE trigger).
const { error } = await supabase
.from('webhook_deliveries')
.update({
status: 'dead',
attempts: outcome && 'attempts' in outcome ? outcome.attempts : undefined,
response_status: outcome && 'responseStatus' in outcome ? outcome.responseStatus : null,
response_body: outcome && 'responseBody' in outcome ? outcome.responseBody : null,
response_headers: outcome && 'responseHeaders' in outcome ? outcome.responseHeaders : null,
error: reason,
})
.eq('id', id)
if (error) log.warn('mark dead update failed', { id, code: error.code })
}
async function disableWebhook(
supabase: SupabaseClient,
webhookId: string,
companyId: string,
reason: string,
): Promise<void> {
// Snapshot before the disable so the audit entry can record the prior
// state. Service-role read; bypasses RLS.
//
// `webhooks` has NO user_id column: the legacy automation_webhooks table
// never had one and webhooks_v2 (20260515170000) didn't add one. Tenancy
// lives on company_id (NOT NULL) and the only actor attribution is
// created_by_api_key_id. Selecting a phantom column makes PostgREST fail
// the ENTIRE read (42703), which previously left `prior` null and the
// audit row unscoped, so the read error is checked and logged here rather
// than silently degrading to a null snapshot.
const { data: prior, error: priorErr } = await supabase
.from('webhooks')
.select('company_id, name, active, disabled_at, disabled_reason')
.eq('id', webhookId)
.maybeSingle()
if (priorErr) {
log.warn('webhook prior-state snapshot failed', { webhookId, code: priorErr.code })
}
const { error } = await supabase
.from('webhooks')
.update({
disabled_at: new Date().toISOString(),
disabled_reason: reason,
active: false,
})
.eq('id', webhookId)
if (error) {
log.warn('webhook auto-disable failed', { webhookId, code: error.code })
return
}
// V16 security event log. Auto-disable is a privileged action taken by
// the dispatcher cron itself, not by a human caller or a credential:
// user_id and actor_id stay null and the attribution is carried by
// actor_type/actor_label instead ('system' is in the
// audit_log_actor_type_check allowlist; same vocabulary the
// commit_journal_entry audit row uses, migration 20260619120000).
//
// company_id comes from the webhook row the dispatch loop already
// loaded and tenancy-checked against the delivery: the audit row keeps
// its company scope even when the prior-state snapshot read fails, so
// it stays visible to the per-company audit trail (getAuditLog filters
// on company_id).
//
// The entry is written UNCONDITIONALLY (even when prior is null)
// because the SECURITY_EVENT must produce a durable record
// (A.8.15 / V16.1.1 / CC7.2).
//
// The reason discriminates between the three auto-disable paths
// (http_410_gone / redirect_blocked / url_unsafe:<class>) so SIEM
// tooling can alert on systematic patterns.
const p = prior as {
company_id: string | null
name: string
active: boolean
disabled_at: string | null
disabled_reason: string | null
} | null
const { error: auditErr } = await supabase.from('audit_log').insert({
user_id: null,
company_id: companyId,
action: 'SECURITY_EVENT',
table_name: 'webhooks',
record_id: webhookId,
actor_id: null,
actor_type: 'system',
actor_label: 'webhook-dispatcher',
description: p
? `Webhook auto-disabled by dispatcher: ${reason} (was "${p.name}")`
: `Webhook auto-disabled by dispatcher: ${reason} (prior snapshot unavailable)`,
old_state: p
? { active: p.active, disabled_at: p.disabled_at, disabled_reason: p.disabled_reason }
: null,
new_state: { active: false, disabled_reason: reason, disabled_at: new Date().toISOString() },
})
if (auditErr) {
// Escalated from warn to error on purpose. Webhook registrations are
// NOT statutory räkenskapsinformation (BFL/BFNAR do not reach them:
// see the classification in migration 20260515170000), so a dropped
// row here is an operational audit-integrity gap, not a legal one,
// and the dispatch cycle must not abort over it. But warn is
// suppressed in prod-adjacent noise filters and never reaches the
// observability sink, which made this a silent gap; error always
// forwards (lib/logger.ts) so the missing SECURITY_EVENT is alertable
// (CC7.2).
log.error(
'audit_log insert failed for webhook auto-disable',
new Error(auditErr.message ?? 'audit_log insert failed'),
{ webhookId, companyId, reason, code: auditErr.code },
)
}
}
// ──────────────────────────────────────────────────────────────────────
// HTTP attempt
// ──────────────────────────────────────────────────────────────────────
type DeliveredOutcome = {
kind: 'delivered'
attempts: number
responseStatus: number
responseBody: string | null
responseHeaders: Record<string, string> | null
}
type FailedOutcome = {
kind: 'failed'
attempts: number
responseStatus: number | null
responseBody: string | null
responseHeaders: Record<string, string> | null
error: string
}
type DeadOutcome = {
kind: 'dead'
reason: string
disableWebhook: boolean
attempts: number
responseStatus: number | null
responseBody: string | null
responseHeaders: Record<string, string> | null
error?: string
}
type AttemptOutcome = DeliveredOutcome | FailedOutcome | DeadOutcome
async function attemptDelivery(args: {
delivery: DueDelivery
webhook: WebhookForDelivery
pinnedFetchImpl: typeof pinnedHttpsFetch
now: Date
}): Promise<AttemptOutcome> {
const { delivery, webhook, pinnedFetchImpl, now } = args
const attempts = delivery.attempts + 1
const requestId = `whdel_${delivery.id}`
const body = JSON.stringify({
id: delivery.id,
type: delivery.event_type,
api_version: delivery.api_version,
created: Math.floor(now.getTime() / 1000),
data: { object: delivery.payload },
previous_attributes: delivery.previous_attributes,
})
const { header } = signPayload({
body,
secret: webhook.secret,
timestamp: Math.floor(now.getTime() / 1000),
})
// pinnedHttpsFetch performs DNS validation AND opens the socket against
// the validated IP in a single call. The previous shape (separate
// validateWebhookUrl + fetch calls) left a DNS-rebinding window between
// the two, closed here. SNI + Host header continue to carry the
// original hostname so receiver-side TLS + vhost routing still work.
const result = await pinnedFetchImpl(webhook.webhook_url, {
method: 'POST',
headers: {
'Content-Type': 'application/json',
'X-Gnubok-Signature': header,
'X-Gnubok-Event': delivery.event_type,
'X-Gnubok-Delivery': delivery.id,
'X-Gnubok-Api-Version': delivery.api_version,
'X-Request-Id': requestId,
'User-Agent': 'gnubok-webhook/1',
},
body,
timeoutMs: REQUEST_TIMEOUT_MS,
maxResponseBytes: MAX_RESPONSE_BODY_BYTES,
})
switch (result.kind) {
case 'unsafe_url':
return {
kind: 'dead',
reason: `url_unsafe:${result.reason}`,
disableWebhook: true,
attempts,
responseStatus: null,
responseBody: null,
responseHeaders: null,
error: result.detail,
}
case 'redirect_blocked':
return {
kind: 'dead',
reason: 'redirect_blocked',
disableWebhook: true,
attempts,
responseStatus: result.status,
responseBody: null,
responseHeaders: null,
error: truncateError(result.detail),
}
case 'timeout':
case 'transport_error':
return {
kind: 'failed',
attempts,
responseStatus: null,
responseBody: null,
responseHeaders: null,
error: truncateError(result.detail),
}
case 'ok': {
const responseHeaders = filterResponseHeaders(result.headers)
const responseBody = isSafeContentType(result.headers['content-type'] ?? '')
? result.body
: null
// HTTP 410: receiver explicitly asks us to stop. Auto-disable.
if (result.status === 410) {
return {
kind: 'dead',
reason: 'http_410_gone',
disableWebhook: true,
attempts,
responseStatus: 410,
responseBody,
responseHeaders,
}
}
if (result.status >= 200 && result.status < 300) {
return {
kind: 'delivered',
attempts,
responseStatus: result.status,
responseBody,
responseHeaders,
}
}
return {
kind: 'failed',
attempts,
responseStatus: result.status,
responseBody,
responseHeaders,
error: `HTTP ${result.status}`,
}
}
}
}
function truncateError(message: string): string {
return message.length > 500 ? `${message.slice(0, 497)}...` : message
}
// Content-Type prefixes for which we persist response_body verbatim. Other
// types (text/html error pages, application/octet-stream, ...) get dropped
// because they routinely echo PII back from receiver-side error renderers
// (Art.32(1)(b), A.8.12). A null body is just as useful for debugging
// when the operator can see the response_status and response_headers.
const SAFE_BODY_CONTENT_TYPE_PREFIXES = ['text/plain', 'application/json']
function isSafeContentType(contentType: string): boolean {
const lower = contentType.toLowerCase()
return SAFE_BODY_CONTENT_TYPE_PREFIXES.some((p) => lower.startsWith(p))
}
// Allowlist for response_headers persistence. Receiver-side headers like
// Set-Cookie, Authorization, WWW-Authenticate, internal tracing, and
// vendor x-* headers can carry credentials or sensitive identifiers; we
// don't need them for delivery diagnostics. (CC7.2 / Art.32(1)(b))
//
// 'server' is deliberately NOT in the allowlist (A.8.12): it carries no
// diagnostic value but routinely leaks receiver infrastructure version
// strings (nginx/1.21.6, Apache/2.4.41, ...) into a multi-tenant audit
// table.
const SAFE_RESPONSE_HEADERS = new Set([
'content-type',
'content-length',
'date',
'x-request-id',
'cf-ray',
])
function filterResponseHeaders(headers: Record<string, string>): Record<string, string> {
const obj: Record<string, string> = {}
for (const [k, v] of Object.entries(headers)) {
if (SAFE_RESPONSE_HEADERS.has(k.toLowerCase())) {
obj[k] = v
}
}
return obj
}
export const __TESTING__ = {
RETRY_BACKOFF_SECONDS,
MAX_ATTEMPTS,
REQUEST_TIMEOUT_MS,
MAX_RESPONSE_BODY_BYTES,
}