diff --git a/changelog.d/st-gje8.md b/changelog.d/st-gje8.md new file mode 100644 index 00000000..608dd72c --- /dev/null +++ b/changelog.d/st-gje8.md @@ -0,0 +1,6 @@ +### Added + +- The W3C Basic HTTP Event I/O Processor, `Statifier.Send.BasicHTTP`: register it in `:send_types` under `http://www.w3.org/TR/scxml/#BasicHTTPEventProcessor` and `basichttp` as `{Statifier.Send.BasicHTTP, base_url: ...}`, and both `_ioprocessors` keys carry one location, the base URL and the session id. A registration without `:base_url` is refused when the session starts. Sends POST through an injectable `Statifier.Send.BasicHTTP.Transport` (OTP `:httpc` by default; no new dependency); a delayed send is the processor's timer and a `` stops it. `decode/1` turns a request into an event for a host's own front, and a failed delivery raises `error.communication` on the sender. +- Every Basic HTTP POST carries the send's dedup key in an `scxml-send-key` header, so a receiver that deduplicates on it delivers each send once; `decode/1` takes the header's value as `:send_key` and refuses a malformed one, and sets no event field from it. +- A `:send_types` value may now be `{module, opts}`: the options reach the processor's callbacks under the plan context's `:opts` key and are recorded as strings. A processor may implement the optional `ioprocessors_entry/2`, which receives the session id and the options. +- `Statifier.Testing.Case.test_scxml/5` takes a `:send_types` option. diff --git a/docs/adr/0075-basichttp-event-io-processor.md b/docs/adr/0075-basichttp-event-io-processor.md index 1ec2b217..f07e3857 100644 --- a/docs/adr/0075-basichttp-event-io-processor.md +++ b/docs/adr/0075-basichttp-event-io-processor.md @@ -361,3 +361,66 @@ nothing: - [ADR-0051](0051-invoke-handlers-are-registered-per-session.md) (decision 4, the planning and performing split) - [ADR-0057](0057-recording-identity-and-serialization.md) (decision 5, registrations recorded as strings) - [ADR-0054](0054-durable-timers-consume-the-effect-vocabulary.md) (the processor-owned timer and its cancellation key) + +### Amendment 2026-09-30: every POST carries the send's dedup key, and the receiver deduplicates + +Status: proposed (2026-09-30) - amends decision 4 (the outbound mapping) +and decision 5 (the inbound decoder) by addition; every other decision, +and the record's own Status above, are unchanged. The header, +at-least-once delivery and deduplication by the receiver were ruled by +the operator, 2026-09-30; what the decoder does with the header is this +record's. + +[ADR-0069](0069-host-registered-send-types.md) decision 4 binds every +registered processor: "A processor MUST be idempotent on the ADR-0054 +decision 3 dedup key's components read off the effect", because "after a +crash and retry, a host may perform the same effect more than once." +Decision 8 point d above has the processor make one attempt per +`perform/2` and keep no memory between calls, so a host that performs the +same instruction twice POSTs twice. Neither decision 4 nor decision 5 said +how the MUST is met. This Amendment says it. + +**The processor is at-least-once, and the receiver deduplicates.** Every +POST the processor makes, immediate or delayed, whatever its body, +carries the send's dedup key in one request header: + +- **Name:** `scxml-send-key`. +- **Value:** the eight components of + [ADR-0054](0054-durable-timers-consume-the-effect-vocabulary.md) + decision 3's deduplication key, as that record and ADR-0059 order them, + joined by `/`: the session scope, `send_id`, `macrostep`, `microstep`, + `round`, `c_index`, `owner`, `ordinal`. + - The session scope is the plan context's `session_id` (spec 5.10's + `_sessionid` for a live session, a host's own scope for a + process-less host), percent-encoded. + - `send_id` is percent-encoded. Percent-encoding here escapes every + byte outside RFC 3986's unreserved set (`A-Z a-z 0-9 - . _ ~`), so + neither field can carry a `/`. + - `macrostep`, `microstep`, `round`, `c_index` and `ordinal` are + decimal integers. + - `owner` is spelled `onentry.S.B`, `onexit.S.B` or `finalize.S.B` with + its state and block indexes, or `transition.T` with its transition + index. + - A component the effect does not carry is the empty string. + +Every component is a deterministic counter or a static position, so a +re-performed instruction sends a byte-identical value. A receiver that +enqueues a request only when it has not already enqueued one with the same +`scxml-send-key` delivers each send once: for such a receiver ADR-0069's +MUST holds end to end. A receiver that ignores the header sees +at-least-once delivery. The processor itself still keeps no memory across +`perform/2` calls. + +**What the decoder does with it.** `decode/1`'s request map takes the +header's value under an optional `:send_key` key, and sets no event field +from it: an inbound event's `sendid` stays unset, as it is for a request +without the header. A value that is not eight `/`-separated fields whose +second field percent-decodes to UTF-8 is +`{:error, {:malformed_send_key, value}}`, answered 400 by decision 5's +status rule; an absent one changes nothing. The decoder does no +deduplicating: it is pure and remembers nothing. A front deduplicates on +the `scxml-send-key` header's value itself, and a request it has already +enqueued is answered 204 again with nothing enqueued. This repository's +loopback front, which lives only as long as one test run, does not +deduplicate; statifier_router's durable front, which must survive a +restart, is the one that will. diff --git a/lib/mix/statifier/basic_http_front.ex b/lib/mix/statifier/basic_http_front.ex new file mode 100644 index 00000000..26afdbb6 --- /dev/null +++ b/lib/mix/statifier/basic_http_front.ex @@ -0,0 +1,128 @@ +defmodule Mix.Statifier.BasicHTTPFront do + @moduledoc """ + The loopback inbound front this repository's own runs deliver Basic HTTP + sends through (ADR-0075), on OTP's `:inets` httpd, with no dependency. It + is repository tooling, not part of the package. + + `start/0` binds a free port on 127.0.0.1 and answers the base URL a + `Statifier.Send.BasicHTTP` registration takes as its `:base_url`, so a + session's `_ioprocessors` location is `base_url <> "/" <> session_id`. + A request to that location is resolved to the live session registered + under the id in `Statifier.Registry`, decoded by + `Statifier.Send.BasicHTTP.decode/1`, and enqueued on the session as an + external event, with the request's `scxml-send-key` header handed to the + decoder. This front does not deduplicate on that header (ADR-0075's + Amendment of 2026-09-30 leaves that to a front that outlives a restart). + The status rule is the decoder's (ADR-0075 decision 5): + + - 204 once the event is enqueued, before it is processed; + - 405 with `Allow: POST` for any other method; + - 400 for a request that forms no event; + - 404 for a path that names no live session. + """ + + alias Statifier.Send.BasicHTTP + alias Statifier.Session + + require Record + + Record.defrecordp(:mod, Record.extract(:mod, from_lib: "inets/include/httpd.hrl")) + + # `:inets` is not an application this package lists (ADR-0075 decision 6), + # so dialyzer's PLT does not see httpd; the calls into it are kept to the + # two functions named here. + @dialyzer {:nowarn_function, start: 0, stop: 1} + + @prefix "/basichttp" + + @typedoc "A running front: the httpd process and the base URL it answers at." + @type t :: %{pid: pid(), base_url: String.t()} + + @doc """ + Starts a front on a free loopback port. Starts `:inets` first. + """ + @spec start() :: {:ok, t()} | {:error, term()} + def start do + root = String.to_charlist(System.tmp_dir!()) + + with {:ok, _apps} <- Application.ensure_all_started(:inets), + {:ok, pid} <- + :inets.start(:httpd, + port: 0, + bind_address: {127, 0, 0, 1}, + server_name: ~c"statifier-basichttp", + server_root: root, + document_root: root, + modules: [__MODULE__] + ) do + [port: port] = :httpd.info(pid, [:port]) + {:ok, %{pid: pid, base_url: "http://127.0.0.1:#{port}#{@prefix}"}} + end + end + + @doc "Stops a front `start/0` started." + @spec stop(front :: t()) :: :ok | {:error, term()} + def stop(%{pid: pid}), do: :inets.stop(:httpd, pid) + + @doc """ + The httpd callback, `do/1` (a name only `unquote/1` can define): + answers one request by the status rule in the moduledoc, through + `respond/1`. + """ + @spec unquote(:do)(mod_data :: tuple()) :: {:proceed, list()} + defdelegate unquote(:do)(mod_data), to: __MODULE__, as: :respond + + @doc "Answers one httpd request by the status rule in the moduledoc." + @spec respond(mod_data :: tuple()) :: {:proceed, list()} + def respond(mod_data) do + %URI{path: path, query: query} = mod_data |> mod(:request_uri) |> to_string() |> URI.parse() + + headers = mod(mod_data, :parsed_header) + + request = %{ + method: mod_data |> mod(:method) |> to_string(), + content_type: header(headers, ~c"content-type"), + send_key: header(headers, ~c"scxml-send-key"), + body: mod_data |> mod(:entity_body) |> IO.iodata_to_binary(), + query: query + } + + {:proceed, [{:response, {:response, head(answer(path, request)), []}}]} + end + + @spec answer(path :: String.t() | nil, request :: BasicHTTP.request()) :: + 204 | 400 | 404 | 405 + defp answer(@prefix <> "/" <> session_id, request) do + with {:ok, pid} <- whereis(session_id), + {:ok, event} <- BasicHTTP.decode(request) do + :ok = Session.send_event(pid, event) + 204 + else + :no_session -> 404 + {:error, {:method_not_allowed, _method}} -> 405 + {:error, _reason} -> 400 + end + end + + defp answer(_path, _request), do: 404 + + @spec head(status :: 204 | 400 | 404 | 405) :: keyword() + defp head(405), do: [code: 405, allow: ~c"POST", content_length: ~c"0"] + defp head(status), do: [code: status, content_length: ~c"0"] + + @spec whereis(session_id :: String.t()) :: {:ok, pid()} | :no_session + defp whereis(session_id) do + case Registry.lookup(Statifier.Registry, session_id) do + [{pid, _value}] -> {:ok, pid} + [] -> :no_session + end + end + + @spec header(headers :: [{charlist(), charlist()}], name :: charlist()) :: String.t() | nil + defp header(headers, name) do + case List.keyfind(headers, name, 0) do + {_name, value} -> to_string(value) + nil -> nil + end + end +end diff --git a/lib/statifier/effect/send.ex b/lib/statifier/effect/send.ex index 224b465c..3c956752 100644 --- a/lib/statifier/effect/send.ex +++ b/lib/statifier/effect/send.ex @@ -43,9 +43,9 @@ defmodule Statifier.Effect.Send do different case entirely: it is an argument failure that discards the whole `` (ADR-0036), never a value that reaches `data` at all. This settles the `#_internal`, same-session, and `#_scxml_` routes - that exist today; a future external-wire processor (BasicHTTP or - otherwise) owns its own `:undefined` encoding at its own boundary, not - here. + that exist today; a registered external-wire processor, such as + `Statifier.Send.BasicHTTP` (ADR-0075), owns its own `:undefined` + encoding at its own boundary, not here. """ alias Statifier.Machine.Content diff --git a/lib/statifier/evaluator/system_variables.ex b/lib/statifier/evaluator/system_variables.ex index 66d9dc32..6fd2271e 100644 --- a/lib/statifier/evaluator/system_variables.ex +++ b/lib/statifier/evaluator/system_variables.ex @@ -63,7 +63,13 @@ defmodule Statifier.Evaluator.SystemVariables do `"location"`, is always there. A session that registers send types (ADR-0069) supports each of them too, so `_ioprocessors` also carries one entry per registered type string, whose value the type's processor - supplied (`Statifier.Send.Types.from_send_types/1` read it). With + supplied. A processor that exports the optional + `c:Statifier.Send.Processor.ioprocessors_entry/2` is asked here, with + the type string and a context carrying `session_id` and the + registration's `opts`, so its entry can address this one session + (ADR-0075 decision 3); a processor that exports only + `c:Statifier.Send.Processor.ioprocessors_entry/1` gets the entry + `Statifier.Send.Types.from_send_types/1` read, exactly as before. With `send_types` `nil`, which is what a session registering nothing carries, the map holds the SCXML entry alone, exactly as before registered types existed. A registered entry never replaces the SCXML entry: a set naming @@ -108,16 +114,41 @@ defmodule Statifier.Evaluator.SystemVariables do "_event" => :undefined, "_ioprocessors" => Map.put( - registered_entries(send_types), + registered_entries(send_types, session_id), @scxml_event_processor, %{"location" => scxml_location(session_id)} ) } end - @spec registered_entries(send_types :: Types.t() | nil) :: %{String.t() => map()} - defp registered_entries(nil), do: %{} - defp registered_entries(%Types{entries: entries}), do: entries + @spec registered_entries(send_types :: Types.t() | nil, session_id :: String.t()) :: + %{String.t() => map()} + defp registered_entries(nil, _session_id), do: %{} + + defp registered_entries(%Types{entries: entries, processors: processors}, session_id) do + Map.new(entries, fn {type, entry} -> + case Map.fetch(processors, type) do + {:ok, {module, opts}} -> {type, session_entry(module, type, opts, session_id, entry)} + :error -> {type, entry} + end + end) + end + + # ADR-0075 decision 3: a module that exports `/2` is asked for its entry + # with the session id; one that exports only `/1` keeps the entry the + # registered set was built with. + @spec session_entry( + module :: module(), + type :: String.t(), + opts :: keyword(), + session_id :: String.t(), + entry :: map() + ) :: map() + defp session_entry(module, type, opts, session_id, entry) do + if function_exported?(module, :ioprocessors_entry, 2), + do: Types.session_entry!(module, type, %{session_id: session_id, opts: opts}), + else: entry + end @doc """ `_event`'s value for `event` - spec 5.10.1's fields, read straight off diff --git a/lib/statifier/replay.ex b/lib/statifier/replay.ex index c0c56eca..6500ac6f 100644 --- a/lib/statifier/replay.ex +++ b/lib/statifier/replay.ex @@ -171,7 +171,7 @@ defmodule Statifier.Replay do # plan context's `send_processors` for the same reason # `invoke_handlers` is, and the live holds `Statifier.Session` # keeps for its cancel routing, kept here by the same rule. - send_processors: %{String.t() => module()}, + send_processors: %{String.t() => SendTypes.registration()}, held_sends: Effects.held_sends() } end diff --git a/lib/statifier/send/basic_http.ex b/lib/statifier/send/basic_http.ex new file mode 100644 index 00000000..448aa5d7 --- /dev/null +++ b/lib/statifier/send/basic_http.ex @@ -0,0 +1,474 @@ +defmodule Statifier.Send.BasicHTTP do + @moduledoc """ + The W3C Basic HTTP Event I/O Processor (SCXML appendix C.2), a + `Statifier.Send.Processor` a host registers like any other send type + (ADR-0069, ADR-0075). + + ## Registering it + + Register it under the spec's processor URI and its short form + `basichttp`, with the base URL the host's own front answers at: + + base = "https://example.org/scxml" + + Statifier.Session.start_link(chart, + send_types: %{ + "http://www.w3.org/TR/scxml/#BasicHTTPEventProcessor" => + {Statifier.Send.BasicHTTP, base_url: base}, + "basichttp" => {Statifier.Send.BasicHTTP, base_url: base} + } + ) + + Neither string is a built-in spelling, so the registration redirects no + built-in send (ADR-0075 decision 2). The options: + + - `:base_url` (required) - the address the host's front answers at. + The session's `_ioprocessors` carries an entry under each registered + string, both holding the same `"location"`: this URL, `/`, and the + session's `_sessionid` (C.2.3, ADR-0075 decision 3). A registration + without it is refused when the session starts, with an + `ArgumentError` naming the option. + - `:transport` - a `Statifier.Send.BasicHTTP.Transport` module the + POSTs go through. Default `Statifier.Send.BasicHTTP.Transport.Httpc`, + on OTP's `:httpc`; this package adds no dependency for it + (ADR-0075 decision 6). + + ## Outbound (C.2.2) + + `deliver/3` plans and `perform/2` POSTs, the split every processor has + (`Statifier.Send.Processor`). The mapping (ADR-0075 decision 4): + + - `event` becomes the form parameter `_scxmleventname`, and each + `namelist` entry and `` a form parameter, in an + `application/x-www-form-urlencoded` body. + - A `` child is the body, sent as `text/plain`; when the send + also names an `event`, `_scxmleventname` travels as a query parameter + of the target URL. + - The two are told apart by `data`'s shape: a map is form-encoded, + `:undefined` sends `_scxmleventname` alone as a form body, and any + other value is the body. A `` that evaluates to a map + is therefore form-encoded. + - A send with neither `target` nor `targetexpr` raises C.2.2's + `error.communication` on the sender's internal queue, carrying the + send id, and makes no request. + + A parameter value is written as text: a string as it is, a number or a + boolean as its literal, `nil` as `null`, `:undefined` as the empty + string. Any other value (a list, a map) is written with `inspect/1`; its + encoding is not decided yet. + + **One attempt, and a miss reaches the sender.** `perform/2` makes one + POST. On a transport error, or a status outside 2xx, it reports the miss + through `Statifier.Session.failed_send/3` to the sending session it finds + in `Statifier.Registry` under the plan context's `session_id`, so the + sender sees C.1's `error.communication` carrying the send id, and it + returns `{:error, reason}`. When no live session is registered under that + id it returns `{:error, reason}` only, and the dead-letter rule of + `failed_send/3`'s documentation is the host's (ADR-0075 decision 8, point + d). + + **At-least-once, deduplicated by the receiver** (ADR-0075's Amendment of + 2026-09-30). The processor keeps no memory across `perform/2` calls, so + a host that performs the same instruction twice POSTs twice. Every POST + therefore carries the send's ADR-0054 decision 3 dedup key in the + `scxml-send-key` header: eight fields joined by `/`, in the record's + order - the session scope (the plan context's `session_id`), the send + id, `macrostep`, `microstep`, `round`, `c_index`, `owner` and `ordinal`. + The session scope and the send id are percent-encoded (every byte outside + RFC 3986's unreserved set), the counters are decimal, and `owner` is + spelled `onentry.S.B`, `onexit.S.B`, `transition.T` or `finalize.S.B` + with its indexes. A receiver that deduplicates on the header sees each + send once, which is ADR-0069's idempotency MUST end to end; a receiver + that ignores it sees at-least-once delivery. + + `perform/2` runs in the process that performs the instruction, which for + `Statifier.Session` is the sending session, so a slow location holds that + session for the length of the request; the default transport bounds each + request with a timeout. + + **A delayed send is this processor's timer** (ADR-0069 decision 4). A + `` is held by a timer process `perform/2` starts, and + `cancel/2` plans the cancellation of every timer held under the send id + (spec 6.3). The timers are kept in the dictionary of the process that + performs the instructions, so a host that performs them itself performs + a delayed send and its cancel in one process. At fire time the timer + POSTs only if that process is a `Statifier.Session` still running: a + session that has stopped, or halted, discards the send (spec 6.2). + + ## Inbound (C.2.1) + + `decode/1` turns one HTTP request into a `%Statifier.Event{}`. It is pure + and knows no session: a front resolves the location to a session (or to + an execution, for a durable host), calls it, and enqueues the event as + an external event. The status rule a front applies (ADR-0075 decision 5): + + - `{:ok, event}` - answer 204 once the event is enqueued, before it is + processed (a front that has already enqueued a request carrying the + same `scxml-send-key` answers 204 again and enqueues nothing); + - `{:error, {:method_not_allowed, method}}` - answer 405 with + `Allow: POST`; + - any other `{:error, _}` - answer 400; + - a location that names no session the front can reach - the front's + own 404. + """ + + @behaviour Statifier.Send.Processor + + alias Statifier.{Effect, Event, EventData, Session} + alias Statifier.Send.BasicHTTP.Transport + + @uri "http://www.w3.org/TR/scxml/#BasicHTTPEventProcessor" + @event_name_param "_scxmleventname" + @form "application/x-www-form-urlencoded" + @send_key_header "scxml-send-key" + + @typedoc """ + What `decode/1` is handed: the request's method, its content type (`nil` + when the request carries none), its body, its query string (`nil` when + the URL has none), and optionally the value of its `scxml-send-key` + header (`nil` or absent when the request carries none). + """ + @type request :: %{ + required(:method) => String.t(), + required(:content_type) => String.t() | nil, + required(:body) => binary(), + required(:query) => String.t() | nil, + optional(:send_key) => String.t() | nil + } + + @typedoc "Why `decode/1` could not form an event from a request." + @type decode_error :: + {:method_not_allowed, String.t()} + | {:not_utf8, :query | :body} + | {:malformed_send_key, String.t()} + + @typep post :: %{ + url: String.t(), + headers: [{String.t(), String.t()}], + body: binary(), + transport: module(), + send: Effect.Send.t() | Effect.SendDelayed.t() + } + + @doc """ + The `_ioprocessors` entry for `type`: a `"location"` that is the + registration's `:base_url`, `/`, and the session's id (C.2.3, ADR-0075 + decision 3). Raises `ArgumentError` when the registration carries no + `:base_url`, which refuses the session's start. + """ + @impl Statifier.Send.Processor + @spec ioprocessors_entry( + type :: String.t(), + context :: Statifier.Send.Processor.entry_context() + ) :: + map() + def ioprocessors_entry(type, %{session_id: session_id, opts: opts}) do + case Keyword.fetch(opts, :base_url) do + {:ok, base_url} when is_binary(base_url) -> + %{"location" => base_url <> "/" <> session_id} + + _missing -> + raise ArgumentError, + "#{inspect(__MODULE__)} registered for #{inspect(type)} needs a :base_url " <> + "option, the address its _ioprocessors location is built from (C.2.3)" + end + end + + @doc """ + Plans one send (see the moduledoc's "Outbound"). Pure. + """ + @impl Statifier.Send.Processor + @spec deliver( + send :: Effect.Send.t() | Effect.SendDelayed.t(), + event :: Event.t(), + ctx :: Statifier.Send.Processor.ctx() + ) :: {:ok, [Statifier.Send.Processor.instruction()]} + def deliver(%{target: nil} = send, _event, _ctx) do + {:ok, + [ + {:raise, :platform, "error.communication", {:content, send.c_index, send.owner}, + sendid: send.send_id} + ]} + end + + def deliver(%Effect.SendDelayed{} = send, _event, ctx), + do: {:ok, [{:handler, __MODULE__, {:post_after, send.delay_ms, post(send, ctx)}}]} + + def deliver(%Effect.Send{} = send, _event, ctx), + do: {:ok, [{:handler, __MODULE__, {:post, post(send, ctx)}}]} + + @doc """ + Plans the cancellation of every delayed send this processor holds under + `cancel.send_id`. Pure. + """ + @impl Statifier.Send.Processor + @spec cancel(cancel :: Effect.Cancel.t(), ctx :: Statifier.Send.Processor.ctx()) :: + {:ok, [Statifier.Send.Processor.instruction()]} + def cancel(%Effect.Cancel{send_id: send_id}, _ctx), + do: {:ok, [{:handler, __MODULE__, {:cancel, send_id}}]} + + @doc """ + Performs one instruction `deliver/3` or `cancel/2` planned: a POST, a + delayed POST's timer, or the cancellation of the timers held under a + send id (see the moduledoc's "Outbound"). + """ + @impl Statifier.Send.Processor + @spec perform(payload :: term(), ctx :: Statifier.Send.Processor.ctx()) :: + :ok | {:error, term()} + def perform({:post, post}, ctx), do: post_now(post, ctx) + + def perform({:post_after, delay_ms, post}, ctx) do + owner = self() + key = {__MODULE__, post.send.send_id} + # ADR-0075 decision 9 / ADR-0069 decision 4: the processor owns the delay. + timer = spawn(fn -> hold(owner, delay_ms, post, ctx) end) + Process.put(key, [timer | Enum.filter(Process.get(key, []), &Process.alive?/1)]) + :ok + end + + def perform({:cancel, send_id}, _ctx) do + {__MODULE__, send_id} |> Process.delete() |> List.wrap() |> Enum.each(&send(&1, :cancel)) + :ok + end + + @doc """ + Decodes one HTTP request into an external event (C.2.1, ADR-0075 + decision 5). Pure. + + - The method must be POST; any other is + `{:error, {:method_not_allowed, method}}`. + - The event name is the first `_scxmleventname` found, the query string + before a form body, else `HTTP.` and the method in upper case + (`HTTP.POST`). + - A form body's other parameters, with the query string's, become + `_event.data`, each value through `Statifier.EventData`'s text rung + (a predicator literal, else the string), so `2` reads as the number 2. + A body of any other content type becomes `_event.data` through the + same text rung, and the query string then contributes the event name + only. + - `origintype` is the processor URI. + - `:send_key`, the `scxml-send-key` header's value, sets no event + field: the front deduplicates on the header's value itself, and the + event's `sendid` stays unset. A value that is not the header's eight + fields is `{:error, {:malformed_send_key, value}}`. + + A query string or a body that is not UTF-8 once decoded forms no + datamodel string, and is `{:error, {:not_utf8, :query | :body}}`. + """ + @spec decode(request :: request()) :: {:ok, Event.t()} | {:error, decode_error()} + def decode(%{method: method} = request) do + with :ok <- post_only(method), + {:ok, query} <- pairs(Map.get(request, :query), :query), + {:ok, body, text} <- body(request), + :ok <- check_send_key(Map.get(request, :send_key)) do + {names, params} = Enum.split_with(query ++ body, &match?({@event_name_param, _value}, &1)) + + name = + case names do + [{_key, name} | _rest] -> name + [] -> "HTTP." <> String.upcase(method) + end + + {:ok, Event.external(name, data: data(params, text), origintype: @uri)} + end + end + + # Checks an `scxml-send-key` value's shape (ADR-0075's Amendment of + # 2026-09-30); the decoder sets no event field from it. + @spec check_send_key(send_key :: String.t() | nil) :: :ok | {:error, decode_error()} + defp check_send_key(nil), do: :ok + + defp check_send_key(send_key) do + with [_scope, send_id, _macro, _micro, _round, _c_index, _owner, _ordinal] <- + String.split(send_key, "/"), + true <- send_id |> URI.decode() |> String.valid?() do + :ok + else + _malformed -> {:error, {:malformed_send_key, send_key}} + end + end + + @spec post_only(method :: String.t()) :: :ok | {:error, decode_error()} + defp post_only(method) do + if String.upcase(method) == "POST", + do: :ok, + else: {:error, {:method_not_allowed, method}} + end + + # A form body's pairs, or a text body; `text` is `nil` for a form. + @spec body(request :: request()) :: + {:ok, [{String.t(), String.t()}], String.t() | nil} | {:error, decode_error()} + defp body(%{body: body} = request) do + cond do + form?(Map.get(request, :content_type)) -> + with {:ok, pairs} <- pairs(body, :body), do: {:ok, pairs, nil} + + String.valid?(body) -> + {:ok, [], body} + + true -> + {:error, {:not_utf8, :body}} + end + end + + @spec form?(content_type :: String.t() | nil) :: boolean() + defp form?(nil), do: false + defp form?(content_type), do: content_type |> String.downcase() |> String.starts_with?(@form) + + @spec data(params :: [{String.t(), String.t()}], text :: String.t() | nil) :: term() + defp data(params, nil), + do: EventData.coerce({:params, Enum.map(params, fn {k, v} -> {k, text(v)} end)}) + + defp data(_params, text), do: text(text) + + @spec text(value :: String.t() | nil) :: term() + defp text(nil), do: :undefined + defp text(value), do: EventData.coerce({:text, value}) + + @spec pairs(encoded :: String.t() | nil, part :: :query | :body) :: + {:ok, [{String.t(), String.t()}]} | {:error, decode_error()} + defp pairs(nil, _part), do: {:ok, []} + + defp pairs(encoded, part) do + pairs = encoded |> URI.query_decoder(:www_form) |> Enum.to_list() + + if Enum.all?(pairs, fn {key, value} -> String.valid?(key) and String.valid?(value) end), + do: {:ok, pairs}, + else: {:error, {:not_utf8, part}} + end + + # The request one send becomes (ADR-0075 decision 4). + @spec post(send :: Effect.Send.t() | Effect.SendDelayed.t(), ctx :: map()) :: post() + defp post(send, ctx) do + named = if send.event, do: [{@event_name_param, send.event}], else: [] + + {url, content_type, body} = + case send.data do + data when is_map(data) -> + {send.target, @form, form(named ++ Enum.map(data, fn {k, v} -> {k, encode(v)} end))} + + :undefined -> + {send.target, @form, form(named)} + + content -> + {with_query(send.target, named), "text/plain", encode(content)} + end + + %{ + url: url, + headers: [{"content-type", content_type}, {@send_key_header, send_key(send, ctx)}], + body: body, + transport: ctx |> Map.get(:opts, []) |> Keyword.get(:transport, Transport.Httpc), + send: send + } + end + + # The ADR-0054 decision 3 dedup key, spelled as ADR-0075's Amendment of + # 2026-09-30 states (see the moduledoc's "At-least-once"). + @spec send_key(send :: Effect.Send.t() | Effect.SendDelayed.t(), ctx :: map()) :: String.t() + defp send_key(send, %{session_id: session_id}) do + Enum.join( + [ + escape(session_id), + escape(send.send_id), + field(send.macrostep), + field(send.microstep), + field(send.round), + field(send.c_index), + owner(send.owner), + field(send.ordinal) + ], + "/" + ) + end + + @spec escape(value :: String.t() | nil) :: String.t() + defp escape(nil), do: "" + defp escape(value), do: URI.encode(value, &URI.char_unreserved?/1) + + @spec field(value :: non_neg_integer() | nil) :: String.t() + defp field(nil), do: "" + defp field(value), do: Integer.to_string(value) + + @spec owner(owner :: Statifier.Machine.Content.owner() | nil) :: String.t() + defp owner({kind, state, block}), do: "#{kind}.#{state}.#{block}" + defp owner({:transition, transition}), do: "transition.#{transition}" + defp owner(nil), do: "" + + @spec form(pairs :: [{String.t(), String.t()}]) :: String.t() + defp form(pairs), do: URI.encode_query(pairs, :www_form) + + @spec with_query(url :: String.t(), pairs :: [{String.t(), String.t()}]) :: String.t() + defp with_query(url, []), do: url + + defp with_query(url, pairs) do + separator = if String.contains?(url, "?"), do: "&", else: "?" + url <> separator <> form(pairs) + end + + @spec encode(value :: term()) :: String.t() + defp encode(value) when is_binary(value), do: value + defp encode(nil), do: "null" + defp encode(:undefined), do: "" + defp encode(value) when is_number(value) or is_boolean(value), do: to_string(value) + defp encode(value), do: inspect(value) + + # One attempt; a miss goes to the sender through `failed_send/3` + # (ADR-0075 decision 8, point d). + @spec post_now(post :: post(), ctx :: map()) :: :ok | {:error, term()} + defp post_now(post, ctx) do + case post.transport.post(post.url, post.headers, post.body) do + {:ok, status} when status in 200..299 -> :ok + {:ok, status} -> report(post.send, ctx, {:http_status, status}) + {:error, reason} -> report(post.send, ctx, reason) + end + end + + @spec report(send :: Effect.Send.t() | Effect.SendDelayed.t(), ctx :: map(), reason :: term()) :: + {:error, term()} + defp report(send, %{session_id: session_id}, reason) do + case whereis(session_id) do + nil -> :ok + pid -> Session.failed_send(pid, send, reason: reason) + end + + {:error, reason} + end + + # `Registry.lookup/2` raises `ArgumentError` when `Statifier.Registry` + # itself is not running, which is the "no live session registered" case. + @spec whereis(session_id :: String.t()) :: pid() | nil + defp whereis(session_id) do + case Registry.lookup(Statifier.Registry, session_id) do + [{pid, _value}] -> pid + [] -> nil + end + rescue + ArgumentError -> nil + end + + # A delayed send's timer process: POSTs after `delay_ms` unless it is + # cancelled first or its owner stops, and at fire time only while its + # owner is a session still running (spec 6.2's discard at termination). + @spec hold(owner :: pid(), delay_ms :: non_neg_integer(), post :: post(), ctx :: map()) :: + :ok | {:error, term()} + defp hold(owner, delay_ms, post, ctx) do + # ADR-0075 decision 9: the timer ends when the session that owns it does. + ref = Process.monitor(owner) + + # ADR-0075 decision 9: a cancel, the owner's end, or the delay, first. + receive do + :cancel -> :ok + {:DOWN, ^ref, :process, ^owner, _reason} -> :ok + after + delay_ms -> if running?(owner), do: post_now(post, ctx), else: :ok + end + end + + @spec running?(owner :: pid()) :: boolean() + defp running?(owner) do + match?(%{status: :running}, Session.status(owner)) + catch + :exit, _reason -> false + end +end diff --git a/lib/statifier/send/basic_http/transport.ex b/lib/statifier/send/basic_http/transport.ex new file mode 100644 index 00000000..d996439e --- /dev/null +++ b/lib/statifier/send/basic_http/transport.ex @@ -0,0 +1,36 @@ +defmodule Statifier.Send.BasicHTTP.Transport do + @moduledoc """ + How `Statifier.Send.BasicHTTP` POSTs a message (ADR-0075 decision 6): one + callback that sends a body with headers to a URL and answers the status + code, or an error when no status came back. + + The default is `Statifier.Send.BasicHTTP.Transport.Httpc`, on OTP's + `:httpc`, which this package always compiles and adds no dependency for. + A host injects any other module through the registration's `:transport` + option. An adapter makes one attempt: retrying is not the transport's + job, and a status outside 2xx is answered as `{:ok, status}`, not as an + error, so the processor reports it as a miss. + + A host that uses `req` writes an adapter of its own: + + defmodule MyApp.ReqTransport do + @behaviour Statifier.Send.BasicHTTP.Transport + + @impl true + def post(url, headers, body) do + case Req.post(url, headers: headers, body: body, retry: false) do + {:ok, %Req.Response{status: status}} -> {:ok, status} + {:error, exception} -> {:error, exception} + end + end + end + """ + + @doc """ + POSTs `body` with `headers` (lower-case names, `content-type` among + them) to `url`, once. Answers `{:ok, status}` for any HTTP status, or + `{:error, reason}` when the request got no response. + """ + @callback post(url :: String.t(), headers :: [{String.t(), String.t()}], body :: binary()) :: + {:ok, 100..599} | {:error, term()} +end diff --git a/lib/statifier/send/basic_http/transport/httpc.ex b/lib/statifier/send/basic_http/transport/httpc.ex new file mode 100644 index 00000000..4bd74455 --- /dev/null +++ b/lib/statifier/send/basic_http/transport/httpc.ex @@ -0,0 +1,115 @@ +defmodule Statifier.Send.BasicHTTP.Transport.Httpc do + @moduledoc """ + The default `Statifier.Send.BasicHTTP.Transport`, on OTP's `:httpc` + (ADR-0075 decision 6). No dependency: `:inets`, `:ssl` and `:public_key` + ship with OTP. + + This package's application list names none of them, so a host that + registers nothing starts nothing. The adapter starts `:inets` and `:ssl` + itself on its first POST, and checks that `:ssl` and `:public_key` can + be loaded; when they cannot, the POST answers `{:error, reason}` and the + processor reports the miss. + + For an `https` URL it verifies the peer (`verify: :verify_peer`) against + the system CA store (`:public_key.cacerts_get/0`) and checks the host + name. Each request is bounded by a five-second connect timeout and a + five-second request timeout, so a slow location holds the sending + session no longer than that. + """ + + @behaviour Statifier.Send.BasicHTTP.Transport + + # `:inets`, `:ssl` and `:public_key` are applications this package does + # not list (ADR-0075 decision 6: a host that registers nothing starts + # nothing), so neither the compiler nor dialyzer's PLT sees their + # modules. The calls into them are resolved at run time, after + # `ensure_started/0` has started `:inets` and `:ssl`, and they are kept + # to the two functions named here so no other function's warnings are + # silenced. + @compile {:no_warn_undefined, [:public_key]} + @dialyzer {:nowarn_function, request: 2, ssl: 1} + + @timeout_ms 5_000 + + @doc """ + POSTs `body` with `headers` to `url` through `:httpc`, once (see + `c:Statifier.Send.BasicHTTP.Transport.post/3`). + """ + @impl Statifier.Send.BasicHTTP.Transport + @spec post(url :: String.t(), headers :: [{String.t(), String.t()}], body :: binary()) :: + {:ok, 100..599} | {:error, term()} + def post(url, headers, body) do + with :ok <- ensure_started() do + {content_type, rest} = content_type(headers) + + request = + {String.to_charlist(url), Enum.map(rest, &charlist_header/1), + String.to_charlist(content_type), body} + + http_options = [timeout: @timeout_ms, connect_timeout: @timeout_ms] ++ ssl(url) + + request(request, http_options) + end + end + + @spec request(request :: tuple(), http_options :: keyword()) :: + {:ok, 100..599} | {:error, term()} + defp request(request, http_options) do + case :httpc.request(:post, request, http_options, body_format: :binary) do + {:ok, {{_version, status, _reason}, _headers, _body}} -> {:ok, status} + {:error, reason} -> {:error, reason} + end + end + + @spec ensure_started() :: :ok | {:error, term()} + defp ensure_started do + with {:ok, _inets} <- Application.ensure_all_started(:inets), + {:ok, _ssl} <- Application.ensure_all_started(:ssl), + do: ensure_loaded([:ssl, :public_key]) + end + + # `:httpc` calls into `:public_key` on every request, `http:` included, + # and an application that reports itself started can still have modules + # the code path cannot load (a Mix task whose path was pruned, for one). + # A module that cannot load answers `{:error, {:not_loadable, module, + # reason}}`, so the processor reports a miss instead of crashing the + # sending session. Internal; public only so its refusal can be tested + # with a module that does not exist. + @doc false + @spec ensure_loaded(modules :: [module()]) :: :ok | {:error, term()} + def ensure_loaded(modules) do + Enum.reduce_while(modules, :ok, fn module, :ok -> + case Code.ensure_loaded(module) do + {:module, ^module} -> {:cont, :ok} + {:error, reason} -> {:halt, {:error, {:not_loadable, module, reason}}} + end + end) + end + + @spec content_type(headers :: [{String.t(), String.t()}]) :: + {String.t(), [{String.t(), String.t()}]} + defp content_type(headers) do + case List.keytake(headers, "content-type", 0) do + {{_name, content_type}, rest} -> {content_type, rest} + nil -> {"application/octet-stream", headers} + end + end + + @spec charlist_header(header :: {String.t(), String.t()}) :: {charlist(), charlist()} + defp charlist_header({name, value}), do: {String.to_charlist(name), String.to_charlist(value)} + + @spec ssl(url :: String.t()) :: keyword() + defp ssl("https:" <> _rest) do + [ + ssl: [ + verify: :verify_peer, + cacerts: :public_key.cacerts_get(), + customize_hostname_check: [ + match_fun: :public_key.pkix_verify_hostname_match_fun(:https) + ] + ] + ] + end + + defp ssl(_url), do: [] +end diff --git a/lib/statifier/send/processor.ex b/lib/statifier/send/processor.ex index ea8bcc93..2da7afed 100644 --- a/lib/statifier/send/processor.ex +++ b/lib/statifier/send/processor.ex @@ -93,13 +93,32 @@ defmodule Statifier.Send.Processor do `Statifier.MachineState.new/2` writes it, once, when the session starts; `Statifier.Send.Types.from_send_types/1` says how it reads after a resume. + An entry that must address one session, such as a location a receiver + POSTs to, cannot be built from the type alone. A processor that + implements the optional `c:Statifier.Send.Processor.ioprocessors_entry/2` + is asked for its entry with the type string and a context carrying the + session id and the registration's options, when the session starts + (ADR-0075 decision 3); it is asked instead of `/1`, and a processor that + implements only `/1` is asked as before. + + ## Registration options + + A `:send_types` value is a bare module or `{module, opts}` (ADR-0075 + decision 8, point b). The options reach + `c:Statifier.Send.Processor.ioprocessors_entry/2`'s context as `:opts` + and, for a `{module, opts}` registration only, the plan context + `deliver/3` and `cancel/2` receive, under `:opts`. A bare-module + registration's plan context carries no `:opts` key, exactly as before. + ## `ctx` The plan context `Statifier.Session.Effects.plan/2` threads through its fold, handed over unchanged: a plain map carrying `session_id` (spec 5.10's `_sessionid`) and no pid, no `%MachineState{}` and no session struct, so a processor cannot reach into the session through it. A key - added to it later is additive for every processor already written. + added to it later is additive for every processor already written: + `:opts`, a `{module, opts}` registration's options, is one (see + "Registration options" above). """ alias Statifier.{Effect, Event} @@ -150,5 +169,25 @@ defmodule Statifier.Send.Processor do """ @callback ioprocessors_entry(type :: String.t()) :: map() - @optional_callbacks perform: 2, ioprocessors_entry: 1 + @typedoc """ + What `c:Statifier.Send.Processor.ioprocessors_entry/2` is handed beside the + type string: the session's `_sessionid` and the registration's options + (`[]` for a bare-module registration). + """ + @type entry_context :: %{session_id: String.t(), opts: keyword()} + + @doc """ + The value of this processor's `_ioprocessors` entry for the registered + type string `type` in the session `context` names (ADR-0075 decision 3): + for example a `"location"` built from the session id and a base URL in + the registration's options. Called once, when the session starts, and + asked instead of `c:Statifier.Send.Processor.ioprocessors_entry/1` when + a processor implements both. Pure and deterministic, and string-keyed at + every level, as `c:Statifier.Send.Processor.ioprocessors_entry/1` is. A + processor that cannot build its entry from `context` raises + `ArgumentError`, which refuses the session's start. Optional. + """ + @callback ioprocessors_entry(type :: String.t(), context :: entry_context()) :: map() + + @optional_callbacks perform: 2, ioprocessors_entry: 1, ioprocessors_entry: 2 end diff --git a/lib/statifier/send/target.ex b/lib/statifier/send/target.ex index bffa1e7a..7588c5be 100644 --- a/lib/statifier/send/target.ex +++ b/lib/statifier/send/target.ex @@ -78,8 +78,11 @@ defmodule Statifier.Send.Target do ("Processors MAY define short form notations") - a MAY this codebase chose to honor, not a MUST - and the processor's own type URI (`SystemVariables.scxml_event_processor/0`) is the long form. Anything - else is unsupported: this engine implements only the SCXML Event I/O - Processor (see the plan's "What We're NOT Doing" on BasicHTTP). + else is not built in: the SCXML Event I/O Processor is the one processor + this engine delivers for itself. Any other type, the Basic HTTP Event I/O + Processor (`Statifier.Send.BasicHTTP`, ADR-0075) included, is supported + only when a host registers it (`Statifier.Send.Types.classify/2`, + ADR-0069). """ @spec supported_type?(type :: String.t() | nil) :: boolean() def supported_type?(nil), do: true diff --git a/lib/statifier/send/types.ex b/lib/statifier/send/types.ex index 15469626..7f04bb9f 100644 --- a/lib/statifier/send/types.ex +++ b/lib/statifier/send/types.ex @@ -24,9 +24,12 @@ defmodule Statifier.Send.Types do Beside the set, a registered set carries each type's `_ioprocessors` entry (spec 5.10), the value its processor supplies through the optional - `c:Statifier.Send.Processor.ioprocessors_entry/1` callback, so + `c:Statifier.Send.Processor.ioprocessors_entry/1` callback, and each + type's module and registration options, so `Statifier.MachineState.new/2` can write the entries from the same - stamp it classifies against. + stamp it classifies against, asking a processor that implements the + optional `c:Statifier.Send.Processor.ioprocessors_entry/2` once the + session id is known (ADR-0075 decision 3). `unsupported_sends/2` is the pure pre-start check of ADR-0069 decision 3. It lives here rather than in `Statifier.Validator`, because @@ -40,9 +43,30 @@ defmodule Statifier.Send.Types do alias Statifier.Parser.Location alias Statifier.Send.Target - defstruct types: MapSet.new(), entries: %{} + defstruct types: MapSet.new(), entries: %{}, processors: %{} - @type t :: %__MODULE__{types: MapSet.t(String.t()), entries: %{String.t() => map()}} + @typedoc """ + One `:send_types` value (ADR-0075 decision 8, point b): a bare + `Statifier.Send.Processor` module, or the module with its registration + options. The options reach the processor's + `c:Statifier.Send.Processor.ioprocessors_entry/2` context and, for a + `{module, opts}` registration only, the plan context its callbacks + receive, under `:opts`. + """ + @type registration :: module() | {module(), keyword()} + + @typedoc """ + The registered set. `types` is the set `classify/2` answers against; + `entries` holds each type's `_ioprocessors` value as + `c:Statifier.Send.Processor.ioprocessors_entry/1` returned it (or an + empty map); `processors` holds each type's module and options, a bare + module's options being `[]`. + """ + @type t :: %__MODULE__{ + types: MapSet.t(String.t()), + entries: %{String.t() => map()}, + processors: %{String.t() => {module(), keyword()}} + } @typedoc """ What `classify/2` answers for one resolved ``: @@ -60,7 +84,8 @@ defmodule Statifier.Send.Types do @doc """ Builds the registered set from a `:send_types` map - (`%{type_string => module}`), derived from the map's own keys rather than + (`%{type_string => registration}`, a `t:registration/0` being a module + or `{module, opts}`), derived from the map's own keys rather than declared beside it - the `` counterpart of `Statifier.Invoke.Types.from_handlers/1`. @@ -79,20 +104,45 @@ defmodule Statifier.Send.Types do empty map when the module does not export it. Raises `ArgumentError` when a returned value is not a map, or holds an atom key other than `true` or `false` at any level, because every datamodel key is a string. - `Statifier.Evaluator.SystemVariables.initial/3` says when the entries are + `processors` keeps each type's module and options (ADR-0075 decision 3), + so `Statifier.Evaluator.SystemVariables.initial/3` can ask a module that + exports `c:Statifier.Send.Processor.ioprocessors_entry/2` for its entry + once the session id is known; that function says when the entries are written and how they read after a resume. """ - @spec from_send_types(send_types :: %{optional(String.t()) => module()}) :: t() | nil + @spec from_send_types(send_types :: %{optional(String.t()) => registration()}) :: t() | nil def from_send_types(send_types) when is_map(send_types) and map_size(send_types) == 0, do: nil def from_send_types(send_types) when is_map(send_types) do + processors = Map.new(send_types, fn {type, registration} -> {type, split(registration)} end) + %__MODULE__{ types: send_types |> Map.keys() |> MapSet.new(), - entries: Map.new(send_types, fn {type, module} -> {type, entry!(module, type)} end) + entries: + Map.new(processors, fn {type, {module, _opts}} -> {type, entry!(module, type)} end), + processors: processors } end + # The module and options of one `t:registration/0`: a bare module's + # options are `[]`. Callable across the library's own modules (the + # planner reads registrations through it) but not part of its public API, + # hence `@doc false`. + @doc false + @spec split(registration :: registration()) :: {module(), keyword()} + def split({module, opts}) when is_atom(module) and is_list(opts), do: {module, opts} + def split(module) when is_atom(module), do: {module, []} + + # The value `module` returns from `ioprocessors_entry/2` for `type` and + # `context`, checked as `from_send_types/1` checks a `/1` entry. + # Internal: `Statifier.Evaluator.SystemVariables.initial/3` is its one + # caller (ADR-0075 decision 3), hence `@doc false`. + @doc false + @spec session_entry!(module :: module(), type :: String.t(), context :: map()) :: map() + def session_entry!(module, type, context), + do: checked!(module.ioprocessors_entry(type, context), module, type) + # The processor's own `_ioprocessors` value for `type`, checked at the one # constructor so a value that reaches the datamodel is string-keyed by # construction, as `Statifier.MachineState`'s datamodel invariant needs. @@ -103,6 +153,11 @@ defmodule Statifier.Send.Types do do: module.ioprocessors_entry(type), else: %{} + checked!(entry, module, type) + end + + @spec checked!(entry :: term(), module :: module(), type :: String.t()) :: map() + defp checked!(entry, module, type) do unless is_map(entry) and string_keyed?(entry) do raise ArgumentError, "#{inspect(module)}.ioprocessors_entry(#{inspect(type)}) must return a map " <> diff --git a/lib/statifier/session.ex b/lib/statifier/session.ex index ff81058a..5f755a91 100644 --- a/lib/statifier/session.ex +++ b/lib/statifier/session.ex @@ -412,7 +412,7 @@ defmodule Statifier.Session do # `%{}`. `init/1` derives the `%MachineState{}` `send_types` stamp # from this same map's keys through # `Statifier.Send.Types.from_send_types/1`, the one constructor. - send_types: %{String.t() => module()}, + send_types: %{String.t() => SendTypes.registration()}, # ADR-0069 decision 4's cancel routing: which registered types' # processors hold a delayed send under each send id, so a # `` naming it reaches them. Kept by the `{:notify, _}` @@ -557,9 +557,12 @@ defmodule Statifier.Session do subtrees their own. Handlers descend independently of `:invoke_source`, which ADR-0038 leaves to its own option, and independently of `:inherit_observers`, which is an observation knob. - - `:send_types` - a `%{type_string => module}` map of the Event I/O - Processor types this session registers for `` - (ADR-0069 decision 2). Default `%{}`, which registers nothing: only + - `:send_types` - a `%{type_string => registration}` map of the Event + I/O Processor types this session registers for `` + (ADR-0069 decision 2), each registration a `Statifier.Send.Processor` + module or `{module, opts}` (ADR-0075 decision 8, point b; see + `Statifier.Send.Processor`'s "Registration options"). Default `%{}`, + which registers nothing: only the built-in types (the attribute absent, `"scxml"`, and the SCXML Event I/O Processor URI) are supported, exactly as before. The `%MachineState{}` `send_types` snapshot this session's core is stamped @@ -1264,7 +1267,7 @@ defmodule Statifier.Session do session_id :: String.t(), invoked_by :: {pid(), String.t()} | nil, invoke_handlers :: %{String.t() => module()}, - send_types :: %{String.t() => module()} + send_types :: %{String.t() => SendTypes.registration()} ) :: {MachineState.t(), [Effect.t()], keyword(), :initialize | :resume, binary() | nil} defp boot(:fresh, machine, opts, session_id, invoked_by, invoke_handlers, send_types) do machine_opts = diff --git a/lib/statifier/session/effects.ex b/lib/statifier/session/effects.ex index eea86025..ba729834 100644 --- a/lib/statifier/session/effects.ex +++ b/lib/statifier/session/effects.ex @@ -57,6 +57,10 @@ defmodule Statifier.Session.Effects do instructions are spliced into the plan in place of any delivery. The target is never parsed. A delayed send of a registered type is planned the same way, with no `{:schedule, ...}`: the processor owns the timer. + A `:send_processors` value is a module or `{module, opts}` (ADR-0075 + decision 8, point b); for `{module, opts}` the context the processor's + `deliver/3` and `cancel/2` receive carries the options under an added + `:opts` key, and for a bare module it is the plan context unchanged. A `` always plans `{:cancel_timers, send_id}` for the library's own timers, as before. When the plan context's `:held_sends` map, or a @@ -97,7 +101,8 @@ defmodule Statifier.Session.Effects do optional keys for registered send types (ADR-0069): `send_types`, the registered set read off `%MachineState{}` exactly as `invoke_types` is; `send_processors`, the session's `:send_types` map from type to - `Statifier.Send.Processor` module; and `held_sends`, the live + `Statifier.Send.Processor` registration (a module or `{module, opts}`); + and `held_sends`, the live `send_id => [type]` map of delayed sends a processor holds. A context without them plans exactly as a session with no registered send type does. @@ -192,7 +197,7 @@ defmodule Statifier.Session.Effects do required(:invoke_handlers) => %{String.t() => module()}, required(:invocation_types) => %{String.t() => String.t()}, optional(:send_types) => Types.t() | nil, - optional(:send_processors) => %{String.t() => module()}, + optional(:send_processors) => %{String.t() => Types.registration()}, optional(:held_sends) => held_sends() } @@ -306,7 +311,8 @@ defmodule Statifier.Session.Effects do defp plan_one({:cancel, %Cancel{send_id: send_id} = cancel} = effect, context) do processors = Enum.flat_map(Map.get(held(context), send_id, []), fn type -> - {:ok, instructions} = processor_for(type, context).cancel(cancel, context) + {module, ctx} = processor_for(type, context) + {:ok, instructions} = module.cancel(cancel, ctx) instructions end) @@ -408,19 +414,26 @@ defmodule Statifier.Session.Effects do # planned for it - no delivery, no `{:schedule, ...}`, no target parse. @spec hand_off(send :: Send.t() | SendDelayed.t(), context :: context()) :: [instruction()] defp hand_off(send, %{session_id: session_id} = context) do - {:ok, instructions} = - processor_for(send.type, context).deliver(send, delivered_event(send, session_id), context) - + {module, ctx} = processor_for(send.type, context) + {:ok, instructions} = module.deliver(send, delivered_event(send, session_id), ctx) instructions end - # The processor module registered for `type`. `Statifier.Session` and - # `Statifier.Replay` derive `:send_types` and `:send_processors` from one - # map, so a registered type always has a module; a context that declares a - # set without the map is the caller's error and raises here. - @spec processor_for(type :: String.t(), context :: context()) :: module() - defp processor_for(type, context), - do: context |> Map.get(:send_processors, %{}) |> Map.fetch!(type) + # The processor module registered for `type`, with the context its + # planning callbacks receive. `Statifier.Session` and `Statifier.Replay` + # derive `:send_types` and `:send_processors` from one map, so a + # registered type always has a module; a context that declares a set + # without the map is the caller's error and raises here. A `{module, + # opts}` registration adds its options to the context under `:opts`; a + # bare module's context is the plan context unchanged (ADR-0075 decision + # 8, point b). + @spec processor_for(type :: String.t(), context :: context()) :: {module(), map()} + defp processor_for(type, context) do + case context |> Map.get(:send_processors, %{}) |> Map.fetch!(type) do + {module, opts} -> {module, Map.put(context, :opts, opts)} + module -> {module, context} + end + end # An ``'s own routing (see moduledoc's "`` routing" # section). Unlike `plan_send/3`, there is no target to check - `` diff --git a/lib/statifier/session/recording.ex b/lib/statifier/session/recording.ex index 27115ccf..5643738c 100644 --- a/lib/statifier/session/recording.ex +++ b/lib/statifier/session/recording.ex @@ -135,6 +135,15 @@ defmodule Statifier.Session.Recording do - into `{:error, {:unknown_handler_modules, names}}`, sorted, so a host learns the whole set of modules it needs to load in one round trip. + A `{module, opts}` `:send_types` registration (ADR-0075 decision 8, + point b) crosses the same way, with its options as strings too: the + module is written as its name, each option key as its name, and each + option value that is an atom other than `true`, `false` and `nil` as + `{:atom, name}`; any other option value is written as it is. + `from_binary/1` resolves the module, the keys and the atom values back + with `String.to_existing_atom/1`, and a name that does not resolve joins + the same `{:unknown_handler_modules, names}` list. + What the codec does not, and cannot, verify: that a resolved handler module's planning callbacks (ADR-0051 decision 4) behave the way they did when the recording was made. Replay's determinism depends on that @@ -227,14 +236,16 @@ defmodule Statifier.Session.Recording do context is built from this recorded map, not from an empty one. `:send_types` is `Statifier.Session.start_link/2`'s own - `%{type_string => module}` map (ADR-0069), kept as the map rather than + `%{type_string => registration}` map (ADR-0069; a registration is a + module or `{module, opts}`, ADR-0075), kept as the map rather than the snapshot derived from it, so `Statifier.Replay` re-derives the snapshot through `Statifier.Send.Types.from_send_types/1`, the one constructor. It is kept only when non-empty: a session that registers no send type records exactly the options it recorded before ADR-0069, and an absent key replays as "no declaration". `to_binary/1` writes its - module values as strings, never atoms or code, under ADR-0057 decision - 5's rule for `:invoke_handlers`. + module values, and a registration's options, as strings, never atoms or + code, under ADR-0057 decision 5's rule for `:invoke_handlers` (see the + moduledoc's "The binary contract" section). `opts[:session_id]` should be the id the session actually resolved to (`machine_state.datamodel["_sessionid"]`), not merely whatever the caller @@ -634,12 +645,27 @@ defmodule Statifier.Session.Recording do defp encode_opts(opts) do opts |> Keyword.replace_lazy(:invoke_handlers, &module_names/1) - |> Keyword.replace_lazy(:send_types, &module_names/1) + |> Keyword.replace_lazy(:send_types, ®istration_names/1) end defp module_names(map), do: Map.new(map, fn {type, module} -> {type, Atom.to_string(module)} end) + # ADR-0075 decision 8, point b: a `{module, opts}` registration is written + # with its options as strings too (see the moduledoc's "The binary + # contract" section); a bare module exactly as `module_names/1` writes it. + defp registration_names(map) do + Map.new(map, fn + {type, {module, opts}} -> {type, {Atom.to_string(module), Enum.map(opts, &option_name/1)}} + {type, module} -> {type, Atom.to_string(module)} + end) + end + + defp option_name({key, value}) when is_atom(value) and value not in [true, false, nil], + do: {Atom.to_string(key), {:atom, Atom.to_string(value)}} + + defp option_name({key, value}), do: {Atom.to_string(key), value} + @spec decode_opts(opts :: keyword()) :: {:ok, keyword()} | {:error, {:unknown_handler_modules, [String.t()]}} defp decode_opts(opts) do @@ -647,7 +673,7 @@ defmodule Statifier.Session.Recording do Enum.reduce([:invoke_handlers, :send_types], {opts, []}, fn key, {opts, unknown} -> case Keyword.fetch(opts, key) do {:ok, modules} -> - {resolved, missing} = resolve_modules(modules) + {resolved, missing} = resolve_modules(modules, key == :send_types) {Keyword.put(opts, key, resolved), missing ++ unknown} :error -> @@ -665,16 +691,59 @@ defmodule Statifier.Session.Recording do # returning the resolved map and every name that did not resolve. Shared # by `:invoke_handlers` and `:send_types`, whose unresolved names are # reported together in one `{:unknown_handler_modules, names}` error. - @spec resolve_modules(modules :: map()) :: {map(), [String.t()]} - defp resolve_modules(modules) do + @spec resolve_modules(modules :: map(), options? :: boolean()) :: {map(), [String.t()]} + defp resolve_modules(modules, options?) do Enum.reduce(modules, {%{}, []}, fn {type, name}, {resolved, unknown} -> - case existing_atom(name) do - {:ok, module} -> {Map.put(resolved, type, module), unknown} - :error -> {resolved, [handler_name(name) | unknown]} + case resolve_registration(name, options?) do + {:ok, registration} -> {Map.put(resolved, type, registration), unknown} + {:error, missing} -> {resolved, missing ++ unknown} end end) end + # One recorded value: a module name, or (for `:send_types` only, ADR-0075) + # a `{module name, options}` pair written by `registration_names/1`. + @spec resolve_registration(name :: term(), options? :: boolean()) :: + {:ok, term()} | {:error, [String.t()]} + defp resolve_registration({name, opts}, true) when is_list(opts) do + {options, missing} = Enum.map_reduce(opts, [], &resolve_option/2) + + case {existing_atom(name), missing} do + {{:ok, module}, []} -> {:ok, {module, options}} + {{:ok, _module}, missing} -> {:error, missing} + {:error, missing} -> {:error, [handler_name(name) | missing]} + end + end + + defp resolve_registration(name, _options?) do + case existing_atom(name) do + {:ok, module} -> {:ok, module} + :error -> {:error, [handler_name(name)]} + end + end + + @spec resolve_option(option :: term(), missing :: [String.t()]) :: {term(), [String.t()]} + defp resolve_option({key, {:atom, value}}, missing) do + {resolved_key, missing} = resolve_name(key, missing) + {resolved_value, missing} = resolve_name(value, missing) + {{resolved_key, resolved_value}, missing} + end + + defp resolve_option({key, value}, missing) do + {resolved_key, missing} = resolve_name(key, missing) + {{resolved_key, value}, missing} + end + + defp resolve_option(other, missing), do: {other, [handler_name(other) | missing]} + + @spec resolve_name(name :: term(), missing :: [String.t()]) :: {term(), [String.t()]} + defp resolve_name(name, missing) do + case existing_atom(name) do + {:ok, atom} -> {atom, missing} + :error -> {name, [handler_name(name) | missing]} + end + end + # `String.to_existing_atom/1` has no non-raising variant, so the rescue is # function-level here - the same shape `safe_decode/1` below uses, rather # than a `try` block inline in the reduce. This is not a rescue-to-default diff --git a/lib/statifier/testing/case.ex b/lib/statifier/testing/case.ex index dbec3015..269b41b6 100644 --- a/lib/statifier/testing/case.ex +++ b/lib/statifier/testing/case.ex @@ -149,6 +149,14 @@ defmodule Statifier.Testing.Case do `#{@default_configuration_deadline_ms}`) - upper bound on waiting for a session to reach an expected configuration. Bounds only the wrong answer: a chart that cannot change again exits the poll immediately. + + One more option registers send types, and has no default: + + - `:send_types` - a `:send_types` map (`Statifier.Session.start_link/2`'s + option of that name) the session is started with, for a document whose + ``s go through a registered Event I/O Processor (ADR-0075). A call + that passes it always drives the document through a session; a call + without it starts its session exactly as before. """ @spec test_scxml( xml :: String.t(), @@ -160,7 +168,7 @@ defmodule Statifier.Testing.Case do def test_scxml(xml, description, expected_initial_config, events, opts \\ []) do detected = validate_features!(xml, description) - if session_required?(detected) do + if session_required?(detected) or Keyword.has_key?(opts, :send_types) do drive_through_session(xml, expected_initial_config, events, opts) else drive_synchronously(xml, expected_initial_config, events) @@ -203,7 +211,12 @@ defmodule Statifier.Testing.Case do # discard-on-termination). defp drive_through_session(xml, expected_initial_config, events, opts) do machine = parse_document(xml) - {:ok, session} = Statifier.start_session(machine, subscribers: [self()]) + + {:ok, session} = + Statifier.start_session( + machine, + [subscribers: [self()]] ++ Keyword.take(opts, [:send_types]) + ) try do assert_configuration_eventually(session, expected_initial_config, opts) diff --git a/test/mix/statifier/basic_http_front_test.exs b/test/mix/statifier/basic_http_front_test.exs new file mode 100644 index 00000000..61716bf0 --- /dev/null +++ b/test/mix/statifier/basic_http_front_test.exs @@ -0,0 +1,212 @@ +defmodule Mix.Statifier.BasicHTTPFrontTest do + use ExUnit.Case, async: true + + # ADR-0075 end to end over a real socket: the loopback front answers by + # the decoder's status rule (decision 5), and the default `:httpc` + # transport (decision 6) delivers a session's send to its own location. + + import Statifier.Testing.Case, only: [test_scxml: 5] + + alias Mix.Statifier.BasicHTTPFront + alias Statifier.Send.BasicHTTP + alias Statifier.Send.BasicHTTP.Transport.Httpc + + @uri "http://www.w3.org/TR/scxml/#BasicHTTPEventProcessor" + + setup_all do + {:ok, front} = BasicHTTPFront.start() + on_exit(fn -> BasicHTTPFront.stop(front) end) + %{front: front} + end + + defp send_types(%{base_url: base_url}), + do: %{@uri => {BasicHTTP, base_url: base_url}, "basichttp" => {BasicHTTP, base_url: base_url}} + + defp start_idle!(front) do + {:ok, machine} = + Statifier.compile(""" + + + + """) + + {:ok, session} = Statifier.start_session(machine, send_types: send_types(front)) + on_exit(fn -> if Process.alive?(session), do: Statifier.Session.stop(session) end) + {session, Statifier.Session.snapshot(session).datamodel["_sessionid"]} + end + + # A raw request through `:httpc`, answering the status and the headers. + defp request(method, url, body \\ "") do + request = + if method == :post, + do: {String.to_charlist(url), [], ~c"application/x-www-form-urlencoded", body}, + else: {String.to_charlist(url), []} + + {:ok, {{_version, status, _reason}, headers, _body}} = :httpc.request(method, request, [], []) + {status, headers} + end + + describe "a session's send reaches its own location" do + # sabotage: `BasicHTTPFront.answer/2` answers 204 without + # `Session.send_event/2` -> the event never arrives and the chart + # stays in `s`. Confirmed red and reverted. + test "the event and its data arrive through the front", %{front: front} do + test_scxml( + """ + + + + + + + + + + + + + + """, + "a send to the session's own location", + ["pass"], + [], + send_types: send_types(front) + ) + end + + # sabotage: `decode/1` names every nameless event `"HTTP"` -> the + # transition on `HTTP.POST` never matches. Confirmed red and reverted. + test "a POST with no _scxmleventname arrives as HTTP.POST", %{front: front} do + test_scxml( + """ + + + + + + + + + + + + + + """, + "a nameless send", + ["pass"], + [], + send_types: send_types(front) + ) + end + end + + describe "the dedup key" do + # sabotage: `decode/1` sets the event's `sendid` from the key's second + # field -> the generated send id reaches `_event.sendid`, the first + # transition matches, and the chart reaches `fail`. Confirmed red and + # reverted. + test "a POST whose key names a generated send id delivers an event with no sendid", + %{front: front} do + test_scxml( + """ + + + + + + + + + + + + + """, + "a generated send id is not the event's sendid", + ["pass"], + [], + send_types: send_types(front) + ) + end + + # sabotage: `BasicHTTPFront.respond/1` hands the decoder `send_key: nil` + # -> the malformed header is never checked, the POST is answered 204, + # and the match reddens. Confirmed red and reverted. + test "the front hands the header to the decoder, which refuses a malformed one", + %{front: front} do + {_session, session_id} = start_idle!(front) + url = String.to_charlist(front.base_url <> "/" <> session_id) + request = {url, [{~c"scxml-send-key", ~c"not/a/key"}], ~c"text/plain", "x"} + + assert {:ok, {{_version, 400, _reason}, _headers, _body}} = + :httpc.request(:post, request, [], []) + end + end + + describe "the front's status rule" do + # sabotage: `answer/2` answers 200 after enqueueing -> the status + # reads 200 and the equality reddens. Confirmed red and reverted. + test "204 once the event is enqueued", %{front: front} do + {_session, session_id} = start_idle!(front) + + assert {204, _headers} = + request(:post, front.base_url <> "/" <> session_id, "_scxmleventname=hello") + end + + # sabotage: `head/1`'s 405 clause drops `allow` -> the header is + # missing and the match reddens. Confirmed red and reverted. + test "405 with Allow: POST for another method", %{front: front} do + {_session, session_id} = start_idle!(front) + + assert {405, headers} = request(:get, front.base_url <> "/" <> session_id) + assert {~c"allow", ~c"POST"} in headers + end + + # sabotage: `answer/2`'s `{:error, _}` arm answers 204 -> the status + # reads 204 and the equality reddens. Confirmed red and reverted. + test "400 for a request that forms no event", %{front: front} do + {_session, session_id} = start_idle!(front) + + assert {400, _headers} = request(:post, front.base_url <> "/" <> session_id, "a=%FF") + end + + # sabotage: `answer/2` answers 204 for `:no_session` -> the unknown id + # reads 204 and the first match reddens. Confirmed red and reverted. + test "404 for a location that names no live session", %{front: front} do + assert {404, _headers} = request(:post, front.base_url <> "/sess_nobody", "") + assert {404, _headers} = request(:post, String.replace(front.base_url, "/basichttp", "/x")) + end + end + + describe "the default :httpc transport" do + # sabotage: `Httpc.post/3` answers `{:ok, 204}` without requesting -> + # the unknown-session POST reads 204, not the front's 404. Confirmed red + # and reverted. + test "answers the status the front answered", %{front: front} do + assert Httpc.post(front.base_url <> "/sess_nobody", [{"content-type", "text/plain"}], "x") == + {:ok, 404} + end + + # sabotage: `request/2`'s `{:error, reason}` arm answers `{:ok, 204}` + # -> the refused connection reads as a status and the first match + # reddens. Confirmed red and reverted. + test "answers an error when nothing listens, over http and https" do + assert {:error, _reason} = Httpc.post("http://127.0.0.1:1/x", [], "") + assert {:error, _reason} = Httpc.post("https://127.0.0.1:1/x", [], "") + end + + # sabotage: `Httpc.ensure_loaded/1` answers `:ok` without loading -> + # the missing module reads as loadable and the equality reddens. + # Confirmed red and reverted. + test "a module it needs that cannot be loaded is an error, not a crash" do + assert Httpc.ensure_loaded([:ssl, :public_key]) == :ok + + assert Httpc.ensure_loaded([:ssl, :no_such_module_for_this_test]) == + {:error, {:not_loadable, :no_such_module_for_this_test, :nofile}} + end + end +end diff --git a/test/statifier/send/basic_http_session_test.exs b/test/statifier/send/basic_http_session_test.exs new file mode 100644 index 00000000..8ac8cce3 --- /dev/null +++ b/test/statifier/send/basic_http_session_test.exs @@ -0,0 +1,255 @@ +defmodule Statifier.Send.BasicHTTPSessionTest do + use ExUnit.Case, async: false + + # ADR-0075 through a live `Statifier.Session`: the two `_ioprocessors` + # keys and their one location (decision 3), the refusal of a registration + # without `:base_url` (decision 8, point b), a send's POST, C.2.2's + # `error.communication` for a missing target, a failed delivery reaching + # the sender through `failed_send/3` (decision 8, point d), and a delayed + # send and its cancel (decision 9). Every POST goes through + # `Statifier.BasicHTTPTestTransport`, which this test process registers + # under its own name to receive them, so the tests are `async: false`. + + import Statifier.Testing.Case, only: [test_scxml: 5] + + alias Statifier.{BasicHTTPTestTransport, Session} + alias Statifier.Send.BasicHTTP + + @uri "http://www.w3.org/TR/scxml/#BasicHTTPEventProcessor" + @base_url "http://front.test/basichttp" + @opts [base_url: @base_url, transport: BasicHTTPTestTransport] + @send_types %{@uri => {BasicHTTP, @opts}, "basichttp" => {BasicHTTP, @opts}} + + setup do + Process.register(self(), BasicHTTPTestTransport) + :ok + end + + defp start!(xml, send_types \\ @send_types) do + {:ok, machine} = Statifier.compile(xml) + {:ok, session} = Statifier.start_session(machine, send_types: send_types) + on_exit(fn -> if Process.alive?(session), do: Session.stop(session) end) + session + end + + @idle """ + + + + """ + + describe "the _ioprocessors entry (C.2.3)" do + # sabotage: `SystemVariables.initial/3` keeps the `/1` entry for every + # type (`session_entry/5` returns `entry`) -> both keys read `%{}` and + # the equality reddens. Confirmed red and reverted. + test "both registered keys carry one location: the base URL and the session id" do + session = start!(@idle) + %{datamodel: datamodel} = Session.snapshot(session) + location = %{"location" => @base_url <> "/" <> datamodel["_sessionid"]} + + assert Map.take(datamodel["_ioprocessors"], [@uri, "basichttp"]) == %{ + @uri => location, + "basichttp" => location + } + end + + # sabotage: `ioprocessors_entry/2` falls back to a default base URL + # instead of raising -> the session starts and the match reddens. + # Confirmed red and reverted. + test "a registration without :base_url is refused when the session starts" do + {:ok, machine} = Statifier.compile(@idle) + + ExUnit.CaptureLog.capture_log(fn -> + assert {:error, {%ArgumentError{message: message}, _stack}} = + Statifier.start_session(machine, send_types: %{"basichttp" => BasicHTTP}) + + assert message =~ ":base_url" + end) + end + end + + describe "an outbound send" do + # sabotage: `perform/2`'s `{:post, _}` clause returns `:ok` without + # calling the transport -> no POST reaches this process and + # `assert_receive` reddens. Confirmed red and reverted. + test "POSTs the event name and parameters to the target" do + start!(""" + + + + + + + + + + + """) + + assert_receive {:basichttp_post, "http://sink.test/in", headers, body} + + assert {"content-type", "application/x-www-form-urlencoded"} in headers + assert URI.decode_query(body) == %{"_scxmleventname" => "ping", "Var1" => "2", "p" => "x"} + end + + # sabotage: `deliver/3`'s nil-target clause raises `error.execution` + # instead -> the chart takes the catch-all to `fail` and the + # configuration assertion reddens. Confirmed red and reverted. + test "a send with no target raises error.communication with the send id" do + test_scxml( + """ + + + + + + + + + + """, + "no target", + ["pass"], + [], + send_types: @send_types + ) + end + + for {what, target} <- [ + {"a transport error", "http://sink.test/answer/error"}, + {"a status outside 2xx", "http://sink.test/answer/500"} + ] do + # sabotage: `post_now/2` answers `:ok` for every transport answer -> + # no miss is reported, no error.communication is raised, and the + # chart never reaches `pass`. Confirmed red and reverted. + test "#{what} reaches the sender as error.communication through failed_send/3" do + test_scxml( + """ + + + + + + + + + + + + """, + unquote(what), + ["pass"], + [], + send_types: @send_types + ) + end + end + + # sabotage: `report/3` returns `:ok` when no session is registered -> + # the miss is swallowed and the equality reddens. Confirmed red and + # reverted. + test "with no live session under the sender's id, a miss is returned only" do + post = %{ + url: "http://sink.test/answer/503", + headers: [], + body: "", + transport: BasicHTTPTestTransport, + send: %Statifier.Effect.Send{ + send_id: "gone", + event: "e", + macrostep: 1, + microstep: 1, + round: 0 + } + } + + assert BasicHTTP.perform({:post, post}, %{session_id: "sess_nobody"}) == + {:error, {:http_status, 503}} + end + end + + describe "a delayed send and its cancel" do + # sabotage: `hold/4`'s `after` arm discards instead of POSTing -> no + # POST arrives and `assert_receive` reddens. Confirmed red and reverted. + test "a delayed send POSTs once its delay has passed" do + start!(""" + + + + + + + + """) + + refute_received {:basichttp_post, _url, _headers, _body} + + assert_receive {:basichttp_post, "http://sink.test/later", _headers, + "_scxmleventname=later"} + end + + # sabotage: `perform/2`'s `{:cancel, _}` clause returns `:ok` without + # messaging the timers -> the POST fires and `refute_receive` reddens. + # Confirmed red and reverted. + test "a cancel before the delay passes stops the POST" do + start!(""" + + + + + + + + + """) + + refute_receive {:basichttp_post, _url, _headers, _body}, 300 + end + + # sabotage: `hold/4` POSTs at fire time without the `running?/1` check + # -> the halted session's send fires and `refute_receive` reddens. + # Confirmed red and reverted. + test "a session that has halted discards a delayed send it still holds" do + start!(""" + + + + + + + + + + """) + + refute_receive {:basichttp_post, _url, _headers, _body}, 250 + end + + # sabotage: `hold/4`'s `{:DOWN, ...}` arm is deleted -> the timer + # outlives its stopped session, and `Process.alive?/1` reads true. + # Confirmed red and reverted. + test "a timer ends when the session holding it stops" do + session = + start!(""" + + + + + + + + """) + + # A call is answered only after the session has performed its start. + _status = Session.status(session) + {:dictionary, dictionary} = Process.info(session, :dictionary) + {_key, [timer]} = List.keyfind(dictionary, {BasicHTTP, "t"}, 0) + ref = Process.monitor(timer) + Session.stop(session) + assert_receive {:DOWN, ^ref, :process, ^timer, _reason} + end + end +end diff --git a/test/statifier/send/basic_http_test.exs b/test/statifier/send/basic_http_test.exs new file mode 100644 index 00000000..ac474570 --- /dev/null +++ b/test/statifier/send/basic_http_test.exs @@ -0,0 +1,331 @@ +defmodule Statifier.Send.BasicHTTPTest do + use ExUnit.Case, async: true + + # ADR-0075: the Basic HTTP Event I/O Processor's pure halves - the + # outbound mapping `deliver/3` plans (decision 4), the `_ioprocessors` + # entry (decision 3), the cancel plan, and the inbound decoder + # (decision 5). The performing half runs in + # `Statifier.Send.BasicHTTPSessionTest`. + + alias Statifier.Effect.{Cancel, Send, SendDelayed} + alias Statifier.Event + alias Statifier.Send.BasicHTTP + alias Statifier.Send.BasicHTTP.Transport.Httpc + + @uri "http://www.w3.org/TR/scxml/#BasicHTTPEventProcessor" + @target "http://127.0.0.1:1/basichttp/sess_receiver" + @form "application/x-www-form-urlencoded" + + defp send_effect(fields) do + struct!( + %Send{ + event: "ping", + type: @uri, + target: @target, + data: :undefined, + send_id: "send_1", + c_index: 3, + owner: {:onentry, 2, 0}, + macrostep: 1, + microstep: 1, + round: 0, + ordinal: 1 + }, + fields + ) + end + + defp ctx(opts \\ nil) do + base = %{session_id: "sess_sender"} + if opts, do: Map.put(base, :opts, opts), else: base + end + + # The one `{:post, request}` payload `deliver/3` plans for `effect`. + defp planned(effect, ctx \\ ctx()) do + assert {:ok, [{:handler, BasicHTTP, {:post, post}}]} = + BasicHTTP.deliver(effect, Event.external("ignored"), ctx) + + post + end + + defp content_type(post), do: post.headers |> List.keyfind("content-type", 0) |> elem(1) + defp send_key(post), do: post.headers |> List.keyfind("scxml-send-key", 0) |> elem(1) + + describe "deliver/3: the outbound mapping (C.2.2)" do + # sabotage: `post/2` leaves `_scxmleventname` out of `named` -> the body + # carries the params alone and the equality reddens. Confirmed red and + # reverted. + test "event, namelist and params are form parameters in a POST body" do + post = planned(send_effect(data: %{"Var1" => 2, "name" => "two words"})) + + assert post.url == @target + assert content_type(post) == @form + + assert URI.decode_query(post.body) == %{ + "_scxmleventname" => "ping", + "Var1" => "2", + "name" => "two words" + } + end + + # sabotage: the `:undefined` arm of `post/2` sends `text/plain` with an + # empty body -> the content type and the decoded body both differ and + # the match reddens. Confirmed red and reverted. + test "a send with no parameters and no content sends the event name alone" do + post = planned(send_effect(data: :undefined)) + + assert content_type(post) == @form + assert post.body == "_scxmleventname=ping" + end + + # sabotage: `with_query/2` returns the URL unchanged -> the event name is + # lost and the URL equality reddens. Confirmed red and reverted. + test "a content body is the body, and the event name goes in the query string" do + post = planned(send_effect(data: "some content")) + + assert post.url == @target <> "?_scxmleventname=ping" + assert content_type(post) == "text/plain" + assert post.body == "some content" + end + + # sabotage: `with_query/2` always uses `?` -> the URL carries two `?` + # and the equality reddens. Confirmed red and reverted. + test "the event name joins a query string the target already carries" do + post = planned(send_effect(target: @target <> "?a=1", data: "x")) + + assert post.url == @target <> "?a=1&_scxmleventname=ping" + end + + # sabotage: `post/2` always names the event, with an empty name when the + # send has none -> the target gains a query and the equality reddens. + # Confirmed red and reverted. + test "a content body without an event name leaves the target alone" do + post = planned(send_effect(event: nil, data: 42)) + + assert post.url == @target + assert post.body == "42" + end + + # sabotage: `encode/1`'s `:undefined` clause returns `"undefined"` -> + # the decoded `unbound` value reads "undefined" and the equality + # reddens. Confirmed red and reverted. + test "parameter values are written as text" do + post = + planned( + send_effect( + data: %{"null" => nil, "unbound" => :undefined, "flag" => true, "list" => [1, 2]} + ) + ) + + assert URI.decode_query(post.body) == %{ + "_scxmleventname" => "ping", + "null" => "null", + "unbound" => "", + "flag" => "true", + "list" => "[1, 2]" + } + end + + # sabotage: the `deliver/3` clause for a nil target is deleted -> a + # `{:handler, ...}` POST is planned instead of the raise and the match + # reddens. Confirmed red and reverted. + test "a send with no target plans error.communication and no request" do + effect = send_effect(target: nil) + + assert BasicHTTP.deliver(effect, Event.external("ping"), ctx()) == + {:ok, + [ + {:raise, :platform, "error.communication", {:content, 3, {:onentry, 2, 0}}, + sendid: "send_1"} + ]} + end + + # sabotage: `post/2` reads the transport from `ctx` without the + # `:opts` key (always the default) -> the injected module is lost and + # the first equality reddens. Confirmed red and reverted. + test "the transport is the registration's :transport option, else the httpc adapter" do + assert planned(send_effect([]), ctx(transport: SomeTransport)).transport == SomeTransport + assert planned(send_effect([]), ctx(base_url: "x")).transport == Httpc + assert planned(send_effect([]), ctx()).transport == Httpc + end + + # sabotage: the `SendDelayed` clause of `deliver/3` plans `{:post, _}` + # -> no delay travels with the payload and the match reddens. Confirmed + # red and reverted. + test "a delayed send plans its delay with the request" do + delayed = struct!(SendDelayed, Map.from_struct(send_effect([])) |> Map.put(:delay_ms, 250)) + + assert {:ok, [{:handler, BasicHTTP, {:post_after, 250, %{url: @target, send: ^delayed}}}]} = + BasicHTTP.deliver(delayed, Event.external("ping"), ctx()) + end + + # sabotage: `send_key/2` leaves out `ordinal` -> the value has seven + # fields and both equalities redden. Confirmed red and reverted. + test "every POST carries the send's dedup key, form body or content body" do + key = "sess_sender/send_1/1/1/0/3/onentry.2.0/1" + + assert send_key(planned(send_effect(data: %{"a" => 1}))) == key + assert send_key(planned(send_effect(data: "some content"))) == key + end + + # sabotage: `escape/1` returns the value unencoded -> the send id's + # `/` and space survive and the equality reddens. Confirmed red and + # reverted. + test "the key escapes the session scope and send id and spells a transition owner" do + effect = + send_effect( + send_id: "a/b c", + owner: {:transition, 4}, + macrostep: 7, + microstep: 2, + round: 1, + ordinal: 9 + ) + + assert send_key(planned(effect, %{session_id: "sess/x"})) == + "sess%2Fx/a%2Fb%20c/7/2/1/3/transition.4/9" + end + + # sabotage: the `SendDelayed` clause of `deliver/3` builds its request + # with `headers: []` -> no key travels and the membership assertion + # reddens. Confirmed red and reverted. + test "a delayed send's POST carries its key too" do + delayed = + struct!( + SendDelayed, + Map.from_struct(send_effect(owner: {:onexit, 5, 1}, ordinal: 4)) + |> Map.put(:delay_ms, 250) + ) + + assert {:ok, [{:handler, BasicHTTP, {:post_after, 250, post}}]} = + BasicHTTP.deliver(delayed, Event.external("ping"), ctx()) + + assert {"scxml-send-key", "sess_sender/send_1/1/1/0/3/onexit.5.1/4"} in post.headers + end + + # sabotage: `cancel/2` returns `{:ok, []}` -> no cancel instruction is + # planned and the equality reddens. Confirmed red and reverted. + test "a cancel plans the cancellation of the timers held under its send id" do + assert BasicHTTP.cancel( + %Cancel{send_id: "send_1", macrostep: 1, microstep: 1, round: 0, ordinal: 2}, + ctx() + ) == + {:ok, [{:handler, BasicHTTP, {:cancel, "send_1"}}]} + end + end + + describe "ioprocessors_entry/2 (C.2.3)" do + # sabotage: the location drops the `/` between the base URL and the + # session id -> the equality reddens. Confirmed red and reverted. + test "the location is the base URL, a slash, and the session id" do + assert BasicHTTP.ioprocessors_entry(@uri, %{ + session_id: "sess_1", + opts: [base_url: "http://h/b"] + }) == + %{"location" => "http://h/b/sess_1"} + end + + # sabotage: the missing-option clause returns `%{}` instead of raising + # -> `assert_raise` reddens. Confirmed red and reverted. + test "a registration without :base_url is refused" do + assert_raise ArgumentError, ~r/needs a :base_url option/, fn -> + BasicHTTP.ioprocessors_entry("basichttp", %{session_id: "sess_1", opts: []}) + end + end + end + + describe "decode/1: the inbound half (C.2.1)" do + defp request(fields), + do: Map.merge(%{method: "POST", content_type: @form, body: "", query: nil}, fields) + + # sabotage: `decode/1` reads the name from the last `_scxmleventname` + # instead of the first -> the body's name wins and the equality + # reddens. Confirmed red and reverted. + test "the event name is the first _scxmleventname, the query string before the body" do + assert {:ok, %Event{name: "from.query"}} = + BasicHTTP.decode( + request(%{ + query: "_scxmleventname=from.query", + body: "_scxmleventname=from.body&_scxmleventname=again" + }) + ) + + assert {:ok, %Event{name: "from.body"}} = + BasicHTTP.decode(request(%{body: "_scxmleventname=from.body&_scxmleventname=x"})) + end + + # sabotage: the no-name arm returns `"HTTP." <> method` without + # upcasing -> a lower-case `post` method yields `HTTP.post` and the + # match reddens. Confirmed red and reverted. + test "without _scxmleventname the event is named for the method" do + assert {:ok, %Event{name: "HTTP.POST", data: :undefined}} = + BasicHTTP.decode(request(%{method: "post"})) + end + + # sabotage: `data/2` skips `text/1` and keeps each value a string -> + # `Var1` reads "2" and the equality reddens. Confirmed red and reverted. + test "a form body's other parameters become the data, each through the text rung" do + assert {:ok, event} = + BasicHTTP.decode( + request(%{query: "q=yes", body: "_scxmleventname=e&Var1=2&name=two+words"}) + ) + + assert event.data == %{"q" => "yes", "Var1" => 2, "name" => "two words"} + assert event.type == :external + assert event.origintype == @uri + end + + # sabotage: `form?/1` answers true for every content type -> the text + # body is read as a form, `data` is a one-key map, and the equality + # reddens. Confirmed red and reverted. + test "a body of another content type becomes the data through the text rung" do + assert {:ok, %Event{name: "e", data: 42}} = + BasicHTTP.decode( + request(%{content_type: "text/plain", body: " 42 ", query: "_scxmleventname=e"}) + ) + + assert {:ok, %Event{name: "HTTP.POST", data: "plain words"}} = + BasicHTTP.decode(request(%{content_type: nil, body: "plain words"})) + end + + # sabotage: `decode/1` sets the event's `sendid` from the key's second + # field -> the generated `send_3` reaches `sendid` and the first match + # reddens. Confirmed red and reverted. + test "a well-formed scxml-send-key sets no event field: sendid stays nil" do + assert {:ok, %Event{sendid: nil}} = + BasicHTTP.decode(request(%{send_key: "sess_1/send_3/1/1/0/3/onentry.2.0/1"})) + + assert {:ok, %Event{sendid: nil}} = BasicHTTP.decode(request(%{send_key: nil})) + assert {:ok, %Event{sendid: nil}} = BasicHTTP.decode(request(%{})) + end + + # sabotage: `check_send_key/1` accepts any field count (the `with` pattern + # matches `[_scope, send_id | _rest]`) -> the short key decodes and + # the equality reddens. Confirmed red and reverted. + test "a malformed scxml-send-key is refused" do + assert BasicHTTP.decode(request(%{send_key: "sess_1/x/1"})) == + {:error, {:malformed_send_key, "sess_1/x/1"}} + + assert BasicHTTP.decode(request(%{send_key: "s/%FF/1/1/0/3/transition.1/1"})) == + {:error, {:malformed_send_key, "s/%FF/1/1/0/3/transition.1/1"}} + end + + # sabotage: `post_only/1` answers `:ok` for every method -> a GET + # decodes and the match reddens. Confirmed red and reverted. + test "a method other than POST is refused as method_not_allowed" do + assert BasicHTTP.decode(request(%{method: "GET"})) == + {:error, {:method_not_allowed, "GET"}} + end + + # sabotage: `pairs/2` checks the values only -> the query string's + # non-UTF-8 key decodes and the second equality reddens. Confirmed red + # and reverted. + test "a query string or a body that is not UTF-8 once decoded is refused" do + assert BasicHTTP.decode(request(%{body: "a=%FF"})) == {:error, {:not_utf8, :body}} + assert BasicHTTP.decode(request(%{query: "%FF=1"})) == {:error, {:not_utf8, :query}} + + assert BasicHTTP.decode(request(%{content_type: "text/plain", body: <<0xFF>>})) == + {:error, {:not_utf8, :body}} + end + end +end diff --git a/test/statifier/send/registration_options_test.exs b/test/statifier/send/registration_options_test.exs new file mode 100644 index 00000000..2c414f87 --- /dev/null +++ b/test/statifier/send/registration_options_test.exs @@ -0,0 +1,240 @@ +defmodule Statifier.Send.RegistrationOptionsTest do + use ExUnit.Case, async: true + + # ADR-0075 decisions 3 and 8 (point b): a `:send_types` value may be + # `{module, opts}`. The registered set keeps each type's module and + # options, `_ioprocessors` asks the optional `ioprocessors_entry/2` with + # the session id, the planner hands the options to a `{module, opts}` + # registration's callbacks under `:opts` (a bare module's context is + # unchanged), the recording writes the options as strings, and + # `test_scxml/5` takes a registration. + + import Statifier.Testing.Case, only: [test_scxml: 5] + + alias Statifier.Effect.{Cancel, Send, SendDelayed} + alias Statifier.Evaluator.SystemVariables + alias Statifier.Send.Types + alias Statifier.Session.{Effects, Recording} + + defmodule Echo do + @moduledoc false + @behaviour Statifier.Send.Processor + + @impl Statifier.Send.Processor + def deliver(_effect, _event, ctx), do: {:ok, [{:handler, __MODULE__, {:deliver, ctx}}]} + + @impl Statifier.Send.Processor + def cancel(_cancel, ctx), do: {:ok, [{:handler, __MODULE__, {:cancel, ctx}}]} + + @impl Statifier.Send.Processor + def ioprocessors_entry(type), do: %{"type" => type} + + @impl Statifier.Send.Processor + def ioprocessors_entry("myapp:atom", _context), do: %{"k" => %{bad: 1}} + + def ioprocessors_entry(type, %{session_id: session_id, opts: opts}), + do: %{"at" => "#{Keyword.get(opts, :prefix, "none")}/#{type}/#{session_id}"} + end + + defmodule OnlyOne do + @moduledoc false + @behaviour Statifier.Send.Processor + + @impl Statifier.Send.Processor + def deliver(_effect, _event, _ctx), do: {:ok, []} + + @impl Statifier.Send.Processor + def cancel(_cancel, _ctx), do: {:ok, []} + + @impl Statifier.Send.Processor + def ioprocessors_entry(type), do: %{"only" => type} + end + + defp machine do + {:ok, machine} = + Statifier.compile(""" + + + + """) + + machine + end + + describe "the registered set" do + # sabotage: `from_send_types/1` stores `processors: %{}` -> the equality + # reddens. Confirmed red and reverted. + test "keeps each type's module and options, a bare module's options []" do + assert %Types{processors: processors} = + Types.from_send_types(%{"myapp:a" => {Echo, prefix: "p"}, "myapp:b" => Echo}) + + assert processors == %{"myapp:a" => {Echo, [prefix: "p"]}, "myapp:b" => {Echo, []}} + end + end + + describe "_ioprocessors through ioprocessors_entry/2" do + # sabotage: `session_entry/5` passes `opts: []` whatever the registration + # carries -> the prefix reads "none" and the equality reddens. Confirmed + # red and reverted. + test "a module exporting /2 is asked with the session id and the options" do + types = Types.from_send_types(%{"myapp:a" => {Echo, prefix: "p"}, "myapp:b" => Echo}) + entries = SystemVariables.initial(machine(), "sess_x", types)["_ioprocessors"] + + assert entries["myapp:a"] == %{"at" => "p/myapp:a/sess_x"} + assert entries["myapp:b"] == %{"at" => "none/myapp:b/sess_x"} + end + + # sabotage: `session_entry/5` answers `%{}` for a module without `/2` + # instead of its `/1` entry -> the entry reads `%{}` and the equality + # reddens. Confirmed red and reverted. + test "a module exporting only /1 gets its /1 entry" do + types = Types.from_send_types(%{"myapp:one" => {OnlyOne, x: 1}}) + + assert SystemVariables.initial(machine(), "sess_x", types)["_ioprocessors"]["myapp:one"] == + %{"only" => "myapp:one"} + end + + # sabotage: `Types.session_entry!/3` returns the entry unchecked -> the + # atom-keyed map is accepted and `assert_raise` reddens. Confirmed red + # and reverted. + test "an entry from /2 must be string-keyed" do + types = Types.from_send_types(%{"myapp:atom" => Echo}) + + assert_raise ArgumentError, ~r/string-keyed at every level/, fn -> + SystemVariables.initial(machine(), "sess_x", types) + end + end + end + + describe "the plan context a processor's callbacks receive" do + defp context(send_types) do + %{ + session_id: "sess_x", + invoke_types: nil, + invoke_handlers: %{}, + invocation_types: %{}, + send_types: Types.from_send_types(send_types), + send_processors: send_types + } + end + + defp delayed(type), + do: %SendDelayed{ + event: "e", + type: type, + target: "t", + send_id: "id", + delay_ms: 10, + macrostep: 1, + microstep: 1, + round: 0, + ordinal: 1 + } + + defp cancel, + do: %Cancel{send_id: "id", macrostep: 1, microstep: 2, round: 0, ordinal: 2} + + # sabotage: `processor_for/2` drops the options (answers `{module, + # context}` for a pair too) -> `deliver/3` sees no `:opts` and the match + # reddens. Confirmed red and reverted. + test "a {module, opts} registration adds :opts to deliver/3's and cancel/2's context" do + ctx = context(%{"myapp:a" => {Echo, prefix: "p"}}) + + assert [ + {:notify, _send}, + {:handler, Echo, {:deliver, %{opts: [prefix: "p"]}}}, + {:notify, _cancel}, + {:cancel_timers, "id"}, + {:handler, Echo, {:cancel, %{opts: [prefix: "p"]}}} + ] = Effects.plan([{:send_delayed, delayed("myapp:a")}, {:cancel, cancel()}], ctx) + end + + # sabotage: `processor_for/2` puts `opts: []` for a bare module too -> + # the context gains a key and the equality reddens. Confirmed red and + # reverted. + test "a bare-module registration's context is the plan context unchanged" do + ctx = context(%{"myapp:b" => Echo}) + send = %Send{event: "e", type: "myapp:b", target: "t", macrostep: 1, microstep: 1, round: 0} + + assert [{:notify, _send}, {:handler, Echo, {:deliver, received}}] = + Effects.plan([{:send, send}], ctx) + + assert received == ctx + end + end + + describe "the recording" do + defp recorded(send_types) do + {:ok, machine} = + Statifier.compile(""" + + + + """) + + Recording.new(machine, send_types: send_types) + end + + # sabotage: `option_name/1`'s atom clause writes the atom itself -> the + # blob carries the module atom, and the match on the written options + # reddens. Confirmed red and reverted. + test "writes a registration's options as strings and reads them back" do + send_types = %{"myapp:a" => {Echo, prefix: "p", via: OnlyOne, on: true}, "myapp:b" => Echo} + assert {:ok, blob} = Recording.to_binary(recorded(send_types)) + + {_tag, _version, _chart, opts, _entries, _anchor} = :erlang.binary_to_term(blob) + written = Keyword.fetch!(opts, :send_types)["myapp:a"] + refute match?({Echo, _opts}, written) + assert {_name, [{"prefix", "p"}, {"via", {:atom, _via}}, {"on", true}]} = written + + assert {:ok, decoded} = Recording.from_binary(blob) + assert Keyword.fetch!(Recording.opts(decoded), :send_types) == send_types + end + + # sabotage: `resolve_option/2` keeps an unresolved key without adding it + # to `missing` -> the blob decodes and the match reddens. Confirmed red + # and reverted. + test "an option name this node does not know is reported like a module" do + {:ok, blob} = Recording.to_binary(recorded(%{"myapp:a" => {Echo, prefix: "p"}})) + {tag, version, chart, opts, entries, anchor} = :erlang.binary_to_term(blob) + + doctored = + Keyword.put(opts, :send_types, %{ + "myapp:a" => {Atom.to_string(Echo), [{"no_such_option_name_x", "v"}]} + }) + + blob = :erlang.term_to_binary({tag, version, chart, doctored, entries, anchor}) + + assert Recording.from_binary(blob) == + {:error, {:unknown_handler_modules, ["no_such_option_name_x"]}} + end + end + + describe "Statifier.Testing.Case.test_scxml/5's :send_types option" do + # sabotage: `drive_through_session/4` drops the `:send_types` option -> + # the session registers nothing, the send is unsupported, and the chart + # reaches `fail` on error.execution. Confirmed red and reverted. + test "starts the session with the registration" do + test_scxml( + """ + + + + + + + + + + + + + """, + "a registered send", + ["pass"], + [], + send_types: %{"myapp:only" => OnlyOne} + ) + end + end +end diff --git a/test/support/basic_http_test_transport.ex b/test/support/basic_http_test_transport.ex new file mode 100644 index 00000000..aa2a8e55 --- /dev/null +++ b/test/support/basic_http_test_transport.ex @@ -0,0 +1,28 @@ +defmodule Statifier.BasicHTTPTestTransport do + @moduledoc """ + A test-only `Statifier.Send.BasicHTTP.Transport` that makes no request. + It reports each POST to the process registered under this module's name, + when there is one, as `{:basichttp_post, url, headers, body}`, and answers + from the URL's path: `/answer/` answers `{:ok, status}`, + `/answer/error` answers `{:error, :econnrefused}`, and anything else + `{:ok, 204}`. + """ + + @behaviour Statifier.Send.BasicHTTP.Transport + + @impl Statifier.Send.BasicHTTP.Transport + @spec post(url :: String.t(), headers :: [{String.t(), String.t()}], body :: binary()) :: + {:ok, 100..599} | {:error, term()} + def post(url, headers, body) do + case Process.whereis(__MODULE__) do + nil -> :ok + pid -> send(pid, {:basichttp_post, url, headers, body}) + end + + case URI.parse(url).path do + "/answer/error" -> {:error, :econnrefused} + "/answer/" <> status -> {:ok, String.to_integer(status)} + _other -> {:ok, 204} + end + end +end