StatifierRouter

CI Hex.pm Version Hex Downloads Hex Docs License

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

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.