onepipeline 0.38.1

Execute a task DAG over oneagentgraph and onevcs, merging their event streams into one.
Documentation
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
645
646
647
648
649
650
651
652
653
654
655
656
657
658
659
660
661
662
663
664
665
666
667
668
669
670
671
672
673
674
675
676
677
678
679
680
681
682
683
684
685
686
687
688
689
690
691
692
693
694
695
696
697
698
699
700
701
702
703
704
705
706
707
708
709
710
711
712
713
714
715
716
717
718
719
720
721
722
723
724
725
726
727
728
729
730
731
732
733
734
735
736
737
738
739
740
741
742
743
744
745
746
747
748
749
750
751
752
753
754
755
756
757
758
759
760
761
762
763
764
765
766
767
768
769
770
771
772
773
774
775
776
777
778
779
780
781
782
783
784
785
786
787
788
789
790
791
792
793
794
795
796
797
798
799
800
801
802
803
804
805
806
807
808
809
810
811
812
813
814
815
816
817
818
819
820
821
822
823
824
825
826
827
828
829
830
831
832
833
834
835
836
837
838
839
840
841
842
843
844
845
846
847
848
849
850
851
852
853
854
855
856
857
858
859
860
861
862
863
864
865
866
867
868
869
870
871
872
873
874
875
876
877
878
879
880
881
882
883
884
885
886
887
888
889
890
891
892
893
894
895
896
897
898
899
900
901
902
903
904
905
906
907
908
909
910
911
912
913
914
915
916
917
918
919
920
921
922
923
924
925
926
927
928
929
930
931
932
933
934
935
936
937
938
939
940
941
942
943
944
945
946
947
948
949
950
951
952
953
954
955
956
957
958
959
960
961
962
//! The executor seam.
//!
//! An [`Executor`] is *where* a node's dispatch runs. v1 ships [`LocalExecutor`]
//! only — it supports both workspace variants — while the trait and the
//! [rules grammar](crate::rules) are shaped so a dispatch-server executor over a
//! WebSocket, and a Kubernetes one, drop in behind the same interface. That is
//! what decouples where a dispatch runs from the caller that asked for it.
//!
//! Two of the request's fields are a sibling library's types, so this seam is
//! also where the cross-repo wiring is proven at compile time: the agent-graph
//! config comes from `oneagentgraph` and the repository session from `onevcs`.
//! The contract first named those types `ResolvedGraphRef` and `SessionSpec`,
//! which neither sibling exports; it now names `ConfigRef` and `SessionRequest`,
//! which they do. Divergences 1 and 2 in
//! [`docs/contract-divergences.md`](../../../docs/contract-divergences.md)
//! record the ruling.

// llmlint: ignore-file[invalid_states_unrepresentable] every shape in this module is the
// one `docs/contract.md` declares in its own Rust block, character for character, and
// narrowing any of them is interface drift. That covers `Executor::name -> &str` (an
// `ExecutorName` newtype is a public item the contract does not name; the rules file
// validates the name against the declared executors), `Capabilities.vcs_sessions: bool`
// (written as `{ vcs_sessions: bool, ... }`), and `CapacityReport.load1: f64` (written as
// `{ slots_free, load1, mem_free_bytes }`, where the probe already refuses a negative or
// NaN load by never producing one).

use std::path::{Path, PathBuf};
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::Arc;

use oneagentgraph::config::ConfigRef;
use onevcs::SessionRequest;

use crate::agentgraph::{Ending, Environment, GraphOutput, GraphRun, Launch};
use crate::controls::{NodeControls, WORKER_MEMBER};
use crate::error::{Error, Result};
use crate::event::{Envelope, Labels};

/// Where a node's dispatch runs.
pub trait Executor {
    /// The name the [rules](crate::rules) file selects this executor by.
    fn name(&self) -> &str;
    /// What this executor can do.
    fn capabilities(&self) -> Capabilities;
    /// What it currently has free.
    fn capacity(&self) -> CapacityReport;
    /// Start one dispatch.
    fn dispatch(&self, req: DispatchRequest) -> Result<Box<dyn DispatchHandle>>;
}

/// What an [`Executor`] can do.
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Hash)]
pub struct Capabilities {
    /// Whether it can open a `onevcs` session — that is, whether it accepts
    /// [`WorkspaceSpec::VcsSession`] as well as [`WorkspaceSpec::Path`].
    pub vcs_sessions: bool,
}

/// What an [`Executor`] currently has free.
#[derive(Debug, Clone, Copy, Default, PartialEq)]
pub struct CapacityReport {
    /// How many more dispatches it will accept.
    pub slots_free: u32,
    /// Its one-minute load average.
    pub load1: f64,
    /// Its free memory, in bytes.
    pub mem_free_bytes: u64,
}

/// One dispatch, as an [`Executor`] is asked for it.
#[derive(Debug, Clone, PartialEq)]
pub struct DispatchRequest {
    /// The content-addressed node-scope agent-graph config, an `oneagentgraph`
    /// type.
    pub graph: ConfigRef,
    /// The task prose.
    pub task: String,
    /// Where in the run this dispatch sits. The reserved keys are `run_id`,
    /// `node`, `step`, and `persona`.
    pub labels: Labels,
    /// The per-node controls this dispatch runs under.
    ///
    /// Carried on the request rather than on the labels: a label is what an
    /// envelope is stamped with and what a `node_label` rule selects on, while a
    /// control changes the agent graph's own effective configuration. `persona`
    /// is both, and is the label, which is why it is not here.
    pub controls: NodeControls,
    /// The workspace to run in.
    pub workspace: WorkspaceSpec,
    /// Raised to stop the dispatch cooperatively.
    pub cancel: CancellationToken,
}

/// The workspace a dispatch runs in.
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum WorkspaceSpec {
    /// A directory that already exists on the machine running the dispatch.
    Path(PathBuf),
    /// A `onevcs` session the machine running the dispatch opens *there* — the
    /// clone, worktree, and branch are cut where the work happens, not shipped
    /// to it.
    VcsSession(SessionRequest),
}

/// The cooperative cancellation signal a [`DispatchRequest`] carries.
///
/// Shared rather than copied: the engine's loop raises it on one side while the
/// dispatch observes it on the other, which is what makes a `drop`, a `retry`,
/// or a `stop` end in-flight work without killing it.
#[derive(Debug, Clone, Default)]
pub struct CancellationToken(Arc<AtomicBool>);

impl CancellationToken {
    /// A signal nobody has raised.
    pub fn new() -> Self {
        Self::default()
    }

    /// Raise it.
    pub fn cancel(&self) {
        self.0.store(true, Ordering::SeqCst);
    }

    /// Whether it has been raised.
    pub fn is_cancelled(&self) -> bool {
        self.0.load(Ordering::SeqCst)
    }
}

impl PartialEq for CancellationToken {
    fn eq(&self, other: &Self) -> bool {
        self.is_cancelled() == other.is_cancelled()
    }
}

/// A started dispatch.
pub trait DispatchHandle {
    /// The envelope NDJSON it produces, relayed from wherever it runs.
    fn events(&mut self) -> EventStream;
    /// Block until it settles.
    fn wait(&mut self) -> Result<DispatchOutcome>;
    /// Stop it.
    fn cancel(&self, mode: CancelMode);
}

/// A dispatch's relayed event stream.
///
/// A boxed iterator rather than a newtype: the contract names `EventStream` as
/// `events`' return type and nothing else about it, and a newtype would need
/// constructors and accessors the contract does not name.
pub type EventStream = Box<dyn Iterator<Item = Result<Envelope>> + Send>;

/// How a dispatch is stopped.
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
pub enum CancelMode {
    /// Raise the cancellation signal and let the dispatch preserve its work.
    Cooperative,
    /// Terminate it.
    Kill,
}

/// How a dispatch settled.
///
/// Everything a caller cannot recover from the relayed event stream: whether the
/// dispatch succeeded, and — because the machine running the dispatch is the one
/// that opened the session — the session it left open for its node to publish.
/// `docs/contract.md` declares these four; divergence 3 in
/// [the divergence record](../../../docs/contract-divergences.md) is the ruling
/// that put them there, and `#[non_exhaustive]` keeps a fifth additive.
#[derive(Debug, Clone, PartialEq, Eq, Hash, Default)]
#[non_exhaustive]
pub struct DispatchOutcome {
    /// Whether the dispatch completed successfully.
    pub succeeded: bool,
    /// What it said when it did not.
    pub detail: String,
    /// The `onevcs` session token, when the workspace was a session.
    pub session: Option<String>,
    /// The branch that session has checked out.
    pub branch: Option<String>,
}

/// The executor that runs a dispatch on this machine.
///
/// The only one v1 ships, and the only one that supports both
/// [`WorkspaceSpec`] variants.
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Hash)]
pub struct LocalExecutor;

impl Executor for LocalExecutor {
    fn name(&self) -> &str {
        "local"
    }

    fn capabilities(&self) -> Capabilities {
        // The one capability the contract states for this executor: it supports
        // both workspace variants, because the machine running the dispatch is
        // this one.
        Capabilities { vcs_sessions: true }
    }

    fn capacity(&self) -> CapacityReport {
        let load1 = load_average().unwrap_or(0.0);
        let cores = std::thread::available_parallelism()
            .map(std::num::NonZeroUsize::get)
            .unwrap_or(1);
        // Every unreadable input resolves toward "has capacity": refusing to
        // dispatch on numbers nobody could measure would stall a healthy host.
        let busy = load1.ceil().max(0.0);
        let busy = if busy.is_finite() { busy as u64 } else { 0 };
        CapacityReport {
            slots_free: u32::try_from(u64::try_from(cores).unwrap_or(1).saturating_sub(busy))
                .unwrap_or(u32::MAX),
            load1,
            mem_free_bytes: available_memory().unwrap_or(u64::MAX),
        }
    }

    fn dispatch(&self, req: DispatchRequest) -> Result<Box<dyn DispatchHandle>> {
        // The run's launch record, read once at the last responsible moment: the
        // labels are what identify the run, and the overrides, the source filter
        // and the dispatch-env hook below are all what that launch decided.
        let launched = launched_with(&req.labels)?;
        // Relayed: this dispatch is read turn by turn into the merged store.
        let node_sets = node_sets(launched.as_ref(), &req.labels, &req.controls)?;
        // The dispatch-env hook, **before** anything of the launch begins — the
        // session below included, so a launch the hook refuses cuts nothing —
        // and never for a dispatch built outside a run, which has no launch to
        // name one. What it adds is this child's alone: it goes into `env` below
        // and nowhere else, and this process's own environment is untouched.
        let added = match (
            &launched,
            req.labels.run_id.as_deref(),
            req.labels.node.as_deref(),
        ) {
            (Some(record), Some(run), Some(node)) => {
                let paths = crate::ledger::RunPaths::under(&crate::ledger::runs_root(), run);
                crate::dispatchenv::run_and_check(&crate::dispatchenv::Launching {
                    paths: &paths,
                    record,
                    node,
                    graph: &req.graph,
                    sets: &node_sets,
                    own: &own_dispatch_variables(&req.workspace),
                })?
            }
            _ => Vec::new(),
        };
        // `WorkspaceSpec::VcsSession` means the machine running the dispatch
        // opens the session *there* — the clone, worktree, and branch are cut
        // where the work happens rather than shipped to it. This executor is
        // that machine, so it opens the session itself and runs in the worktree
        // `onevcs` hands back.
        let (dir, session) = match &req.workspace {
            WorkspaceSpec::Path(path) => (path.clone(), None),
            WorkspaceSpec::VcsSession(request) => {
                let session = crate::vcs::session_open(request)?;
                remember_worktree(&session);
                (session.worktree.clone(), Some(session))
            }
        };
        // The session this dispatch works in: the one just opened, or — for every
        // later dispatch of the same node, which names the worktree rather than
        // asking for a session of its own — the one this executor opened there.
        let token = session
            .as_ref()
            .map(|session| session.token.0.clone())
            .or_else(|| session_of_worktree(&dir));
        // Every node-scope launch a run starts is one of that run's
        // `oneagentgraph` sources, so it carries the same source filter the
        // observer graph does.
        let filters = launched.map(|record| record.filters).unwrap_or_default();
        // The hook's additions first and this crate's own keys after them, so a
        // hook cannot move where a dispatch keeps its scratch or which run it
        // belongs to: a later pair of the same name is the one the child gets.
        let mut env = added;
        env.extend(prepare_dispatch_env(&req.labels, token.as_deref())?);
        let mut run = GraphRun::start(&Launch {
            graph: &req.graph.0,
            task: &req.task,
            dir: &dir,
            labels: &req.labels,
            env: &env,
            // The scratch directory in `env` is this dispatch's and no other's,
            // so the launch is given a process to hold it in: the library
            // backend has nowhere per-launch to put a pair, and two dispatches
            // sharing one driver would read and overwrite each other's.
            environment: Environment::PerLaunch,
            sets: &node_sets,
            filter: filters.agentgraph.as_ref(),
            output: GraphOutput::Relayed,
            // The dispatch settles on the terminal envelope this launch relays,
            // which is the answer the graph gives before its own final teardown
            // — so the launch is held for the sibling's write of its record,
            // which its `history` lists the run by, and not for the reap after.
            ending: Ending::Announced,
        })?;
        // The run's registry of what it is running, and where. Recorded here
        // because this is the layer that knows: the executor is *where a
        // dispatch runs*, so the process the work is in is its answer to give
        // and nobody else's — an executor that ran the dispatch on another
        // machine would have no local process to name, and would say so by
        // recording nothing.
        //
        // A dispatch this run cannot register does not run. The registry is the
        // only record of where the work is, so an unregistered dispatch is a
        // process no view will show and no `stop` will reach — work that can only
        // be found by a person reading a process table, on a run whose own
        // records say it has nothing running. So the graph that has just started
        // is taken back down and the failure is the caller's: a dispatch that
        // could not start is an outcome this seam already has, and the engine
        // retries it and settles the node saying so.
        let claim = match register_dispatch(&req.labels, run.process()) {
            Ok(claim) => claim,
            Err(refusal) => {
                // Ended and collected, not merely signalled: what this returns
                // to the caller is that the dispatch is not running, and a
                // process nobody has waited on is a zombie — which answers a
                // liveness probe as alive and would leave the very row an
                // operator would go looking for.
                run.cancel();
                let _ = run.wait();
                return Err(refusal);
            }
        };
        Ok(Box::new(LocalDispatch {
            run,
            cancel: req.cancel,
            labels: req.labels,
            session,
            _claim: claim,
        }))
    }
}

/// Where one dispatch may write whatever it likes.
///
/// An **absolute** path to a directory this crate created, that exists and is
/// writable before the dispatch's first turn, that is unique to that dispatch —
/// a retry, a requeue and a resumed pin of the same node each get their own — and
/// that nothing here removes while the dispatch is running. Nothing more is
/// promised: the spelling below is not a contract and no consumer may derive one
/// path from another.
///
/// Divergence 48 in
/// [the divergence record](../../../docs/contract-divergences.md) is why, and the
/// proposal this answers.
pub(crate) const NODE_SCRATCH_DIR_ENV: &str = "ONEPIPELINE_NODE_SCRATCH_DIR";

/// The sessions this executor opened, by the worktree each handed back.
///
/// A lifecycle node's first dispatch asks for a session and every later one —
/// its remaining steps and its drafting dispatch — names the worktree that
/// session opened, because a second session on the same branch would reclaim
/// the first. The token is not on that request: [`WorkspaceSpec::Path`] is a
/// directory and the contract fixes it as one. So the executor that opened the
/// session is what remembers which one, which is the same fact
/// [`WorkspaceSpec::VcsSession`] states — the machine running the dispatch is
/// the one that opened the session there. A path nothing here opened answers
/// nothing, which is every direct node's dispatch.
fn opened_worktrees() -> std::sync::MutexGuard<'static, std::collections::BTreeMap<PathBuf, String>>
{
    static OPENED: std::sync::Mutex<std::collections::BTreeMap<PathBuf, String>> =
        std::sync::Mutex::new(std::collections::BTreeMap::new());
    // A poisoned lock holds a map a panicking thread was mid-insert into, which
    // is still a map of sessions this process opened.
    OPENED
        .lock()
        .unwrap_or_else(std::sync::PoisonError::into_inner)
}

fn remember_worktree(session: &onevcs::Session) {
    opened_worktrees().insert(session.worktree.clone(), session.token.0.clone());
}

fn session_of_worktree(worktree: &Path) -> Option<String> {
    opened_worktrees().get(worktree).cloned()
}

/// The names of the variables [`prepare_dispatch_env`] will set on a dispatch in
/// `workspace`, known before any of their values are.
///
/// What the dispatch-env hook's check is told is present by name: the pairs
/// themselves are composed only once the launch goes ahead — the scratch
/// directory is made, and a session's token exists once the session is open —
/// and a config that sources one of this crate's own variables through
/// `env_from` is a launch that will have it. The session is the one variable
/// not every dispatch carries, and whether this one will is decided here the
/// way [`dispatch`](LocalExecutor::dispatch) decides it: a session it opens, or
/// one already opened on the worktree it names.
fn own_dispatch_variables(workspace: &WorkspaceSpec) -> Vec<&'static str> {
    let mut names = vec![
        crate::agentgraph::RUN_ID_ENV,
        crate::agentgraph::RUNS_DIR_ENV,
        crate::channel::ASKER_ENV,
        NODE_SCRATCH_DIR_ENV,
    ];
    let in_session = match workspace {
        WorkspaceSpec::VcsSession(_) => true,
        WorkspaceSpec::Path(path) => session_of_worktree(path).is_some(),
    };
    if in_session {
        names.push(crate::agentgraph::SESSION_ENV);
    }
    names
}

/// Compose what every dispatch this executor makes carries in its own
/// environment, **making** the scratch directory one of those pairs names.
///
/// The **run id** is what the operator's `ask-manager` wrapper addresses a
/// manager by, and a dispatch outside a run carries none for the same reason it
/// registers nothing. The **runs root** goes beside it, absolute, so a dispatch
/// that runs `onepipeline transcript` reads this run's store from wherever it
/// is working. The **session** is the `onevcs` token whose worktree the dispatch
/// runs in, so a worker can address its own session; a dispatch in no session
/// carries none. The **scratch directory** is this dispatch's alone, which
/// is why the launch below declares [`Environment::PerLaunch`]: the pair has to
/// live somewhere no sibling dispatch can read or overwrite, and that is a
/// process rather than a map. The **asker** is that same uniqueness read as an
/// identity: the wrapper above asks through a succession of the host bus server
/// listeners, and this is what tells them they are serving one side that is
/// still waiting rather than a series of sides that have each gone.
///
/// # Errors
///
/// [`Error::Ledger`] where the scratch directory cannot be made: a promised
/// directory that is not there would fail the agent's writes one at a time, and
/// those failures read as the agent's own work going wrong. And where the runs
/// root cannot be resolved to an absolute path, for the reason stated at that
/// call.
fn prepare_dispatch_env(labels: &Labels, session: Option<&str>) -> Result<Vec<(String, String)>> {
    let mut env: Vec<(String, String)> = labels
        .run_id
        .iter()
        .map(|run| (crate::agentgraph::RUN_ID_ENV.to_string(), run.clone()))
        .collect();
    // Absolute, because the dispatch does not run where this process was
    // started: the default root is the relative `runs`, which from inside a
    // worktree names a directory that is not there. A resolution that failed —
    // this process's own working directory unreadable, which is what `absolute`
    // consults for a relative root — is **refused** rather than fallen back
    // from, because exporting the relative root anyway hands every dispatch a
    // path the contract promises is absolute and that resolves, in the worktree,
    // somewhere else or nowhere.
    let root = crate::ledger::runs_root();
    let root = std::path::absolute(&root).map_err(|source| Error::Ledger {
        path: root.clone(),
        source,
    })?;
    env.push((
        crate::agentgraph::RUNS_DIR_ENV.to_string(),
        root.display().to_string(),
    ));
    if let Some(token) = session {
        env.push((crate::agentgraph::SESSION_ENV.to_string(), token.to_owned()));
    }
    let scratch = make_node_scratch_dir(labels)?.display().to_string();
    // The asker's name is the scratch directory's own path rather than a second
    // thing minted beside it: what has to be true of it is that every session
    // this dispatch serves through carries the same value and no other dispatch
    // carries it, and that is exactly what the directory above already is. It is
    // read as an opaque word and never as a path — see `channel::ASKER_ENV`.
    env.push((crate::channel::ASKER_ENV.to_string(), scratch.clone()));
    env.push((NODE_SCRATCH_DIR_ENV.to_string(), scratch));
    Ok(env)
}

/// Make this dispatch's own scratch directory, and answer where it is.
///
/// Named for the making, because that is the whole of it: this mints a number no
/// caller has had, creates a directory under it, and answers the path — so two
/// calls with identical arguments answer differently, and neither answer existed
/// before the call.
///
/// Under the run's own directory, so a run's scratch is thrown away with the run.
///
/// Uniqueness is the directory's *creation*, not its name: `create_dir` refuses
/// one that is already there, which a name minted from a pid and a counter would
/// not — a host reissues pids and a counter starts again in every process.
// llmlint: ignore-block[changed_behavior_has_e2e] no command reaches the two arms below:
// every dispatch a run makes carries its id, and a `scratch` that will not be created
// needs a run directory that exists holding a file by that name. Both are driven against
// the real filesystem by
// `tests::every_dispatch_is_given_a_directory_of_its_own_and_no_two_share_one`.
fn make_node_scratch_dir(labels: &Labels) -> Result<PathBuf> {
    /// Enough numbers that walking past every directory a run has already made is
    /// never the reason a dispatch fails, and few enough that a base directory
    /// nothing can be created in fails rather than spinning.
    const TRIES: u64 = 4096;
    static MINTED: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0);

    let base = match labels.run_id.as_deref() {
        Some(run) => crate::ledger::RunPaths::under(&crate::ledger::runs_root(), run)
            .dir
            .join("scratch"),
        None => std::env::temp_dir().join("onepipeline-scratch"),
    };
    let ledger = |path: &Path| {
        let path = path.to_path_buf();
        move |source: std::io::Error| Error::Ledger { path, source }
    };
    std::fs::create_dir_all(&base).map_err(ledger(&base))?;
    let pid = crate::sys::pid();
    for _ in 0..TRIES {
        let at = base.join(format!(
            "{pid}-{}",
            MINTED.fetch_add(1, std::sync::atomic::Ordering::Relaxed)
        ));
        match std::fs::create_dir(&at) {
            // Absolute, because the value is read by a program whose working
            // directory is its own business: a relative runs root — the default
            // is one — would name a different place from the workspace a
            // dispatch runs in.
            Ok(()) => return std::fs::canonicalize(&at).map_err(ledger(&at)),
            Err(error) if error.kind() == std::io::ErrorKind::AlreadyExists => {}
            Err(error) => return Err(ledger(&at)(error)),
        }
    }
    Err(Error::Ledger {
        path: base.clone(),
        source: std::io::Error::new(
            std::io::ErrorKind::AlreadyExists,
            format!(
                "no scratch directory under {} could be created",
                base.display()
            ),
        ),
    })
} // llmlint: ignore-end[changed_behavior_has_e2e]

/// Record this dispatch in its run's registry, and hold the entry open.
///
/// `process` is the graph run's own, where the graph is a process this crate
/// started; a graph running **in this process** is recorded as this process,
/// which is the true answer to where that dispatch's work is and the one a
/// teardown would have to aim at.
///
/// A dispatch outside a run records nothing and is not refused for it: the
/// contract's own example and the seam's tests carry no `run_id`, and there is no
/// registry for a run that does not exist. So is one whose node the labels do not
/// name — an entry that could not say which node it belonged to would be a pid an
/// operator could not act on. Every dispatch a *run* makes carries both.
fn register_dispatch(
    labels: &Labels,
    process: Option<u32>,
) -> Result<Option<crate::ledger::DispatchClaim>> {
    let (Some(run), Some(node)) = (labels.run_id.as_deref(), labels.node.as_deref()) else {
        return Ok(None);
    };
    let paths = crate::ledger::RunPaths::under(&crate::ledger::runs_root(), run);
    crate::ledger::claim_dispatch(&paths, node, process.unwrap_or_else(crate::sys::pid)).map(Some)
}

/// The overrides one dispatch's graph launch carries, in the order they apply.
///
/// The run's opaque node-scope overrides are read at the last responsible
/// moment — the labels already identify the launch ledger for every local
/// dispatch — and the node's own settings are applied *after* them: an operator's
/// `--node-set` is run-wide, and a control the plan wrote against one node is the
/// more specific of the two.
///
/// **None of them for the drafting dispatch.** `--node-set` is forwarded to every
/// *node-scope* launch, and the persona override names `members.worker`, which is
/// the member of the node-scope graph this crate composes: a run's pr-author
/// graph is the operator's whole statement about how a change request is
/// drafted, and it declares its own members under its own names. Composing
/// either onto it refuses the launch — `this graph has no worker` — which is a
/// drafting dispatch that could never start.
///
/// The **persona** is what tells the two apart, and it can be: `pr-author` is
/// this crate's own, and a plan naming it for a node or a step is refused where
/// the plan is read — see [`RESERVED_PERSONA`](crate::graph::RESERVED_PERSONA) —
/// so a dispatch arriving here under it is the drafting one and nothing else. A
/// second condition on the graph would not narrow that: an operator may point
/// `--pr-author-graph` at the same document a node dispatches under, and then
/// the graph says nothing about which dispatch this is.
fn node_sets(
    launched: Option<&crate::ledger::LaunchRecord>,
    labels: &Labels,
    controls: &NodeControls,
) -> Result<Vec<String>> {
    if labels.persona.as_deref() == Some(crate::lifecycle::PR_AUTHOR_PERSONA) {
        return Ok(Vec::new());
    }
    let mut sets = launched.map_or_else(Vec::new, |record| record.node_sets.clone());
    if let Some(persona) = &labels.persona {
        sets.push(format!("members.{WORKER_MEMBER}.persona={persona}"));
    }
    // A control this build cannot apply refuses the launch here as well as at
    // validation, so no path composes a launch that drops one on the floor.
    sets.extend(controls.overrides().map_err(Error::Invalid)?);
    Ok(sets)
}

/// The launch record of the run this dispatch belongs to, when it belongs to one.
///
/// A dispatch built outside a run — the contract's own example, and the seam's
/// tests — carries no `run_id` and so has no launch to read: it takes the
/// defaults rather than being refused, because nothing about it is wrong.
fn launched_with(labels: &Labels) -> Result<Option<crate::ledger::LaunchRecord>> {
    let Some(run) = labels.run_id.as_deref() else {
        return Ok(None);
    };
    let paths = crate::ledger::RunPaths::under(&crate::ledger::runs_root(), run);
    crate::ledger::read_json::<crate::ledger::LaunchRecord>(&paths.launch()).map(Some)
}

/// One dispatch running on this machine.
#[derive(Debug)]
struct LocalDispatch {
    run: GraphRun,
    cancel: CancellationToken,
    labels: Labels,
    session: Option<onevcs::Session>,
    /// This dispatch's entry in the run's registry, held for exactly as long as
    /// the dispatch is: dropping the handle — settled, failed, cancelled,
    /// retried — takes the entry with it, so the registry holds live dispatches
    /// and nothing else.
    ///
    /// Underscored because nothing reads it and nothing should: what it does, it
    /// does by existing and then not.
    _claim: Option<crate::ledger::DispatchClaim>,
}

impl DispatchHandle for LocalDispatch {
    fn events(&mut self) -> EventStream {
        let opened = self.session.as_ref().map(|session| {
            // The opened session is `onevcs`'s own contribution to the merged
            // stream: without it a lifecycle node's branch would appear in the
            // ledger with nothing saying where it came from.
            Ok(crate::vcs::session_opened_event(session, &self.labels))
        });
        match opened {
            Some(event) => Box::new(std::iter::once(event).chain(self.run.events())),
            None => self.run.events(),
        }
    }

    fn wait(&mut self) -> Result<DispatchOutcome> {
        let settled = self.run.wait()?;
        Ok(DispatchOutcome {
            succeeded: settled.succeeded(),
            detail: settled.stderr.trim().to_string(),
            session: self.session.as_ref().map(|s| s.token.0.clone()),
            branch: self.session.as_ref().map(|s| s.branch.clone()),
        })
    }

    /// Stop this dispatch, as far as the mode asks.
    ///
    /// Both modes raise the cooperative signal, because both are the caller
    /// changing its mind. `Kill` additionally tears the graph run down —
    /// `GraphRun::cancel` acts on either backend and reaps the process tree —
    /// which is what a `Cooperative` stop deliberately does not do: the engine
    /// asks the live turn to commit and end first, and escalates to this only
    /// when the dispatch has not exited by its deadline.
    fn cancel(&self, mode: CancelMode) {
        self.cancel.cancel();
        if mode == CancelMode::Kill {
            self.run.cancel();
        }
    }
}

/// This host's one-minute load average, where it can be read.
fn load_average() -> Option<f64> {
    let text = std::fs::read_to_string("/proc/loadavg").ok()?;
    text.split_whitespace()
        .next()?
        .parse::<f64>()
        .ok()
        .filter(|value| value.is_finite() && *value >= 0.0)
}

/// This host's available memory in bytes, where it can be read.
fn available_memory() -> Option<u64> {
    let text = std::fs::read_to_string("/proc/meminfo").ok()?;
    for line in text.lines() {
        if let Some(rest) = line.strip_prefix("MemAvailable:") {
            let kib = rest.split_whitespace().next()?.parse::<u64>().ok()?;
            return kib.checked_mul(1024);
        }
    }
    None
}

#[cfg(test)]
mod tests {
    use super::*;

    #[test]
    fn the_local_executor_is_named_and_capable_of_both_workspaces() {
        let executor = LocalExecutor;
        assert_eq!(executor.name(), "local");
        assert!(executor.capabilities().vcs_sessions);
    }

    #[test]
    fn the_capacity_probe_reports_finite_numbers_on_any_host() {
        let report = LocalExecutor.capacity();
        assert!(
            report.load1.is_finite() && report.load1 >= 0.0,
            "{report:?}"
        );
        assert!(report.mem_free_bytes > 0, "{report:?}");
    }

    #[test]
    fn a_cancellation_signal_is_shared_between_the_two_sides() {
        let token = CancellationToken::new();
        let observer = token.clone();
        assert!(!observer.is_cancelled());
        token.cancel();
        assert!(
            observer.is_cancelled(),
            "the signal did not reach the dispatch"
        );
        assert_eq!(token, observer);
        assert_ne!(CancellationToken::new(), observer);
    }

    #[test]
    fn a_dispatch_request_carries_both_siblings_types() {
        // The seam's whole point: this fails to compile if either sibling's
        // vocabulary drifts out from under it.
        let request = DispatchRequest {
            graph: ConfigRef("./graphs/node-scope.yaml".into()),
            task: "## What\ndo it".into(),
            labels: Labels::default(),
            controls: NodeControls::default(),
            workspace: WorkspaceSpec::VcsSession(SessionRequest {
                repo: "owner/repo".into(),
                branch: None,
                base: None,
                execution_checkout: None,
            }),
            cancel: CancellationToken::new(),
        };
        assert!(matches!(request.workspace, WorkspaceSpec::VcsSession(_)));
        assert_eq!(request.graph.0, "./graphs/node-scope.yaml");
    }

    /// Serialises the tests below that set `RUNS_DIR_ENV`, which belongs to the
    /// whole process rather than to the test that set it.
    ///
    /// nextest gives each test its own process; plain `cargo test` runs a
    /// module's tests as *threads of one process*, where both tests below would
    /// otherwise read whichever value the other set last. The lock costs nothing
    /// under nextest and makes both runners say the same thing.
    static RUNS_DIR: std::sync::Mutex<()> = std::sync::Mutex::new(());

    /// Held for the length of a test that sets `RUNS_DIR_ENV`. A poisoned lock
    /// is recovered rather than propagated: the test that panicked holding it
    /// has already failed, and refusing to run the next one would report a
    /// second failure belonging to nobody.
    fn runs_dir_lock() -> std::sync::MutexGuard<'static, ()> {
        RUNS_DIR
            .lock()
            .unwrap_or_else(|poisoned| poisoned.into_inner())
    }

    /// A temporary root nobody else has, keyed on the *test* rather than on the
    /// process it runs in. Two tests running as threads share a pid; the counter
    /// is what they do not share.
    fn scratch_root(what: &str) -> PathBuf {
        static NTH: std::sync::atomic::AtomicU32 = std::sync::atomic::AtomicU32::new(0);
        let root = std::env::temp_dir().join(format!(
            "onepipeline-{what}-{}-{}",
            crate::sys::pid(),
            NTH.fetch_add(1, std::sync::atomic::Ordering::SeqCst)
        ));
        let _ = std::fs::remove_dir_all(&root);
        root
    }

    /// The drafting dispatch takes the graph the launch named as it was written.
    ///
    /// Neither half of the node-scope composition is a statement about it: the
    /// persona override names a member only the node-scope graph has, and
    /// `--node-set` is forwarded to node-scope launches. Both are dropped on the
    /// persona alone, which is a name a plan may not claim — so this is the one
    /// dispatch that reaches it.
    #[test]
    fn the_drafting_dispatch_composes_nothing_onto_the_graph_the_launch_named() {
        let _runs_dir = runs_dir_lock();
        let root = scratch_root("drafting");
        let paths = crate::ledger::RunPaths::under(&root, "demo");
        paths.create().expect("the run directory");
        let record = r#"{"run_id":"demo","plan":"p.json","node_graph":"./node.yaml",
            "pr_author_graph":"./author.yaml","launcher":"l","session":"s","pid":1,
            "host":"h","started_at":"now","heartbeat_interval":1,
            "node_sets":["members.worker.model=m"]}"#;
        std::fs::write(paths.launch(), record).expect("the launch record is written");
        std::env::set_var(crate::ledger::RUNS_DIR_ENV, &root);

        let sets = |persona: &str| {
            let labels = Labels {
                run_id: Some("demo".into()),
                persona: Some(persona.into()),
                ..Labels::default()
            };
            let launched = launched_with(&labels).expect("the launch record is readable");
            node_sets(launched.as_ref(), &labels, &NodeControls::default())
                .expect("the overrides compose")
        };
        assert!(
            sets(crate::lifecycle::PR_AUTHOR_PERSONA).is_empty(),
            "the drafting dispatch was given a member this graph never declared"
        );
        // The node's own work, under the same run and the same record.
        assert_eq!(
            sets("engineer"),
            vec![
                "members.worker.model=m".to_string(),
                "members.worker.persona=engineer".to_string(),
            ]
        );
        std::env::remove_var(crate::ledger::RUNS_DIR_ENV);
        let _ = std::fs::remove_dir_all(&root);
    }

    /// Every dispatch is given a directory of its own, and no two are given one.
    ///
    /// The end-to-end halves are `scratch::a_dispatch_is_given_an_absolute_writable_directory_of_its_own`
    /// and `scratch::every_dispatch_of_one_node_is_given_its_own_directory_and_none_is_taken_away`,
    /// which read the value out of a real dispatch's own environment. What is held
    /// here is the promise itself, against the real filesystem: two dispatches of
    /// one node — the pair a retry produces, and the pair that agree on every name
    /// a path could have been derived from — and a run root that cannot hold a
    /// scratch directory at all.
    #[test]
    fn every_dispatch_is_given_a_directory_of_its_own_and_no_two_share_one() {
        let _runs_dir = runs_dir_lock();
        let root = scratch_root("scratch");
        std::env::set_var(crate::ledger::RUNS_DIR_ENV, &root);
        let labels = Labels {
            run_id: Some("demo".into()),
            node: Some("build".into()),
            ..Labels::default()
        };

        let scratch = |labels: &Labels| {
            let env =
                prepare_dispatch_env(labels, None).expect("the dispatch's environment is composed");
            let (_, value) = env
                .iter()
                .find(|(key, _)| key == NODE_SCRATCH_DIR_ENV)
                .expect("every dispatch carries a scratch directory")
                .clone();
            PathBuf::from(value)
        };

        // The same node, twice, which is what a retry is.
        let first = scratch(&labels);
        let second = scratch(&labels);
        assert_ne!(
            first, second,
            "a node asked again was handed the directory its first attempt had"
        );
        for at in [&first, &second] {
            assert!(at.is_absolute(), "{} is not absolute", at.display());
            assert!(at.is_dir(), "{} was not created", at.display());
            std::fs::write(at.join("written"), "by the dispatch")
                .unwrap_or_else(|error| panic!("{} is not writable: {error}", at.display()));
        }
        // And the first is untouched by the second, which is the whole of what
        // "unique to that dispatch" buys.
        assert!(first.join("written").is_file());

        // A dispatch outside a run has no run directory to sit in and is given one
        // anyway: the contract's own example carries no `run_id`.
        assert!(scratch(&Labels::default()).is_dir());

        // A run root that is a file holds no scratch directory, and the dispatch is
        // refused rather than handed a path to nothing.
        let blocked = root.join("blocked");
        std::fs::write(&blocked, "not a directory").expect("the blocking file is written");
        std::env::set_var(crate::ledger::RUNS_DIR_ENV, &blocked);
        assert!(matches!(
            prepare_dispatch_env(&labels, None),
            Err(Error::Ledger { .. })
        ));

        std::env::remove_var(crate::ledger::RUNS_DIR_ENV);
        let _ = std::fs::remove_dir_all(&root);
    }

    /// A runs root that cannot be made absolute refuses the dispatch.
    ///
    /// `ONEPIPELINE_RUNS_DIR` is promised absolute, and a dispatch reads it from
    /// a worktree rather than from where this process started — so a relative
    /// root exported anyway names a different directory there, or none. The one
    /// way the resolution fails for a relative root is this process's own
    /// working directory being unreadable, which is what `std::path::absolute`
    /// consults; that is induced here by removing it.
    ///
    /// The working directory belongs to the whole process, so this holds the
    /// same lock the root-setting tests do and puts it back before releasing it.
    /// nextest gives each test its own process and pays nothing for either.
    #[test]
    #[cfg(unix)]
    fn a_runs_root_that_cannot_be_made_absolute_refuses_the_dispatch() {
        let _runs_dir = runs_dir_lock();
        let labels = Labels {
            run_id: Some("demo".into()),
            node: Some("build".into()),
            ..Labels::default()
        };
        // Relative, which is the only kind whose resolution consults the working
        // directory — and the shipped default is one.
        std::env::set_var(crate::ledger::RUNS_DIR_ENV, "runs");
        let here = std::env::current_dir().expect("this process has a working directory");
        let gone = scratch_root("cwd-gone");
        std::fs::create_dir_all(&gone).expect("the directory to stand in is made");
        std::env::set_current_dir(&gone).expect("this process can stand in it");
        std::fs::remove_dir(&gone).expect("and it can be taken away underneath");

        let refused = prepare_dispatch_env(&labels, None);

        std::env::set_current_dir(&here).expect("the working directory is put back");
        std::env::remove_var(crate::ledger::RUNS_DIR_ENV);

        match refused {
            Err(Error::Ledger { path, .. }) => assert_eq!(
                path,
                PathBuf::from("runs"),
                "the refusal does not name the root that could not be resolved"
            ),
            other => panic!(
                "a runs root that cannot be made absolute was accepted rather than refused: \
                 {other:?}"
            ),
        }
    }

    #[test]
    fn a_dispatch_with_no_run_still_carries_its_nodes_own_controls() {
        // No `run_id`, so there is no launch record to read: the node's own
        // budget is what the launch must still carry, because a dispatch that
        // dropped it here would run to the base config's default instead.
        let sets = node_sets(
            None,
            &Labels {
                persona: Some("engineer".into()),
                ..Labels::default()
            },
            &NodeControls {
                max_turns: std::num::NonZeroU32::new(45),
            },
        )
        .expect("both are appliable");
        assert_eq!(
            sets,
            vec![
                "members.worker.persona=engineer".to_string(),
                "members.worker.max_turns=45".to_string(),
            ],
            "the node's own control must apply after the run-wide ones"
        );
    }
}