lgwks_bot 2.2.0

Capability-gated automation bots on a change-detecting ECS schedule: Observe, Evaluate, Execute, and Query, with an async runtime facade.
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
963
964
965
966
967
968
969
970
971
972
973
974
975
976
977
978
979
980
981
982
983
984
985
986
987
988
989
990
991
992
993
994
995
996
997
998
999
1000
1001
1002
1003
1004
1005
1006
1007
1008
1009
1010
1011
1012
1013
1014
1015
1016
1017
1018
1019
1020
1021
1022
1023
1024
1025
1026
1027
1028
1029
1030
1031
1032
1033
1034
1035
1036
1037
1038
1039
1040
1041
1042
1043
1044
1045
1046
1047
1048
1049
1050
1051
1052
1053
1054
1055
1056
1057
1058
//! The environment broker: registering environments, fencing the generations
//! they go through, and handing out authority for one exact attempt.
//!
//! [`crate::journal`] records what happened. This decides what may happen at
//! all. A durable record of a dispatch is worth little if the dispatch was
//! aimed at an environment that had already been replaced: the record would be
//! accurate and the effect would still have landed somewhere nobody expected.
//!
//! # Fencing is a comparison, not a flag
//!
//! An [`EffectKey`] names an environment and the generation it was prepared
//! against. Authority is minted only when that generation is the one the broker
//! currently holds, so a command prepared before a replacement cannot execute
//! after it. Two failure directions are distinguished rather than collapsed:
//!
//! - [`BrokerError::Superseded`] when the key names an *older* generation. This
//!   is the real case, a command that was correct and is now stale.
//! - [`BrokerError::NeverIssued`] when the key names a *newer* one. The broker
//!   never minted it, so the key was not produced by this broker at all. Folded
//!   into one "generation mismatch" error these would print identically while
//!   calling for opposite investigations.
//!
//! Minting is one instant of that comparison, and one instant is not a
//! guarantee. [`Broker::revalidate`] runs the identical check again at the
//! handoff, because a replacement landing between the mint and the handoff would
//! otherwise leave a legitimately minted warrant dispatchable.
//!
//! # What the broker owns
//!
//! The environment identity itself comes from the host, because the host is what
//! has the entropy and what actually created the process or container. What the
//! broker owns is the environment's lifetime and its generation: once registered
//! an environment is this broker's to replace and to close, and a closed
//! environment mints no further authority. That is the resource contract this
//! module declares, and it is why registration is explicit rather than a
//! counter the broker invents.
//!
//! [`Broker::authorize`]: crate::broker::Broker::authorize
//! [`Broker::revalidate`]: crate::broker::Broker::revalidate
//! [`BrokerError::NeverIssued`]: crate::broker::BrokerError::NeverIssued
//! [`BrokerError::Superseded`]: crate::broker::BrokerError::Superseded
//! [`DurableAck`]: crate::journal::DurableAck
//! [`EffectJournal`]: crate::journal::EffectJournal
//! [`EffectKey`]: crate::effect::EffectKey

use core::fmt;
use core::num::NonZeroU64;
use std::collections::HashMap;

use crate::effect::{EffectKey, EnvironmentEpoch, EnvironmentId};
use crate::journal::{DurableAck, EffectEvent, EffectJournal, JournalError, JournalPosition};

/// What the broker knows about one environment.
#[derive(Debug, Clone, Copy)]
struct Environment {
    /// The generation currently authorized to receive commands.
    epoch: EnvironmentEpoch,
    /// Whether this broker still owns the environment. A closed environment
    /// mints no further authority, which is the cleanup half of the resource
    /// contract.
    open: bool,
}

/// What the broker can refuse.
#[derive(Debug)]
#[non_exhaustive]
pub enum BrokerError {
    /// The environment was already registered with this broker.
    ///
    /// Refused rather than treated as idempotent, because the two cases a
    /// silent re-register would merge are "the host is setting this up" and
    /// "the host lost track of what it already set up", and only the first
    /// should reset nothing.
    AlreadyRegistered {
        /// The environment that was already known.
        id: EnvironmentId,
    },
    /// The broker has never heard of the environment.
    UnknownEnvironment {
        /// The environment named by the key.
        id: EnvironmentId,
    },
    /// The environment has been closed.
    Closed {
        /// The environment that was closed.
        id: EnvironmentId,
    },
    /// The command names a generation that has been replaced.
    ///
    /// A command prepared against an older generation is refused here rather
    /// than dispatched, which is the whole point of fencing.
    Superseded {
        /// The environment the command was prepared for.
        environment: EnvironmentId,
        /// The generation the command names.
        presented: EnvironmentEpoch,
        /// The generation the broker currently holds.
        current: EnvironmentEpoch,
    },
    /// The command names a generation this broker never issued.
    ///
    /// Distinct from [`Self::Superseded`] on purpose: a stale command means one
    /// investigation and a generation from nowhere means another, and one error
    /// for both would print the same sentence for each.
    NeverIssued {
        /// The environment the command was prepared for.
        environment: EnvironmentId,
        /// The generation the command names.
        presented: EnvironmentEpoch,
        /// The generation the broker currently holds.
        current: EnvironmentEpoch,
    },
    /// The journal names an environment this broker was not asked to adopt.
    ///
    /// Its own arm rather than a fold into a read failure: adopting the wrong
    /// environment would fence against a generation counter that has nothing to
    /// do with the one the previous owner moved, which is a silent downgrade
    /// wearing a successful return.
    ForeignEnvironment {
        /// The environment the caller asked to adopt.
        asked: EnvironmentId,
        /// The environment the journal's events name.
        named: EnvironmentId,
    },
    /// There is nothing in the journal to take over.
    ///
    /// An empty history cannot say which generation the previous owner was at, so
    /// claiming one would be a guess dressed as a fence. A caller that wants a
    /// fresh environment over an empty journal wants [`Broker::register`], and
    /// the difference between the two is exactly this refusal.
    NothingToAdopt {
        /// The environment with no committed history.
        id: EnvironmentId,
    },
    /// The journal could not be read, so its generation is unknown.
    ///
    /// An error is never absence (INV-BOT-7): a journal that cannot answer is not
    /// a journal with nothing in it, and adopting on the strength of an unread
    /// file is how a takeover silently becomes a fresh start.
    Journal(crate::journal::JournalError),
    /// The environment has no generations left.
    ///
    /// Reachable only after 2^64 replacements. Named rather than wrapped,
    /// because a wrapped generation would re-authorize commands the broker
    /// already fenced.
    Exhausted {
        /// The environment whose generation space is spent.
        id: EnvironmentId,
    },
}

impl fmt::Display for BrokerError {
    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
        // The scrutinee is `*self` so each pattern's type is the enum's own type
        // rather than a reference to it. `clippy::pattern_type_mismatch` is
        // forbidden in this workspace.
        match *self {
            Self::AlreadyRegistered { id } => {
                write!(f, "environment {id} is already registered with this broker")
            }
            Self::UnknownEnvironment { id } => {
                write!(f, "this broker has no environment {id}")
            }
            Self::Closed { id } => {
                write!(f, "environment {id} is closed and mints no authority")
            }
            Self::Superseded {
                environment,
                presented,
                current,
            } => write!(
                f,
                "a command for {environment} at generation {presented} names a \
                 generation that has been replaced; the broker now holds {current}"
            ),
            Self::NeverIssued {
                environment,
                presented,
                current,
            } => write!(
                f,
                "a command for {environment} names generation {presented}, which \
                 this broker never issued; it holds {current}"
            ),
            Self::ForeignEnvironment { asked, named } => write!(
                f,
                "the journal describes environment {named}, not the {asked} being adopted"
            ),
            Self::NothingToAdopt { id } => write!(
                f,
                "environment {id} has no committed history to take over; register it instead"
            ),
            Self::Journal(ref cause) => {
                write!(f, "the journal being adopted could not be read: {cause}")
            }
            Self::Exhausted { id } => {
                write!(f, "environment {id} has no generations left")
            }
        }
    }
}

impl std::error::Error for BrokerError {}

/// Authority to hand one exact attempt to the environment it was prepared for.
///
/// Minted only by [`Broker::authorize`], and deliberately not [`Clone`]. The
/// seal is that the fields are private and there is no public constructor, so a
/// warrant cannot be assembled from an environment id and a generation a caller
/// chose.
///
/// Not being `Clone` stops a warrant being copied; it does not stop one being
/// used twice through a shared reference, and it is not what makes a double
/// dispatch impossible. The journal's ordering ladder is what does that. This
/// type stops the warrant from being *scattered*.
#[derive(Debug)]
pub struct Authority {
    /// The environment the attempt may be handed to.
    environment: EnvironmentId,
    /// The generation it was authorized against.
    epoch: EnvironmentEpoch,
}

impl Authority {
    /// Seal a warrant. Crate-private on purpose: this is the seal.
    /// [`Broker::authorize`] is the only caller, so a warrant can never name a
    /// generation the broker did not check.
    pub(crate) const fn new(environment: EnvironmentId, epoch: EnvironmentEpoch) -> Self {
        Self { environment, epoch }
    }

    /// The environment the attempt may be handed to.
    #[must_use]
    pub const fn environment(&self) -> EnvironmentId {
        self.environment
    }

    /// The generation it was authorized against.
    #[must_use]
    pub const fn epoch(&self) -> EnvironmentEpoch {
        self.epoch
    }
}

impl fmt::Display for Authority {
    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
        write!(
            f,
            "authority for {} at generation {}",
            self.environment, self.epoch
        )
    }
}

/// The environments a run may act on, and the generation each is at.
///
/// One broker per run. It holds no native handle itself: the adopter that
/// actually owns the process, socket or input seat lives behind whatever
/// [`EffectJournal`] adapter and dispatch path the host supplies. What this owns
/// is the fence, because fencing is a comparison and a comparison needs one
/// place that holds the current generation.
#[derive(Debug, Default)]
pub struct Broker {
    /// The known environments, by identity.
    ///
    /// Bounded by the environments one run created, not by traffic: an entry is
    /// added by [`Broker::register`] or [`Broker::adopt`], both of which the host
    /// calls once per environment it owns, and neither is a cache that refills on
    /// its own. There is no eviction, and that is the point — [`Broker::close`]
    /// marks an entry closed rather than removing it, so a warrant for a closed
    /// environment is refused as `Closed` instead of looking like a warrant for
    /// an environment this broker never heard of. A host that creates and
    /// discards environments without bound holds one small entry per creation
    /// for the life of the broker, which is the run's own resource count and is
    /// the host's to keep bounded.
    environments: HashMap<EnvironmentId, Environment>,
}

impl Broker {
    /// A broker that owns no environment.
    #[must_use]
    pub fn new() -> Self {
        Self {
            environments: HashMap::new(),
        }
    }

    /// Take ownership of an environment at its first generation.
    ///
    /// The identity comes from the host, which is what created the process or
    /// container and therefore holds the entropy; the broker owns the
    /// environment's lifetime from here.
    ///
    /// # Errors
    ///
    /// [`BrokerError::AlreadyRegistered`] when this broker already owns `id`.
    pub fn register(&mut self, id: EnvironmentId) -> Result<EnvironmentEpoch, BrokerError> {
        if self.environments.contains_key(&id) {
            let refusal = Err(BrokerError::AlreadyRegistered { id });
            lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "register: returning an error to the caller");
            return refusal;
        }
        // The first generation is one, not zero: the counter is a `NonZeroU64`
        // precisely so that a zeroed or truncated field cannot be read as a
        // generation the broker issued.
        let epoch = EnvironmentEpoch::new(NonZeroU64::MIN);
        self.environments
            .insert(id, Environment { epoch, open: true });
        Ok(epoch)
    }

    /// Replace the environment, invalidating every command prepared against the
    /// generation it was at.
    ///
    /// This is what a restart of the underlying process looks like from the
    /// broker's side: same identity, new generation, and every warrant already
    /// handed out is now stale.
    ///
    /// # Errors
    ///
    /// [`BrokerError::UnknownEnvironment`], [`BrokerError::Closed`] or
    /// [`BrokerError::Exhausted`].
    pub fn replace(&mut self, id: EnvironmentId) -> Result<EnvironmentEpoch, BrokerError> {
        let environment = self
            .environments
            .get_mut(&id)
            .ok_or(BrokerError::UnknownEnvironment { id })?;
        if !environment.open {
            let refusal = Err(BrokerError::Closed { id });
            lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "replace: returning an error to the caller");
            return refusal;
        }
        let next = environment
            .epoch
            .checked_next()
            .ok_or(BrokerError::Exhausted { id })?;
        environment.epoch = next;
        Ok(next)
    }

    /// Take ownership of this environment at the generation a journal that
    /// already exists was last written at, and move past it.
    ///
    /// The honest form of a takeover, and the only one that means anything across
    /// a process boundary. [`Self::register`] starts an environment at
    /// generation 1 whatever is on the disk, so a process that adopts a journal
    /// another worker wrote would mint warrants for a generation that worker had
    /// already been replaced past — and two processes would then disagree about
    /// which commands are current while each was internally consistent. This
    /// reads the generation the journal's own committed history was written at
    /// and claims the one after it, so the new owner's every warrant is fenced
    /// against everything the previous owner ever prepared.
    ///
    /// A journal with no committed event, or one whose events name another
    /// environment, is a refusal rather than a generation of 1: there is nothing
    /// here to take over, and claiming generation 1 for it would be the same
    /// silent downgrade [`Self::register`] is.
    ///
    /// # Errors
    ///
    /// [`BrokerError::UnknownEnvironment`] when this broker already owns `id`,
    /// and [`JournalError`] when the journal cannot be read or does not describe
    /// `id`.
    pub fn adopt(
        &mut self,
        id: EnvironmentId,
        journal: &dyn crate::journal::EffectJournal,
    ) -> Result<EnvironmentEpoch, BrokerError> {
        if self.environments.contains_key(&id) {
            let refusal = Err(BrokerError::AlreadyRegistered { id });
            lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "adopt: returning an error to the caller");
            return refusal;
        }
        let mut highest = None;
        for event in journal.committed().map_err(BrokerError::Journal)? {
            let key = event.key();
            if key.environment() != id {
                let refusal = Err(BrokerError::ForeignEnvironment {
                    asked: id,
                    named: key.environment(),
                });
                lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "adopt: returning an error to the caller");
                return refusal;
            }
            highest = Some(match highest {
                None => key.epoch(),
                Some(current) if key.epoch() > current => key.epoch(),
                Some(current) => current,
            });
        }
        let Some(highest) = highest else {
            let refusal = Err(BrokerError::NothingToAdopt { id });
            lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "adopt: returning an error to the caller");
            return refusal;
        };
        let next = highest
            .checked_next()
            .ok_or(BrokerError::Exhausted { id })?;
        self.environments.insert(
            id,
            Environment {
                epoch: next,
                open: true,
            },
        );
        Ok(next)
    }

    /// Release the environment. Idempotent, because a cleanup path that has to
    /// know whether it already ran is a cleanup path that will sometimes not run.
    ///
    /// # Errors
    ///
    /// [`BrokerError::UnknownEnvironment`] when this broker never owned `id`.
    pub fn close(&mut self, id: EnvironmentId) -> Result<(), BrokerError> {
        let environment = self
            .environments
            .get_mut(&id)
            .ok_or(BrokerError::UnknownEnvironment { id })?;
        environment.open = false;
        Ok(())
    }

    /// The generation this broker currently holds for `id`.
    #[must_use]
    pub fn epoch(&self, id: EnvironmentId) -> Option<EnvironmentEpoch> {
        self.environments
            .get(&id)
            .map(|environment| environment.epoch)
    }

    /// Whether this broker still owns `id` and will mint authority for it.
    #[must_use]
    pub fn is_open(&self, id: EnvironmentId) -> bool {
        self.environments
            .get(&id)
            .is_some_and(|environment| environment.open)
    }

    /// The generation check, run both when a warrant is minted and when it is
    /// presented again.
    ///
    /// One function because the two callers must agree exactly. A second copy
    /// of this comparison that drifted would mean a warrant the broker would
    /// refuse to mint and would accept at the boundary, which is the failure
    /// fencing exists to prevent.
    fn check_generation(
        &self,
        id: EnvironmentId,
        presented: EnvironmentEpoch,
    ) -> Result<EnvironmentEpoch, BrokerError> {
        let environment = self
            .environments
            .get(&id)
            .ok_or(BrokerError::UnknownEnvironment { id })?;
        if !environment.open {
            let refusal = Err(BrokerError::Closed { id });
            lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "check_generation: returning an error to the caller");
            return refusal;
        }
        let current = environment.epoch;
        if presented.get() > current.get() {
            let refusal = Err(BrokerError::NeverIssued {
                environment: id,
                presented,
                current,
            });
            lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "check_generation: returning an error to the caller");
            return refusal;
        }
        if presented != current {
            let refusal = Err(BrokerError::Superseded {
                environment: id,
                presented,
                current,
            });
            lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "check_generation: returning an error to the caller");
            return refusal;
        }
        Ok(current)
    }

    /// Mint authority for `key`, if and only if the generation it names is the
    /// one this broker currently holds.
    ///
    /// This is a check at one instant. It is not sufficient on its own: see
    /// [`Broker::revalidate`] for why the handoff path has to ask again.
    ///
    /// # Errors
    ///
    /// [`BrokerError::UnknownEnvironment`], [`BrokerError::Closed`],
    /// [`BrokerError::Superseded`] or [`BrokerError::NeverIssued`].
    pub fn authorize(&self, key: EffectKey) -> Result<Authority, BrokerError> {
        let id = key.environment();
        let current = self.check_generation(id, key.epoch())?;
        Ok(Authority::new(id, current))
    }

    /// Re-check a warrant that was minted earlier, immediately before the
    /// handoff it authorizes.
    ///
    /// Authorizing and handing over are two instants, and a replacement landing
    /// between them would leave a warrant that was minted legitimately and is
    /// now stale. RQ-006 says replacement invalidates *all* old commands, and a
    /// check that only runs at mint time does not deliver that. So the handoff
    /// path presents its warrant here, and a warrant whose generation has been
    /// replaced is refused even though nothing was wrong when it was minted.
    ///
    /// # Errors
    ///
    /// [`BrokerError::UnknownEnvironment`], [`BrokerError::Closed`],
    /// [`BrokerError::Superseded`] or [`BrokerError::NeverIssued`].
    pub fn revalidate(&self, authority: &Authority) -> Result<(), BrokerError> {
        self.check_generation(authority.environment(), authority.epoch())?;
        Ok(())
    }
}

/// What went wrong on the one path that may reach an external system.
///
/// One enum rather than a tuple of two, so a caller cannot match the broker
/// refusal and forget the journal one.
#[derive(Debug)]
#[non_exhaustive]
pub enum DispatchError {
    /// The environment refused to authorize the attempt.
    Broker(BrokerError),
    /// The journal refused to record it.
    Journal(JournalError),
}

impl fmt::Display for DispatchError {
    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
        match *self {
            Self::Broker(ref cause) => write!(f, "dispatch refused by the broker: {cause}"),
            Self::Journal(ref cause) => write!(f, "dispatch refused by the journal: {cause}"),
        }
    }
}

impl std::error::Error for DispatchError {
    fn source(&self) -> Option<&(dyn std::error::Error + 'static)> {
        match *self {
            Self::Broker(ref cause) => Some(cause),
            Self::Journal(ref cause) => Some(cause),
        }
    }
}

impl From<BrokerError> for DispatchError {
    fn from(cause: BrokerError) -> Self {
        Self::Broker(cause)
    }
}

impl From<JournalError> for DispatchError {
    fn from(cause: JournalError) -> Self {
        Self::Journal(cause)
    }
}

/// An attempt that has been authorized and durably recorded, and may now be
/// handed to its environment.
///
/// Holding one of these is the only way to reach the handoff, which is how the
/// required order stops being a convention: authority is checked first, the
/// `DispatchPrepared` append lands second, and only then does anything exist to
/// hand over.
///
/// The warrant it carries was checked when it was minted, which is an instant
/// and not a guarantee. The handoff path must present it to
/// [`Broker::revalidate`] immediately before handing over, because a replacement
/// landing in between would otherwise leave this preparation dispatchable.
#[derive(Debug)]
pub(crate) struct Prepared {
    /// Authority for the environment generation the attempt was prepared
    /// against.
    authority: Authority,
    /// The committed journal position for the `DispatchPrepared` append.
    ack: DurableAck,
}

impl Prepared {
    /// Spend the preparation, taking the authority and the acknowledgment.
    ///
    /// This is the only way to read it. Borrowing accessors for the two fields
    /// existed while the type was public and nothing outside this crate ever
    /// called them; the one caller, the dispatch path in `ecs`, wants both
    /// halves and consumes the preparation to get them.
    #[must_use]
    pub(crate) fn into_parts(self) -> (Authority, DurableAck) {
        (self.authority, self.ack)
    }
}

/// Authorize `key` against the current environment generation, then append
/// `DispatchPrepared` for it.
///
/// The order is the durable half of RQ-007: authority is obtained and the
/// attempt recorded *before* any bytes exist to hand over, so a crash after this
/// returns leaves a committed `DispatchPrepared` with no outcome, which recovery
/// reads back as an unknown rather than as a never-sent.
///
/// A refusal from either half leaves the journal unchanged. In particular a
/// superseded key never reaches the append, so a fenced command cannot become a
/// recorded attempt.
///
/// # Errors
///
/// [`DispatchError::Broker`] when the environment refuses, and
/// [`DispatchError::Journal`] when the append is refused. The append is also
/// where a second `DispatchPrepared` for one attempt is refused, so a retry that
/// reuses an `AttemptId` fails here rather than dispatching twice.
pub(crate) async fn prepare_dispatch(
    broker: &Broker,
    journal: &mut dyn EffectJournal,
    expected_tail: JournalPosition,
    key: EffectKey,
) -> Result<Prepared, DispatchError> {
    let authority = broker.authorize(key)?;
    let ack = journal
        .compare_and_append_async(expected_tail, &EffectEvent::DispatchPrepared { key })
        .await?;
    Ok(Prepared { authority, ack })
}

#[cfg(test)]
mod tests {
    use super::*;
    use crate::effect::{ActionDigest, ActionId, AttemptId, EnvironmentId, FlowRevision, RunId};
    use crate::journal::{EffectEvidence, MemoryJournal};
    use std::collections::HashMap;

    const RUN: &str = "0102030405060708090a0b0c0d0e0f10";
    const ACTION: &str = "1112131415161718191a1b1c1d1e1f20";
    const ENV: &str = "2122232425262728292a2b2c2d2e2f30";
    const OTHER_ENV: &str = "5152535455565758595a5b5c5d5e5f60";
    const FLOW_HEX: &str = "000102030405060708090a0b0c0d0e0f101112131415161718191a1b1c1d1e1f";
    const DIGEST_HEX: &str = "f0f1f2f3f4f5f6f7f8f9fafbfcfdfeffe0e1e2e3e4e5e6e7e8e9eaebecedeeef";

    /// The crate refuses a panicking path anywhere, tests included, so a test
    /// that needs a parsed value returns `Result` and propagates.
    type TestResult = Result<(), Box<dyn std::error::Error>>;

    /// A key for the shared run and action, on the named environment and
    /// generation.
    fn key_on(env: &str, epoch: &str) -> Result<EffectKey, Box<dyn std::error::Error>> {
        let key_run = RunId::from_hex(RUN)?;
        let key_action = ActionId::from_hex(ACTION)?;
        let key_attempt = AttemptId::from_decimal("1")?;
        let key_flow_revision = FlowRevision::from_tagged("blake3_256", FLOW_HEX)?;
        let key_digest = ActionDigest::from_tagged("blake3_256", DIGEST_HEX)?;
        let key_environment = EnvironmentId::from_hex(env)?;
        let key_epoch = EnvironmentEpoch::from_decimal(epoch)?;
        Ok(
            crate::effect::EffectIdentity::new(key_run, key_environment, key_flow_revision).key(
                key_action,
                key_attempt,
                key_digest,
                key_epoch,
            ),
        )
    }

    /// The environment id most of these tests use.
    fn env() -> Result<EnvironmentId, Box<dyn std::error::Error>> {
        Ok(EnvironmentId::from_hex(ENV)?)
    }

    /// A key on the shared environment at the named generation.
    fn key(epoch: &str) -> Result<EffectKey, Box<dyn std::error::Error>> {
        key_on(ENV, epoch)
    }

    /// Walk an attempt to the point where a dispatch may be prepared.
    fn admitted(
        journal: &mut MemoryJournal,
        key: EffectKey,
    ) -> Result<(), Box<dyn std::error::Error>> {
        let tail = journal.tail();
        journal.compare_and_append(tail, &EffectEvent::IntentAdmitted { key })?;
        Ok(())
    }

    /// The shared subject of the dispatch tests: an environment the broker will
    /// authorize, the key carrying its current generation, and a journal that
    /// has already admitted that attempt.
    ///
    /// One struct rather than a four-tuple, because the three facts have to be
    /// the *same* facts: a test that assembled its own could admit one key and
    /// authorize another, and every assertion after that would be about a
    /// situation the module cannot produce.
    struct AdmittedEnvironment {
        /// The broker, holding the environment open at `generation`.
        broker: Broker,
        /// The generation the broker issued, which `key` carries.
        generation: EnvironmentEpoch,
        /// The admitted attempt, on that environment and at that generation.
        key: EffectKey,
        /// A journal with that attempt already admitted.
        journal: MemoryJournal,
    }

    /// A broker holding the shared environment at its first generation, that
    /// generation's key, and a journal with the attempt already admitted.
    fn admitted_environment() -> Result<AdmittedEnvironment, Box<dyn std::error::Error>> {
        let mut broker = Broker::new();
        let generation = broker.register(env()?)?;
        let key = key(&generation.get().to_string())?;
        let mut journal = MemoryJournal::new();
        admitted(&mut journal, key)?;
        Ok(AdmittedEnvironment {
            broker,
            generation,
            key,
            journal,
        })
    }

    #[test]
    fn register_takes_the_first_generation() -> TestResult {
        let mut broker = Broker::new();
        let epoch = broker.register(env()?)?;
        assert_eq!(broker.epoch(env()?), Some(epoch));
        assert!(broker.is_open(env()?));
        Ok(())
    }

    #[test]
    fn a_second_register_of_one_environment_is_refused() -> TestResult {
        let mut broker = Broker::new();
        let first = broker.register(env()?)?;
        let refused = broker.register(env()?);
        assert!(
            matches!(refused, Err(BrokerError::AlreadyRegistered { .. })),
            "a double register must not silently reset anything"
        );
        assert_eq!(broker.epoch(env()?), Some(first));
        Ok(())
    }

    #[test]
    fn a_current_generation_authorizes() -> TestResult {
        let mut broker = Broker::new();
        let first = broker.register(env()?)?;
        let authority = broker.authorize(key(&first.get().to_string())?)?;
        assert_eq!(authority.environment(), env()?);
        assert_eq!(authority.epoch(), first);
        Ok(())
    }

    #[test]
    fn an_unknown_environment_is_refused() -> TestResult {
        let broker = Broker::new();
        let refused = broker.authorize(key("1")?);
        assert!(matches!(
            refused,
            Err(BrokerError::UnknownEnvironment { .. })
        ));
        Ok(())
    }

    #[test]
    fn replacement_supersedes_a_command_prepared_before_it() -> TestResult {
        let mut broker = Broker::new();
        let first = broker.register(env()?)?;
        let stale = key(&first.get().to_string())?;
        assert!(broker.authorize(stale).is_ok());

        let second = broker.replace(env()?)?;
        assert_ne!(second, first);

        match broker.authorize(stale) {
            Err(BrokerError::Superseded {
                environment,
                presented,
                current,
            }) => {
                assert_eq!(environment, env()?);
                assert_eq!(presented, first);
                assert_eq!(current, second);
            }
            Err(other) => return Err(format!("expected a superseded refusal, got {other}").into()),
            Ok(authority) => {
                return Err(format!(
                    "a command from a replaced generation must not be authorized, got {authority}"
                )
                .into());
            }
        }
        Ok(())
    }

    #[test]
    fn a_generation_the_broker_never_issued_is_not_reported_as_stale() -> TestResult {
        let mut broker = Broker::new();
        let first = broker.register(env()?)?;
        let ahead = EnvironmentEpoch::new(NonZeroU64::new(first.get() + 1).ok_or("increment")?);
        let forged = key(&ahead.get().to_string())?;

        match broker.authorize(forged) {
            Err(BrokerError::NeverIssued { current, .. }) => assert_eq!(current, first),
            Err(other) => {
                return Err(format!("expected a never-issued refusal, got {other}").into());
            }
            Ok(authority) => {
                return Err(format!(
                    "a generation the broker never issued must not authorize, got {authority}"
                )
                .into());
            }
        }
        Ok(())
    }

    #[test]
    fn a_warrant_minted_before_a_replacement_is_refused_at_the_boundary() -> TestResult {
        let mut broker = Broker::new();
        let first = broker.register(env()?)?;
        let warrant = broker.authorize(key(&first.get().to_string())?)?;
        assert!(broker.revalidate(&warrant).is_ok());

        // The warrant was minted legitimately. A replacement between the mint
        // and the handoff is what makes it stale, and the boundary check is the
        // only thing that catches it.
        let second = broker.replace(env()?)?;

        match broker.revalidate(&warrant) {
            Err(BrokerError::Superseded {
                presented, current, ..
            }) => {
                assert_eq!(presented, first);
                assert_eq!(current, second);
            }
            Err(other) => return Err(format!("expected a superseded refusal, got {other}").into()),
            Ok(()) => {
                return Err("a warrant from before the replacement must not revalidate".into());
            }
        }
        Ok(())
    }

    #[test]
    fn a_warrant_is_refused_at_the_boundary_once_the_environment_closes() -> TestResult {
        let mut broker = Broker::new();
        let first = broker.register(env()?)?;
        let warrant = broker.authorize(key(&first.get().to_string())?)?;
        broker.close(env()?)?;
        assert!(matches!(
            broker.revalidate(&warrant),
            Err(BrokerError::Closed { .. })
        ));
        Ok(())
    }

    #[test]
    fn a_closed_environment_mints_nothing() -> TestResult {
        let mut broker = Broker::new();
        let first = broker.register(env()?)?;
        let live = key(&first.get().to_string())?;
        broker.close(env()?)?;

        assert!(!broker.is_open(env()?));
        assert!(matches!(
            broker.authorize(live),
            Err(BrokerError::Closed { .. })
        ));
        assert!(matches!(
            broker.replace(env()?),
            Err(BrokerError::Closed { .. })
        ));
        Ok(())
    }

    #[test]
    fn closing_twice_is_allowed() -> TestResult {
        let mut broker = Broker::new();
        broker.register(env()?)?;
        broker.close(env()?)?;
        broker.close(env()?)?;
        Ok(())
    }

    #[test]
    fn closing_an_unknown_environment_is_refused() -> TestResult {
        let mut broker = Broker::new();
        assert!(matches!(
            broker.close(env()?),
            Err(BrokerError::UnknownEnvironment { .. })
        ));
        Ok(())
    }

    #[test]
    fn an_exhausted_generation_space_is_named_rather_than_wrapped() -> TestResult {
        let id = env()?;
        let last = EnvironmentEpoch::new(NonZeroU64::MAX);
        let mut broker = Broker {
            environments: HashMap::from([(
                id,
                Environment {
                    epoch: last,
                    open: true,
                },
            )]),
        };
        assert!(matches!(
            broker.replace(id),
            Err(BrokerError::Exhausted { .. })
        ));
        assert_eq!(broker.epoch(id), Some(last));
        Ok(())
    }

    #[test]
    fn one_broker_does_not_authorize_another_environments_key() -> TestResult {
        let mut broker = Broker::new();
        broker.register(env()?)?;
        let other = key_on(OTHER_ENV, "1")?;
        assert!(matches!(
            broker.authorize(other),
            Err(BrokerError::UnknownEnvironment { .. })
        ));
        Ok(())
    }

    /// The synchronous shape of [`prepare_dispatch`], for a test with no
    /// executor of its own.
    ///
    /// One definition rather than a `block_on` at each call site: these tests
    /// are about what the order records and refuses, and the await is not part
    /// of either question.
    fn prepare(
        broker: &Broker,
        journal: &mut MemoryJournal,
        tail: JournalPosition,
        key: EffectKey,
    ) -> Result<Prepared, DispatchError> {
        lgwks_std::task::block_on(prepare_dispatch(broker, journal, tail, key))
    }

    #[test]
    fn prepare_dispatch_records_the_attempt_it_authorized() -> TestResult {
        let AdmittedEnvironment {
            broker,
            generation,
            key,
            mut journal,
        } = admitted_environment()?;

        let tail = journal.tail();
        let (authority, ack) = prepare(&broker, &mut journal, tail, key)?.into_parts();
        assert_eq!(authority.epoch(), generation);
        assert_eq!(ack.position(), journal.tail());
        assert_eq!(journal.committed().len(), 2);
        Ok(())
    }

    #[test]
    fn a_superseded_key_never_reaches_the_journal() -> TestResult {
        let AdmittedEnvironment {
            mut broker,
            key: stale,
            mut journal,
            ..
        } = admitted_environment()?;
        let before = journal.tail();

        broker.replace(env()?)?;

        let refused = prepare(&broker, &mut journal, before, stale);
        assert!(
            matches!(
                refused,
                Err(DispatchError::Broker(BrokerError::Superseded { .. }))
            ),
            "a fenced command must be refused before it can be recorded"
        );
        assert_eq!(journal.tail(), before);
        assert_eq!(journal.committed().len(), 1);
        Ok(())
    }

    #[test]
    fn a_second_preparation_of_one_attempt_is_refused_by_the_ladder() -> TestResult {
        let AdmittedEnvironment {
            broker,
            key,
            mut journal,
            ..
        } = admitted_environment()?;

        let tail = journal.tail();
        let first_prepare = prepare(&broker, &mut journal, tail, key)?;
        let spent = first_prepare.into_parts();

        let tail = journal.tail();
        let refused = prepare(&broker, &mut journal, tail, key);
        assert!(
            matches!(
                refused,
                Err(DispatchError::Journal(JournalError::OutOfOrder { .. }))
            ),
            "the environment still authorizes, and the journal is what refuses the resend"
        );
        assert_eq!(spent.1.position(), journal.tail());
        assert_eq!(journal.committed().len(), 2);
        Ok(())
    }

    #[test]
    fn an_unadmitted_attempt_is_refused_even_with_authority() -> TestResult {
        let mut broker = Broker::new();
        let first = broker.register(env()?)?;
        let key = key(&first.get().to_string())?;
        let mut journal = MemoryJournal::new();

        // The environment would authorize this attempt: the generation is
        // current. The order still refuses it, because nothing was admitted.
        assert!(broker.authorize(key).is_ok());
        let tail = journal.tail();
        let refused = prepare(&broker, &mut journal, tail, key);
        assert!(matches!(
            refused,
            Err(DispatchError::Journal(JournalError::OutOfOrder { .. }))
        ));
        assert!(journal.committed().is_empty());
        Ok(())
    }

    #[test]
    fn a_stale_controller_cannot_append_after_recovery() -> TestResult {
        let AdmittedEnvironment {
            broker,
            key,
            mut journal,
            ..
        } = admitted_environment()?;

        // A second controller reads the tail before the first appends.
        let stale_tail = journal.tail();
        let tail = journal.tail();
        prepare(&broker, &mut journal, tail, key)?;

        let late = journal.compare_and_append(
            stale_tail,
            &EffectEvent::OutcomeObserved {
                key,
                evidence: EffectEvidence::Applied,
            },
        );
        assert!(
            matches!(late, Err(JournalError::TailMismatch { .. })),
            "the tail check is what stops the second controller writing over the first"
        );
        Ok(())
    }
}