fed-sx-types Phase 6: DefineTrigger verb + trigger_registry
Some checks failed
Test, Build, and Deploy / test-build-deploy (push) Failing after 53s
Some checks failed
Test, Build, and Deploy / test-build-deploy (push) Failing after 53s
The trigger declaration layer (fed-sx-triggers-loop.md Phase 1): bind an
activity-type to a durable flow so an arriving activity can fan out into
a business flow.
- next/genesis/activity-types/define_trigger.sx — the DefineTrigger verb
(DefineActivity form, nested-get schema). :object carries
:activity-type, :flow-name, optional :guard / :actor-scope.
- next/kernel/trigger_registry.erl — pure core + registered gen_server,
mirroring peer_actors/peer_types. Keyed by activity-type, multiple
specs per type fire independently. Spec = {TriggerCid, FlowName,
Guard, ActorScope}. Hydrates on start from a fold over DefineTrigger
activities (restart-safe, same content-addressing as define_registry).
Manifest activity-types 7->8 (total bundle 38->39); the four bootstrap
count suites + genesis_parse bumped, and bootstrap_load's internal
timeout raised (the larger bundle's double cid:to_string was truncating).
Tests: define_trigger.sh (6), trigger_registry.sh (17). lib/erlang
771/771 + next/flow 34/34 untouched.
Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
This commit is contained in:
33
next/genesis/activity-types/define_trigger.sx
Normal file
33
next/genesis/activity-types/define_trigger.sx
Normal file
@@ -0,0 +1,33 @@
|
|||||||
|
;; next/genesis/activity-types/define_trigger.sx
|
||||||
|
;;
|
||||||
|
;; Bootstrap definition of the DefineTrigger verb per
|
||||||
|
;; plans/agent-briefings/fed-sx-triggers-loop.md (Phase 1) and
|
||||||
|
;; plans/fed-sx-design.md §13. Read as data by the bundler
|
||||||
|
;; (bootstrap.erl) — never evaluated as code.
|
||||||
|
;;
|
||||||
|
;; DefineTrigger binds an activity-type to a flow. When a matching
|
||||||
|
;; activity is appended to the log, the kernel's trigger fan-out
|
||||||
|
;; (pipeline.erl, post-append) looks the type up in the trigger
|
||||||
|
;; registry and starts the named flow with the activity as input.
|
||||||
|
;; The activity's :object is the binding record:
|
||||||
|
;; {:activity-type "Create" ;; the verb to fire on
|
||||||
|
;; :flow-name "blog-publish-digest"
|
||||||
|
;; :guard <optional predicate> ;; discriminator
|
||||||
|
;; :actor-scope <optional actor id>} ;; default: any
|
||||||
|
;;
|
||||||
|
;; The schema validates the *activity* shape: :object present with
|
||||||
|
;; string :activity-type and :flow-name. The optional :guard lets one
|
||||||
|
;; type bind to multiple flows with discriminators; it is resolved to
|
||||||
|
;; an Erlang predicate at registration time (trigger_registry), not
|
||||||
|
;; carried in the pure-predicate schema here. Schema bodies use nested
|
||||||
|
;; `get` (not keyword-threading) so the predicate is evaluatable.
|
||||||
|
(DefineActivity
|
||||||
|
:name "DefineTrigger"
|
||||||
|
:doc "Bind an activity-type to a flow. :object carries :activity-type, :flow-name, and optional :guard and :actor-scope."
|
||||||
|
:schema (fn
|
||||||
|
(act)
|
||||||
|
(and
|
||||||
|
(not (nil? (get act :object)))
|
||||||
|
(string? (get (get act :object) :activity-type))
|
||||||
|
(string? (get (get act :object) :flow-name))))
|
||||||
|
:semantics (fn (state act) state))
|
||||||
@@ -24,7 +24,8 @@
|
|||||||
"activity-types/announce.sx"
|
"activity-types/announce.sx"
|
||||||
"activity-types/endorse.sx"
|
"activity-types/endorse.sx"
|
||||||
"activity-types/define_type.sx"
|
"activity-types/define_type.sx"
|
||||||
"activity-types/subtype_of.sx")
|
"activity-types/subtype_of.sx"
|
||||||
|
"activity-types/define_trigger.sx")
|
||||||
:object-types ("object-types/sx-artifact.sx"
|
:object-types ("object-types/sx-artifact.sx"
|
||||||
"object-types/note.sx"
|
"object-types/note.sx"
|
||||||
"object-types/tombstone.sx"
|
"object-types/tombstone.sx"
|
||||||
|
|||||||
180
next/kernel/trigger_registry.erl
Normal file
180
next/kernel/trigger_registry.erl
Normal file
@@ -0,0 +1,180 @@
|
|||||||
|
-module(trigger_registry).
|
||||||
|
-export([new/0, add/3, remove/2, lookup/2, all/1, fold/2, fold_fn/0,
|
||||||
|
mk_spec/4, spec_cid/1, spec_flow_name/1, spec_guard/1,
|
||||||
|
spec_actor_scope/1,
|
||||||
|
start_link/0, start_link/1, stop/0,
|
||||||
|
add/2, remove/1, lookup/1, all_triggers/0]).
|
||||||
|
-export([init/1, handle_call/3, handle_cast/2, handle_info/2]).
|
||||||
|
-behaviour(gen_server).
|
||||||
|
|
||||||
|
%% Trigger registry — binds activity-types to durable flows
|
||||||
|
%% (plans/agent-briefings/fed-sx-triggers-loop.md, Phase 1). When an
|
||||||
|
%% activity is appended, the kernel's post-append fan-out
|
||||||
|
%% (pipeline.erl, Phase 2) looks the activity's type up here and starts
|
||||||
|
%% each registered flow. Mirrors the peer_actors / peer_types shape: a
|
||||||
|
%% pure-functional core plus a registered gen_server, hydrated on start
|
||||||
|
%% from a fold over DefineTrigger activities.
|
||||||
|
%%
|
||||||
|
%% State shape (pure-functional):
|
||||||
|
%% [{ActivityType, [Spec, ...]}, ...]
|
||||||
|
%% Multiple triggers may bind the same activity-type; they fire
|
||||||
|
%% independently. A Spec is a 4-tuple:
|
||||||
|
%% {TriggerCid, FlowName, Guard, ActorScope}
|
||||||
|
%% TriggerCid — content-address of the DefineTrigger activity
|
||||||
|
%% (dedup + audit); `undefined` if not yet addressed.
|
||||||
|
%% FlowName — the flow_store-registered flow to start.
|
||||||
|
%% Guard — fun ((Activity, ActorState) -> bool) | undefined.
|
||||||
|
%% Lets one type bind multiple flows with
|
||||||
|
%% discriminators ("only Articles in :newsletter").
|
||||||
|
%% Resolved to a fun at registration; not carried over
|
||||||
|
%% the wire (term_codec can't encode funs).
|
||||||
|
%% ActorScope — an actor id the trigger is scoped to, or `any`.
|
||||||
|
|
||||||
|
%% ── Spec constructor / accessors ────────────────────────────────
|
||||||
|
|
||||||
|
mk_spec(TriggerCid, FlowName, Guard, ActorScope) ->
|
||||||
|
{TriggerCid, FlowName, Guard, ActorScope}.
|
||||||
|
|
||||||
|
spec_cid({Cid, _, _, _}) -> Cid.
|
||||||
|
spec_flow_name({_, FlowName, _, _}) -> FlowName.
|
||||||
|
spec_guard({_, _, Guard, _}) -> Guard.
|
||||||
|
spec_actor_scope({_, _, _, Scope}) -> Scope.
|
||||||
|
|
||||||
|
%% ── Pure-functional API ─────────────────────────────────────────
|
||||||
|
|
||||||
|
new() -> [].
|
||||||
|
|
||||||
|
%% add(ActivityType, Spec, State) — append Spec to ActivityType's list.
|
||||||
|
add(ActivityType, Spec, State) ->
|
||||||
|
Existing = lookup(ActivityType, State),
|
||||||
|
set_keyed(ActivityType, append1(Existing, Spec), State).
|
||||||
|
|
||||||
|
%% remove(TriggerCid, State) — drop every spec carrying TriggerCid,
|
||||||
|
%% across all activity-types; empties are pruned.
|
||||||
|
remove(TriggerCid, State) ->
|
||||||
|
prune([{T, drop_cid(TriggerCid, Specs)} || {T, Specs} <- State]).
|
||||||
|
|
||||||
|
%% lookup(ActivityType, State) — the specs bound to ActivityType ([] if
|
||||||
|
%% none).
|
||||||
|
lookup(ActivityType, State) ->
|
||||||
|
case find_keyed(ActivityType, State) of
|
||||||
|
{ok, Specs} -> Specs;
|
||||||
|
not_found -> []
|
||||||
|
end.
|
||||||
|
|
||||||
|
all(State) -> State.
|
||||||
|
|
||||||
|
%% ── Hydration fold ──────────────────────────────────────────────
|
||||||
|
%%
|
||||||
|
%% fold(Activity, State) — register the binding carried by a
|
||||||
|
%% DefineTrigger activity. Replaying the actor log through this fold
|
||||||
|
%% rebuilds the registry after a restart (same content-addressing
|
||||||
|
%% discipline as define_registry). A non-DefineTrigger activity passes
|
||||||
|
%% through untouched.
|
||||||
|
|
||||||
|
fold(Activity, State) ->
|
||||||
|
case envelope:get_field(type, Activity) of
|
||||||
|
{ok, define_trigger} -> fold_trigger(Activity, State);
|
||||||
|
_ -> State
|
||||||
|
end.
|
||||||
|
|
||||||
|
fold_trigger(Activity, State) ->
|
||||||
|
case envelope:get_field(object, Activity) of
|
||||||
|
{ok, Obj} ->
|
||||||
|
case binding_of(Activity, Obj) of
|
||||||
|
{ok, AType, Spec} -> add(AType, Spec, State);
|
||||||
|
not_a_binding -> State
|
||||||
|
end;
|
||||||
|
_ -> State
|
||||||
|
end.
|
||||||
|
|
||||||
|
binding_of(Activity, Obj) ->
|
||||||
|
case envelope:get_field(activity_type, Obj) of
|
||||||
|
{ok, AType} ->
|
||||||
|
case envelope:get_field(flow_name, Obj) of
|
||||||
|
{ok, FlowName} ->
|
||||||
|
Guard = field_or(guard, Obj, undefined),
|
||||||
|
Scope = field_or(actor_scope, Obj, any),
|
||||||
|
Cid = field_or(id, Activity, undefined),
|
||||||
|
{ok, AType, mk_spec(Cid, FlowName, Guard, Scope)};
|
||||||
|
_ -> not_a_binding
|
||||||
|
end;
|
||||||
|
_ -> not_a_binding
|
||||||
|
end.
|
||||||
|
|
||||||
|
%% fold_fn/0 — a 2-arity fun the projection scheduler can plant.
|
||||||
|
fold_fn() ->
|
||||||
|
fun (Activity, State) -> fold(Activity, State) end.
|
||||||
|
|
||||||
|
%% ── gen_server wrapper ──────────────────────────────────────────
|
||||||
|
|
||||||
|
start_link() ->
|
||||||
|
start_link([]).
|
||||||
|
|
||||||
|
start_link(InitialState) ->
|
||||||
|
Pid = gen_server:start_link(trigger_registry, [InitialState]),
|
||||||
|
erlang:register(trigger_registry, Pid),
|
||||||
|
Pid.
|
||||||
|
|
||||||
|
stop() ->
|
||||||
|
R = gen_server:call(trigger_registry, '$gen_stop'),
|
||||||
|
erlang:unregister(trigger_registry),
|
||||||
|
R.
|
||||||
|
|
||||||
|
add(ActivityType, Spec) ->
|
||||||
|
gen_server:call(trigger_registry, {add, ActivityType, Spec}).
|
||||||
|
|
||||||
|
remove(TriggerCid) ->
|
||||||
|
gen_server:call(trigger_registry, {remove, TriggerCid}).
|
||||||
|
|
||||||
|
lookup(ActivityType) ->
|
||||||
|
gen_server:call(trigger_registry, {lookup, ActivityType}).
|
||||||
|
|
||||||
|
all_triggers() ->
|
||||||
|
gen_server:call(trigger_registry, all_triggers).
|
||||||
|
|
||||||
|
init([InitialState]) ->
|
||||||
|
{ok, InitialState}.
|
||||||
|
|
||||||
|
handle_call({add, ActivityType, Spec}, _From, State) ->
|
||||||
|
{reply, ok, add(ActivityType, Spec, State)};
|
||||||
|
handle_call({remove, TriggerCid}, _From, State) ->
|
||||||
|
{reply, ok, remove(TriggerCid, State)};
|
||||||
|
handle_call({lookup, ActivityType}, _From, State) ->
|
||||||
|
{reply, lookup(ActivityType, State), State};
|
||||||
|
handle_call(all_triggers, _From, State) ->
|
||||||
|
{reply, State, State}.
|
||||||
|
|
||||||
|
handle_cast(_, S) -> {noreply, S}.
|
||||||
|
|
||||||
|
handle_info(_, S) -> {noreply, S}.
|
||||||
|
|
||||||
|
%% ── helpers ─────────────────────────────────────────────────────
|
||||||
|
|
||||||
|
field_or(Key, Proplist, Default) ->
|
||||||
|
case envelope:get_field(Key, Proplist) of
|
||||||
|
{ok, V} -> V;
|
||||||
|
_ -> Default
|
||||||
|
end.
|
||||||
|
|
||||||
|
drop_cid(_, []) -> [];
|
||||||
|
drop_cid(Cid, [Spec | Rest]) ->
|
||||||
|
case spec_cid(Spec) of
|
||||||
|
Cid -> drop_cid(Cid, Rest);
|
||||||
|
_ -> [Spec | drop_cid(Cid, Rest)]
|
||||||
|
end.
|
||||||
|
|
||||||
|
prune([]) -> [];
|
||||||
|
prune([{_, []} | Rest]) -> prune(Rest);
|
||||||
|
prune([P | Rest]) -> [P | prune(Rest)].
|
||||||
|
|
||||||
|
append1([], X) -> [X];
|
||||||
|
append1([H | T], X) -> [H | append1(T, X)].
|
||||||
|
|
||||||
|
find_keyed(_, []) -> not_found;
|
||||||
|
find_keyed(K, [{K, V} | _]) -> {ok, V};
|
||||||
|
find_keyed(K, [_ | Rest]) -> find_keyed(K, Rest).
|
||||||
|
|
||||||
|
set_keyed(K, V, []) -> [{K, V}];
|
||||||
|
set_keyed(K, V, [{K, _} | Rest]) -> [{K, V} | Rest];
|
||||||
|
set_keyed(K, V, [P | Rest]) -> [P | set_keyed(K, V, Rest)].
|
||||||
@@ -79,7 +79,7 @@ cat > "$TMPFILE" <<'EPOCHS'
|
|||||||
(eval "(get (erlang-eval-ast \"R = bootstrap:read_genesis(), {ok, S1} = bootstrap:load_genesis(R), {ok, S2} = bootstrap:load_genesis(R), cid:to_string(S1) =:= cid:to_string(S2)\") :name)")
|
(eval "(get (erlang-eval-ast \"R = bootstrap:read_genesis(), {ok, S1} = bootstrap:load_genesis(R), {ok, S2} = bootstrap:load_genesis(R), cid:to_string(S1) =:= cid:to_string(S2)\") :name)")
|
||||||
EPOCHS
|
EPOCHS
|
||||||
|
|
||||||
OUTPUT=$(timeout 300 "$SX_SERVER" < "$TMPFILE" 2>/dev/null)
|
OUTPUT=$(timeout 590 "$SX_SERVER" < "$TMPFILE" 2>/dev/null)
|
||||||
|
|
||||||
check() {
|
check() {
|
||||||
local epoch="$1" desc="$2" expected="$3"
|
local epoch="$1" desc="$2" expected="$3"
|
||||||
@@ -106,7 +106,7 @@ check 10 "strip suffix create.sx -> create" "true"
|
|||||||
check 11 "strip suffix hello unchanged" "true"
|
check 11 "strip suffix hello unchanged" "true"
|
||||||
check 12 "strip suffix .sx -> empty" "true"
|
check 12 "strip suffix .sx -> empty" "true"
|
||||||
check 13 "load_genesis rejects bad shape" "ok"
|
check 13 "load_genesis rejects bad shape" "ok"
|
||||||
check 20 "loaded activity_types count = 7" "7"
|
check 20 "loaded activity_types count = 8" "8"
|
||||||
check 21 "loaded object_types count = 13" "13"
|
check 21 "loaded object_types count = 13" "13"
|
||||||
check 22 "loaded projections count = 7" "7"
|
check 22 "loaded projections count = 7" "7"
|
||||||
check 23 "loaded validators count = 3" "3"
|
check 23 "loaded validators count = 3" "3"
|
||||||
|
|||||||
@@ -99,8 +99,8 @@ check() {
|
|||||||
check 2 "gen_server loaded" "gen_server"
|
check 2 "gen_server loaded" "gen_server"
|
||||||
check 3 "registry loaded" "registry"
|
check 3 "registry loaded" "registry"
|
||||||
check 4 "bootstrap loaded" "bootstrap"
|
check 4 "bootstrap loaded" "bootstrap"
|
||||||
check 10 "populate returns total 38" "38"
|
check 10 "populate returns total 39" "39"
|
||||||
check 20 "activity_types count = 7" "7"
|
check 20 "activity_types count = 8" "8"
|
||||||
check 21 "object_types count = 13" "13"
|
check 21 "object_types count = 13" "13"
|
||||||
check 22 "projections count = 7" "7"
|
check 22 "projections count = 7" "7"
|
||||||
check 23 "validators count = 3" "3"
|
check 23 "validators count = 3" "3"
|
||||||
|
|||||||
@@ -102,7 +102,7 @@ check 10 "sections/0 length" "7"
|
|||||||
check 11 "ends_with_sx create.sx" "true"
|
check 11 "ends_with_sx create.sx" "true"
|
||||||
check 12 "ends_with_sx hello" "false"
|
check 12 "ends_with_sx hello" "false"
|
||||||
check 13 "ends_with_sx empty" "false"
|
check 13 "ends_with_sx empty" "false"
|
||||||
check 20 "section activity_types count" "7"
|
check 20 "section activity_types count" "8"
|
||||||
check 21 "section object_types count" "13"
|
check 21 "section object_types count" "13"
|
||||||
check 22 "section projections count" "7"
|
check 22 "section projections count" "7"
|
||||||
check 23 "section validators count" "3"
|
check 23 "section validators count" "3"
|
||||||
@@ -111,7 +111,7 @@ check 25 "section sig_suites count" "2"
|
|||||||
check 26 "section audience count" "3"
|
check 26 "section audience count" "3"
|
||||||
check 30 "read_genesis returns 7 sections" "7"
|
check 30 "read_genesis returns 7 sections" "7"
|
||||||
check 31 "first section name" "activity_types"
|
check 31 "first section name" "activity_types"
|
||||||
check 32 "first section entry count" "7"
|
check 32 "first section entry count" "8"
|
||||||
|
|
||||||
TOTAL=$((PASS+FAIL))
|
TOTAL=$((PASS+FAIL))
|
||||||
if [ $FAIL -eq 0 ]; then
|
if [ $FAIL -eq 0 ]; then
|
||||||
|
|||||||
@@ -121,10 +121,10 @@ check() {
|
|||||||
|
|
||||||
check 10 "bootstrap module loaded" "bootstrap"
|
check 10 "bootstrap module loaded" "bootstrap"
|
||||||
check 20 "whereis(nx_kernel) is Pid" "true"
|
check 20 "whereis(nx_kernel) is Pid" "true"
|
||||||
check 21 "activity_types count = 7" "7"
|
check 21 "activity_types count = 8" "8"
|
||||||
check 22 "object_types count = 13" "13"
|
check 22 "object_types count = 13" "13"
|
||||||
check 23 "projections count = 7" "7"
|
check 23 "projections count = 7" "7"
|
||||||
check 24 "total entries = 38" "38"
|
check 24 "total entries = 39" "39"
|
||||||
check 25 "fresh log_tip = 0" "0"
|
check 25 "fresh log_tip = 0" "0"
|
||||||
check 26 "publish advances tip to 1" "1"
|
check 26 "publish advances tip to 1" "1"
|
||||||
check 27 "actor_id = alice" "true"
|
check 27 "actor_id = alice" "true"
|
||||||
|
|||||||
99
next/tests/define_trigger.sh
Executable file
99
next/tests/define_trigger.sh
Executable file
@@ -0,0 +1,99 @@
|
|||||||
|
#!/usr/bin/env bash
|
||||||
|
# next/tests/define_trigger.sh — fed-sx triggers Phase 1 (verb).
|
||||||
|
#
|
||||||
|
# The DefineTrigger genesis verb
|
||||||
|
# (next/genesis/activity-types/define_trigger.sx) binds an activity-type
|
||||||
|
# to a flow. This suite confirms it parses with the expected
|
||||||
|
# DefineActivity head + :name, that its :schema accepts a well-formed
|
||||||
|
# binding and rejects malformed ones, and that a DefineTrigger envelope
|
||||||
|
# round-trips through term_codec.
|
||||||
|
|
||||||
|
set -uo pipefail
|
||||||
|
cd "$(git rev-parse --show-toplevel)"
|
||||||
|
|
||||||
|
SX_SERVER="${SX_SERVER:-hosts/ocaml/_build/default/bin/sx_server.exe}"
|
||||||
|
if [ ! -x "$SX_SERVER" ]; then
|
||||||
|
SX_SERVER="/root/rose-ash/hosts/ocaml/_build/default/bin/sx_server.exe"
|
||||||
|
fi
|
||||||
|
if [ ! -x "$SX_SERVER" ]; then
|
||||||
|
echo "ERROR: sx_server.exe not found." >&2
|
||||||
|
exit 1
|
||||||
|
fi
|
||||||
|
|
||||||
|
VERBOSE="${1:-}"
|
||||||
|
PASS=0; FAIL=0; ERRORS=""
|
||||||
|
TMPFILE=$(mktemp); trap "rm -f $TMPFILE" EXIT
|
||||||
|
|
||||||
|
SCH='(eval-expr (get (apply dict (rest (parse (file-read \"next/genesis/activity-types/define_trigger.sx\")))) :schema))'
|
||||||
|
|
||||||
|
cat > "$TMPFILE" <<EPOCHS
|
||||||
|
(epoch 1)
|
||||||
|
(load "lib/erlang/tokenizer.sx")
|
||||||
|
(load "lib/erlang/parser.sx")
|
||||||
|
(load "lib/erlang/parser-core.sx")
|
||||||
|
(load "lib/erlang/parser-expr.sx")
|
||||||
|
(load "lib/erlang/parser-module.sx")
|
||||||
|
(load "lib/erlang/transpile.sx")
|
||||||
|
(load "lib/erlang/runtime.sx")
|
||||||
|
(load "lib/erlang/vm/dispatcher.sx")
|
||||||
|
(epoch 2)
|
||||||
|
(eval "(get (erlang-load-module (file-read \"next/kernel/term_codec.erl\")) :name)")
|
||||||
|
|
||||||
|
;; ── parse / shape ──────────────────────────────────────────
|
||||||
|
(epoch 10)
|
||||||
|
(eval "(first (parse (file-read \"next/genesis/activity-types/define_trigger.sx\")))")
|
||||||
|
(epoch 11)
|
||||||
|
(eval "(get (apply dict (rest (parse (file-read \"next/genesis/activity-types/define_trigger.sx\")))) :name)")
|
||||||
|
|
||||||
|
;; ── schema accept / reject ─────────────────────────────────
|
||||||
|
;; valid binding: string :activity-type + :flow-name -> true
|
||||||
|
(epoch 20)
|
||||||
|
(eval "(define sch ${SCH}) (sch (dict :object (dict :activity-type \"Create\" :flow-name \"blog-publish-digest\")))")
|
||||||
|
;; reject: missing :activity-type -> false
|
||||||
|
(epoch 21)
|
||||||
|
(eval "(define sch ${SCH}) (sch (dict :object (dict :flow-name \"f\")))")
|
||||||
|
;; reject: missing :flow-name -> false
|
||||||
|
(epoch 22)
|
||||||
|
(eval "(define sch ${SCH}) (sch (dict :object (dict :activity-type \"Create\")))")
|
||||||
|
|
||||||
|
;; ── envelope round-trip through term_codec ─────────────────
|
||||||
|
(epoch 30)
|
||||||
|
(eval "(get (erlang-eval-ast \"A = [{type, define_trigger}, {actor, alice}, {object, [{activity_type, create}, {flow_name, blog_publish_digest}]}], {ok, D, _} = term_codec:decode(term_codec:encode(A)), D =:= A\") :name)")
|
||||||
|
EPOCHS
|
||||||
|
|
||||||
|
OUTPUT=$(timeout 180 "$SX_SERVER" < "$TMPFILE" 2>/dev/null)
|
||||||
|
|
||||||
|
check() {
|
||||||
|
local epoch="$1" desc="$2" expected="$3"
|
||||||
|
local actual
|
||||||
|
actual=$(echo "$OUTPUT" | awk -v e="$epoch" '
|
||||||
|
$0 ~ "^\\(ok-len " e " " { getline; print; exit }
|
||||||
|
$0 ~ "^\\(ok " e " " { print; exit }
|
||||||
|
$0 ~ "^\\(error " e " " { print; exit }
|
||||||
|
')
|
||||||
|
[ -z "$actual" ] && actual="<no output for epoch $epoch>"
|
||||||
|
if echo "$actual" | grep -qF -- "$expected"; then
|
||||||
|
PASS=$((PASS+1))
|
||||||
|
[ "$VERBOSE" = "-v" ] && echo " ok $desc"
|
||||||
|
else
|
||||||
|
FAIL=$((FAIL+1))
|
||||||
|
ERRORS+=" FAIL [$desc] (epoch $epoch) expected: $expected | actual: $actual
|
||||||
|
"
|
||||||
|
fi
|
||||||
|
}
|
||||||
|
|
||||||
|
check 10 "define_trigger.sx head form" "DefineActivity"
|
||||||
|
check 11 "define_trigger.sx name" "DefineTrigger"
|
||||||
|
check 20 "schema accepts valid binding" "true"
|
||||||
|
check 21 "schema rejects missing type" "false"
|
||||||
|
check 22 "schema rejects missing flow-name" "false"
|
||||||
|
check 30 "DefineTrigger envelope round-trips" "true"
|
||||||
|
|
||||||
|
TOTAL=$((PASS+FAIL))
|
||||||
|
if [ $FAIL -eq 0 ]; then
|
||||||
|
echo "ok $PASS/$TOTAL next/tests/define_trigger.sh passed"
|
||||||
|
else
|
||||||
|
echo "FAIL $PASS/$TOTAL passed, $FAIL failed:"
|
||||||
|
echo "$ERRORS"
|
||||||
|
fi
|
||||||
|
[ $FAIL -eq 0 ]
|
||||||
@@ -56,6 +56,10 @@ cat > "$TMPFILE" <<'EPOCHS'
|
|||||||
(eval "(first (parse (file-read \"next/genesis/activity-types/subtype_of.sx\")))")
|
(eval "(first (parse (file-read \"next/genesis/activity-types/subtype_of.sx\")))")
|
||||||
(epoch 204)
|
(epoch 204)
|
||||||
(eval "(get (apply dict (rest (parse (file-read \"next/genesis/activity-types/subtype_of.sx\")))) :name)")
|
(eval "(get (apply dict (rest (parse (file-read \"next/genesis/activity-types/subtype_of.sx\")))) :name)")
|
||||||
|
(epoch 205)
|
||||||
|
(eval "(first (parse (file-read \"next/genesis/activity-types/define_trigger.sx\")))")
|
||||||
|
(epoch 206)
|
||||||
|
(eval "(get (apply dict (rest (parse (file-read \"next/genesis/activity-types/define_trigger.sx\")))) :name)")
|
||||||
(epoch 19)
|
(epoch 19)
|
||||||
(eval "(len (get (apply dict (rest (parse (file-read \"next/genesis/manifest.sx\")))) :activity-types))")
|
(eval "(len (get (apply dict (rest (parse (file-read \"next/genesis/manifest.sx\")))) :activity-types))")
|
||||||
(epoch 30)
|
(epoch 30)
|
||||||
@@ -192,7 +196,9 @@ check 201 "define_type.sx head form" "DefineActivity"
|
|||||||
check 202 "define_type.sx name" "DefineType"
|
check 202 "define_type.sx name" "DefineType"
|
||||||
check 203 "subtype_of.sx head form" "DefineActivity"
|
check 203 "subtype_of.sx head form" "DefineActivity"
|
||||||
check 204 "subtype_of.sx name" "SubtypeOf"
|
check 204 "subtype_of.sx name" "SubtypeOf"
|
||||||
check 19 "manifest has 7 activity-types" "7"
|
check 205 "define_trigger.sx head form" "DefineActivity"
|
||||||
|
check 206 "define_trigger.sx name" "DefineTrigger"
|
||||||
|
check 19 "manifest has 8 activity-types" "8"
|
||||||
check 30 "sx-artifact.sx head form" "DefineObject"
|
check 30 "sx-artifact.sx head form" "DefineObject"
|
||||||
check 31 "sx-artifact.sx name" "SXArtifact"
|
check 31 "sx-artifact.sx name" "SXArtifact"
|
||||||
check 32 "note.sx name" "Note"
|
check 32 "note.sx name" "Note"
|
||||||
|
|||||||
143
next/tests/trigger_registry.sh
Executable file
143
next/tests/trigger_registry.sh
Executable file
@@ -0,0 +1,143 @@
|
|||||||
|
#!/usr/bin/env bash
|
||||||
|
# next/tests/trigger_registry.sh — fed-sx triggers Phase 1 (registry).
|
||||||
|
#
|
||||||
|
# trigger_registry binds activity-types to durable flows. The kernel's
|
||||||
|
# post-append fan-out (Phase 2) looks an arriving activity's type up
|
||||||
|
# here and starts each registered flow. Mirrors peer_actors / peer_types:
|
||||||
|
# a pure core + a gen_server, hydrated from a fold over DefineTrigger
|
||||||
|
# activities.
|
||||||
|
|
||||||
|
set -uo pipefail
|
||||||
|
cd "$(git rev-parse --show-toplevel)"
|
||||||
|
|
||||||
|
SX_SERVER="${SX_SERVER:-hosts/ocaml/_build/default/bin/sx_server.exe}"
|
||||||
|
if [ ! -x "$SX_SERVER" ]; then
|
||||||
|
SX_SERVER="/root/rose-ash/hosts/ocaml/_build/default/bin/sx_server.exe"
|
||||||
|
fi
|
||||||
|
if [ ! -x "$SX_SERVER" ]; then
|
||||||
|
echo "ERROR: sx_server.exe not found." >&2
|
||||||
|
exit 1
|
||||||
|
fi
|
||||||
|
|
||||||
|
VERBOSE="${1:-}"
|
||||||
|
PASS=0; FAIL=0; ERRORS=""
|
||||||
|
TMPFILE=$(mktemp); trap "rm -f $TMPFILE" EXIT
|
||||||
|
|
||||||
|
# Spec1/Spec2 bind activity-type `create`. TrigAct/TrigAct2 are
|
||||||
|
# DefineTrigger activities the fold hydrates from.
|
||||||
|
SETUP='S1 = trigger_registry:mk_spec(<<99,49>>, flow_a, undefined, any), S2 = trigger_registry:mk_spec(<<99,50>>, flow_b, undefined, any), TrigAct = [{type, define_trigger}, {actor, alice}, {id, <<99,49>>}, {object, [{activity_type, create}, {flow_name, flow_a}]}], TrigAct2 = [{type, define_trigger}, {actor, alice}, {id, <<99,50>>}, {object, [{activity_type, follow}, {flow_name, flow_c}]}], Note = [{type, note}, {actor, alice}, {object, [{content, hi}]}],'
|
||||||
|
|
||||||
|
cat > "$TMPFILE" <<EPOCHS
|
||||||
|
(epoch 1)
|
||||||
|
(load "lib/erlang/tokenizer.sx")
|
||||||
|
(load "lib/erlang/parser.sx")
|
||||||
|
(load "lib/erlang/parser-core.sx")
|
||||||
|
(load "lib/erlang/parser-expr.sx")
|
||||||
|
(load "lib/erlang/parser-module.sx")
|
||||||
|
(load "lib/erlang/transpile.sx")
|
||||||
|
(load "lib/erlang/runtime.sx")
|
||||||
|
(load "lib/erlang/vm/dispatcher.sx")
|
||||||
|
(epoch 2)
|
||||||
|
(eval "(er-load-gen-server!)")
|
||||||
|
(epoch 3)
|
||||||
|
(eval "(get (erlang-load-module (file-read \"next/kernel/envelope.erl\")) :name)")
|
||||||
|
(epoch 4)
|
||||||
|
(eval "(get (erlang-load-module (file-read \"next/kernel/trigger_registry.erl\")) :name)")
|
||||||
|
|
||||||
|
;; ── pure core ──────────────────────────────────────────────
|
||||||
|
(epoch 10)
|
||||||
|
(eval "(get (erlang-eval-ast \"trigger_registry:new() =:= []\") :name)")
|
||||||
|
;; add + lookup round-trip
|
||||||
|
(epoch 11)
|
||||||
|
(eval "(get (erlang-eval-ast \"${SETUP} St = trigger_registry:add(create, S1, trigger_registry:new()), trigger_registry:lookup(create, St) =:= [S1]\") :name)")
|
||||||
|
;; lookup with no match -> []
|
||||||
|
(epoch 12)
|
||||||
|
(eval "(get (erlang-eval-ast \"${SETUP} trigger_registry:lookup(create, trigger_registry:new()) =:= []\") :name)")
|
||||||
|
;; multi-bind: two specs on the same activity-type, both returned in order
|
||||||
|
(epoch 13)
|
||||||
|
(eval "(get (erlang-eval-ast \"${SETUP} St = trigger_registry:add(create, S2, trigger_registry:add(create, S1, trigger_registry:new())), trigger_registry:lookup(create, St) =:= [S1, S2]\") :name)")
|
||||||
|
;; remove by trigger cid
|
||||||
|
(epoch 14)
|
||||||
|
(eval "(get (erlang-eval-ast \"${SETUP} St = trigger_registry:add(create, S2, trigger_registry:add(create, S1, trigger_registry:new())), trigger_registry:lookup(create, trigger_registry:remove(<<99,49>>, St)) =:= [S2]\") :name)")
|
||||||
|
;; remove last spec for a type prunes the type
|
||||||
|
(epoch 15)
|
||||||
|
(eval "(get (erlang-eval-ast \"${SETUP} St = trigger_registry:add(create, S1, trigger_registry:new()), trigger_registry:remove(<<99,49>>, St) =:= []\") :name)")
|
||||||
|
;; spec accessors
|
||||||
|
(epoch 16)
|
||||||
|
(eval "(get (erlang-eval-ast \"${SETUP} {trigger_registry:spec_cid(S1), trigger_registry:spec_flow_name(S1), trigger_registry:spec_guard(S1), trigger_registry:spec_actor_scope(S1)} =:= {<<99,49>>, flow_a, undefined, any}\") :name)")
|
||||||
|
|
||||||
|
;; ── hydration fold ─────────────────────────────────────────
|
||||||
|
;; a DefineTrigger activity registers its binding
|
||||||
|
(epoch 20)
|
||||||
|
(eval "(get (erlang-eval-ast \"${SETUP} St = trigger_registry:fold(TrigAct, trigger_registry:new()), trigger_registry:lookup(create, St) =:= [trigger_registry:mk_spec(<<99,49>>, flow_a, undefined, any)]\") :name)")
|
||||||
|
;; a non-trigger activity passes through untouched
|
||||||
|
(epoch 21)
|
||||||
|
(eval "(get (erlang-eval-ast \"${SETUP} trigger_registry:fold(Note, trigger_registry:new()) =:= []\") :name)")
|
||||||
|
;; folding several Trigger activities rebuilds the whole registry
|
||||||
|
(epoch 22)
|
||||||
|
(eval "(get (erlang-eval-ast \"${SETUP} St = trigger_registry:fold(TrigAct2, trigger_registry:fold(TrigAct, trigger_registry:new())), {trigger_registry:lookup(create, St), trigger_registry:lookup(follow, St)} =:= {[trigger_registry:mk_spec(<<99,49>>, flow_a, undefined, any)], [trigger_registry:mk_spec(<<99,50>>, flow_c, undefined, any)]}\") :name)")
|
||||||
|
;; fold_fn/0 is a 2-arity fun
|
||||||
|
(epoch 23)
|
||||||
|
(eval "(get (erlang-eval-ast \"is_function(trigger_registry:fold_fn(), 2)\") :name)")
|
||||||
|
|
||||||
|
;; ── gen_server ─────────────────────────────────────────────
|
||||||
|
(epoch 30)
|
||||||
|
(eval "(get (erlang-eval-ast \"${SETUP} trigger_registry:start_link(), trigger_registry:add(create, S1), trigger_registry:lookup(create) =:= [S1]\") :name)")
|
||||||
|
(epoch 31)
|
||||||
|
(eval "(get (erlang-eval-ast \"trigger_registry:start_link(), trigger_registry:lookup(create) =:= []\") :name)")
|
||||||
|
(epoch 32)
|
||||||
|
(eval "(get (erlang-eval-ast \"${SETUP} trigger_registry:start_link(), trigger_registry:add(create, S1), trigger_registry:add(create, S2), trigger_registry:remove(<<99,49>>), trigger_registry:lookup(create) =:= [S2]\") :name)")
|
||||||
|
(epoch 33)
|
||||||
|
(eval "(get (erlang-eval-ast \"${SETUP} trigger_registry:start_link(), trigger_registry:add(create, S1), trigger_registry:add(follow, S2), trigger_registry:all_triggers() =:= [{create, [S1]}, {follow, [S2]}]\") :name)")
|
||||||
|
;; start_link/1 pre-populates from a hydrated state
|
||||||
|
(epoch 34)
|
||||||
|
(eval "(get (erlang-eval-ast \"${SETUP} St = trigger_registry:fold(TrigAct, trigger_registry:new()), trigger_registry:start_link(St), trigger_registry:lookup(create) =:= [trigger_registry:mk_spec(<<99,49>>, flow_a, undefined, any)]\") :name)")
|
||||||
|
EPOCHS
|
||||||
|
|
||||||
|
OUTPUT=$(timeout 300 "$SX_SERVER" < "$TMPFILE" 2>/dev/null)
|
||||||
|
|
||||||
|
check() {
|
||||||
|
local epoch="$1" desc="$2" expected="$3"
|
||||||
|
local actual
|
||||||
|
actual=$(echo "$OUTPUT" | awk -v e="$epoch" '
|
||||||
|
$0 ~ "^\\(ok-len " e " " { getline; print; exit }
|
||||||
|
$0 ~ "^\\(ok " e " " { print; exit }
|
||||||
|
$0 ~ "^\\(error " e " " { print; exit }
|
||||||
|
')
|
||||||
|
[ -z "$actual" ] && actual="<no output for epoch $epoch>"
|
||||||
|
if echo "$actual" | grep -qF -- "$expected"; then
|
||||||
|
PASS=$((PASS+1))
|
||||||
|
[ "$VERBOSE" = "-v" ] && echo " ok $desc"
|
||||||
|
else
|
||||||
|
FAIL=$((FAIL+1))
|
||||||
|
ERRORS+=" FAIL [$desc] (epoch $epoch) expected: $expected | actual: $actual
|
||||||
|
"
|
||||||
|
fi
|
||||||
|
}
|
||||||
|
|
||||||
|
check 4 "trigger_registry module loaded" "trigger_registry"
|
||||||
|
check 10 "new/0 -> []" "true"
|
||||||
|
check 11 "add + lookup round-trip" "true"
|
||||||
|
check 12 "lookup no match -> []" "true"
|
||||||
|
check 13 "multi-bind same type, ordered" "true"
|
||||||
|
check 14 "remove by trigger cid" "true"
|
||||||
|
check 15 "remove last prunes the type" "true"
|
||||||
|
check 16 "spec accessors" "true"
|
||||||
|
check 20 "fold registers a binding" "true"
|
||||||
|
check 21 "fold non-trigger passes through" "true"
|
||||||
|
check 22 "fold hydration rebuilds registry" "true"
|
||||||
|
check 23 "fold_fn/0 is fun/2" "true"
|
||||||
|
check 30 "gen_server add + lookup" "true"
|
||||||
|
check 31 "gen_server lookup no match -> []" "true"
|
||||||
|
check 32 "gen_server remove" "true"
|
||||||
|
check 33 "gen_server all_triggers" "true"
|
||||||
|
check 34 "start_link/1 pre-populates" "true"
|
||||||
|
|
||||||
|
TOTAL=$((PASS+FAIL))
|
||||||
|
if [ $FAIL -eq 0 ]; then
|
||||||
|
echo "ok $PASS/$TOTAL next/tests/trigger_registry.sh passed"
|
||||||
|
else
|
||||||
|
echo "FAIL $PASS/$TOTAL passed, $FAIL failed:"
|
||||||
|
echo "$ERRORS"
|
||||||
|
fi
|
||||||
|
[ $FAIL -eq 0 ]
|
||||||
Reference in New Issue
Block a user