Current section
Files
Jump to
Current section
Files
README.md
# AshHooks[](https://hex.pm/packages/ash_hooks)Webhooks for [Ash Framework](https://ash-hq.org) — **inbound** (receive,verify per-provider signatures, deduplicate, invoke the provider handlerand record its outcome) and **outbound** (sign, deliver, retry, track).> **Status: v0.1.0 — inbound + outbound complete.** Not yet shipped:> retention/TTL cleanup hooks (ledger and delivery rows accumulate until> you clean them — ADR-0005 names the hooks as a floor; tracked for a> following release).> Inbound: per-provider signature verification (`AshHooks.Provider`> behaviour; ComplyCube + HubSpot v3 reference providers), a fenced> unique-ingest ledger with claim/lease fencing and a reaper, and> fail-closed DSL verifiers. Outbound: `AshHooks.dispatch/4` fanout with> per-endpoint isolation, the delivery runtime on Oban> (`use AshHooks.Worker`) with row-owned retry policy, Standard Webhooks> signing (v1/v1a, legacy `:dual` migration mode), and a memory-bounded> native HTTP adapter. Security floors ship in the package (ADR-0005):> secrets as sources only, SSRF guards at registration + send, response> snippets store no body bytes by default. Telemetry events for the> whole lifecycle (see `AshHooks.Telemetry`). Design records:> [#1](https://github.com/baselabs/ash_hooks/issues/1),> ADR-0001–0008.Inbound and outbound are independently consumable: inbound-only applicationspull no queue infrastructure.## Installation```elixirdef deps do [ {:ash_hooks, "~> 0.1.0"}, # outbound delivery only (inbound-only apps need none of these): {:oban, "~> 2.20"} ]end```Elixir ~> 1.15, Ash ~> 3.0. Optional components: **Oban** (~> 2.20) foroutbound delivery — added, migrated, configured, and supervised by yourapp; **Plug/Phoenix** for inbound receipt (the raw-body reader below).Or with [igniter](https://hex.pm/packages/igniter):```mix igniter.install ash_hooks```The installer ATTEMPTS to patch your endpoint's `Plug.Parsers` with a`body_reader` — signature schemes sign the exact wire bytes, and a routerplug cannot recover pre-parser bytes. Review the generated diff; if itcould not locate the call, add it by hand:```elixirplug Plug.Parsers, parsers: [:json], pass: ["*/*"], body_reader: {AshHooks.BodyReader, :read_body, []}, json_decoder: Phoenix.json_library()```By default every parsed request carries a cached copy of its raw body;pass `[only: ["/webhooks"]]` as the reader's third element to scope thatmemory cost to the webhook routes.Your migrations create the tables AND the two UNIQUE INDEXES the dedupguarantees rest on (`[provider, external_event_id | scope]` for inbound,`[endpoint_id, event_uuid]` for outbound) — the[get-started tutorial](https://github.com/baselabs/ash_hooks/blob/main/documentation/tutorials/get-started.md)carries complete runnable shapes.## UsageAttach to a resource and declare sources and events:```elixiruse Ash.Resource, data_layer: AshSqlite.DataLayer, extensions: [AshHooks, AshHooks.InboundDelivery]inbound_delivery do # provider event ids are not globally unique across accounts — the # unique-ingest identity extends by your scope slots; each slot must # be a non-nullable attribute (the DSL verifier enforces both) scope_identity([:account_id])endattributes do attribute(:account_id, :string, allow_nil?: false)endwebhooks do # convention-resolves to AshHooks.Provider.ComplyCube inbound :comply_cube do secret {:app_env, [:my_app, :complycube_secret]} end # convention-resolves to AshHooks.Provider.HubSpotV3 — the vendor's # five-minute replay window applies by default; replay_window_seconds # overrides it in either direction inbound :hub_spot_v3 do secret {:app_env, [:my_app, :hubspot_client_secret]} end outbound :order_paid do signing_mode :standard endend```Inbound (sync mode) — a controller reads the cached raw body and drives thefenced machine:```elixirAshHooks.Ingress.ingest(Ledger, :comply_cube, conn.private[:ash_hooks_raw_body], %{ signature: get_req_header(conn, "complycube-signature") |> List.first(), headers: Map.new(conn.req_headers), scope: %{account_id: connection.account_id}})```HubSpot's v3 scheme signs the method and the full request URI alongsidethe body, so its controller passes both. Build the PUBLIC URI from valuesyou configure (a base URL you control), not blindly from `conn` — behinda TLS-terminating proxy, `conn.host`/port/scheme are the INTERNAL onesand the signature will not match the URI HubSpot signed. Where you doreconstruct from `conn`, remember `conn.query_string` excludes the `?`and `conn.host` excludes a non-default port:```elixirquery = if conn.query_string == "", do: "", else: "?" <> conn.query_stringAshHooks.Ingress.ingest(Ledger, :hub_spot_v3, conn.private[:ash_hooks_raw_body], %{ signature: get_req_header(conn, "x-hubspot-signature-v3") |> List.first(), headers: Map.new(conn.req_headers), method: conn.method, request_uri: "https://" <> conn.host <> conn.request_path <> query})```HubSpot delivers batches — a top-level JSON array of event objects. Theledger stores the array verbatim; a homogeneous batch parses to itssubscription type (`contact.creation` → `:contact_creation`), a mixed batchto `:mixed` (fan out per event in your handler), and an undocumentedsubscription type fails closed into the ledger as `failed_permanent`(`unknown_event_type`) — recorded and auditable.The machine persists the decoded payload (with a digest binding it tothe signed raw bytes) before handling, deduplicates onstorage-level uniqueness (exactly one `:created` per delivery, concurrentor sequential), and fences claims with a monotonic token and an expiringlease. Once the durable row exists, a crash between any two stepsre-drives on redelivery instead of silently dropping, and a stale owner(superseded or expired lease) can never mark. Terminal rows(`:processed` / `:failed_permanent`) are never processed again, but acrash after handler side effects and before the ledger mark re-invokesthe handler on redelivery — durable deduplication with AT-LEAST-ONCEhandler invocation; write handlers idempotent, keyed on the externalevent identity.Handler outcomes land in the ledger (`:processed`,`:failed_retryable`, `:failed_permanent`); expired leases are re-driven by`AshHooks.Ingress.reap/1`.Outbound deliveries are signed per the [Standard Webhooks](https://www.standardwebhooks.com)specification (`webhook-id` / `webhook-timestamp` / `webhook-signature`, `v1`HMAC-SHA256 and `v1a` ed25519 — old+new key rotation on both schemes), soreceivers verify with any conformant library; a `:dual` mode additionallyemits a legacy envelope during receiver migration — it REQUIRES theendpoint to carry a `legacy_secret_ref` (the resolver's base value alwayssigns the Standard Webhooks envelope; legacy slots come only from theendpoint's legacy references).Outbound fanout (landed): declare the Subscription / Endpoint /OutboundDelivery resources on your data layer, point the outbounddeclaration at them, and dispatch:```elixirdefmodule MyApp.WebhookEndpoint do use Ash.Resource, data_layer: AshSqlite.DataLayer, domain: MyApp, extensions: [AshHooks.Endpoint] sqlite do table("webhook_endpoints") repo(MyApp.Repo) end actions do defaults([:read, :create, :update]) endenddefmodule MyApp.WebhookSubscription do use Ash.Resource, data_layer: AshSqlite.DataLayer, domain: MyApp, extensions: [AshHooks.Subscription] sqlite do table("webhook_subscriptions") repo(MyApp.Repo) end actions do defaults([:read, :create]) end subscription do endpoint_resource(MyApp.WebhookEndpoint) endenddefmodule MyApp.OutboundDelivery do use Ash.Resource, data_layer: AshSqlite.DataLayer, domain: MyApp, extensions: [AshHooks.OutboundDelivery] sqlite do table("outbound_deliveries") repo(MyApp.Repo) end actions do defaults([:read]) endend``````elixirwebhooks do outbound :order_paid do subscriptions(MyApp.WebhookSubscription) deliveries(MyApp.OutboundDelivery) endend```With no `enqueue:` configured this persists `:pending` rows — thedurable ledger only; nothing sends until a runtime drives them:```elixir{:ok, event} = AshHooks.Event.new(type: :order_paid, payload: Jason.encode!(order))AshHooks.dispatch(Order, :order_paid, event)```Each matching enabled endpoint gets a durable delivery row unique on`{endpoint_id, event_uuid}` — the same pair the Oban job uniqueness keysuse — carrying the exact payload bytes to sign and the frozen effectivesigning mode. Endpoints store secret REFERENCES only (`whsec_`-shapedliterals are rejected at cast, on every write path); endpoints carry adurable `:enabled | :disabled` state the dispatcher respects. Oneendpoint's enqueue failure records `:enqueue_failed` on its row and neverstops its siblings; a re-dispatch claims the failed row via a CAS andretries the enqueue exactly once per won claim. With no `:enqueue`configured, rows persist `:pending` (`:deferred` results).Delivery runtime (landed): define ONE worker module in your app and wireit as the dispatch enqueuer —```elixirdefmodule MyApp.WebhookDeliveryWorker do use AshHooks.Worker, deliveries: MyApp.OutboundDelivery, endpoints: MyApp.WebhookEndpoint, secret_resolver: {MyApp.Secrets, :webhook_secret}, queue: :webhooksendAshHooks.dispatch(Order, :order_paid, event, enqueue: {MyApp.WebhookDeliveryWorker, :enqueue})```The worker drives `AshHooks.Delivery` (ADR-0008: the ROW owns the retrypolicy — attempts, `next_attempt_at`, the dead-letter ceiling — and Obanis the durable trigger). Sends are Standard-Webhooks signed per the row'smode with the same `webhook-id` on every retry; only 2xx succeeds;redirects are never followed; 410 disables the endpoint durably;408/429 honor `Retry-After` (bounded); 5xx/transport failures back offexponentially with jitter; other 4xx and refused redirects dead-letterimmediately. Response snippets store NO body bytes by default — a fixedsummary of the status and an allowlisted content-type token (ADR-0005'ssnippet amendment); a per-call `snippet_capture: true` in the`AshHooks.Delivery.run/2` config opts one diagnostic run into bodycapture, which persists the `[captured]`-marked body under the packagefloor (NFKC homoglyph folding, a bounded-fixpoint decode chain,separator-tolerant marker patterns, and a ≥16-char union-alphabetentropy rule) with an optional fail-closed `snippet_redactor` callbackahead of it. Machine-written fieldsaccept no action input. SSRF is guarded at registration (the endpoint's`url` type rejects private/loopback/link-local/metadata literals andnon-http schemes on every write path) and re-checked at send with DNSre-resolution. HTTP goes through the `AshHooks.Http` adapter behaviour —the default is `AshHooks.Http.Bounded`, a minimal HTTP/1.1 client whoseEVERY read is capped (headers, and bodies under Content-Length, chunked,and read-to-close framings alike — no response can balloon a worker'smemory); `AshHooks.Http.Httpc` (OTP `:httpc`) is available as analternative, and you can inject your own for tests or proxies. The packagestill compiles and runs Oban-free (CI no-optional leg + the inbound-onlyproof); `use AshHooks.Worker` without Oban on the host failsdeterministically at compile.### Dispatch-time capture (consumer-owned, no package change)`snippet_capture` is deliberately a per-call runtime config key, not aworker-macro knob — but dispatch-time opt-in needs no contract change:bring your own enqueue seam and your own Oban worker driving the publicruntime with the flag merged in.```elixirdefmodule MyApp.CaptureWorker do use Oban.Worker, queue: :webhook_diagnostics def perform(%Oban.Job{args: args}) do AshHooks.Delivery.run(args, deliveries: MyApp.OutboundDelivery, endpoints: MyApp.WebhookEndpoint, secret_resolver: {MyApp.Secrets, :webhook_secret}, snippet_capture: true ) endend# the enqueue seam contract is (delivery, event) -> :ok | {:error, term}AshHooks.dispatch(Order, :order_paid, event, enqueue: fn delivery, _event -> %{ "endpoint_id" => to_string(delivery.endpoint_id), "event_uuid" => delivery.event_uuid } |> MyApp.CaptureWorker.new() |> Oban.insert() |> case do {:ok, _job} -> :ok {:error, reason} -> {:error, reason} end end)```For a one-row diagnostic re-drive of an already-dispatched event, call`AshHooks.Delivery.run/2` directly with `snippet_capture: true` — therow's `{endpoint_id, event_uuid}` args and your config are all it takes.## ObservabilityNine lifecycle events cover the send/receive hot paths — ingressverify/dedup/claim, dispatch enqueue-failure, deliveryattempt/result/backoff/dead-letter/endpoint-disable. (Successfuldispatches and post-claim inbound outcomes are observed on the ledgerrows themselves, not as events.) Events carry ids,integers, fixed-vocabulary atoms, and classified reason strings only —never secrets, bodies, or payloads (ADR-0005), so they are safe to shipto any metrics/APM backend. `:telemetry.execute/3` matches exact eventnames, so consume the surface with one `attach_many` — the full listand a copy-paste handler live in `AshHooks.Telemetry`'s docs and the[get-started tutorial](https://github.com/baselabs/ash_hooks/blob/main/documentation/tutorials/get-started.md).## Design recordsArchitectural decisions live in[`docs/adr/`](https://github.com/baselabs/ash_hooks/tree/main/docs/adr).## LicenseMIT.