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..1cb0abc --- /dev/null +++ b/lib/statifier_examples/hold_desk/desk_post.ex @@ -0,0 +1,157 @@ +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} = decode(args["instruction"]) + + 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) + + # 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 + 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..67ce243 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,11 @@ 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.{DeskTransport, FirstWorkflow, HoldDesk, Repo, RoutedWorkflow} + alias StatifierExamples.HoldDesk.DeskPost alias StatifierPersistence.{Execution, Executions, Storage} alias StatifierRouter.BasicHTTP @@ -20,7 +22,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 +32,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 +65,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 +172,177 @@ 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) - :ok = :telemetry.attach(handler, reentered, forward, nil) - on_exit(fn -> :telemetry.detach(handler) end) + execution_id = requested_hold("hold-0420") + 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 + + # 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. + # 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 +352,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