StatifierRouter
Delivers external events to the right durable statifier execution, creating it when absent, for Elixir developers whose events arrive from other systems, late or twice. Bindings pick the execution, an address table remembers it, dedupe records the second copy instead of delivering it, and one key's events step one at a time.
Why this package
A parcel is scanned at the depot, by the carrier and by the van's handheld, and each system reports it on its own schedule: a scan arrives late, a webhook is retried, the doorstep scan lands before the depot's. Every one of those events belongs to the one durable execution that tracks that parcel, which has to exist from the first scan on. Written by hand, that is a lookup table from parcel to execution, a dedupe table, a lock, and one transaction around all three, for every source.
With this package you declare bindings instead: which source, which events, which key, which document. Each event is routed in one transaction on statifier_persistence: the dedupe claim, the address lookup or the create, the step, and a ledger row naming the outcome. A retried webhook is recorded as a duplicate, a late scan for a finished parcel is recorded as dropped, and two scans of one parcel step its execution one after the other. The package starts no process: the host starts the Broadway pipeline in its own tree and schedules the reapers.
Installation
def deps do
[
{:statifier_router, "~> 0.12.0"}
]
end
The router's tables come from one migration of the host's own, run after
statifier_persistence's
(StatifierPersistence.Ecto.Migrations):
defmodule MyApp.Repo.Migrations.AddStatifierRouter do
use Ecto.Migration
def up, do: StatifierRouter.Migrations.up()
def down, do: StatifierRouter.Migrations.down()
end
Basic usage
Two systems scan the same parcel, one sending its id as a number and one as a
padded string; both bindings trim it to one key, so both reach one execution.
MyApp.Persistence is the host's use StatifierPersistence.Ecto, repo: MyApp.Repo module.
{:ok, machine} = Statifier.compile(File.read!("priv/charts/parcel_delivery.scxml"))
hash = Statifier.Machine.identity(machine).content_hash
{:ok, store} = StatifierPersistence.Storage.new(StatifierPersistence.Storage.Ecto, persistence: MyApp.Persistence)
{:ok, resolver} = StatifierRouter.Resolver.Static.new(%{{"depot_north", "parcel_delivery"} => machine})
{:ok, config} =
StatifierRouter.Config.new(
repo: MyApp.Repo,
store: store,
executor: fn _effect, _context -> :ok end,
resolver: resolver,
chart_resolver: fn ^hash -> {:ok, machine} end,
bindings: [
%{id: "depot_scans", source: "depot_scanners", match: ~s(event.kind == "scan"),
key: "trim(event.parcel_id::string)", document: "parcel_delivery", event: "parcel.scanned"},
%{id: "carrier_scans", source: "carrier", match: ~s(event.kind == "delivered"),
key: "trim(event.parcel_id::string)", document: "parcel_delivery", event: "parcel.delivered"}
]
)
scan = %{scope: "depot_north", source: "depot_scanners", message_id: "scan-881",
data: %{"kind" => "scan", "parcel_id" => 1042771}}
StatifierRouter.route(config, scan)
#=> {:ok, [{:created_and_delivered, "depot_scans", "ex_..."}]}
StatifierRouter.route(config, scan)
#=> {:ok, [{:duplicate, "depot_scans"}]}
StatifierRouter.route(config, %{scope: "depot_north", source: "carrier", message_id: "c-77",
data: %{"kind" => "delivered", "parcel_id" => " 1042771 "}})
#=> {:ok, [{:delivered, "carrier_scans", "ex_..."}]}
The second copy of the scan is recorded and not delivered, and the carrier's event reaches the execution the depot's scan created.
Documentation
- Learn
- Basic usage: two sources, one parcel, one execution, and a duplicate recorded rather than delivered.
- Do
- How to route a producer's messages through the Broadway pipeline: the pipeline in the host's tree, its configuration, and what a failed message does.
- How to bind events from several sources to one execution: bindings, a key program that normalizes, and bindings that differ by scope.
- How to deliver a chart's sends to a sink: routes, a transactional outbox end to end, and a finished execution's hand-off.
- How to take webhooks and form posts: the webhook front, the host's signature check, the message id and the status.
- How to give an execution an HTTP location: the BasicHTTP front, its tokens, and sending from a durable execution.
- How to resolve a document to its chart: the resolver behaviour, the static resolver, and the chart an existing execution keeps.
- How to fit the router into an engine of your own: the create and step hooks, the whole-delivery wrapper, execution ids, send types and a timer queue.
- How to fit the router's tables to a host: the reapers, a host column at a fixed position, and a primary key of the host's own.
- Upgrading a host from 0.6 to 0.12: what a host changes for each minor, the V03 migration, and the opt-in location table.
- Look up
- The configuration: every option, its default and how it is checked.
- The binding: its fields, their defaults and the key rules.
- Routing and its outcomes:
route/3, the outcome vocabulary and what each outcome writes. - The migrations: the versions, their options and the table prefix.
- The conformance corpus: the language-neutral cases and their format.
- The changelog: what changed in each version, with every breaking change marked.
- Understand
- What the router owns, and what it leaves to the host: the pieces, where each lives, and why an address row pins a chart.
- Why one key's events step one at a time: the partitioner, the per-execution lock, the ceiling it sets and the ways around it.
- The decision records: why bindings, addressing, delivery, routes and the fronts are shaped the way they are.
Compatibility
The package needs Elixir 1.18 or later (elixir: "~> 1.18" in mix.exs). It
builds on statifier ~> 2.10, statifier_persistence ~> 0.18,
predicator ~> 9.4, broadway ~> 1.3 and ecto_sql ~> 3.14, and brings no
database driver: the host's repo does. Its CI runs the suite against
Postgres, and the migrations are also run on SQLite.
Until 1.0, the public surface may change between minor releases: a release may
rename modules, callbacks, table columns, telemetry events or error vocabulary
with no compatibility shim. Every such change is recorded in the
changelog
under a bold Breaking heading that says what to do about it, and pinning
to an exact minor, ~> X.Y.0, is the recommended way to take the package
until then.
License
MIT - see LICENSE.