Skip to content

Commit aac19d8

Browse files
feat(webhooks): create subscriptions idempotently
1 parent 47faab8 commit aac19d8

16 files changed

Lines changed: 808 additions & 6 deletions

README.md

Lines changed: 8 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -137,8 +137,7 @@ O envelope versionado e o catálogo inicial de eventos podem ser descobertos em
137137
`GET /v1/events/catalog`. O catálogo referencia o JSON Schema público do envelope
138138
e contém somente metadados estáticos, por isso não exige autenticação. A presença
139139
de um tipo no catálogo não significa que ele já esteja sendo emitido: cada
140-
transição precisa ser conectada ao outbox, e subscriptions, assinatura e entrega
141-
at-least-once serão adicionadas nas próximas slices de F0.038.
140+
transição ainda precisa ser conectada explicitamente ao outbox.
142141

143142
A criação de projeto já persiste `project.created` e `project.version.created`
144143
no outbox, atomicamente com o projeto, sua versão inicial e o registro de
@@ -160,6 +159,13 @@ HTTPS e seu fingerprint; secrets são apenas metadados de versão, fingerprint e
160159
estado, sem `keyRef` ou material criptográfico. Os filtros exatos da subscription
161160
são visíveis porque definem o comportamento contratado da entrega.
162161

162+
Novas subscriptions podem ser criadas por `POST /v1/webhooks/subscriptions`
163+
para um endpoint existente. O command exige `Idempotency-Key`: repetir endpoint
164+
e filtro idênticos com a mesma chave devolve o recurso original, enquanto reutilizar
165+
a chave para outro filtro ou tentar duplicar o filtro com outra chave retorna
166+
conflito explícito. A criação nasce ativa somente quando o endpoint está ativo;
167+
endpoints ainda em challenge produzem uma subscription pendente.
168+
163169
O status de uma subscription pode ser alterado por
164170
`PUT /v1/webhooks/subscriptions/{subscriptionId}/status`. A resposta de consulta
165171
fornece uma revisão opaca que deve ser enviada como `baseRevision`: alterações

TODO.md

Lines changed: 34 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -496,7 +496,7 @@
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.
498498
- [x] Implementar at-least-once, backoff, dead-letter e replay controlado. Evidência F0-031/F0-046: render possui lease/fencing, backoff, checkpoint, dead-letter e retry manual; webhooks possuem outbox/fan-out, claim/lease/fencing, dispatch assinado, transporte DNS-pinado, heartbeat, discovery, coordenação durável de shards, secret provider configurado, entrypoint operacional, backoff, dead-letter e replay idempotente individual ou por evento exato. Replay por intervalo permanece enhancement administrativo separado.
499-
- [ ] Criar UI/API administrativa de status, attempts e rotação de secret. Parcial F0-042/F0-044/F0-047/F0-048/F0-049: API externa lista/lê endpoints, subscriptions e deliveries, executa replay e altera lifecycle de endpoints/subscriptions; UI, criação/challenge e rotação de secret continuam abertas.
499+
- [ ] Criar UI/API administrativa de status, attempts e rotação de secret. Parcial F0-042/F0-044/F0-047/F0-048/F0-049/F0-050: API externa lista/lê endpoints, subscriptions e deliveries, cria subscriptions com filtros exatos, executa replay e altera lifecycle de endpoints/subscriptions; UI, criação/challenge de endpoint e rotação de secret continuam abertas.
500500
- [x] Criar integration tests de duplicação, timeout, assinatura inválida e replay. Evidência F0-035/F0-043: assinatura adulterada, anti-replay durável, deadline absoluto, DNS/rebinding, claim concorrente, lease/fencing, retry/dead-letter e replay administrativo idempotente estão cobertos em contratos, Prisma e HTTP.
501501

502502
### F0.039 — Idempotência e concorrência externa [FR-245]
@@ -539,7 +539,7 @@
539539
### F0.043 — Governança da API [FR-249]
540540

541541
- [ ] Criar administração de clients, scopes, secrets, environments e status.
542-
- [ ] Criar administração de webhooks, subscriptions e delivery diagnostics. Parcial F0-042/F0-044/F0-047/F0-048/F0-049: capabilities de endpoints, subscriptions, deliveries e replay entregam consulta, lifecycle com cascatas e replay workspace-scoped; criação, challenge, rotação e UI continuam abertas.
542+
- [ ] Criar administração de webhooks, subscriptions e delivery diagnostics. Parcial F0-042/F0-044/F0-047/F0-048/F0-049/F0-050: capabilities de endpoints, subscriptions, deliveries e replay entregam consulta, criação idempotente de subscription, lifecycle com cascatas e replay workspace-scoped; criação/challenge de endpoint, rotação e UI continuam abertas.
543543
- [ ] Implementar rate limits, quotas, concurrency e spend budgets por client/workspace.
544544
- [ ] Criar usage e audit queries paginadas com redaction.
545545
- [ ] Criar sandbox isolado com provider fakes e custos simulados.
@@ -3991,7 +3991,7 @@ Limites explícitos desta slice:
39913991

39923992
### Slice F0-049 — Lifecycle e cascatas de endpoints de webhook
39933993

3994-
**Status:** concluído localmente em 15 de julho de 2026; ainda não commitado.
3994+
**Status:** publicado no `main` em 15 de julho de 2026 (`47faab8`); hosted CI `29445313005` aprovada.
39953995

39963996
Entregas:
39973997

@@ -4026,3 +4026,34 @@ Limites explícitos desta slice:
40264026
- criação e alteração de filtros de subscription permanecem abertas;
40274027
- rotação de signing secret e reload dinâmico de provider permanecem abertas;
40284028
- UI administrativa, audit query, métricas e alertas operacionais continuam futuros.
4029+
4030+
### Slice F0-050 — Criação externa idempotente de subscriptions
4031+
4032+
**Status:** concluído localmente em 15 de julho de 2026; ainda não commitado.
4033+
4034+
Entregas:
4035+
4036+
- capability `apollo.webhooks.subscriptions.create` expõe `POST /v1/webhooks/subscriptions` sob `webhooks:admin`;
4037+
- body fechado aceita endpoint, tipos do catálogo e IDs de recurso opcionais, sempre convertidos para filtro canônico exato;
4038+
- `Idempotency-Key` obrigatória é vinculada a workspace, API client, endpoint e hash do filtro;
4039+
- primeira criação retorna 201; repetição idêntica retorna o mesmo recurso com 200 e `replayed: true`;
4040+
- mesma chave com payload diferente retorna `IDEMPOTENCY_PAYLOAD_MISMATCH`; filtro idêntico com outra chave retorna `WEBHOOK_SUBSCRIPTION_ALREADY_EXISTS`;
4041+
- ledger e subscription são gravados atomicamente em transação serializável;
4042+
- endpoint ativo cria subscription ativa; endpoint pendente cria subscription pendente; suspenso ou revogado rejeita o command;
4043+
- resposta reutiliza presenter redigido e não expõe workspace, URL, `filterHash`, secrets ou ledger;
4044+
- ADR-039 formaliza idempotência, unicidade exata, estados iniciais e concorrência.
4045+
4046+
Regressões e evidências locais:
4047+
4048+
- suíte global passa com 108 testes; 40 são contratos de webhook;
4049+
- integração Prisma cobre criação pendente, replay, payload divergente, filtro duplicado e ausência de duplicação;
4050+
- jornada HTTP cobre OpenAPI, 201/200, chave ausente, replay, conflitos distintos e 403 sem scope;
4051+
- contratos públicos passam com 38 capabilities, 46 schemas, 62 exemplos e 34 paths;
4052+
- build registra `POST /v1/webhooks/subscriptions` junto do GET existente;
4053+
- integração de webhook e jornada pública passam no protótipo SQLite.
4054+
4055+
Limites explícitos desta slice:
4056+
4057+
- criação, alteração de URL e challenge público de endpoint continuam abertos;
4058+
- alteração de filtro de subscription permanece um command futuro separado;
4059+
- rotação de signing secret, UI administrativa, audit query, métricas e alertas continuam futuros.
Lines changed: 37 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,37 @@
1+
# ADR-039 — Criação idempotente de subscriptions de webhook
2+
3+
> **Status:** Accepted
4+
>
5+
> **Data:** 15 de julho de 2026
6+
7+
## Contexto
8+
9+
A administração externa já consultava subscriptions e alterava seu lifecycle, mas ainda não conseguia anexar um novo filtro a um endpoint existente. Agentes externos precisam repetir uma chamada após timeout sem criar recursos duplicados, distinguir replay de conflito e receber o mesmo resultado original. A criação também precisa respeitar o estado do endpoint no mesmo instante em que persiste a subscription.
10+
11+
## Decisão
12+
13+
- A capability `apollo.webhooks.subscriptions.create` expõe `POST /v1/webhooks/subscriptions` sob `webhooks:admin` e confirmação humana.
14+
- O body fechado aceita `endpointId`, `eventTypes` e `resourceIds` opcional; tipos de evento devem pertencer ao catálogo público e ambos os filtros são ordenados e hasheados canonicamente.
15+
- `Idempotency-Key` é obrigatório, possui de 1 a 128 caracteres ASCII imprimíveis e é isolado por workspace e API client.
16+
- O fingerprint vincula a chave ao endpoint e ao hash do filtro canônico. A mesma chave e payload devolve a subscription original com HTTP 200 e `replayed: true`; payload diferente devolve 409.
17+
- A primeira criação devolve HTTP 201 e `replayed: false`. O ledger e a subscription são persistidos atomicamente em transação serializável.
18+
- Uma subscription criada em endpoint ativo nasce ativa; em endpoint pendente de verificação nasce pendente. Endpoint suspenso ou revogado rejeita a criação com 409.
19+
- O par endpoint + filtro exato é único. Outra chave tentando reproduzir o mesmo filtro recebe `WEBHOOK_SUBSCRIPTION_ALREADY_EXISTS`, sem adotar silenciosamente um recurso anterior.
20+
- Endpoint, API client e workspace são verificados dentro da transação. Replays reidratam o recurso por workspace e falham de forma fechada se o resultado persistido não existir.
21+
- A resposta pública usa o presenter redigido: não expõe workspace, URL, `filterHash`, secret, ledger ou detalhes da transação.
22+
23+
## Consequências
24+
25+
- Agentes podem repetir chamadas ambíguas de rede com segurança e obter identidade estável.
26+
- A unicidade natural do filtro não substitui o ledger: ela bloqueia duplicatas, enquanto a chave preserva a semântica do request e sua resposta.
27+
- Criações concorrentes convergem pelo ledger ou falham explicitamente por filtro duplicado/serialização, sem estado parcial.
28+
- Alteração de filtro continua sendo uma operação futura separada; não há atualização implícita durante a criação.
29+
30+
## Evidências exigidas
31+
32+
- criação em endpoint ativo e pendente escolhe o estado correto;
33+
- replay retorna o mesmo ID e não cria uma segunda linha;
34+
- mesma chave com payload diferente e filtro duplicado com outra chave retornam conflitos distintos;
35+
- endpoint ausente/inativo, body inválido, chave ausente e falta de scope falham antes de efeitos indevidos;
36+
- OpenAPI declara body, header idempotente e respostas 200/201;
37+
- contratos unitários, Prisma e jornada HTTP exercitam a implementação real.

src/app/v1/webhooks/subscriptions/route.ts

Lines changed: 70 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,9 +1,14 @@
1+
import { randomUUID } from 'node:crypto'
12
import { NextRequest, NextResponse } from 'next/server'
23

34
import { requireScope } from '@/v2/application/authenticate-api-client'
5+
import { createWebhookSubscriptionService } from '@/v2/application/create-webhook-subscription'
46
import { listWebhookSubscriptionsService } from '@/v2/application/list-webhook-administration'
57
import { DomainError } from '@/v2/domain/errors'
6-
import { createWebhookAdministrationQueryRepository } from '@/v2/infrastructure/repository-factory'
8+
import {
9+
createWebhookAdministrationQueryRepository,
10+
createWebhookSubscriptionCreationRepository,
11+
} from '@/v2/infrastructure/repository-factory'
712
import { authenticateExternalRequest } from '@/v2/public-api/authentication'
813
import { publicApiHeaders, resolveRequestId, respondPublicError } from '@/v2/public-api/errors'
914
import { presentSuccess, presentWebhookSubscription } from '@/v2/public-api/presenters'
@@ -33,3 +38,67 @@ export async function GET(request: NextRequest) {
3338
}), { status: 200, headers: publicApiHeaders(requestId) })
3439
} catch (error) { return respondPublicError(error, requestId) }
3540
}
41+
42+
export async function POST(request: NextRequest) {
43+
const requestId = resolveRequestId(request)
44+
try {
45+
const actor = await authenticateExternalRequest(request)
46+
requireScope(actor, 'webhooks:admin')
47+
const idempotencyKey = request.headers.get('idempotency-key')?.trim() ?? ''
48+
let body: Record<string, unknown>
49+
try {
50+
const value = await request.json() as unknown
51+
if (!value || typeof value !== 'object' || Array.isArray(value)) {
52+
throw new DomainError('INVALID_ARGUMENT', 'Request body must be an object')
53+
}
54+
body = value as Record<string, unknown>
55+
} catch (error) {
56+
if (error instanceof DomainError) throw error
57+
throw new DomainError('INVALID_ARGUMENT', 'Request body must be valid JSON')
58+
}
59+
const allowed = new Set(['endpointId', 'eventTypes', 'resourceIds'])
60+
for (const name of Object.keys(body)) {
61+
if (!allowed.has(name)) throw new DomainError('INVALID_ARGUMENT', `${name} is not supported`)
62+
}
63+
if (typeof body.endpointId !== 'string') {
64+
throw new DomainError('INVALID_ARGUMENT', 'endpointId must be a string')
65+
}
66+
if (!Array.isArray(body.eventTypes) || !body.eventTypes.every((value) => typeof value === 'string')) {
67+
throw new DomainError('INVALID_ARGUMENT', 'eventTypes must be an array of strings')
68+
}
69+
if (
70+
body.resourceIds !== undefined &&
71+
(!Array.isArray(body.resourceIds) || !body.resourceIds.every((value) => typeof value === 'string'))
72+
) {
73+
throw new DomainError('INVALID_ARGUMENT', 'resourceIds must be an array of strings')
74+
}
75+
76+
const createSubscription = createWebhookSubscriptionService({
77+
repository: createWebhookSubscriptionCreationRepository(),
78+
clock: () => new Date(),
79+
createId: (kind) => kind === 'webhook-subscription'
80+
? randomUUID()
81+
: `${kind}-${randomUUID()}`,
82+
})
83+
const result = await createSubscription({
84+
workspaceId: actor.workspaceId,
85+
endpointId: body.endpointId,
86+
eventTypes: body.eventTypes,
87+
...(body.resourceIds ? { resourceIds: body.resourceIds as string[] } : {}),
88+
createdByClientId: actor.clientId,
89+
idempotencyKey,
90+
})
91+
return NextResponse.json(
92+
presentSuccess({
93+
subscription: presentWebhookSubscription(result.subscription),
94+
replayed: result.replayed,
95+
}),
96+
{
97+
status: result.replayed ? 200 : 201,
98+
headers: publicApiHeaders(requestId),
99+
},
100+
)
101+
} catch (error) {
102+
return respondPublicError(error, requestId)
103+
}
104+
}
Lines changed: 90 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,90 @@
1+
import { assertDomain } from '../domain/errors.ts'
2+
import { createWebhookSubscription } from '../domain/webhook.ts'
3+
import type {
4+
WebhookSubscriptionCreationRepository,
5+
WebhookSubscriptionCreationResult,
6+
} from './ports/webhook-subscription-creation-repository.ts'
7+
import { calculateVersionHash } from './version-hash.ts'
8+
9+
export type WebhookSubscriptionCreationEntityKind =
10+
| 'webhook-subscription'
11+
| 'idempotency-record'
12+
13+
export interface CreateWebhookSubscriptionRequest {
14+
workspaceId: string
15+
endpointId: string
16+
eventTypes: readonly string[]
17+
resourceIds?: readonly string[]
18+
createdByClientId: string
19+
idempotencyKey: string
20+
idempotencyTtlSeconds?: number
21+
}
22+
23+
export interface CreateWebhookSubscriptionDependencies {
24+
repository: WebhookSubscriptionCreationRepository
25+
clock: () => Date
26+
createId: (kind: WebhookSubscriptionCreationEntityKind) => string
27+
}
28+
29+
const DEFAULT_IDEMPOTENCY_TTL_SECONDS = 24 * 60 * 60
30+
const PRINTABLE_IDEMPOTENCY_KEY = /^[\x21-\x7e]+$/
31+
32+
export function createWebhookSubscriptionService(
33+
dependencies: CreateWebhookSubscriptionDependencies,
34+
) {
35+
return async function execute(
36+
request: CreateWebhookSubscriptionRequest,
37+
): Promise<Readonly<WebhookSubscriptionCreationResult>> {
38+
const workspaceId = request.workspaceId.trim()
39+
const endpointId = request.endpointId.trim()
40+
const clientId = request.createdByClientId.trim()
41+
const idempotencyKey = request.idempotencyKey.trim()
42+
const ttlSeconds = request.idempotencyTtlSeconds ?? DEFAULT_IDEMPOTENCY_TTL_SECONDS
43+
44+
assertDomain(workspaceId.length > 0, 'INVALID_ARGUMENT', 'workspaceId is required')
45+
assertDomain(clientId.length > 0, 'INVALID_ARGUMENT', 'createdByClientId is required')
46+
assertDomain(
47+
idempotencyKey.length >= 1 && idempotencyKey.length <= 128 && PRINTABLE_IDEMPOTENCY_KEY.test(idempotencyKey),
48+
'INVALID_ARGUMENT',
49+
'Idempotency-Key must contain 1 to 128 printable ASCII characters',
50+
)
51+
assertDomain(
52+
Number.isInteger(ttlSeconds) && ttlSeconds >= 60 && ttlSeconds <= 7 * 24 * 60 * 60,
53+
'INVALID_ARGUMENT',
54+
'idempotency ttlSeconds must be between 60 seconds and 7 days',
55+
)
56+
57+
const now = dependencies.clock()
58+
const createdAt = now.toISOString()
59+
const subscription = createWebhookSubscription({
60+
id: dependencies.createId('webhook-subscription'),
61+
workspaceId,
62+
endpointId,
63+
status: 'active',
64+
filter: {
65+
eventTypes: request.eventTypes,
66+
...(request.resourceIds ? { resourceIds: request.resourceIds } : {}),
67+
},
68+
createdByClientId: clientId,
69+
createdAt,
70+
})
71+
const requestFingerprint = calculateVersionHash({
72+
action: 'webhook-subscription-create/v1',
73+
endpointId: subscription.endpointId,
74+
filterHash: subscription.filter.hash,
75+
})
76+
77+
return dependencies.repository.createOrReplay({
78+
subscription,
79+
idempotency: {
80+
id: dependencies.createId('idempotency-record'),
81+
workspaceId,
82+
clientId,
83+
key: idempotencyKey,
84+
requestFingerprint,
85+
requestedAt: createdAt,
86+
expiresAt: new Date(now.getTime() + ttlSeconds * 1000).toISOString(),
87+
},
88+
})
89+
}
90+
}
Lines changed: 27 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,27 @@
1+
import type { WebhookSubscription } from '../../domain/webhook.ts'
2+
3+
export interface WebhookSubscriptionCreationIdempotency {
4+
id: string
5+
workspaceId: string
6+
clientId: string
7+
key: string
8+
requestFingerprint: string
9+
requestedAt: string
10+
expiresAt: string
11+
}
12+
13+
export interface WebhookSubscriptionCreationBundle {
14+
subscription: Readonly<WebhookSubscription>
15+
idempotency: Readonly<WebhookSubscriptionCreationIdempotency>
16+
}
17+
18+
export interface WebhookSubscriptionCreationResult {
19+
subscription: Readonly<WebhookSubscription>
20+
replayed: boolean
21+
}
22+
23+
export interface WebhookSubscriptionCreationRepository {
24+
createOrReplay(
25+
bundle: Readonly<WebhookSubscriptionCreationBundle>,
26+
): Promise<Readonly<WebhookSubscriptionCreationResult>>
27+
}

src/v2/domain/errors.ts

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -33,6 +33,8 @@ export const DOMAIN_ERROR_CODES = [
3333
'WEBHOOK_ENDPOINT_TRANSITION_REJECTED',
3434
'WEBHOOK_ENDPOINT_REVISION_MISMATCH',
3535
'WEBHOOK_SUBSCRIPTION_NOT_FOUND',
36+
'WEBHOOK_SUBSCRIPTION_ALREADY_EXISTS',
37+
'WEBHOOK_SUBSCRIPTION_CREATE_REJECTED',
3638
'WEBHOOK_SUBSCRIPTION_TRANSITION_REJECTED',
3739
'WEBHOOK_SUBSCRIPTION_REVISION_MISMATCH',
3840
'WEBHOOK_DELIVERY_NOT_FOUND',

0 commit comments

Comments
 (0)