Skip to main content

aion_server/worker/
declared_body.rs

1//! Server-side execution of declared action bodies.
2//!
3//! An action whose deployed contract carries an [`ActionBodyContract`] is
4//! executed BY THE SERVER, with no connected worker: the dispatch is
5//! intercepted at the [`ActivityDispatcher`] seam before task-queue routing,
6//! the declared command runs through the worker SDK's own executor
7//! ([`aion_worker::shell::ShellAction`] — argv-element substitution, no
8//! shell, process-group containment), and the result flows back through the
9//! engine's normal completion path. The engine still schedules, records, and
10//! replays the activity exactly as if a worker had served it.
11//!
12//! Actions with no declared body are delegated to the wrapped production
13//! dispatcher unchanged, so remote workers keep working exactly as before.
14//!
15//! # The command is readable while it runs
16//!
17//! The executing activity carries a live transcript seam
18//! ([`ActivityContext::with_transcript`](aion_worker::ActivityContext::with_transcript)),
19//! so every line the command writes to stdout or stderr is published onto the
20//! server's transcript sequencer AS IT ARRIVES — the same stream, envelope, and
21//! cursor reads an agent step's transcript uses (see
22//! [`super::declared_body_transcript`]). The activity's recorded result is
23//! untouched by this: it still carries the command's complete output.
24//!
25//! # The command can be stopped, by its bound and by its run
26//!
27//! Two things end a server-executed command early, and both reach the same
28//! cancellation the worker path already acts on — `SIGTERM` → grace →
29//! `SIGKILL` across the whole process group, with the verdict withheld until
30//! the group has been proven gone.
31//!
32//! The first is the attempt's own deadline. A dispatch carrying an authored
33//! per-attempt timeout (#223) ends its command at that bound HERE, in the
34//! server, where the process is — see [`run_bounded`]. The engine's own
35//! deadline stops the run WAITING and cannot reach a process, which is the
36//! right division of labour for a remote worker and no division at all for a
37//! body the server itself started. A dispatch that authored no bound is
38//! unbounded, exactly as before: the server adds no deadline of its own.
39//!
40//! The second is the run being cancelled. Every executing attempt registers in
41//! [`super::declared_body_cancel::DeclaredCommandAttempts`] for exactly as long
42//! as its command runs, which is how the cancel path reaches an activity no
43//! worker holds and no heartbeat tracks. An attempt that cannot register is
44//! refused rather than run: a command a cancelled run could not stop is the
45//! defect the registration exists to prevent.
46
47use std::collections::BTreeMap;
48use std::sync::{Arc, OnceLock};
49
50use aion::{ActivityDispatch, ActivityDispatcher};
51use aion_package::{ActionBodyContract, ContentHash};
52use aion_worker::shell::ShellAction;
53
54use super::declared_body_ambiguity::{DeclaringVersion, ambiguous_body_refusal};
55use super::declared_body_cancel::DeclaredCommandAttempts;
56use super::declared_body_selection::select_declared_body;
57use super::declared_body_transcript::publish_declared_transcript;
58use super::workspace_root::{WORKSPACE_ROOT_PLACEHOLDER, WorkspaceRoot};
59use crate::activity_publisher::ActivityEventPublisher;
60
61/// What a declared-body lookup found for one `(task_queue, action)` address.
62#[derive(Clone, Debug)]
63pub enum DeclaredBodyLookup {
64    /// No retained contract declares a body for this action — it is a
65    /// requirement on an out-of-band worker and must be delegated.
66    None,
67    /// Exactly one distinct body is declared across every retained package
68    /// version. Safe to execute.
69    Declared(ActionBodyContract),
70    /// Retained package versions declare DIFFERENT bodies for this action.
71    /// Executing one of them would guess which deploy the running workflow
72    /// meant, so the dispatch is refused by name instead.
73    Ambiguous {
74        /// Every retained version that declares a body for this action, in
75        /// catalog order. Carried rather than counted because the refusal has
76        /// to name the versions the operator must retire — a bare count leaves
77        /// them holding a terminal error with no way to act on it.
78        declaring: Vec<DeclaringVersion>,
79    },
80    /// The catalog could not be read. The reader reports why; the dispatch
81    /// is delegated so a readable worker path can still serve it.
82    Unreadable(String),
83}
84
85/// Which run a declared-body lookup is being made for.
86///
87/// A body is a property of the run's own package version, not of the queue, so
88/// the lookup cannot answer correctly without knowing whose dispatch it is —
89/// see [`super::declared_body_selection`].
90#[derive(Clone, Copy, Debug)]
91pub struct DispatchingRun<'a> {
92    /// The workflow the activity belongs to.
93    pub workflow_id: &'a aion_core::WorkflowId,
94    /// The concrete run within that workflow.
95    pub run_id: &'a aion_core::RunId,
96}
97
98/// A reader over the deployed contracts' declared action bodies.
99pub trait DeclaredBodies: Send + Sync {
100    /// Look up the declared body for `action` on `task_queue`, as the run
101    /// issuing the dispatch sees it.
102    fn body_for(
103        &self,
104        task_queue: &str,
105        action: &str,
106        run: DispatchingRun<'_>,
107    ) -> DeclaredBodyLookup;
108}
109
110/// Shared, install-once handle the dispatcher holds from construction and the
111/// boot path fills in once the engine exists.
112///
113/// Mirrors [`super::QueueDeclarationSource`]: the dispatcher is built before
114/// the engine, so the seam it consults is handed over afterwards through a
115/// clone of this handle rather than by rebuilding the dispatcher.
116#[derive(Clone, Default)]
117pub struct DeclaredBodySource {
118    inner: Arc<OnceLock<Arc<dyn DeclaredBodies>>>,
119}
120
121impl std::fmt::Debug for DeclaredBodySource {
122    fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
123        formatter
124            .debug_struct("DeclaredBodySource")
125            .field("installed", &self.inner.get().is_some())
126            .finish()
127    }
128}
129
130impl DeclaredBodySource {
131    /// Install the reader. A second install is ignored and logged: the source
132    /// is process-wide and must not silently change identity.
133    pub fn install(&self, source: Arc<dyn DeclaredBodies>) {
134        if self.inner.set(source).is_err() {
135            tracing::warn!("declared body source already installed; ignoring duplicate set");
136        }
137    }
138
139    /// Look up the declared body, or [`DeclaredBodyLookup::None`] when no
140    /// reader is installed yet.
141    ///
142    /// An uninstalled consult is stated at ERROR before delegating to the
143    /// worker path, never silently: a dispatch can only reach this seam from
144    /// a live run, and a live run's deploy is durable — so "nothing installed
145    /// yet" is a boot-ordering defect, not an empty catalog. This exact
146    /// silence was #266 Defect A: startup recovery replay re-dispatched
147    /// adopted in-flight declared-body activities before
148    /// `install_engine_backed_seams` filled this source, and every one fell
149    /// through here to a queue with no pollers and parked forever. The fix
150    /// (deferred startup recovery) removes the caller; this arm stays loud so
151    /// any future pre-install dispatch path names itself in the log instead
152    /// of stranding runs silently.
153    #[must_use]
154    pub fn body_for(
155        &self,
156        task_queue: &str,
157        action: &str,
158        run: DispatchingRun<'_>,
159    ) -> DeclaredBodyLookup {
160        self.inner.get().map_or_else(
161            || {
162                tracing::error!(
163                    operation = "declared_command_dispatch",
164                    task_queue,
165                    action,
166                    workflow_id = %run.workflow_id,
167                    run_id = %run.run_id,
168                    "declared body source consulted before it was installed; the dispatch \
169                     falls through to the worker path and will park if the queue's only \
170                     service is its declared bodies (#266 boot-ordering defect)"
171                );
172                DeclaredBodyLookup::None
173            },
174            |source| source.body_for(task_queue, action, run),
175        )
176    }
177}
178
179/// Reads declared bodies out of the engine's live workflow catalog.
180pub struct EngineDeclaredBodies {
181    engine: Arc<aion::Engine>,
182}
183
184impl EngineDeclaredBodies {
185    /// Build a reader over `engine`'s catalog.
186    #[must_use]
187    pub const fn new(engine: Arc<aion::Engine>) -> Self {
188        Self { engine }
189    }
190}
191
192impl std::fmt::Debug for EngineDeclaredBodies {
193    fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
194        formatter.write_str("EngineDeclaredBodies")
195    }
196}
197
198impl EngineDeclaredBodies {
199    /// The package version `run` is pinned to, or `None` when the registry
200    /// cannot name it.
201    ///
202    /// Two ways to reach `None`, and both are reported rather than swallowed:
203    /// the run has no handle (it left the registry), or the registry could not
204    /// be read at all. Neither is a reason to guess a body — the caller falls
205    /// back to the queue-wide reading, which refuses on disagreement.
206    fn version_of(&self, run: DispatchingRun<'_>) -> Option<ContentHash> {
207        match self.engine.registry().get(run.workflow_id, run.run_id) {
208            Ok(Some(handle)) => Some(handle.loaded_version().clone()),
209            Ok(None) => {
210                tracing::warn!(
211                    operation = "declared_command_dispatch",
212                    workflow_id = %run.workflow_id,
213                    run_id = %run.run_id,
214                    "no registry handle for the dispatching run; resolving its body \
215                     from the whole queue instead of from its own package version"
216                );
217                None
218            }
219            Err(error) => {
220                tracing::error!(
221                    operation = "declared_command_dispatch",
222                    workflow_id = %run.workflow_id,
223                    run_id = %run.run_id,
224                    %error,
225                    "registry unreadable while resolving the dispatching run's version; \
226                     resolving its body from the whole queue instead"
227                );
228                None
229            }
230        }
231    }
232}
233
234impl DeclaredBodies for EngineDeclaredBodies {
235    fn body_for(
236        &self,
237        task_queue: &str,
238        action: &str,
239        run: DispatchingRun<'_>,
240    ) -> DeclaredBodyLookup {
241        let contracts = match self.engine.worker_contracts_for_queue(task_queue) {
242            Ok(contracts) => contracts,
243            Err(error) => return DeclaredBodyLookup::Unreadable(error.to_string()),
244        };
245        // The RAW retained set is the right input here, unlike worker admission
246        // (see `Engine::worker_contracts_for_queue`): a run pinned to a version
247        // nothing else can reach still has to execute that version's body. What
248        // narrows the answer is the run's own identity, not reachability.
249        select_declared_body(&contracts, action, self.version_of(run).as_ref())
250    }
251}
252
253/// The dispatcher decorator that executes declared bodies at the server.
254///
255/// Wraps the production dispatcher. Consults the declared-body source before
256/// every dispatch; delegates untouched whenever the action carries no body.
257pub struct DeclaredCommandDispatcher {
258    inner: Arc<dyn ActivityDispatcher>,
259    bodies: DeclaredBodySource,
260    attempts: DeclaredCommandAttempts,
261    tokio: tokio::runtime::Handle,
262    workspace_root: WorkspaceRoot,
263    transcript: ActivityEventPublisher,
264}
265
266impl DeclaredCommandDispatcher {
267    /// Wrap `inner`, consulting `bodies` before every dispatch, registering
268    /// every attempt it executes in `attempts` so the run's cancel can reach
269    /// it, expanding `{workspace_root}` in declared commands with the
270    /// server-resolved `workspace_root`, and streaming each executed command's
271    /// output onto `transcript` — the deployment's one transcript sequencer,
272    /// shared with every agent step.
273    ///
274    /// `attempts` is required rather than optional because a dispatcher without
275    /// one would execute commands nothing could stop, which is precisely the
276    /// state this argument exists to end.
277    #[must_use]
278    pub fn new(
279        inner: Arc<dyn ActivityDispatcher>,
280        bodies: DeclaredBodySource,
281        attempts: DeclaredCommandAttempts,
282        tokio: tokio::runtime::Handle,
283        workspace_root: WorkspaceRoot,
284        transcript: ActivityEventPublisher,
285    ) -> Self {
286        Self {
287            inner,
288            bodies,
289            attempts,
290            tokio,
291            workspace_root,
292            transcript,
293        }
294    }
295
296    /// Parse the declared command into the executor's action, with the
297    /// server-resolved `{workspace_root}` already spliced in.
298    ///
299    /// Ratification condition (#139): a body that USES the placeholder is
300    /// refused terminally, by name, when the root cannot resolve to an absolute
301    /// directory that exists — no fallback to cwd, temp, or anything else. A
302    /// body without the placeholder never reaches the resolution at all
303    /// (`expand` returns `Ok(None)` untouched).
304    fn declared_action(
305        &self,
306        request: &ActivityDispatch,
307        command: &str,
308    ) -> Result<ShellAction, String> {
309        let expanded = self.workspace_root.expand(command).map_err(|error| {
310            format!(
311                "terminal:declared body for action `{name}` uses the {placeholder} \
312                 placeholder and cannot dispatch: {error}",
313                name = request.name,
314                placeholder = WORKSPACE_ROOT_PLACEHOLDER,
315            )
316        })?;
317        if let Some(expansion) = &expanded {
318            tracing::info!(
319                operation = "declared_command_dispatch",
320                workflow_id = %request.workflow_id,
321                activity_id = %request.activity_id,
322                activity_name = %request.name,
323                task_queue = %request.task_queue,
324                attempt = request.attempt,
325                workspace_root = %expansion.workspace_root,
326                "expanded the workspace-root placeholder in the declared command"
327            );
328        }
329        let command = expanded
330            .as_ref()
331            .map_or(command, |expansion| expansion.command.as_str());
332        ShellAction::new(command).map_err(|error| {
333            // The AWL checker refuses these at compile time, so reaching this
334            // arm means a defective contract got deployed — name the defect
335            // rather than hiding it behind a generic dispatch failure.
336            format!("terminal:declared command failed to parse at dispatch: {error}")
337        })
338    }
339
340    /// Put the attempt on this server's cancel path BEFORE its command starts.
341    ///
342    /// The returned guard keeps it there for exactly as long as the command
343    /// runs, so an attempt that ended — completed, failed, or unwound — can
344    /// never be signalled afterwards. A registration that cannot be made is a
345    /// command nothing could stop, so the dispatch is refused rather than run:
346    /// an uncancellable command on the operator's machine is the whole defect
347    /// this registration exists to prevent, and starting one to avoid an error
348    /// message would be choosing it.
349    fn join_cancel_path(
350        &self,
351        request: &ActivityDispatch,
352        cancellation: &aion_worker::ActivityCancellationHandle,
353    ) -> Result<super::DeclaredAttemptRegistration, String> {
354        self.attempts
355            .register(
356                super::AttemptKey::new(
357                    request.workflow_id.clone(),
358                    request.run_id.clone(),
359                    request.activity_id.clone(),
360                    request.attempt,
361                ),
362                cancellation.clone(),
363            )
364            .map_err(|error| match error {
365                // A draining refusal is a park, not a failure: same sentinel a
366                // worker dispatch returns mid-drain. Nothing is recorded, the
367                // engine parks the attempt, and the next boot re-dispatches
368                // it — a terminal error here would fail the workflow for the
369                // crime of the operator stopping the server.
370                crate::error::ServerError::DrainingRefusedDeclaredAttempt { .. } => {
371                    tracing::info!(
372                        operation = "declared_command_dispatch",
373                        workflow_id = %request.workflow_id,
374                        activity_id = %request.activity_id,
375                        activity_name = %request.name,
376                        task_queue = %request.task_queue,
377                        attempt = request.attempt,
378                        "declared command parked: this server is draining and starts no new work"
379                    );
380                    aion::PARKED_ACTIVITY_REASON.to_owned()
381                }
382                other => format!(
383                    "terminal:declared body for action `{name}` cannot dispatch: the attempt \
384                     could not join this server's cancel path, and a command a cancelled run \
385                     could not stop must not be started: {other}",
386                    name = request.name,
387                ),
388            })
389    }
390
391    /// Execute one declared command attempt and encode the outcome onto the
392    /// FFI string contract (`retryable:`/`terminal:` on the error side).
393    fn run_declared_command(
394        &self,
395        request: &ActivityDispatch,
396        command: &str,
397    ) -> Result<String, String> {
398        let arguments = decode_arguments(&request.input)?;
399        let action = self.declared_action(request, command)?;
400        // The live transcript seam for this attempt. The context owns the
401        // sending end, so dropping it after the run closes the stream and ends
402        // the pump — which is then awaited, so no observed line is abandoned
403        // unpublished when the command finishes.
404        let (events, drain) = tokio::sync::mpsc::unbounded_channel();
405        let (context, cancellation) = aion_worker::ActivityContext::with_transcript(
406            request.workflow_id.clone(),
407            request.run_id.clone(),
408            request.activity_id.clone(),
409            request.attempt,
410            events,
411        );
412        let registration = self.join_cancel_path(request, &cancellation)?;
413
414        tracing::info!(
415            operation = "declared_command_dispatch",
416            workflow_id = %request.workflow_id,
417            activity_id = %request.activity_id,
418            activity_name = %request.name,
419            task_queue = %request.task_queue,
420            attempt = request.attempt,
421            "executing declared action body at the server"
422        );
423        // All three 2026-08-16 anonymous deaths correlated with workflow
424        // execution and the third died on exactly this path; the breadcrumb
425        // makes the in-flight site a death-note fact, not a log inference.
426        crate::death_note::breadcrumb(&format!(
427            "declared-action start action={} workflow_id={} run_id={} activity_id={} attempt={}",
428            request.name, request.workflow_id, request.run_id, request.activity_id, request.attempt,
429        ));
430
431        // #223: the bound the DISPATCH authored, or `None` when it authored
432        // none. Read through the engine's own decoder so the server cannot
433        // answer "what did this document authorise" differently from the retry
434        // loop, and so an unbounded body stays unbounded — the server invents
435        // no deadline of its own.
436        let bound = aion::activity_timeout_from_config(&request.config);
437        let transcript = self.transcript.clone();
438        let ended = self.tokio.block_on(async move {
439            let pump = tokio::spawn(publish_declared_transcript(transcript, drain));
440            let ended = run_bounded(&action, &arguments, &context, &cancellation, bound).await;
441            // Closing the seam is what ends the pump; the context holds it.
442            drop(context);
443            if let Err(error) = pump.await {
444                tracing::warn!(
445                    %error,
446                    operation = "declared_command_dispatch",
447                    "declared command transcript: the publishing task ended abnormally; some \
448                     output lines may not have been retained"
449                );
450            }
451            ended
452        });
453        // The command is over and its group is gone, so the attempt leaves the
454        // cancel path. Dropped explicitly, here and not earlier: while this
455        // lives, a cancel arriving mid-run still reaches the process.
456        drop(registration);
457
458        encode_end(request, ended)
459    }
460}
461
462/// Encode how the attempt ended onto the FFI string contract.
463///
464/// Three vocabularies, one per honest outcome: the encoded result, the
465/// classified failure the executor produced (`retryable:`/`terminal:`), and the
466/// engine's own `timeout:` reason for an attempt that outlived its authored
467/// bound.
468fn encode_end(request: &ActivityDispatch, ended: AttemptEnd) -> Result<String, String> {
469    let outcome = match ended {
470        AttemptEnd::Ran(outcome) => outcome,
471        AttemptEnd::Expired { bound, ran_anyway } => {
472            if let Some(exit_code) = ran_anyway {
473                // The command reached its own end inside the stopping window.
474                // Its result is discarded — the attempt is already recorded as
475                // having outlived its bound, and answering with a late success
476                // would contradict a terminal the run has already been told
477                // about — but the fact is said, not swallowed.
478                tracing::warn!(
479                    operation = "declared_command_dispatch",
480                    workflow_id = %request.workflow_id,
481                    activity_id = %request.activity_id,
482                    activity_name = %request.name,
483                    attempt = request.attempt,
484                    exit_code,
485                    bound_ms = bound.as_millis(),
486                    "the declared command finished while it was being stopped on its \
487                     authored bound; its result is discarded in favour of the timeout"
488                );
489            }
490            return Err(aion::activity_timeout_reason(bound));
491        }
492    };
493
494    match outcome {
495        Ok(result) => serde_json::to_string(&result)
496            .map_err(|error| format!("terminal:declared command result failed to encode: {error}")),
497        Err(failure) => {
498            let prefix = match failure.classification() {
499                aion_worker::Classification::Retryable => "retryable",
500                aion_worker::Classification::PolicyRefused => "policy_refused",
501                aion_worker::Classification::Terminal => "terminal",
502            };
503            Err(format!("{prefix}:{}", failure.message()))
504        }
505    }
506}
507
508/// How one declared-command attempt ended.
509#[derive(Debug)]
510enum AttemptEnd {
511    /// The command ran to its own end — completed, failed, or was stopped by
512    /// something other than the authored bound.
513    Ran(Result<aion_worker::shell::ShellOutcome, aion_worker::ActivityFailure>),
514    /// The attempt outlived the per-attempt bound its dispatch authored, and
515    /// its process group has been stopped and PROVEN gone.
516    Expired {
517        /// The authored bound that fired, carried so the refusal can name it.
518        bound: std::time::Duration,
519        /// The exit code of a command that reached its own end inside the
520        /// stopping window, when that happened. `None` — the ordinary case —
521        /// means the command was still running when the bound was enforced.
522        ran_anyway: Option<i32>,
523    },
524}
525
526/// Run the declared command, ending it at the bound its dispatch authored.
527///
528/// # Why the server enforces a bound the engine already applies
529///
530/// The engine wraps every attempt in `tokio::time::timeout` at the same
531/// authored bound (`nif_activity_retry_dispatch::deliver_one_attempt`), and is
532/// explicit about what that achieves: "the dispatch future is DROPPED, which
533/// stops this run waiting and nothing more... the worker-side call runs on to
534/// its own end and its result is discarded". For a REMOTE worker that is
535/// someone else's machine and the right division of labour. For a declared body
536/// it is a process tree in the server's own process group hierarchy, on the
537/// operator's machine, with nothing left that could ever stop it — the run has
538/// already moved on.
539///
540/// So the bound is enforced HERE as well, where the process is. On expiry the
541/// activity's cancellation is signalled and the SAME run future is awaited to
542/// its end: [`aion_worker::run_cancellable_command`] does not return until
543/// `SIGTERM` → [`aion_worker::PROCESS_GROUP_TERMINATION_GRACE`] → `SIGKILL` has
544/// been delivered to the whole group and the group has been PROVEN gone. This
545/// therefore returns only once the command is genuinely stopped, and the
546/// termination ladder and its grace are the worker path's, not a second copy.
547///
548/// A dispatch that authored no bound is awaited exactly as before.
549async fn run_bounded(
550    action: &ShellAction,
551    arguments: &BTreeMap<String, serde_json::Value>,
552    context: &aion_worker::ActivityContext,
553    cancellation: &aion_worker::ActivityCancellationHandle,
554    bound: Option<std::time::Duration>,
555) -> AttemptEnd {
556    let run = action.run(arguments, context);
557    let Some(bound) = bound else {
558        return AttemptEnd::Ran(run.await);
559    };
560    tokio::pin!(run);
561    match tokio::time::timeout(bound, &mut run).await {
562        Ok(outcome) => AttemptEnd::Ran(outcome),
563        Err(_elapsed) => {
564            cancellation.cancel();
565            AttemptEnd::Expired {
566                bound,
567                ran_anyway: run.await.ok().map(|outcome| outcome.exit_code),
568            }
569        }
570    }
571}
572
573impl std::fmt::Debug for DeclaredCommandDispatcher {
574    fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
575        formatter
576            .debug_struct("DeclaredCommandDispatcher")
577            .field("bodies", &self.bodies)
578            .finish_non_exhaustive()
579    }
580}
581
582impl ActivityDispatcher for DeclaredCommandDispatcher {
583    fn dispatch(&self, request: ActivityDispatch) -> Result<String, String> {
584        let run = DispatchingRun {
585            workflow_id: &request.workflow_id,
586            run_id: &request.run_id,
587        };
588        match self
589            .bodies
590            .body_for(&request.task_queue, &request.name, run)
591        {
592            DeclaredBodyLookup::None => self.inner.dispatch(request),
593            DeclaredBodyLookup::Unreadable(reason) => {
594                // Delegated, not refused: a catalog read failure must not
595                // strand a queue that live workers could still serve. Loud so
596                // an operator sees a bodied action falling through.
597                tracing::error!(
598                    operation = "declared_command_dispatch",
599                    workflow_id = %request.workflow_id,
600                    activity_name = %request.name,
601                    task_queue = %request.task_queue,
602                    %reason,
603                    "declared-body catalog read failed; delegating to the worker path"
604                );
605                self.inner.dispatch(request)
606            }
607            DeclaredBodyLookup::Ambiguous { declaring } => Err(ambiguous_body_refusal(
608                &request.name,
609                &request.task_queue,
610                &declaring,
611            )),
612            DeclaredBodyLookup::Declared(ActionBodyContract::Run { command }) => {
613                self.run_declared_command(&request, &command)
614            }
615        }
616    }
617}
618
619/// Decode the dispatch's JSON input into the declared action's arguments.
620///
621/// A declared action's parameters are named in its `.awl` declaration, so the
622/// input must be a JSON object; anything else cannot bind to `$name`
623/// references and is refused by shape. Retrying cannot change the input, so
624/// the refusal is terminal.
625fn decode_arguments(input: &str) -> Result<BTreeMap<String, serde_json::Value>, String> {
626    let value: serde_json::Value = serde_json::from_str(input)
627        .map_err(|error| format!("terminal:declared command input is not valid JSON: {error}"))?;
628    match value {
629        serde_json::Value::Object(members) => Ok(members.into_iter().collect()),
630        other => Err(format!(
631            "terminal:declared command input must be a JSON object binding the action's \
632             parameters by name; got {}",
633            json_kind(&other)
634        )),
635    }
636}
637
638/// A JSON value's kind, named for a refusal message.
639const fn json_kind(value: &serde_json::Value) -> &'static str {
640    match value {
641        serde_json::Value::Null => "null",
642        serde_json::Value::Bool(_) => "a boolean",
643        serde_json::Value::Number(_) => "a number",
644        serde_json::Value::String(_) => "a string",
645        serde_json::Value::Array(_) => "an array",
646        serde_json::Value::Object(_) => "an object",
647    }
648}
649
650/// Containment of a server-executed body: the authored per-attempt bound, and
651/// the run's cancellation. Its own file because its subjects are live process
652/// trees rather than dispatcher shapes; it builds them out of the fixtures
653/// [`tests`] shares with it.
654#[cfg(test)]
655#[path = "declared_body_containment_tests.rs"]
656mod declared_body_containment_tests;
657
658#[cfg(test)]
659mod tests {
660    use std::collections::BTreeMap;
661    use std::sync::{Arc, Mutex};
662
663    use aion::{ActivityDispatch, ActivityDispatcher};
664    use aion_core::{ActivityId, RunId, WorkflowId};
665    use aion_package::ActionBodyContract;
666
667    use aion_core::ActivityEventKind;
668    use aion_store::ActivityStreamKey;
669
670    use super::super::workspace_root::{WorkspaceRoot, WorkspaceRootError};
671    use super::{
672        ActivityEventPublisher, DeclaredBodies, DeclaredBodyLookup, DeclaredBodySource,
673        DeclaredCommandAttempts, DeclaredCommandDispatcher, DeclaringVersion, DispatchingRun,
674        decode_arguments,
675    };
676
677    /// What a test returns. Every fallible step is carried rather than
678    /// unwrapped, because the workspace denies panicking accessors in test
679    /// code as firmly as in library code.
680    pub(super) type TestResult = Result<(), Box<dyn std::error::Error>>;
681
682    /// Inner dispatcher that records whether it was reached.
683    struct RecordingInner {
684        reached: Arc<Mutex<Vec<String>>>,
685        reply: Result<String, String>,
686    }
687
688    impl ActivityDispatcher for RecordingInner {
689        fn dispatch(&self, request: ActivityDispatch) -> Result<String, String> {
690            match self.reached.lock() {
691                Ok(mut names) => names.push(request.name),
692                Err(poisoned) => poisoned.into_inner().push(request.name),
693            }
694            self.reply.clone()
695        }
696    }
697
698    struct FixedBodies {
699        lookup: DeclaredBodyLookup,
700    }
701
702    impl DeclaredBodies for FixedBodies {
703        fn body_for(
704            &self,
705            _task_queue: &str,
706            _action: &str,
707            _run: DispatchingRun<'_>,
708        ) -> DeclaredBodyLookup {
709            self.lookup.clone()
710        }
711    }
712
713    /// A reader that records whose dispatch it was asked about.
714    ///
715    /// The selection rule is unit-tested on its own inputs, which proves the
716    /// rule and nothing about the plumbing. This double closes that gap: it
717    /// captures the [`DispatchingRun`] the dispatcher hands over, so the
718    /// identity can be compared against the request it came from.
719    struct RecordingBodies {
720        seen: Arc<Mutex<Vec<(WorkflowId, RunId)>>>,
721    }
722
723    impl DeclaredBodies for RecordingBodies {
724        fn body_for(
725            &self,
726            _task_queue: &str,
727            _action: &str,
728            run: DispatchingRun<'_>,
729        ) -> DeclaredBodyLookup {
730            let observed = (run.workflow_id.clone(), run.run_id.clone());
731            match self.seen.lock() {
732                Ok(mut seen) => seen.push(observed),
733                Err(poisoned) => poisoned.into_inner().push(observed),
734            }
735            DeclaredBodyLookup::None
736        }
737    }
738
739    pub(super) fn request(name: &str, input: &str) -> ActivityDispatch {
740        ActivityDispatch {
741            namespace: "default".to_owned(),
742            task_queue: "shell".to_owned(),
743            node: None,
744            workflow_id: WorkflowId::new_v4(),
745            run_id: RunId::new_v4(),
746            activity_id: ActivityId::from_sequence_position(1),
747            name: name.to_owned(),
748            input: input.to_owned(),
749            config: "{}".to_owned(),
750            attempt: 1,
751            labels: BTreeMap::new(),
752            advisory: false,
753        }
754    }
755
756    fn dispatcher(
757        lookup: DeclaredBodyLookup,
758        reply: Result<String, String>,
759    ) -> (DeclaredCommandDispatcher, Arc<Mutex<Vec<String>>>) {
760        // These tests exercise bodies without the placeholder, so the root's
761        // value is never read; it is an explicit existing directory rather
762        // than a default so nothing here depends on resolution.
763        let (decorated, reached, _transcript) = dispatcher_with_root(
764            lookup,
765            reply,
766            WorkspaceRoot::from_resolution(Ok(std::env::temp_dir())),
767        );
768        (decorated, reached)
769    }
770
771    /// The live-tail buffer these tests give their transcript sequencer. A
772    /// `const` match rather than an unwrap: the workspace denies panicking
773    /// accessors in test code as firmly as in library code.
774    const TRANSCRIPT_CAPACITY: std::num::NonZeroUsize = match std::num::NonZeroUsize::new(64) {
775        Some(capacity) => capacity,
776        None => std::num::NonZeroUsize::MIN,
777    };
778
779    /// A dispatcher whose executing attempts nothing external will signal.
780    ///
781    /// Correct for every test whose subject runs to its own end. A test that
782    /// CANCELS its subject needs the registry the cancel signals through, and
783    /// uses [`dispatcher_with_attempts`] to hold the same instance.
784    pub(super) fn dispatcher_with_root(
785        lookup: DeclaredBodyLookup,
786        reply: Result<String, String>,
787        workspace_root: WorkspaceRoot,
788    ) -> (
789        DeclaredCommandDispatcher,
790        Arc<Mutex<Vec<String>>>,
791        ActivityEventPublisher,
792    ) {
793        dispatcher_with_attempts(
794            lookup,
795            reply,
796            workspace_root,
797            DeclaredCommandAttempts::new(crate::shutdown::DrainState::default()),
798        )
799    }
800
801    pub(super) fn dispatcher_with_attempts(
802        lookup: DeclaredBodyLookup,
803        reply: Result<String, String>,
804        workspace_root: WorkspaceRoot,
805        attempts: DeclaredCommandAttempts,
806    ) -> (
807        DeclaredCommandDispatcher,
808        Arc<Mutex<Vec<String>>>,
809        ActivityEventPublisher,
810    ) {
811        let reached = Arc::new(Mutex::new(Vec::new()));
812        let inner = RecordingInner {
813            reached: Arc::clone(&reached),
814            reply,
815        };
816        let bodies = DeclaredBodySource::default();
817        bodies.install(Arc::new(FixedBodies { lookup }));
818        let store: Arc<dyn aion_store::ObservabilityStore> =
819            Arc::new(aion_store::InMemoryObservabilityStore::default());
820        let transcript = ActivityEventPublisher::new(
821            store,
822            TRANSCRIPT_CAPACITY,
823            crate::activity_publisher::TranscriptBatchPolicy::UNBATCHED,
824        );
825        let decorated = DeclaredCommandDispatcher::new(
826            Arc::new(inner),
827            bodies,
828            attempts,
829            tokio::runtime::Handle::current(),
830            workspace_root,
831            transcript.clone(),
832        );
833        (decorated, reached, transcript)
834    }
835
836    pub(super) fn reached_names(reached: &Arc<Mutex<Vec<String>>>) -> Vec<String> {
837        match reached.lock() {
838            Ok(names) => names.clone(),
839            Err(poisoned) => poisoned.into_inner().clone(),
840        }
841    }
842
843    #[tokio::test(flavor = "multi_thread")]
844    async fn a_bodiless_action_is_delegated_untouched() -> TestResult {
845        let (decorated, reached) =
846            dispatcher(DeclaredBodyLookup::None, Ok("\"worker-served\"".to_owned()));
847        let handle =
848            tokio::task::spawn_blocking(move || decorated.dispatch(request("plain", "{}")));
849        let result = handle.await?;
850        assert_eq!(result, Ok("\"worker-served\"".to_owned()));
851        assert_eq!(reached_names(&reached), vec!["plain".to_owned()]);
852        Ok(())
853    }
854
855    #[tokio::test(flavor = "multi_thread")]
856    async fn a_declared_body_executes_without_touching_the_worker_path() -> TestResult {
857        let (decorated, reached) = dispatcher(
858            DeclaredBodyLookup::Declared(ActionBodyContract::Run {
859                command: "echo {{greeting}}".to_owned(),
860            }),
861            Err("terminal:the worker path must never be reached".to_owned()),
862        );
863        let handle = tokio::task::spawn_blocking(move || {
864            decorated.dispatch(request(
865                "greet",
866                "{\"greeting\":\"hello from the contract\"}",
867            ))
868        });
869        let result = handle.await?;
870        let encoded = result.map_err(|error| format!("declared command failed: {error}"))?;
871        let outcome: serde_json::Value = serde_json::from_str(&encoded)?;
872        assert_eq!(outcome["stdout"], "hello from the contract");
873        assert_eq!(outcome["exit_code"], 0);
874        assert!(
875            reached_names(&reached).is_empty(),
876            "the worker path must not be consulted for a bodied action"
877        );
878        Ok(())
879    }
880
881    /// A draining server PARKS a declared dispatch instead of running it: the
882    /// dispatch returns the park sentinel (the same face a worker dispatch
883    /// wears mid-drain, so the engine records nothing and the next boot
884    /// re-dispatches), the command's process never starts, and the census the
885    /// drain gate waits on registers nothing — work arriving after `stop` can
886    /// neither launch nor hold the gate open.
887    #[tokio::test(flavor = "multi_thread")]
888    async fn a_draining_server_parks_a_declared_dispatch_without_starting_it() -> TestResult {
889        let marker =
890            std::env::temp_dir().join(format!("aion-drain-park-{}", uuid::Uuid::new_v4().simple()));
891        let drain = crate::shutdown::DrainState::default();
892        let attempts = DeclaredCommandAttempts::new(drain.clone());
893        let (decorated, reached, _transcript) = dispatcher_with_attempts(
894            DeclaredBodyLookup::Declared(ActionBodyContract::Run {
895                command: format!("touch {}", marker.display()),
896            }),
897            Err("terminal:the worker path must never be reached".to_owned()),
898            WorkspaceRoot::from_resolution(Ok(std::env::temp_dir())),
899            attempts.clone(),
900        );
901        assert!(drain.begin(), "the first begin() must flip the latch");
902
903        let handle =
904            tokio::task::spawn_blocking(move || decorated.dispatch(request("touch_marker", "{}")));
905        let result = handle.await?;
906
907        assert_eq!(
908            result,
909            Err(aion::PARKED_ACTIVITY_REASON.to_owned()),
910            "a drained-over declared dispatch must wear the park sentinel, not a failure"
911        );
912        assert!(
913            !marker.exists(),
914            "the declared command must never start on a draining server"
915        );
916        assert!(
917            reached_names(&reached).is_empty(),
918            "the park must not fall through to the worker path"
919        );
920        assert!(
921            attempts
922                .executing()
923                .map_err(|error| format!("census read failed: {error}"))?
924                .is_empty(),
925            "a parked dispatch must leave no census entry to hold the drain gate open"
926        );
927        Ok(())
928    }
929
930    /// THE MID-STEP ANSWER: a server-run declared body's output reaches the
931    /// deployment's transcript sequencer as one event per line, on both streams,
932    /// keyed to the dispatch's own `(workflow, activity, attempt)` — the same
933    /// durable stream an agent step's transcript is read from, so every reader
934    /// that already serves transcripts serves this without change.
935    ///
936    /// The completion contract is asserted on the same run: the recorded result
937    /// still carries the command's whole stdout.
938    #[tokio::test(flavor = "multi_thread")]
939    async fn a_declared_body_publishes_its_output_onto_the_transcript() -> TestResult {
940        let (decorated, reached, transcript) = dispatcher_with_root(
941            DeclaredBodyLookup::Declared(ActionBodyContract::Run {
942                command: "sh -c 'echo one; echo two; echo warned >&2'".to_owned(),
943            }),
944            Err("terminal:the worker path must never be reached".to_owned()),
945            WorkspaceRoot::from_resolution(Ok(std::env::temp_dir())),
946        );
947        let dispatch = request("noisy", "{}");
948        let key = ActivityStreamKey::new(
949            dispatch.workflow_id.clone(),
950            dispatch.run_id.clone(),
951            dispatch.activity_id.clone(),
952            dispatch.attempt,
953        );
954
955        let handle = tokio::task::spawn_blocking(move || decorated.dispatch(dispatch));
956        let encoded = handle
957            .await?
958            .map_err(|error| format!("declared command failed: {error}"))?;
959
960        // The replay-authoritative result is untouched by the streaming.
961        let outcome: serde_json::Value = serde_json::from_str(&encoded)?;
962        assert_eq!(outcome["stdout"], "one\ntwo");
963        assert_eq!(outcome["stderr"], "warned");
964        assert!(reached_names(&reached).is_empty());
965
966        // ...and the same output is on the durable transcript, line by line.
967        let retained = transcript.replay_from(&key, 0).await?;
968        let lines = retained
969            .iter()
970            .map(|record| match &record.event.kind {
971                ActivityEventKind::Message { text, .. } => {
972                    (record.event.agent_role.clone(), text.clone())
973                }
974                other => (record.event.agent_role.clone(), format!("{other:?}")),
975            })
976            .collect::<Vec<_>>();
977        assert!(
978            lines.contains(&("command stdout".to_owned(), "one".to_owned()))
979                && lines.contains(&("command stdout".to_owned(), "two".to_owned())),
980            "each stdout line must be its own transcript event: {lines:?}"
981        );
982        assert!(
983            lines.contains(&("command stderr".to_owned(), "warned".to_owned())),
984            "stderr must be on the transcript, labelled by its stream: {lines:?}"
985        );
986        // Sequencing is the publisher's: the durable order is gap-free from 0.
987        let sequences = retained
988            .iter()
989            .map(|record| record.store_seq)
990            .collect::<Vec<_>>();
991        assert_eq!(
992            sequences,
993            (0..u64::try_from(retained.len())?).collect::<Vec<_>>(),
994            "the sequencer assigns a gap-free durable order"
995        );
996        Ok(())
997    }
998
999    #[tokio::test(flavor = "multi_thread")]
1000    async fn a_failing_declared_command_reports_retryable_with_its_stderr() -> TestResult {
1001        let (decorated, _reached) = dispatcher(
1002            DeclaredBodyLookup::Declared(ActionBodyContract::Run {
1003                command: "sh -c 'echo boom >&2; exit 7'".to_owned(),
1004            }),
1005            Ok("unused".to_owned()),
1006        );
1007        let handle =
1008            tokio::task::spawn_blocking(move || decorated.dispatch(request("fails", "{}")));
1009        let Err(error) = handle.await? else {
1010            return Err("a non-zero exit must fail the dispatch".into());
1011        };
1012        assert!(
1013            error.starts_with("retryable:"),
1014            "a non-zero exit is retryable by default: {error}"
1015        );
1016        assert!(
1017            error.contains("boom"),
1018            "stderr must ride the failure: {error}"
1019        );
1020        Ok(())
1021    }
1022
1023    /// The hash the refusal prints must be one the deploy API will accept, or
1024    /// the remedy is a command that cannot run — the exact failure the old
1025    /// "redeploy so one body remains" wording had.
1026    ///
1027    /// The oracle is the deploy API's own parser, not a length or a shape:
1028    /// `EngineDeclaredBodies` renders the version with `ContentHash::to_string`,
1029    /// so this takes a real hash through that rendering, pulls the token back
1030    /// out of the printed command, and parses it the way
1031    /// `decode_version_target` does.
1032    #[test]
1033    fn the_printed_hash_parses_back_as_a_content_hash() -> TestResult {
1034        let version = aion_package::ContentHash::from_bytes([0x5a; 32]);
1035        let routed = aion_package::ContentHash::from_bytes([0xa5; 32]);
1036        let refusal = super::ambiguous_body_refusal(
1037            "find_repositories",
1038            "local",
1039            &[
1040                DeclaringVersion {
1041                    content_hash: version.to_string(),
1042                    workflow_types: vec!["sweeper".to_owned()],
1043                    route_active: false,
1044                    body: 0,
1045                },
1046                DeclaringVersion {
1047                    content_hash: routed.to_string(),
1048                    workflow_types: vec!["sweeper".to_owned()],
1049                    route_active: true,
1050                    body: 1,
1051                },
1052            ],
1053        );
1054        let Some(command) = refusal.split("`aion unload sweeper ").nth(1) else {
1055            return Err(format!("no unload command in the refusal: {refusal}").into());
1056        };
1057        let Some(printed) = command.split('`').next() else {
1058            return Err(format!("the unload command is unterminated: {refusal}").into());
1059        };
1060        let parsed: aion_package::ContentHash = printed.parse()?;
1061        assert_eq!(
1062            parsed, version,
1063            "the printed hash must round-trip to the version it names"
1064        );
1065        Ok(())
1066    }
1067
1068    #[tokio::test(flavor = "multi_thread")]
1069    async fn ambiguous_bodies_refuse_terminally_by_name() -> TestResult {
1070        let superseded = "1111111111111111111111111111111111111111111111111111111111111111";
1071        let routed = "2222222222222222222222222222222222222222222222222222222222222222";
1072        let (decorated, reached) = dispatcher(
1073            DeclaredBodyLookup::Ambiguous {
1074                declaring: vec![
1075                    DeclaringVersion {
1076                        content_hash: superseded.to_owned(),
1077                        workflow_types: vec!["sweeper".to_owned()],
1078                        route_active: false,
1079                        body: 0,
1080                    },
1081                    DeclaringVersion {
1082                        content_hash: routed.to_owned(),
1083                        workflow_types: vec!["sweeper".to_owned()],
1084                        route_active: true,
1085                        body: 1,
1086                    },
1087                ],
1088            },
1089            Ok(String::new()),
1090        );
1091        let handle = tokio::task::spawn_blocking(move || decorated.dispatch(request("torn", "{}")));
1092        let Err(error) = handle.await? else {
1093            return Err("ambiguous bodies must refuse".into());
1094        };
1095        assert!(error.starts_with("terminal:"), "{error}");
1096        assert!(error.contains("torn"), "{error}");
1097        // The refusal must reach the dispatcher carrying an act-on-able remedy,
1098        // not just a count: the operator reads this string and nothing else.
1099        assert!(
1100            error.contains(&format!("`aion unload sweeper {superseded}`")),
1101            "the dispatch refusal must name the version to retire: {error}"
1102        );
1103        assert!(reached_names(&reached).is_empty());
1104        Ok(())
1105    }
1106
1107    #[tokio::test(flavor = "multi_thread")]
1108    async fn a_placeholder_bearing_body_executes_with_the_expanded_root() -> TestResult {
1109        let scratch = tempfile::tempdir()?;
1110        let root = scratch.path().join("clones");
1111        let root_text = root.to_string_lossy().into_owned();
1112        let (decorated, reached, _transcript) = dispatcher_with_root(
1113            DeclaredBodyLookup::Declared(ActionBodyContract::Run {
1114                command: "echo {workspace_root}".to_owned(),
1115            }),
1116            Err("terminal:the worker path must never be reached".to_owned()),
1117            WorkspaceRoot::from_resolution(Ok(root.clone())),
1118        );
1119        let handle =
1120            tokio::task::spawn_blocking(move || decorated.dispatch(request("provision", "{}")));
1121        let result = handle.await?;
1122        let encoded = result.map_err(|error| format!("declared command failed: {error}"))?;
1123        let outcome: serde_json::Value = serde_json::from_str(&encoded)?;
1124        assert_eq!(
1125            outcome["stdout"], root_text,
1126            "the command must observe the server-resolved root as its argv word"
1127        );
1128        assert_eq!(outcome["exit_code"], 0);
1129        assert!(
1130            root.is_dir(),
1131            "dispatching a placeholder-bearing body must create the missing root"
1132        );
1133        assert!(reached_names(&reached).is_empty());
1134        Ok(())
1135    }
1136
1137    #[tokio::test(flavor = "multi_thread")]
1138    async fn a_placeholder_bearing_body_refuses_terminally_when_the_root_is_unresolved()
1139    -> TestResult {
1140        let (decorated, reached, _transcript) = dispatcher_with_root(
1141            DeclaredBodyLookup::Declared(ActionBodyContract::Run {
1142                command: "echo {workspace_root}".to_owned(),
1143            }),
1144            Ok("unused".to_owned()),
1145            WorkspaceRoot::from_resolution(Err(WorkspaceRootError::Unresolvable {
1146                reason: "cannot resolve Aion home: set AION_HOME or HOME".to_owned(),
1147            })),
1148        );
1149        let handle =
1150            tokio::task::spawn_blocking(move || decorated.dispatch(request("provision", "{}")));
1151        let Err(error) = handle.await? else {
1152            return Err("an unresolved root must refuse a placeholder-bearing body".into());
1153        };
1154        assert!(error.starts_with("terminal:"), "{error}");
1155        assert!(
1156            error.contains("provision"),
1157            "the refusal must name the action: {error}"
1158        );
1159        assert!(
1160            error.contains("cannot resolve Aion home"),
1161            "the refusal must carry the resolution failure's reason: {error}"
1162        );
1163        assert!(
1164            reached_names(&reached).is_empty(),
1165            "a refused body must not fall through to the worker path"
1166        );
1167        Ok(())
1168    }
1169
1170    #[tokio::test(flavor = "multi_thread")]
1171    async fn a_shape_changing_root_refuses_terminally_naming_the_action() -> TestResult {
1172        // A `{` in the root would pair with the `{` the command continues
1173        // with, opening a `{{` interpolation neither of them wrote. `$` used
1174        // to sit here and no longer can: it opens nothing now, so a root
1175        // containing one is an ordinary path.
1176        let (decorated, reached, _transcript) = dispatcher_with_root(
1177            DeclaredBodyLookup::Declared(ActionBodyContract::Run {
1178                command: "echo {workspace_root}".to_owned(),
1179            }),
1180            Ok("unused".to_owned()),
1181            WorkspaceRoot::from_resolution(Ok(std::path::PathBuf::from("/absolute/with{brace"))),
1182        );
1183        let handle =
1184            tokio::task::spawn_blocking(move || decorated.dispatch(request("provision", "{}")));
1185        let Err(error) = handle.await? else {
1186            return Err("a shape-changing root must refuse a placeholder-bearing body".into());
1187        };
1188        assert!(error.starts_with("terminal:"), "{error}");
1189        assert!(
1190            error.contains("provision"),
1191            "the refusal must name the action: {error}"
1192        );
1193        assert!(
1194            error.contains("would change the parsed shape"),
1195            "the refusal must carry the shape-changing diagnosis: {error}"
1196        );
1197        assert!(
1198            reached_names(&reached).is_empty(),
1199            "a refused body must not fall through to the worker path"
1200        );
1201        Ok(())
1202    }
1203
1204    #[tokio::test(flavor = "multi_thread")]
1205    async fn an_uncreatable_root_refuses_terminally_naming_the_action() -> TestResult {
1206        // A root beneath a regular file cannot be created by any retry.
1207        let scratch = tempfile::tempdir()?;
1208        let file = scratch.path().join("occupied");
1209        std::fs::write(&file, b"not a directory")?;
1210        let (decorated, reached, _transcript) = dispatcher_with_root(
1211            DeclaredBodyLookup::Declared(ActionBodyContract::Run {
1212                command: "echo {workspace_root}".to_owned(),
1213            }),
1214            Ok("unused".to_owned()),
1215            WorkspaceRoot::from_resolution(Ok(file.join("clones"))),
1216        );
1217        let handle =
1218            tokio::task::spawn_blocking(move || decorated.dispatch(request("provision", "{}")));
1219        let Err(error) = handle.await? else {
1220            return Err("an uncreatable root must refuse a placeholder-bearing body".into());
1221        };
1222        assert!(error.starts_with("terminal:"), "{error}");
1223        assert!(
1224            error.contains("provision"),
1225            "the refusal must name the action: {error}"
1226        );
1227        assert!(
1228            error.contains("could not be created"),
1229            "the refusal must carry the creation-failure diagnosis: {error}"
1230        );
1231        assert!(
1232            reached_names(&reached).is_empty(),
1233            "a refused body must not fall through to the worker path"
1234        );
1235        Ok(())
1236    }
1237
1238    #[tokio::test(flavor = "multi_thread")]
1239    async fn a_body_without_the_placeholder_is_untouched_by_resolution_failure() -> TestResult {
1240        let (decorated, _reached, _transcript) = dispatcher_with_root(
1241            DeclaredBodyLookup::Declared(ActionBodyContract::Run {
1242                command: "echo {{greeting}}".to_owned(),
1243            }),
1244            Ok("unused".to_owned()),
1245            WorkspaceRoot::from_resolution(Err(WorkspaceRootError::Unresolvable {
1246                reason: "cannot resolve Aion home: set AION_HOME or HOME".to_owned(),
1247            })),
1248        );
1249        let handle = tokio::task::spawn_blocking(move || {
1250            decorated.dispatch(request("greet", "{\"greeting\":\"still served\"}"))
1251        });
1252        let result = handle.await?;
1253        let encoded = result.map_err(|error| format!("declared command failed: {error}"))?;
1254        let outcome: serde_json::Value = serde_json::from_str(&encoded)?;
1255        assert_eq!(outcome["stdout"], "still served");
1256        Ok(())
1257    }
1258
1259    #[tokio::test(flavor = "multi_thread")]
1260    async fn an_unreadable_catalog_delegates_to_the_worker_path() -> TestResult {
1261        let (decorated, reached) = dispatcher(
1262            DeclaredBodyLookup::Unreadable("catalog offline".to_owned()),
1263            Ok("\"served anyway\"".to_owned()),
1264        );
1265        let handle =
1266            tokio::task::spawn_blocking(move || decorated.dispatch(request("resilient", "{}")));
1267        let result = handle.await?;
1268        assert_eq!(result, Ok("\"served anyway\"".to_owned()));
1269        assert_eq!(reached_names(&reached), vec!["resilient".to_owned()]);
1270        Ok(())
1271    }
1272
1273    #[test]
1274    fn non_object_input_is_refused_terminally_by_shape() {
1275        for (input, kind) in [
1276            ("[1,2]", "an array"),
1277            ("\"text\"", "a string"),
1278            ("3", "a number"),
1279            ("null", "null"),
1280            ("true", "a boolean"),
1281        ] {
1282            let Err(error) = decode_arguments(input) else {
1283                unreachable_refusal(input);
1284                return;
1285            };
1286            assert!(error.starts_with("terminal:"), "{error}");
1287            assert!(error.contains(kind), "{error} must name {kind}");
1288        }
1289    }
1290
1291    /// Fails the calling test without a panicking accessor.
1292    fn unreachable_refusal(input: &str) {
1293        assert!(
1294            input.is_empty(),
1295            "input `{input}` must have been refused by shape"
1296        );
1297    }
1298
1299    /// The selection rule cannot be right if it is asked about the wrong run.
1300    ///
1301    /// `select_declared_body` is unit-tested on inputs the test itself
1302    /// constructs, which proves the rule and nothing about the plumbing. This
1303    /// asserts the other half: the identity the dispatcher hands the reader is
1304    /// the identity of the dispatch it is serving, not a placeholder and not
1305    /// another run's.
1306    #[tokio::test(flavor = "multi_thread")]
1307    async fn the_reader_is_asked_about_the_run_that_is_dispatching() -> TestResult {
1308        let seen = Arc::new(Mutex::new(Vec::new()));
1309        let bodies = DeclaredBodySource::default();
1310        bodies.install(Arc::new(RecordingBodies {
1311            seen: Arc::clone(&seen),
1312        }));
1313        let reached = Arc::new(Mutex::new(Vec::new()));
1314        let decorated = DeclaredCommandDispatcher::new(
1315            Arc::new(RecordingInner {
1316                reached: Arc::clone(&reached),
1317                reply: Ok("\"worker-served\"".to_owned()),
1318            }),
1319            bodies,
1320            DeclaredCommandAttempts::new(crate::shutdown::DrainState::default()),
1321            tokio::runtime::Handle::current(),
1322            WorkspaceRoot::from_resolution(Ok(std::env::temp_dir())),
1323            ActivityEventPublisher::new(
1324                Arc::new(aion_store::InMemoryObservabilityStore::default()),
1325                TRANSCRIPT_CAPACITY,
1326                crate::activity_publisher::TranscriptBatchPolicy::UNBATCHED,
1327            ),
1328        );
1329
1330        let dispatch = request("plain", "{}");
1331        let expected = (dispatch.workflow_id.clone(), dispatch.run_id.clone());
1332        let handle = tokio::task::spawn_blocking(move || decorated.dispatch(dispatch));
1333        handle
1334            .await?
1335            .map_err(|error| format!("dispatch failed: {error}"))?;
1336
1337        let observed = match seen.lock() {
1338            Ok(observed) => observed.clone(),
1339            Err(poisoned) => poisoned.into_inner().clone(),
1340        };
1341        assert_eq!(
1342            observed,
1343            vec![expected],
1344            "the body reader must be asked about the dispatching run itself"
1345        );
1346        Ok(())
1347    }
1348}