Skip to content

Commit f2a2e86

Browse files
feat(webhooks): orchestrate delivery worker
1 parent 2012560 commit f2a2e86

7 files changed

Lines changed: 463 additions & 17 deletions

TODO.md

Lines changed: 40 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -495,7 +495,7 @@
495495
- [ ] Implementar outbox transacional a partir de domain/workflow transitions. Parcial F0-033: `project.created` e `project.version.created` são persistidos atomicamente com a criação idempotente; demais transitions continuam abertas.
496496
- [x] Modelar endpoint, subscription, secret, filter e delivery attempt. Evidência F0-034: domínios canônicos, registro transacional, cinco tabelas, constraints e regressões de segurança.
497497
- [x] Implementar challenge, assinatura, timestamp e anti-replay. Evidência F0-035/F0-036: challenge durável one-shot, HMAC dos bytes exatos, janela de timestamp, receipt anti-replay e transporte HTTPS pinado com resolução DNS fail-closed.
498-
- [ ] Implementar at-least-once, backoff, dead-letter e replay controlado. Parcial F0-031/F0-039: render possui lease/fencing, backoff, checkpoint, dead-letter e retry manual; webhooks agora possuem fan-out deduplicado, claim/lease, dispatch assinado com transporte DNS-pinado, classificação de resposta, backoff com jitter e dead-letter. Adapter concreto do secret provider e replay administrativo continuam abertos.
498+
- [ ] Implementar at-least-once, backoff, dead-letter e replay controlado. Parcial F0-031/F0-040: render possui lease/fencing, backoff, checkpoint, dead-letter e retry manual; webhooks possuem fan-out, claim/lease, dispatch assinado, transporte DNS-pinado, heartbeat orquestrado, worker round-robin por workspace, backoff e dead-letter. Adapter concreto do secret provider, discovery/rebalance de shards e replay administrativo continuam abertos.
499499
- [ ] Criar UI/API administrativa de status, attempts e rotação de secret.
500500
- [ ] Criar integration tests de duplicação, timeout, assinatura inválida e replay. Parcial F0-035/F0-039: assinatura adulterada, replay durável, timeout absoluto, DNS misto, rebinding, claim concorrente, lease incorreto, reclaim, worker obsoleto, retry e dead-letter estão cobertos. Dispatcher cobre bytes/headers assinados, fingerprint, DNS privado e settlement integrado; falta replay administrativo ponta a ponta.
501501

@@ -3624,7 +3624,7 @@ Limites explícitos desta slice:
36243624

36253625
### Slice F0-039 — Dispatcher assinado de webhook deliveries
36263626

3627-
**Status:** concluído localmente em 14 de julho de 2026; ainda não commitado.
3627+
**Status:** concluído e publicado em 14 de julho de 2026 no commit `2012560`.
36283628

36293629
Entregas:
36303630

@@ -3656,4 +3656,41 @@ Limites explícitos desta slice:
36563656
- ainda não existe loop executável que una claim, heartbeat e dispatch continuamente;
36573657
- replay administrativo, rotação operacional, rate limit/circuit breaker por endpoint e observabilidade continuam abertos;
36583658
- API/UI administrativa de endpoint, delivery e attempts permanece no incremento previsto;
3659-
- hosted CI do dispatcher será registrado após sua publicação no próximo ciclo.
3659+
- hosted CI `29378737291` aprovou PostgreSQL, 86 testes, contratos, API, FFmpeg, Remotion real, build e auditorias.
3660+
3661+
### Slice F0-040 — Runner e loop do worker de webhook deliveries
3662+
3663+
**Status:** concluído localmente em 14 de julho de 2026; ainda não commitado.
3664+
3665+
Entregas:
3666+
3667+
- runner une claim, heartbeat, dispatch e settlement em uma única unidade workspace-scoped;
3668+
- claim ocioso retorna sem criar timer, abrir secret ou executar transporte;
3669+
- heartbeat não sobrepõe renovações e permanece ativo durante abertura da chave, DNS e HTTPS;
3670+
- factory valida heartbeat menor que lease e constrói runner completo com repository, transporte e provider injetado;
3671+
- settlement fenced continua sendo a autoridade final: `stale` vira `lease-lost`, enquanto sucesso persistido não é rebaixado por heartbeat tardio;
3672+
- outcomes expõem apenas workspace, delivery, attempt e status, sem token, URL, chave, payload, assinatura ou resposta;
3673+
- loop recebe shard explícito de 1 a 1.000 workspaces únicos e processa uma delivery por tenant a cada passagem;
3674+
- round-robin impede starvation por workspace com backlog elevado;
3675+
- falha de um tenant é isolada e reportada apenas com workspace ID, sem interromper os demais;
3676+
- polling ocorre somente quando todo o shard está ocioso e pode ser interrompido por `AbortSignal`;
3677+
- desligamento gracioso impede novos claims e deixa a iteração corrente concluir;
3678+
- ADR-029 formaliza lifecycle, justiça entre tenants, callbacks seguros e pré-requisitos do host.
3679+
3680+
Regressões e evidências locais:
3681+
3682+
- suíte global passa com 89 testes; 21 são contratos de webhook;
3683+
- runner longo comprova heartbeat durante dispatch e ausência de renovação sobreposta;
3684+
- perda do heartbeat com settlement stale retorna `lease-lost`;
3685+
- loop comprova isolamento de erro, passagem para o próximo workspace e parada graciosa;
3686+
- callbacks foram verificados sem token ou chave;
3687+
- integração Prisma agora executa retry due → claim → dispatcher assinado → settlement através do runner;
3688+
- typecheck e integração dedicada de webhook em SQLite passam.
3689+
3690+
Limites explícitos desta slice:
3691+
3692+
- o host ainda precisa fornecer provider de secrets e lista autorizada de workspaces do shard;
3693+
- discovery dinâmica, rebalanceamento de shards e autoscaling permanecem no scheduler posterior;
3694+
- ainda não existe entrypoint de produção porque a escolha do secret provider não foi tomada;
3695+
- replay administrativo, rotação operacional, rate limit/circuit breaker e observabilidade continuam abertos;
3696+
- hosted CI do runner será registrado após publicação no próximo ciclo.
Lines changed: 44 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,44 @@
1+
# ADR-029 — Orquestração do worker de webhook deliveries
2+
3+
> **Status:** Accepted
4+
>
5+
> **Data:** 14 de julho de 2026
6+
7+
## Contexto
8+
9+
Claim, lease e dispatcher já existiam como boundaries separados. Sem uma orquestração única, cada host poderia esquecer heartbeat, interpretar `stale` como sucesso, processar tenants sem justiça ou registrar material sensível. O loop também precisa parar sem abandonar uma tentativa que já entrou em rede.
10+
11+
## Decisão
12+
13+
- Um runner executa exatamente uma unidade `claim → heartbeat → dispatch → settlement` para um workspace explícito.
14+
- Claim ocioso retorna `null` e não agenda heartbeat nem abre secret.
15+
- O heartbeat começa somente depois do claim, não permite renovações sobrepostas e continua durante abertura do secret, DNS e HTTPS.
16+
- Intervalo de heartbeat deve ser menor que o lease. O factory valida essa relação antes de iniciar trabalho.
17+
- Erro ou rejeição de heartbeat interrompe novas renovações. O dispatcher ainda precisa concluir pelo fence; resultado `stale` é apresentado como `lease-lost`.
18+
- Settlement bem-sucedido é a autoridade final. Um heartbeat concorrente que perde a corrida porque o settlement já limpou o lease não converte sucesso durável em falha.
19+
- O outcome contém apenas workspace, delivery, attempt e estado seguro. Token, URL, chave, payload, headers e resposta não aparecem em callbacks.
20+
- O loop recebe uma lista explícita de workspaces autorizados ao shard, limitada a 1.000 IDs únicos.
21+
- Cada passagem processa no máximo uma delivery por workspace e segue ordem round-robin, evitando que um tenant com backlog monopolize o processo.
22+
- Erro de uma iteração é reportado somente com workspace ID e não interrompe os demais tenants.
23+
- O loop só dorme quando todos os workspaces estão ociosos. O poll é interrompível por `AbortSignal`.
24+
- Desligamento gracioso impede novos claims e permite que a iteração em andamento termine seu settlement.
25+
- Provider de secrets e lista de workspaces permanecem dependências do host. Não existe descoberta global ou fallback de chave implícito.
26+
27+
## Consequências
28+
29+
- Um host não precisa reproduzir manualmente as regras de lease e fencing.
30+
- Justiça entre tenants é determinística dentro de cada shard.
31+
- Escala horizontal continua segura porque o claim é CAS e workspace-scoped.
32+
- O deployment precisa fornecer o shard de workspaces e um provider confiável. Descoberta dinâmica e rebalanceamento pertencem a um scheduler posterior.
33+
- Métricas e logs podem consumir outcomes seguros sem risco de vazar credenciais.
34+
35+
## Evidências exigidas
36+
37+
- runner ocioso não chama heartbeat nem dispatcher;
38+
- dispatch longo recebe heartbeat durante a rede;
39+
- heartbeat perdido + settlement stale resulta em `lease-lost`;
40+
- sucesso persistido não é rebaixado por heartbeat concorrente tardio;
41+
- erro de um workspace não bloqueia o próximo;
42+
- callbacks e outcomes não contêm token ou secret;
43+
- workspace duplicado, owner inseguro e intervalos inválidos falham antes do loop;
44+
- integração Prisma executa claim → dispatch assinado → settlement pelo runner.

src/v2/application/dispatch-webhook-delivery.ts

Lines changed: 9 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -207,6 +207,15 @@ export function dispatchWebhookDeliveryService(dependencies: {
207207
},
208208
})
209209
if (!settled) return Object.freeze({ status: 'stale' as const })
210+
if (
211+
settled.delivery.status !== 'retry-scheduled' &&
212+
settled.delivery.status !== 'dead-lettered'
213+
) {
214+
throw new DomainError(
215+
'PERSISTENCE_CONFLICT',
216+
'Webhook failure settlement returned an invalid delivery state',
217+
)
218+
}
210219
return Object.freeze({ status: settled.delivery.status, delivery: settled.delivery })
211220
}
212221
}
Lines changed: 208 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,208 @@
1+
import { DomainError, assertDomain } from '../domain/errors.ts'
2+
3+
type ClaimWebhookDelivery = (request: {
4+
workspaceId: string
5+
leaseOwner: string
6+
}) => Promise<Readonly<{
7+
delivery: Readonly<{ id: string }>
8+
attempt: Readonly<{ attemptNumber: number }>
9+
leaseToken: string
10+
}> | null>
11+
12+
type HeartbeatWebhookDelivery = (request: {
13+
workspaceId: string
14+
deliveryId: string
15+
leaseOwner: string
16+
leaseToken: string
17+
attemptNumber: number
18+
}) => Promise<boolean>
19+
20+
type DispatchWebhookDelivery = (request: {
21+
workspaceId: string
22+
deliveryId: string
23+
leaseOwner: string
24+
leaseToken: string
25+
attemptNumber: number
26+
}) => Promise<Readonly<{
27+
status: 'succeeded' | 'retry-scheduled' | 'dead-lettered' | 'stale'
28+
}>>
29+
30+
export interface WebhookDeliveryWorkerOutcome {
31+
workspaceId: string
32+
deliveryId: string
33+
attemptNumber: number
34+
status: 'succeeded' | 'retry-scheduled' | 'dead-lettered' | 'lease-lost'
35+
}
36+
37+
export function runNextWebhookDeliveryService(dependencies: {
38+
claim: ClaimWebhookDelivery
39+
heartbeat: HeartbeatWebhookDelivery
40+
dispatch: DispatchWebhookDelivery
41+
heartbeatIntervalMs?: number
42+
}) {
43+
const heartbeatIntervalMs = dependencies.heartbeatIntervalMs ?? 10_000
44+
assertDomain(
45+
Number.isSafeInteger(heartbeatIntervalMs) &&
46+
heartbeatIntervalMs >= 100 &&
47+
heartbeatIntervalMs <= 60_000,
48+
'INVALID_WEBHOOK',
49+
'Webhook worker heartbeat interval must be between 100 and 60000 milliseconds',
50+
)
51+
52+
return async function runNextWebhookDelivery(request: {
53+
workspaceId: string
54+
leaseOwner: string
55+
}): Promise<Readonly<WebhookDeliveryWorkerOutcome> | null> {
56+
const claimed = await dependencies.claim(request)
57+
if (!claimed) return null
58+
59+
const command = {
60+
workspaceId: request.workspaceId,
61+
deliveryId: claimed.delivery.id,
62+
leaseOwner: request.leaseOwner,
63+
leaseToken: claimed.leaseToken,
64+
attemptNumber: claimed.attempt.attemptNumber,
65+
}
66+
let stopped = false
67+
let leaseLost = false
68+
let timer: ReturnType<typeof setTimeout> | undefined
69+
let renewal: Promise<boolean> | undefined
70+
71+
const heartbeat = async () => {
72+
if (stopped || leaseLost) return false
73+
if (renewal) return renewal
74+
renewal = (async () => {
75+
try {
76+
const renewed = await dependencies.heartbeat(command)
77+
if (!renewed) leaseLost = true
78+
return renewed
79+
} catch {
80+
leaseLost = true
81+
return false
82+
} finally {
83+
renewal = undefined
84+
}
85+
})()
86+
return renewal
87+
}
88+
const scheduleHeartbeat = () => {
89+
if (stopped || leaseLost) return
90+
timer = setTimeout(async () => {
91+
await heartbeat()
92+
scheduleHeartbeat()
93+
}, heartbeatIntervalMs)
94+
timer.unref?.()
95+
}
96+
const stopHeartbeat = async () => {
97+
stopped = true
98+
if (timer) clearTimeout(timer)
99+
if (renewal) await renewal
100+
}
101+
102+
try {
103+
scheduleHeartbeat()
104+
const dispatched = await dependencies.dispatch(command)
105+
if (dispatched.status === 'stale') {
106+
return Object.freeze({
107+
workspaceId: command.workspaceId,
108+
deliveryId: command.deliveryId,
109+
attemptNumber: command.attemptNumber,
110+
status: 'lease-lost' as const,
111+
})
112+
}
113+
return Object.freeze({
114+
workspaceId: command.workspaceId,
115+
deliveryId: command.deliveryId,
116+
attemptNumber: command.attemptNumber,
117+
status: dispatched.status,
118+
})
119+
} finally {
120+
await stopHeartbeat()
121+
}
122+
}
123+
}
124+
125+
const SAFE_ID_PATTERN = /^[A-Za-z0-9][A-Za-z0-9._:-]{2,127}$/
126+
127+
function normalizeWorkspaceIds(values: readonly string[]): readonly string[] {
128+
const normalized = values.map((value) => value.trim())
129+
assertDomain(
130+
normalized.length >= 1 &&
131+
normalized.length <= 1_000 &&
132+
normalized.every((value) => SAFE_ID_PATTERN.test(value)) &&
133+
new Set(normalized).size === normalized.length,
134+
'INVALID_WEBHOOK',
135+
'Webhook worker requires 1 to 1000 unique workspace IDs',
136+
)
137+
return Object.freeze(normalized)
138+
}
139+
140+
async function waitForPoll(delayMs: number, signal: AbortSignal): Promise<void> {
141+
if (signal.aborted) return
142+
await new Promise<void>((resolve) => {
143+
const timer = setTimeout(done, delayMs)
144+
timer.unref?.()
145+
signal.addEventListener('abort', done, { once: true })
146+
function done() {
147+
clearTimeout(timer)
148+
signal.removeEventListener('abort', done)
149+
resolve()
150+
}
151+
})
152+
}
153+
154+
export async function runWebhookDeliveryWorkerLoop(dependencies: {
155+
runNext: (request: {
156+
workspaceId: string
157+
leaseOwner: string
158+
}) => Promise<Readonly<WebhookDeliveryWorkerOutcome> | null>
159+
workspaceIds: readonly string[]
160+
leaseOwner: string
161+
signal: AbortSignal
162+
pollIntervalMs?: number
163+
onOutcome?: (outcome: Readonly<WebhookDeliveryWorkerOutcome>) => void
164+
onIterationError?: (event: Readonly<{ workspaceId: string }>) => void
165+
wait?: (delayMs: number, signal: AbortSignal) => Promise<void>
166+
}): Promise<void> {
167+
const workspaceIds = normalizeWorkspaceIds(dependencies.workspaceIds)
168+
const leaseOwner = dependencies.leaseOwner.trim()
169+
const pollIntervalMs = dependencies.pollIntervalMs ?? 1_000
170+
assertDomain(
171+
SAFE_ID_PATTERN.test(leaseOwner),
172+
'INVALID_WEBHOOK',
173+
'Webhook worker lease owner is invalid',
174+
)
175+
assertDomain(
176+
Number.isSafeInteger(pollIntervalMs) && pollIntervalMs >= 100 && pollIntervalMs <= 60_000,
177+
'INVALID_WEBHOOK',
178+
'Webhook worker poll interval must be between 100 and 60000 milliseconds',
179+
)
180+
const wait = dependencies.wait ?? waitForPoll
181+
182+
while (!dependencies.signal.aborted) {
183+
let processed = false
184+
for (const workspaceId of workspaceIds) {
185+
if (dependencies.signal.aborted) break
186+
try {
187+
const outcome = await dependencies.runNext({ workspaceId, leaseOwner })
188+
if (outcome) {
189+
processed = true
190+
try {
191+
dependencies.onOutcome?.(outcome)
192+
} catch {
193+
// Observability callbacks cannot control delivery execution.
194+
}
195+
}
196+
} catch {
197+
try {
198+
dependencies.onIterationError?.(Object.freeze({ workspaceId }))
199+
} catch {
200+
// Observability callbacks cannot stop other workspaces.
201+
}
202+
}
203+
}
204+
if (!processed && !dependencies.signal.aborted) {
205+
await wait(pollIntervalMs, dependencies.signal)
206+
}
207+
}
208+
}

src/v2/infrastructure/repository-factory.ts

Lines changed: 47 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -10,6 +10,7 @@ import {
1010
settleWebhookDeliveryService,
1111
} from '../application/manage-webhook-delivery.ts'
1212
import { dispatchWebhookDeliveryService } from '../application/dispatch-webhook-delivery.ts'
13+
import { runNextWebhookDeliveryService } from '../application/run-webhook-delivery-worker.ts'
1314
import { materializeAuthorizedRenderInputService } from '../application/materialize-authorized-render-input.ts'
1415
import { renderAuthorizedInputService } from '../application/render-authorized-input.ts'
1516
import { runNextPublicOperationService } from '../application/run-public-operation-worker.ts'
@@ -154,6 +155,52 @@ export function createWebhookDeliveryWorker(
154155
})
155156
}
156157

158+
export function createWebhookDeliveryRunner(
159+
secrets: WebhookSigningSecretProvider,
160+
environment: NodeJS.ProcessEnv = process.env,
161+
clock: () => Date = () => new Date(),
162+
) {
163+
const configuredLease = Number(environment.APOLLO_V2_WEBHOOK_DELIVERY_LEASE_MS)
164+
const configuredHeartbeat = Number(environment.APOLLO_V2_WEBHOOK_HEARTBEAT_MS)
165+
const configuredTimeout = Number(environment.APOLLO_V2_WEBHOOK_DELIVERY_TIMEOUT_MS)
166+
const configuredRetryBase = Number(environment.APOLLO_V2_WEBHOOK_RETRY_BASE_MS)
167+
const configuredRetryMax = Number(environment.APOLLO_V2_WEBHOOK_RETRY_MAX_MS)
168+
const leaseDurationMs = Number.isSafeInteger(configuredLease) && configuredLease > 0
169+
? configuredLease
170+
: 30_000
171+
const heartbeatIntervalMs = Number.isSafeInteger(configuredHeartbeat) && configuredHeartbeat > 0
172+
? configuredHeartbeat
173+
: 10_000
174+
if (heartbeatIntervalMs >= leaseDurationMs) {
175+
throw new DomainError(
176+
'INVALID_WEBHOOK',
177+
'Webhook heartbeat interval must be shorter than its lease',
178+
)
179+
}
180+
const repository = createWebhookDeliveryRepository()
181+
return runNextWebhookDeliveryService({
182+
claim: claimNextWebhookDeliveryService({ repository, clock, leaseDurationMs }),
183+
heartbeat: heartbeatWebhookDeliveryService({ repository, clock, leaseDurationMs }),
184+
dispatch: dispatchWebhookDeliveryService({
185+
repository,
186+
secrets,
187+
transport: new SafeWebhookDeliveryTransport({
188+
...(Number.isSafeInteger(configuredTimeout) && configuredTimeout > 0
189+
? { timeoutMs: configuredTimeout }
190+
: {}),
191+
}),
192+
clock,
193+
...(Number.isSafeInteger(configuredRetryBase) && configuredRetryBase > 0
194+
? { retryBaseDelayMs: configuredRetryBase }
195+
: {}),
196+
...(Number.isSafeInteger(configuredRetryMax) && configuredRetryMax > 0
197+
? { retryMaxDelayMs: configuredRetryMax }
198+
: {}),
199+
}),
200+
heartbeatIntervalMs,
201+
})
202+
}
203+
157204
export function createWebhookFanoutMaterializer(
158205
clock: () => Date = () => new Date(),
159206
) {

0 commit comments

Comments
 (0)