From 28f07b7b587b5321400cfb79945d38426004e4fa Mon Sep 17 00:00:00 2001 From: JohnnyT Date: Fri, 2 Oct 2026 03:52:20 -0600 Subject: [PATCH 1/2] Posts to the hold desk after the commit The hold desk's executor no longer makes its BasicHTTP POST inside the delivery's transaction. It plans the send with the router's processor and inserts a HoldDesk.DeskPost job on this app's Oban, in the new desk_posts queue: the job commits with the step that sent, a delivery that rolls back takes it with it, and it is unique on the send's dedup key, so a redriven step inserts one job. The job makes the POST after the commit, so a slow desk holds no transaction and no SQLite write lock. A POST the desk refuses is retried; once the third has failed, the job delivers error.communication back into the hold through StatifierRouter.Delivery.deliver_event/4 with create: :never over the row from Addresses.by_execution/2, as statifier_router's ADR-0002 Amendment on the outbound BasicHTTP send recommends. A failure that reaches no execution is cancelled and kept as the dead letter. The controller tests drain the queue, a new test shows the desk is called after the delivery returned and outside any transaction, and the refused-desk test reads error.communication from the execution's input log. The guide says why the example performs after the commit. Refs: se-g9t4 --- config/config.exs | 6 +- docs/guides/basichttp-front.md | 74 +++++--- lib/statifier_examples/hold_desk.ex | 50 +++--- lib/statifier_examples/hold_desk/desk_post.ex | 149 ++++++++++++++++ .../basic_http_controller_test.exs | 165 ++++++++++++++---- test/support/desk_transport.ex | 9 +- 6 files changed, 376 insertions(+), 77 deletions(-) create mode 100644 lib/statifier_examples/hold_desk/desk_post.ex diff --git a/config/config.exs b/config/config.exs index a29c348..d0ad034 100644 --- a/config/config.exs +++ b/config/config.exs @@ -43,6 +43,9 @@ config :statifier_examples, StatifierExamples.Repo, # package leaves to the host: `parcel_notices` carries the hand-off its # one route makes, and `router_maintenance` the reapers. The recipe # checks that both reapers are on this crontab. +# +# `desk_posts` carries the hold desk's BasicHTTP POSTs, each made after the +# delivery that planned it has committed (`StatifierExamples.HoldDesk.DeskPost`). config :statifier_examples, Oban, repo: StatifierExamples.Repo, engine: Oban.Engines.Lite, @@ -51,7 +54,8 @@ config :statifier_examples, Oban, statifier_timers: 5, statifier_invocations: 5, parcel_notices: 1, - router_maintenance: 1 + router_maintenance: 1, + desk_posts: 1 ], plugins: [ {Oban.Plugins.Cron, diff --git a/docs/guides/basichttp-front.md b/docs/guides/basichttp-front.md index 5667e9f..bb649c1 100644 --- a/docs/guides/basichttp-front.md +++ b/docs/guides/basichttp-front.md @@ -9,7 +9,8 @@ path every routed event takes. This guide walks the pieces as this app wires them: `StatifierExamples.HoldDesk` (the router configuration, the chart's -resolver and the executor), `priv/library/hold_desk.scxml` (the chart), +resolver and the executor), `StatifierExamples.HoldDesk.DeskPost` (the job +that makes the outbound POST), `priv/library/hold_desk.scxml` (the chart), `StatifierExamplesWeb.BasicHTTPController` (the front's action) and `test/statifier_examples_web/controllers/basic_http_controller_test.exs`, which drives the whole of it through the controller. @@ -35,7 +36,7 @@ answers 404. | Package | Version | What this guide uses it for | |---|---|---| -| `statifier_router` | 0.9.2 | the `:basichttp` key, the location, `StatifierRouter.BasicHTTP.Front`, the location table | +| `statifier_router` | 0.9.2 | the `:basichttp` key, the location, `StatifierRouter.BasicHTTP.Front`, the location table, `deliver_event/4` for a failed send | | `statifier` | 2.10.0 | the Basic HTTP Event I/O Processor and its decoder | | `statifier_persistence` | 0.24.0 | the chart registry, the execution, the input log | @@ -87,22 +88,53 @@ gives every router table. A durable execution has no session to perform its sends: the step hands each effect to the configuration's executor. `StatifierExamples.HoldDesk.execute/2` -plans a BasicHTTP send with `StatifierRouter.BasicHTTP.deliver/3` and -performs what it planned with the processor's `perform/2`, which POSTs a -form body to the desk with the send's `scxml-send-key` header. The POST -goes through the configuration's `:transport`: statifier's default, -on OTP's `:httpc`, in the dev app, and a transport under test that hands -the POST back to the test instead of sending it. - -The POST is made inside the delivery's transaction, before the step -commits. A desk that does not answer 2xx, or does not answer at all, is a -failed send: `statifier_persistence` enters `error.communication`, -carrying the send id, into the execution in the same step, and the chart -takes it from `waiting` to its other final state, `desk_unreached`. A send -with no target, which statifier plans as an `error.communication` raise and -no request, fails without a POST and enters the execution the same way. A -delayed BasicHTTP send is refused the same way, because its timer would -live in the delivering process rather than in the database. +plans a BasicHTTP send with `StatifierRouter.BasicHTTP.deliver/3`, which +answers the POST to make - a form body for the desk, with the send's +`scxml-send-key` header - and makes none. + +The executor runs inside the delivery's transaction, so it does not make +the POST there. It inserts a `StatifierExamples.HoldDesk.DeskPost` job on +this app's own Oban, in the `desk_posts` queue. The job writes through the +same repo, so it commits with the step that sent the POST and a delivery +that rolls back takes the job with it: the jobs table is the outbox. The +job is unique on the send's dedup key written out, so a step that is +driven again, and re-emits the same send, inserts no second job. + +The job performs the POST after the delivery has committed, with the +processor's `perform/2`, through the configuration's `:transport`: +statifier's default, on OTP's `:httpc`, in the dev app, and a transport +under test that hands the POST back to the test instead of sending it. +This is why the example performs after the commit, as `statifier_router` +recommends in `docs/adr/0002-addressing.md`, the Amendment of 2026-10-02 +on a durable execution's outbound BasicHTTP send. +Made inside the delivery, the POST would keep the transaction, the +execution's lock and SQLite's single write lock held for as long as the +desk took to answer, and would leave even for a step that then rolled +back. Made from the job, a slow desk holds none of them, and no POST +leaves for a step that never committed. + +A desk that does not answer 2xx, or does not answer at all, is retried: +the job makes the POST up to three times. A send that still fails comes +back into the execution through `StatifierRouter.Delivery.deliver_event/4`, +the one way back in the router's record names, with `create: :never` over +the hold's address row from `StatifierRouter.Addresses.by_execution/2`: an +external `error.communication` event carrying the send's id, delivered +in a step of its own under the plan id `desk_post_failure`. The chart +takes it from `waiting` to its other final state, `desk_unreached`. A +failure that reaches no execution - the hold has finished, or its address +row is gone - is cancelled, which keeps the job and its reason in the jobs +table as the dead letter. + +A send with no target, which statifier plans as an `error.communication` +raise and no request, plans no job: the executor fails it at once, and +`statifier_persistence` enters `error.communication` into the execution +in the same step. A delayed BasicHTTP send is refused the same way, +because its timer would live in the delivering process rather than in the +database. + +The job's arguments carry the planned POST as it was planned, body +included, so the `reply_to` location the body hands the desk is written to +the jobs table with it, and stays there until the host prunes the job. ## The front @@ -128,8 +160,10 @@ body of a request under `/basichttp` for the action to hand on. ## Driving it -The controller test routes a hold request, reads the `hold.placed` POST the -desk was sent, takes `reply_to` from it and POSTs +The controller test routes a hold request, drains the `desk_posts` queue +with `Oban.drain_queue/2` (the suite runs Oban with `testing: :manual`), +reads the `hold.placed` POST the desk was sent, takes `reply_to` from it +and POSTs `_scxmleventname=copy.shelved` at that location through the endpoint: ```elixir diff --git a/lib/statifier_examples/hold_desk.ex b/lib/statifier_examples/hold_desk.ex index 3906b61..f558cf2 100644 --- a/lib/statifier_examples/hold_desk.ex +++ b/lib/statifier_examples/hold_desk.ex @@ -30,17 +30,23 @@ defmodule StatifierExamples.HoldDesk do processor's type strings, which is what lets the chart's `` pass the engine's type check and read `_ioprocessors['basichttp']`. A durable execution has no session to perform the send, so the effect - reaches `execute/2`, which plans it with the router's processor and - performs what it planned, through the configuration's transport. That - runs inside the delivery's transaction: one POST to the desk, made - before the step commits. A POST the desk does not answer with a 2xx is - a failed send, which `statifier_persistence` enters into the execution - as `error.communication`, and the chart ends the hold unreached. A - send with no target, which statifier plans as an `error.communication` - raise and no request, fails without a POST and enters the execution the - same way. A delayed BasicHTTP send, whose timer would live in this - process rather than in the database, is refused, and enters the - execution the same way. + reaches `execute/2`, which plans it with the router's processor. That + runs inside the delivery's transaction, so the POST it planned is not + made there: it is handed to a `StatifierExamples.HoldDesk.DeskPost` job, + inserted in the same transaction and performed after the delivery + commits, through the configuration's transport. A desk that is slow to + answer then holds no transaction and no SQLite write lock, and no POST + leaves for a step that rolled back, as `statifier_router` recommends + (its ADR-0002, the Amendment on the outbound BasicHTTP send). A POST + the desk does not answer with a 2xx is retried, and a send that still + fails comes back into the execution as `error.communication`, delivered + by the job, and the chart ends the hold unreached. A send with no + target, which statifier plans as an `error.communication` raise and no + request, plans no job: it fails at once, and `statifier_persistence` + enters `error.communication` into the execution in the same step. A + delayed BasicHTTP send, whose timer would live in this process rather + than in the database, is refused, and enters the execution the same + way. """ @behaviour StatifierRouter.Resolver @@ -49,6 +55,7 @@ defmodule StatifierExamples.HoldDesk do alias Statifier.Machine alias Statifier.Send.Event, as: SendEvent alias StatifierExamples.FirstWorkflow + alias StatifierExamples.HoldDesk.DeskPost alias StatifierExamples.RoutedWorkflow.Stepper alias StatifierPersistence.Storage alias StatifierRouter.{BasicHTTP, Config} @@ -163,9 +170,10 @@ defmodule StatifierExamples.HoldDesk do @doc """ The executor every create and step hands its effects to. A BasicHTTP - `` is planned with `StatifierRouter.BasicHTTP.deliver/3` and each - instruction it plans is performed with the processor's `perform/2`; a - send with no target, which statifier plans as an `error.communication` + `` is planned with `StatifierRouter.BasicHTTP.deliver/3`, and each + POST it plans is handed to a `StatifierExamples.HoldDesk.DeskPost` job, + inserted in the delivery's transaction and performed after it commits; + a send with no target, which statifier plans as an `error.communication` raise and no request, fails as `{:basichttp_send_without_target, send_id}`, and a delayed one is refused as `{:delayed_basichttp_send, send_id}`. Every other effect is passed. @@ -177,7 +185,7 @@ defmodule StatifierExamples.HoldDesk do ctx = %{session_id: execution_id, opts: basichttp()} event = SendEvent.build(send, execution_id) {:ok, instructions} = BasicHTTP.deliver(send, event, ctx) - perform(instructions, ctx, send) + enqueue(instructions, ctx, send) end def execute({:send_delayed, %SendDelayed{type: type} = send}, _context) @@ -186,13 +194,13 @@ defmodule StatifierExamples.HoldDesk do def execute(_effect, _context), do: :ok - @spec perform([term()], map(), Send.t()) :: :ok | {:error, term()} - defp perform(instructions, ctx, send) do + @spec enqueue([term()], map(), Send.t()) :: :ok | {:error, term()} + defp enqueue(instructions, ctx, send) do Enum.reduce_while(instructions, :ok, fn - {:handler, module, payload}, :ok -> - case module.perform(payload, ctx) do - :ok -> {:cont, :ok} - {:error, _reason} = error -> {:halt, error} + {:handler, _module, payload}, :ok -> + case payload |> DeskPost.new(ctx, send) |> Oban.insert() do + {:ok, _job} -> {:cont, :ok} + {:error, reason} -> {:halt, {:error, reason}} end # Statifier plans a send with no target as a raise of diff --git a/lib/statifier_examples/hold_desk/desk_post.ex b/lib/statifier_examples/hold_desk/desk_post.ex new file mode 100644 index 0000000..d361d7d --- /dev/null +++ b/lib/statifier_examples/hold_desk/desk_post.ex @@ -0,0 +1,149 @@ +defmodule StatifierExamples.HoldDesk.DeskPost do + @moduledoc """ + The job that tells a branch desk a hold was placed: one BasicHTTP POST, + made after the delivery that planned it has committed. + + `StatifierExamples.HoldDesk.execute/2` plans the hold's `` with `StatifierRouter.BasicHTTP.deliver/3` and + inserts one of these jobs for the POST it planned, through `new/3`. The + executor runs inside the delivery's transaction, and this app's Oban + writes through the same repo, so the job commits with the step that + sent and a delivery that rolls back takes the job with it. The job + carries the planned instruction and the plan context as they were + planned, written out with `:erlang.term_to_binary/1`, since the + instruction holds a struct and a module that job arguments cannot. + + The job is unique on `key`, the send's dedup key written out: a + redriven step re-emits the same send with the same fields and inserts + nothing new. + + `perform/1` makes the POST with `StatifierRouter.BasicHTTP.perform/2`. + It runs outside every delivery, so a slow desk holds no transaction, no + execution lock and no SQLite write lock while it answers. A POST the + desk does not take is retried; once the third POST has failed too, the + job delivers `error.communication`, carrying the send's id, back into + the hold through `StatifierRouter.Delivery.deliver_event/4`, over the + hold's address row, with `create: :never`. That is the only way back in + `statifier_router` sanctions for a send performed after the commit (its + ADR-0002, the Amendment on the outbound BasicHTTP send). A delivery + that does not settle is retried on the attempts left, without posting + again; a failure that reaches no execution - a hold already finished, or an + address row already reaped - is cancelled, which keeps the job and its + reason in the jobs table as the dead letter. + """ + + use Oban.Worker, + queue: :desk_posts, + max_attempts: 5, + unique: [keys: [:key], period: :infinity] + + # The attempts that POST; the ones after them only deliver the failure. + @post_attempts 3 + + alias Statifier.Effect.Send + alias StatifierExamples.HoldDesk + alias StatifierRouter.{Addresses, BasicHTTP, Delivery} + + # The plan the failure is delivered under. The id is this app's own and + # is neither `execution`, `basichttp` nor a binding's, so the ledger and + # dedupe rows it writes are told apart from theirs; the horizon is the + # router's default dedupe horizon, 72 hours. + @failure_plan %{ + id: "desk_post_failure", + create: :never, + dedupe: %{by: :message_id, horizon_ms: 259_200_000} + } + + @doc """ + The job for one planned instruction's `payload`, sent by `send` from the + execution `ctx.session_id` names. + """ + @spec new(term(), map(), Send.t()) :: Oban.Job.changeset() + def new(payload, %{session_id: execution_id} = ctx, %Send{} = send) do + new(%{ + "execution_id" => execution_id, + "send_id" => send.send_id, + "key" => key(execution_id, send), + "instruction" => Base.encode64(:erlang.term_to_binary({payload, ctx})) + }) + end + + @impl Oban.Worker + def perform(%Oban.Job{args: args, attempt: attempt}) when attempt <= @post_attempts do + {payload, ctx} = args["instruction"] |> Base.decode64!() |> :erlang.binary_to_term([:safe]) + + case BasicHTTP.perform(payload, ctx) do + :ok -> :ok + {:error, reason} when attempt < @post_attempts -> {:error, reason} + {:error, reason} -> failed(args, reason) + end + end + + # Every POST attempt failed and the failure's delivery did not settle: + # deliver it again, without posting again. + def perform(%Oban.Job{args: args}), do: failed(args, :desk_post_attempts_spent) + + # The POST has failed for good: tell the hold, or keep a dead letter. + @spec failed(map(), term()) :: :ok | {:error, term()} | {:cancel, term()} + defp failed(%{"execution_id" => execution_id} = args, reason) do + config = HoldDesk.config() + + case Addresses.by_execution(config, execution_id) do + nil -> + {:cancel, {:desk_unreached_without_execution, reason}} + + row -> + event = + Statifier.Event.external("error.communication", + sendid: args["send_id"], + data: %{"reason" => inspect(reason)} + ) + + config + |> Delivery.deliver_event(Map.put(@failure_plan, :document, row.document), row.key, %{ + event: event, + message_id: args["key"], + scope: row.scope, + now: DateTime.utc_now() + }) + |> settled(reason) + end + end + + @spec settled(StatifierRouter.outcome() | {:error, term()}, term()) :: + :ok | {:error, term()} | {:cancel, term()} + defp settled({:delivered, _plan, _execution_id}, _reason), do: :ok + defp settled({:duplicate, _plan}, _reason), do: :ok + defp settled({:dropped, _plan, why}, reason), do: {:cancel, {:desk_unreached, why, reason}} + defp settled({:error, _reason} = error, _reason_sent), do: error + + # The send's dedup key, the fields the `scxml-send-key` header carries, + # each nil written as "-". + @spec key(String.t(), Send.t()) :: String.t() + defp key(execution_id, send) do + Enum.map_join( + [ + execution_id, + send.send_id, + send.macrostep, + send.microstep, + send.round, + send.c_index, + owner(send.owner), + send.ordinal + ], + "/", + &field/1 + ) + end + + @spec field(String.t() | non_neg_integer() | nil) :: String.t() + defp field(nil), do: "-" + defp field(value) when is_integer(value), do: Integer.to_string(value) + defp field(value), do: URI.encode_www_form(value) + + @spec owner(Send.owner() | nil) :: String.t() | nil + defp owner({kind, state, block}), do: "#{kind}.#{state}.#{block}" + defp owner({:transition, transition}), do: "transition.#{transition}" + defp owner(nil), do: nil +end diff --git a/test/statifier_examples_web/controllers/basic_http_controller_test.exs b/test/statifier_examples_web/controllers/basic_http_controller_test.exs index d1ad8c6..ae975b9 100644 --- a/test/statifier_examples_web/controllers/basic_http_controller_test.exs +++ b/test/statifier_examples_web/controllers/basic_http_controller_test.exs @@ -1,9 +1,10 @@ defmodule StatifierExamplesWeb.BasicHTTPControllerTest do use StatifierExamplesWeb.ConnCase, async: false + import Ecto.Query, only: [from: 2] import ExUnit.CaptureLog - alias StatifierExamples.{FirstWorkflow, HoldDesk, RoutedWorkflow} + alias StatifierExamples.{FirstWorkflow, HoldDesk, Repo, RoutedWorkflow} alias StatifierPersistence.{Execution, Executions, Storage} alias StatifierRouter.BasicHTTP @@ -20,7 +21,9 @@ defmodule StatifierExamplesWeb.BasicHTTPControllerTest do :ok end - defp placed_hold(hold_id \\ "hold-0417") do + # The hold request's delivery returns before the desk is told: the POST + # is a job, performed when the test drains the desk's queue. + defp requested_hold(hold_id) do assert {:ok, [{:created_and_delivered, "hold_requests", execution_id}]} = HoldDesk.request(@scope, %{ "hold_id" => hold_id, @@ -28,6 +31,15 @@ defmodule StatifierExamplesWeb.BasicHTTPControllerTest do "desk" => @desk }) + execution_id + end + + defp drain_desk_posts, + do: Oban.drain_queue(queue: :desk_posts, with_scheduled: true, with_recursion: true) + + defp placed_hold(hold_id \\ "hold-0417") do + execution_id = requested_hold(hold_id) + assert %{success: 1, failure: 0} = drain_desk_posts() assert_received {:desk_post, @desk, headers, body} {execution_id, headers, URI.decode_query(body)} end @@ -52,8 +64,8 @@ defmodule StatifierExamplesWeb.BasicHTTPControllerTest do # chart's basichttp send was an unsupported type, no desk_post arrived # and the route answered an error, red; restored, green. # sabotage: execute/2's BasicHTTP clause made to answer :ok without - # performing -> assert_received {:desk_post, ...} failed, red; restored, - # green. + # inserting the job -> nothing drained and no desk_post arrived, red; + # restored, green. test "the hold tells the desk it was placed, with its own location to answer at" do {execution_id, headers, params} = placed_hold() @@ -159,47 +171,135 @@ defmodule StatifierExamplesWeb.BasicHTTPControllerTest do Executions.inputs(FirstWorkflow.store(), execution_id) end - # The chart has two finals, and a finished execution keeps no - # configuration to read the one it took from; the error.communication - # the refused POST re-entered is what names desk_unreached, since only - # that event leads there. - # sabotage: execute/2 made to answer :ok whatever perform/2 answered -> - # the execution stayed active in waiting and the location still took a - # POST, red; restored, green. - # sabotage: perform/3 made to continue past a failed POST -> no - # error.communication was re-entered, red; restored, green. - test "a desk that refuses the POST ends the hold unreached", %{conn: conn} do - reentered = [:statifier_persistence, :execution, :step, :reentered] - handler = "hold-desk-reentered-#{System.unique_integer([:positive])}" + # The delivery returns, and so has committed, before the desk is called; + # the desk is then called with no transaction open in the process that + # calls it, so however long it takes to answer it holds no delivery. + # sabotage: the executor made to perform the POST itself instead of + # inserting the job -> the desk was called before the delivery returned, + # red; restored, green. + # sabotage: the job's POST wrapped in a Repo transaction -> the desk was + # called inside one, red; restored, green. + test "a slow desk holds no delivery open: the POST is made after the commit" do test_pid = self() - forward = fn _event, _measurements, metadata, nil -> - send(test_pid, {:reentered, metadata}) - end + Process.put(:desk_answering, fn -> + send(test_pid, {:desk_answering, Repo.in_transaction?()}) + end) + + execution_id = requested_hold("hold-0420") - :ok = :telemetry.attach(handler, reentered, forward, nil) - on_exit(fn -> :telemetry.detach(handler) end) + refute_received {:desk_answering, _in_transaction} + refute_received {:desk_post, _url, _headers, _body} + assert status!(execution_id) == :active + assert [%{"execution_id" => ^execution_id}] = desk_post_args() + + assert %{success: 1, failure: 0} = drain_desk_posts() + assert_received {:desk_answering, false} + assert_received {:desk_post, @desk, _headers, _body} + end + # The chart has two finals, and a finished execution keeps no + # configuration to read the one it took from; the error.communication + # the job delivered is what names desk_unreached, since only that event + # leads there. It arrives as a delivered external event, in a step of + # its own, so it is in the execution's input log under the send's id. + # sabotage: the job's last failed POST made to cancel instead of + # delivering error.communication -> the hold stayed waiting, red; + # restored, green. + test "a desk that refuses the POST ends the hold unreached", %{conn: conn} do Process.put(:desk_status, 503) - {execution_id, _headers, %{"reply_to" => location}} = placed_hold("hold-0418") + execution_id = requested_hold("hold-0418") + + assert %{success: 1, failure: 2} = drain_desk_posts() + assert_received {:desk_post, @desk, _headers, body} + %{"reply_to" => location} = URI.decode_query(body) - assert_received {:reentered, - %{execution_id: ^execution_id, name: "error.communication", opts: opts}} + assert {:ok, + [ + %{event: %{name: "hold.requested"}}, + %{event: %{name: "error.communication", sendid: sendid}} + ]} = Executions.inputs(FirstWorkflow.store(), execution_id) - assert is_binary(opts[:sendid]) + assert is_binary(sendid) assert status!(execution_id) == :completed assert conn |> post_event(location, "_scxmleventname=copy.shelved") |> response(404) end + # sabotage: the job's unique option removed -> two jobs, red; restored, + # green. + test "a redriven send inserts one job, keyed on the send's dedup key" do + send = hold_send("https://riverside.example/holds-desk") + + assert :ok = HoldDesk.execute({:send, send}, %{execution_id: "ex_hold_0421"}) + assert :ok = HoldDesk.execute({:send, send}, %{execution_id: "ex_hold_0421"}) + + assert [%{"key" => "ex_hold_0421/send_1/1/0/0/0/transition.0/0", "send_id" => "send_1"}] = + desk_post_args() + end + + # The executor runs inside the delivery's transaction, and the job is + # inserted through the same repo, so a delivery that rolls back takes + # the job with it and no POST is made for it. + # sabotage: the insert moved to a Task outside the transaction -> the + # sandbox refused the Task a connection, so every hold test errored + # before any assertion; a sandboxed suite has one connection and cannot + # show this test red. Restored, green. + test "a delivery that rolls back takes its desk post with it" do + assert {:error, :rolled_back} = + Repo.transaction(fn -> + assert :ok = + HoldDesk.execute({:send, hold_send(@desk)}, %{ + execution_id: "ex_hold_0423" + }) + + assert [_job] = desk_post_args() + Repo.rollback(:rolled_back) + end) + + assert desk_post_args() == [] + assert %{success: 0} = drain_desk_posts() + refute_received {:desk_post, _url, _headers, _body} + end + + # A failed POST whose hold has finished, or whose execution has no + # address row, reaches no execution: the job is cancelled, which keeps it + # in the jobs table with its reason. + # sabotage: a missing address row made to answer :ok -> one cancel + # short, red; restored, green. + # sabotage: a dropped delivery made to answer :ok -> one cancel short, + # red; restored, green. + test "a failed POST that reaches no hold is kept as a dead letter", %{conn: conn} do + {execution_id, _headers, %{"reply_to" => location}} = placed_hold("hold-0422") + assert conn |> post_event(location, "_scxmleventname=copy.shelved") |> response(204) + + Process.put(:desk_status, 503) + send = hold_send(@desk) + assert :ok = HoldDesk.execute({:send, send}, %{execution_id: execution_id}) + assert :ok = HoldDesk.execute({:send, send}, %{execution_id: "ex_hold_nowhere"}) + + assert %{success: 0, failure: 4, cancelled: 2} = drain_desk_posts() + assert status!(execution_id) == :completed + end + # Statifier plans a send with no target as an error.communication raise - # and no request; the executor fails the send instead of performing it. - # sabotage: perform/3's raise clause made to continue -> execute/2 + # and no request; the executor fails the send instead of performing it, + # and plans no job. + # sabotage: enqueue/3's raise clause made to continue -> execute/2 # answered :ok, red; restored, green. test "a send with no target posts nothing and fails, naming the send" do - no_target = %Statifier.Effect.Send{ + assert HoldDesk.execute({:send, hold_send(nil)}, %{execution_id: "ex_hold_0419"}) == + {:error, {:basichttp_send_without_target, "send_1"}} + + assert desk_post_args() == [] + assert %{success: 0} = drain_desk_posts() + refute_received {:desk_post, _url, _headers, _body} + end + + defp hold_send(target) do + %Statifier.Effect.Send{ type: "basichttp", event: "hold.placed", - target: nil, + target: target, data: %{"hold_id" => "hold-0419"}, send_id: "send_1", c_index: 0, @@ -209,11 +309,10 @@ defmodule StatifierExamplesWeb.BasicHTTPControllerTest do round: 0, ordinal: 0 } + end - assert HoldDesk.execute({:send, no_target}, %{execution_id: "ex_hold_0419"}) == - {:error, {:basichttp_send_without_target, "send_1"}} - - refute_received {:desk_post, _url, _headers, _body} + defp desk_post_args do + Repo.all(from(job in Oban.Job, where: job.queue == "desk_posts", select: job.args)) end # sabotage: :basichttp added to RoutedWorkflow's configuration -> red; diff --git a/test/support/desk_transport.ex b/test/support/desk_transport.ex index 295482f..faa632a 100644 --- a/test/support/desk_transport.ex +++ b/test/support/desk_transport.ex @@ -2,15 +2,20 @@ defmodule StatifierExamples.DeskTransport do @moduledoc """ The test transport for `StatifierExamples.HoldDesk`'s outbound BasicHTTP POSTs: it sends `{:desk_post, url, headers, body}` to the - process performing the POST, which is the test process that routed the - hold request, and answers the status the process put under + process performing the POST, which is the test process that drained the + desk's job queue, and answers the status the process put under `:desk_status`, 204 as a branch desk would when it put none. + + A test that needs to see the moment the desk is called puts a + zero-arity function under `:desk_answering`; it is called first, while + the desk is still answering. """ @behaviour Statifier.Send.BasicHTTP.Transport @impl true def post(url, headers, body) do + if answering = Process.get(:desk_answering), do: answering.() send(self(), {:desk_post, url, headers, body}) {:ok, Process.get(:desk_status, 204)} end From f72a3f3caa8c4a668e7df78547d1f511128a2bf4 Mon Sep 17 00:00:00 2001 From: JohnnyT Date: Fri, 2 Oct 2026 04:04:45 -0600 Subject: [PATCH 2/2] Reads a desk post job without :safe DeskPost decoded its instruction with :erlang.binary_to_term/2 and :safe, which refuses an atom the node has not created yet. A job queued before a restart can name one that no module has loaded since (a send's owner kind, a struct field), so its POST attempts raised and the hold ended unreached with no POST made. The job now reads its instruction with :erlang.binary_to_term/1, as the router README's recipe does; the jobs table is written only through this app's repo. A new test writes an instruction naming an atom no code creates and shows the job still posts. Refs: se-g9t4 --- lib/statifier_examples/hold_desk/desk_post.ex | 10 ++++- .../basic_http_controller_test.exs | 45 ++++++++++++++++++- 2 files changed, 53 insertions(+), 2 deletions(-) diff --git a/lib/statifier_examples/hold_desk/desk_post.ex b/lib/statifier_examples/hold_desk/desk_post.ex index d361d7d..1cb0abc 100644 --- a/lib/statifier_examples/hold_desk/desk_post.ex +++ b/lib/statifier_examples/hold_desk/desk_post.ex @@ -70,7 +70,7 @@ defmodule StatifierExamples.HoldDesk.DeskPost do @impl Oban.Worker def perform(%Oban.Job{args: args, attempt: attempt}) when attempt <= @post_attempts do - {payload, ctx} = args["instruction"] |> Base.decode64!() |> :erlang.binary_to_term([:safe]) + {payload, ctx} = decode(args["instruction"]) case BasicHTTP.perform(payload, ctx) do :ok -> :ok @@ -83,6 +83,14 @@ defmodule StatifierExamples.HoldDesk.DeskPost do # deliver it again, without posting again. def perform(%Oban.Job{args: args}), do: failed(args, :desk_post_attempts_spent) + # Read back without :safe, as the router README's recipe reads its own: + # :safe refuses an atom this node has not created yet, and a job queued + # before a restart can name one (a send's owner kind, a struct field) + # that no module has loaded since. The jobs table is written only + # through this app's repo, by `new/3`. + @spec decode(String.t()) :: {term(), map()} + defp decode(instruction), do: instruction |> Base.decode64!() |> :erlang.binary_to_term() + # The POST has failed for good: tell the hold, or keep a dead letter. @spec failed(map(), term()) :: :ok | {:error, term()} | {:cancel, term()} defp failed(%{"execution_id" => execution_id} = args, reason) do diff --git a/test/statifier_examples_web/controllers/basic_http_controller_test.exs b/test/statifier_examples_web/controllers/basic_http_controller_test.exs index ae975b9..67ce243 100644 --- a/test/statifier_examples_web/controllers/basic_http_controller_test.exs +++ b/test/statifier_examples_web/controllers/basic_http_controller_test.exs @@ -4,7 +4,8 @@ defmodule StatifierExamplesWeb.BasicHTTPControllerTest do import Ecto.Query, only: [from: 2] import ExUnit.CaptureLog - alias StatifierExamples.{FirstWorkflow, HoldDesk, Repo, RoutedWorkflow} + alias StatifierExamples.{DeskTransport, FirstWorkflow, HoldDesk, Repo, RoutedWorkflow} + alias StatifierExamples.HoldDesk.DeskPost alias StatifierPersistence.{Execution, Executions, Storage} alias StatifierRouter.BasicHTTP @@ -237,6 +238,48 @@ defmodule StatifierExamplesWeb.BasicHTTPControllerTest do desk_post_args() end + # A job queued before a restart can name an atom this node has not + # created since. The instruction here is written with a placeholder atom + # whose bytes are then swapped for a name of the same length that no + # code creates, so the job decodes it as a node meeting it for the first + # time would. + # sabotage: the job's decode given [:safe] -> the instruction was refused + # and no POST was made, red; restored, green. + test "a job queued before a restart posts, whatever atoms this node has" do + payload_ctx = + {{:post, + %{ + url: @desk, + headers: [{"content-type", @form}], + body: "_scxmleventname=hold.placed", + transport: DeskTransport, + send: hold_send(@desk) + }}, %{session_id: "ex_hold_0424", opts: [], restart_marker: :zz_desk_placeholder_atom}} + + unseen = "zz_desk_" <> Base.encode16(:crypto.strong_rand_bytes(8), case: :lower) + assert byte_size(unseen) == byte_size("zz_desk_placeholder_atom") + + instruction = + payload_ctx + |> :erlang.term_to_binary() + |> :binary.replace("zz_desk_placeholder_atom", unseen) + + assert_raise ArgumentError, fn -> :erlang.binary_to_term(instruction, [:safe]) end + + {:ok, _job} = + %{ + "execution_id" => "ex_hold_0424", + "send_id" => "send_1", + "key" => "ex_hold_0424/send_1/1/0/0/0/transition.0/0", + "instruction" => Base.encode64(instruction) + } + |> DeskPost.new() + |> Oban.insert() + + assert %{success: 1, failure: 0} = drain_desk_posts() + assert_received {:desk_post, @desk, _headers, "_scxmleventname=hold.placed"} + end + # The executor runs inside the delivery's transaction, and the job is # inserted through the same repo, so a delivery that rolls back takes # the job with it and no POST is made for it.