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
//! The durable per-run control ledger: what a repair has already consumed.
//!
//! # Why this is a second store rather than more step records
//!
//! A step record answers "what did this step return?" and is keyed by
//! `(run, step key)`; it is written once and replayed verbatim, because a
//! completed step's value is a fact, not a state. The three facts this ledger
//! holds are the opposite of that. The root **spend** and **attempt budget** are
//! cumulative counters that a repair must *not* reset, the repair **epoch** is
//! advanced by every applied ticket, and an **applied ticket** must be refused a
//! second time rather than replayed. All four are read-modify-write across a
//! resume, which is not what a step record is for and cannot be expressed by
//! replaying one. So they get their own file, its own chain, and its own record.
//!
//! The frame grammar ([`journal::frame`]) and the storage-owner thread
//! ([`journal::owner`]) are the ones the step store already uses, exactly as
//! INV-BOT-51 requires: one answer in this crate to "what is a frame" and to
//! "which thread reaches the disk", not one per store.
//!
//! # The budget is a root budget, and a repair spends it
//!
//! `attempts` and `spend` are the *root* counters. They are charged on every run
//! attempt and are never reduced by a repair, a replay or a fresh resume — so a
//! run that has spent its budget stays spent even after an authorized repair,
//! which is T13's "authorized repair is distinct and root budgets remain
//! charged". A run whose budget is already spent refuses another attempt with a
//! typed [`LeaseRefusal::BudgetSpent`] rather than resetting, so a permanent
//! refusal plus repeated `NotApplied` reaches a finite typed refusal instead of
//! unbounded retries.
//!
//! # One ordered step decides and writes
//!
//! Whether a repair is valid (is the epoch current, was this ticket already
//! applied, is the budget spent) and the act that makes it true (charge the
//! budget, advance the epoch, record the ticket) are decided and performed in one
//! step on the ledger's own thread. A caller that read the state, decided, and
//! charged later could decide against bytes another repair has already moved; the
//! single ordered step is what makes the check and the write un-overtakable, the
//! same property the step store's length fence rests on.

use std::collections::HashMap;
use std::fmt;
use std::fs::{File, OpenOptions};
use std::io::{Read, Seek, SeekFrom, Write};
use std::path::{Path, PathBuf};

use crate::journal::frame::SaturatingFrom;
use std::sync::{Arc, Mutex};

use lgwks_std::hash::{Digest, Hasher};
use lgwks_std::wire::{WireError, from_bytes, to_bytes};

use crate::effect::RunId;
use crate::journal::frame::{self, Cursor, HEAD_BYTES};
use crate::journal::owner::{self, Stage, StorageOwner, SubmitError};

use super::store::{StoreError, StoreLimitKind};

/// The magic at the head of every ledger file.
///
/// Distinct from the step store's magic so a ledger is never opened as a step
/// store or the reverse: the two record different things and reading one as the
/// other would hand a caller records it did not write. The trailing byte is this
/// format's version.
const LEDGER_MAGIC: &[u8; 17] = b"lgwks-runledger\x00\x01";

/// The most archived bytes one ledger entry may occupy.
///
/// A ledger entry is a counter update or an epoch/ticket note — small by
/// construction — so this ceiling exists only to bound a hostile or runaway
/// value before any allocation, and is far below the step store's own record
/// ceiling.
pub const MAX_LEDGER_RECORD_BYTES: usize = 16 * 1024;

/// The most ledger entries one run's chain may hold.
///
/// A run's ledger is a handful of budget charges and repair notes. Sixty-four
/// thousand is generous enough that no honest run reaches it and small enough
/// that a runaway loop is refused rather than filling the disk.
pub const MAX_LEDGER_RECORDS_PER_RUN: u64 = 65_536;

/// The chain digest a ledger starts from, hashed through the same framing as any
/// entry so the first head is a real chain step rather than a constant every
/// file shares.
fn genesis_head() -> Digest {
    let mut hasher = Hasher::new();
    hasher.write_framed(b"lgwks-runledger/genesis");
    hasher.finalize()
}

/// One durable ledger entry: the run it belongs to and the control state it
/// leaves behind.
///
/// The counters are cumulative and absolute — the writer reads the current state,
/// computes the new one, and stores the result — so replaying the whole chain
/// reproduces exactly the counters and the epoch a reader would see live. That is
/// what makes the ledger idempotent under a crash that lost the acknowledgment but
/// kept the bytes: the entry is the fact, and reading it back is the truth.
#[derive(
    Debug, Clone, lgwks_std::wire::Archive, lgwks_std::wire::Serialize, lgwks_std::wire::Deserialize,
)]
#[rkyv(
    attr(non_exhaustive),
    crate = lgwks_std::wire::rkyv,
    compare(PartialEq),
    derive(Debug)
)]
struct Entry {
    /// The run this entry belongs to.
    #[rkyv(attr(doc = "The run this entry belongs to."))]
    run: RunId,
    /// The tenant that owns the run; the only tenant allowed to read it.
    #[rkyv(attr(doc = "The tenant that owns the run."))]
    tenant: String,
    /// The root attempts this run has been charged, in total.
    #[rkyv(attr(doc = "The root attempts this run has been charged, in total."))]
    attempts: u64,
    /// The root spend this run has been charged, in total.
    #[rkyv(attr(doc = "The root spend this run has been charged, in total."))]
    spend: u64,
    /// The repair epoch this entry leaves the run at.
    #[rkyv(attr(doc = "The repair epoch this entry leaves the run at."))]
    epoch: u64,
    /// The identity of the repair ticket this entry applied, when it applied one.
    #[rkyv(attr(doc = "The repair ticket identity this entry applied, if any."))]
    applied: Option<Vec<u8>>,
}

impl Entry {
    /// This entry's chain head over the head before it.
    ///
    /// Framed exactly as the step store frames its record, so a ledger frame and
    /// a step frame cannot hash the same byte string: the run, the tenant and
    /// the applied ticket are each length-framed, and the counters follow.
    fn head_from(&self, previous: &Digest) -> Digest {
        let mut hasher = Hasher::new();
        hasher.write_framed(previous.as_bytes());
        hasher.write_framed(self.run.id().to_hex().as_bytes());
        hasher.write_framed(self.tenant.as_bytes());
        hasher.write_framed(&self.attempts.to_le_bytes());
        hasher.write_framed(&self.spend.to_le_bytes());
        hasher.write_framed(&self.epoch.to_le_bytes());
        hasher.write_framed(match self.applied {
            Some(ref ticket) => ticket.as_slice(),
            None => &[],
        });
        hasher.finalize()
    }
}

/// The run-keyed index one ledger holds, folded in the storage owner's ordered
/// step alongside each write.
#[derive(Debug)]
struct Index {
    /// Per run, its live control state.
    runs: HashMap<RunId, Control>,
    /// The file length every indexed entry accounts for.
    committed: u64,
    /// The chain head the next entry follows.
    tail: Digest,
}

/// One run's live control state, as replay and live updates agree on it.
///
/// The four facts a repair has to be decided against, read as a value rather than
/// through four accessors: a caller checking a budget wants the whole state at
/// once, and a half-read state is exactly the disagreement this ledger exists to
/// prevent.
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct Control {
    /// The tenant that owns the run.
    tenant: String,
    /// The root attempts charged so far, never reset by a repair or resume.
    attempts: u64,
    /// The root spend charged so far, never reset by a repair or resume.
    spend: u64,
    /// The current repair epoch. Zero is the epoch of a run no repair has
    /// touched, which is the epoch every ticket minted before any repair names.
    epoch: u64,
    /// The identities of repair tickets already applied to this run, so the same
    /// ticket delivered twice is recognized and refused the second time.
    applied: Vec<Vec<u8>>,
}

impl Control {
    /// The control state for a run this ledger has never seen, owned by
    /// `tenant`.
    fn fresh(tenant: &str) -> Self {
        Self {
            tenant: tenant.to_owned(),
            attempts: 0,
            spend: 0,
            epoch: 0,
            applied: Vec::new(),
        }
    }

    /// The tenant that owns the run.
    #[must_use]
    pub fn tenant(&self) -> &str {
        &self.tenant
    }

    /// The root attempts charged so far, never reset by a repair or a resume.
    #[must_use]
    pub const fn attempts(&self) -> u64 {
        self.attempts
    }

    /// The root spend charged so far, never reset by a repair or a resume.
    #[must_use]
    pub const fn spend(&self) -> u64 {
        self.spend
    }

    /// The current repair epoch. Zero for a run no repair has touched.
    #[must_use]
    pub const fn epoch(&self) -> u64 {
        self.epoch
    }

    /// How many repair tickets have been applied to this run.
    #[must_use]
    pub fn applied(&self) -> usize {
        self.applied.len()
    }
}

/// The repair ticket a charge is carrying, as the ledger sees it.
///
/// An identity and the epoch it was minted at. The identity is what makes a
/// duplicate recognizable and the epoch is what makes a stale ticket stale; both
/// are the ticket's facts, carried rather than re-derived here, because the
/// ledger cannot know what a caller's ticket "means" — only whether it has seen
/// it before.
///
/// Crate-private, because it is the ledger's own vocabulary: a caller holds a
/// [`RepairTicket`] and never a stamp, and a stamp a caller
/// could rewrite would be a way to forge "the same ticket" or "a fresh epoch"
/// without going through the ticket's content-derived identity.
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) struct TicketStamp {
    /// The identity of the ticket, used only for duplicate recognition.
    pub(crate) identity: Vec<u8>,
    /// The repair epoch the ticket was minted at.
    pub(crate) epoch: u64,
}

/// Why a lease could not be taken. Every refusal charges nothing and mints no
/// epoch, so a refused repair leaves the run's authority and budget exactly where
/// they were.
#[derive(Debug, Clone, PartialEq, Eq)]
#[non_exhaustive]
pub enum LeaseRefusal {
    /// The run belongs to a tenant other than the one asking.
    ForeignTenant {
        /// The tenant that owns the run.
        owner: String,
        /// The tenant that asked.
        asked: String,
    },
    /// The repair epoch on the ticket is not the run's current epoch, so the
    /// ticket was minted against a different repair state. An older ticket is
    /// stale — replaying it would re-apply an earlier authority; a newer one was
    /// minted against a state this run is not in. Both are refused by the same
    /// check, so "stale" names the disagreement rather than the direction.
    StaleEpoch {
        /// The epoch the run is at.
        current: u64,
        /// The epoch the ticket names.
        offered: u64,
    },
    /// The same repair ticket has already been applied to this run. Delivering a
    /// ticket twice must apply its effect once.
    AlreadyApplied,
    /// The root budget this run carries is already spent; a repair does not
    /// refill it.
    BudgetSpent {
        /// The attempts that would have been charged.
        attempts: u64,
        /// The largest attempt count admitted.
        max_attempts: u64,
    },
}

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

impl fmt::Display for LeaseRefusal {
    /// The refusal, naming both sides of the disagreement.
    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
        match *self {
            Self::ForeignTenant {
                ref owner,
                ref asked,
            } => write!(
                formatter,
                "run belongs to tenant {owner:?}, not {asked:?}; refusing to read its repair ledger"
            ),
            Self::StaleEpoch { current, offered } => write!(
                formatter,
                "repair epoch {offered} is not this run's current epoch {current}; \
                 refusing a stale repair ticket"
            ),
            Self::AlreadyApplied => formatter.write_str(
                "this repair ticket was already applied to this run; refusing to apply it twice",
            ),
            Self::BudgetSpent {
                attempts,
                max_attempts,
            } => write!(
                formatter,
                "root budget spent: {attempts} attempts against a ceiling of {max_attempts}; \
                 an authorized repair does not refill it"
            ),
        }
    }
}

/// One durable, file-backed set of per-run repair controls.
///
/// Cloning shares the file and the index through an `Arc`, so a host's handles
/// see one ledger rather than one per clone — the same rule the step store's
/// admission budget follows, and for the same reason: two writers over one chain
/// would fork it.
#[derive(Clone)]
pub struct RunLedger {
    /// Where the ledger lives and the index every clone shares.
    inner: Arc<Inner>,
}

/// The state one ledger owns, shared by every clone of its handle.
struct Inner {
    /// The thread that holds the file and performs every ordered step. Its
    /// answer is the control state a charge leaves behind, so a caller reads the
    /// result of its own decision rather than a second lookup that could race
    /// another charge.
    owner: StorageOwner<Arc<Mutex<Index>>, Control>,
    /// Where the file is, for diagnostics.
    path: PathBuf,
    /// The index every clone reads and the owner folds into.
    index: Arc<Mutex<Index>>,
}

impl RunLedger {
    /// Open (or create) the ledger at `path`.
    ///
    /// The file is created with a header if it does not exist and replayed if it
    /// does, exactly as the step store opens its own file. A torn final entry is
    /// dropped (it was never anyone's answer) and every earlier entry survives;
    /// a committed entry whose head does not follow is refused rather than
    /// trimmed, because those bytes were acknowledged.
    ///
    /// # Errors
    ///
    /// [`StoreError`] when the device refuses, when `path` is not a ledger, or
    /// when a committed entry does not follow from the ones before it.
    pub fn open(path: impl Into<PathBuf>) -> Result<Self, StoreError> {
        let path = path.into();
        let existed = path.exists();
        let mut file = OpenOptions::new()
            .read(true)
            .write(true)
            .create(true)
            .truncate(false)
            .open(&path)
            .map_err(StoreError::storage)?;
        if !existed || file.metadata().map_err(StoreError::storage)?.len() == 0 {
            file.write_all(LEDGER_MAGIC)
                .and_then(|()| file.sync_all())
                .map_err(StoreError::storage)?;
        }
        let index = replay(&mut file)?;
        if file.metadata().map_err(StoreError::storage)?.len() != index.committed {
            file.set_len(index.committed).map_err(StoreError::storage)?;
            file.sync_all().map_err(StoreError::storage)?;
        }
        let index = Arc::new(Mutex::new(index));
        let owner =
            StorageOwner::spawn(file, Arc::clone(&index), false).map_err(StoreError::storage)?;
        Ok(Self {
            inner: Arc::new(Inner { owner, path, index }),
        })
    }

    /// Open a ledger under `dir`, named by the tenant that owns it, so two
    /// tenants pointed at one directory never share a file.
    ///
    /// # Errors
    ///
    /// As [`RunLedger::open`].
    pub fn open_in(dir: &Path, tenant: &str) -> Result<Self, StoreError> {
        std::fs::create_dir_all(dir).map_err(StoreError::storage)?;
        Self::open(dir.join(format!("{tenant}.runledger")))
    }

    /// Where this ledger lives.
    #[must_use]
    pub fn path(&self) -> &Path {
        &self.inner.path
    }

    /// The control state `run` currently holds, or `None` when this ledger has
    /// never seen it.
    ///
    /// A run the ledger has *never seen* is `None`, not a zeroed control state:
    /// the difference between "no budget charged" and "no run here" is the
    /// difference between a first attempt and somebody else's run, and a caller
    /// must be able to tell them apart before it charges anything.
    #[must_use]
    pub fn control(&self, run: RunId) -> Option<Control> {
        owner::lock(&self.inner.index).runs.get(&run).cloned()
    }

    /// Perform one ordered step on `run`: decide the ticket against the current
    /// state and, if it is admissible, charge the budget and apply it.
    ///
    /// What a ticket means for this charge — `Some` authorizes a repair, `None`
    /// is a plain resume — is [`Charge`]'s, and the two ceilings are
    /// [`Ceilings`]'.
    ///
    /// # Errors
    ///
    /// [`CommitError::Refused`] for every logical refusal — the four
    /// [`LeaseRefusal`] arms — charging nothing, and [`CommitError::Store`] when
    /// the ledger's device refuses the write.
    ///
    /// The future is `Send`: this is awaited on the host's own path, which the
    /// engine may drive on a multi-threaded runtime whatever the task body is, so
    /// the answer travels the owner's concrete awaiting future rather than the
    /// non-`Send` box a step body would want.
    pub(crate) async fn charge(
        &self,
        tenant: &str,
        run: RunId,
        ticket: Option<&TicketStamp>,
        spend: u64,
        max_attempts: u64,
        max_spend: u64,
    ) -> Result<Control, CommitError> {
        // The two named values, built here so every caller reaches the ordered step
        // through them: the tenant and the ticket are borrowed for the length of a
        // call, and the ordered step owns them.
        let charge = Charge {
            tenant: tenant.to_owned(),
            run,
            ticket: ticket.cloned(),
            cost: spend,
        };
        let ceilings = Ceilings {
            attempts: max_attempts,
            spend: max_spend,
        };
        self.inner
            .owner
            .enqueue_awaiting(move |file, shared| charge_on_owner(file, shared, &charge, ceilings))
            .await
            .map_err(classify)
    }
}

/// Tell a logical refusal from a device one, on the way out of the owner.
///
/// The owner's one reply channel is an `io::Error`, so a decision made on its
/// thread travels as one — but a refusal that reached the caller as a string would
/// be exactly the flattening this crate's typed refusals exist to prevent: a
/// caller could not tell "this ticket was already applied" from "the device
/// refused". So the refusal is carried as the error itself (it implements
/// [`std::error::Error`]) and recovered here by downcast, and only a genuinely
/// unrecognised error becomes a [`StoreError`].
fn classify(cause: SubmitError) -> CommitError {
    let is_ours = matches!(
        cause,
        SubmitError::Device(ref error)
            if error
                .get_ref()
                .and_then(<dyn std::error::Error + Send + Sync>::downcast_ref::<LeaseRefusal>)
                .is_some()
    );
    if !is_ours {
        return CommitError::Store(StoreError::from(cause));
    }
    match cause {
        SubmitError::Device(error) => match error.into_inner() {
            Some(inner) => match inner.downcast::<LeaseRefusal>() {
                Ok(refusal) => CommitError::Refused(*refusal),
                Err(other) => CommitError::Store(StoreError::storage(std::io::Error::other(other))),
            },
            None => CommitError::Store(StoreError::storage(std::io::Error::other(
                "the ledger owner refused a charge without naming why",
            ))),
        },
        other => CommitError::Store(StoreError::from(other)),
    }
}

/// The one ordered step: decide the charge against the ledger's own state, write
/// the entry, fold it into the index, and answer the state it left behind.
///
/// The decision and the write are here together rather than in two functions so
/// no caller can decide against one state and write against another: the whole
/// run holds the index lock, and a second charge cannot enter until this one has
/// written and answered. That is what makes "this ticket was not applied twice"
/// and "the epoch did not move twice" facts about bytes rather than about the
/// order two threads happened to run in.
///
/// # Errors
///
/// Whatever the device or the ledger's own check reports, carried as the
/// device's own error so the owner's one reply channel serves this store the way
/// it serves the step store. A logical refusal is decided before any byte moves,
/// so a refused charge leaves the file byte-identical.
///
/// The step performs its own `sync_all` and fold under the index lock before it
/// returns, so it answers [`Stage::Committed`]: it owes the batch's flush nothing,
/// because the one ordered step that makes "applied once" a fact about bytes is
/// the same step that flushed it. A later member of a shared batch must decide
/// against the fold this charge left, which is why the fold cannot be deferred to
/// a batch settle the way a step record's can.
fn charge_on_owner(
    file: &mut File,
    shared: &Arc<Mutex<Index>>,
    charge: &Charge,
    ceilings: Ceilings,
) -> std::io::Result<Stage<Control, Arc<Mutex<Index>>>> {
    let mut index = owner::lock(shared);
    let (entry, next) = decide_under(&index, charge, ceilings).map_err(std::io::Error::other)?;
    write_entry(file, &mut index, &entry)
        .map_err(|error| std::io::Error::other(error.to_string()))?;
    Ok(Stage::Committed(next))
}

/// One attempt's charge against a run's root budget.
///
/// A named value rather than four parameters, because these are one request: the
/// tenant that owns the run, which run, the repair authorizing this attempt, and
/// what it costs. Passed positionally, a call site could hand the cost as the
/// ticket and nothing in the types would notice.
#[derive(Clone, Debug)]
pub(crate) struct Charge {
    /// The tenant that owns the run. Checked before any counter moves, because
    /// charging another tenant's run is how a shared ledger becomes a cross-tenant
    /// write.
    pub(crate) tenant: String,
    /// The run this attempt charges.
    pub(crate) run: RunId,
    /// The repair authorizing this attempt, or `None` for a plain resume.
    pub(crate) ticket: Option<TicketStamp>,
    /// What this attempt costs against the root budget.
    ///
    /// Charged whether or not a ticket was carried: a repair is an attempt, and an
    /// authorized repair does not refund the run for the attempts that led to it.
    pub(crate) cost: u64,
}

/// The two ceilings a charge is refused against.
///
/// One value rather than two parameters, so the attempt ceiling and the spend
/// ceiling travel together: a caller that set one and left the other at some
/// default would be admitting a budget nobody declared.
#[derive(Clone, Copy, Debug)]
pub(crate) struct Ceilings {
    /// The most attempts one run may charge.
    pub(crate) attempts: u64,
    /// The most root-budget spend one run may charge.
    pub(crate) spend: u64,
}

/// The decision, over the ledger's own state: what charging `run` would leave, and
/// the entry that would say so.
///
/// Returns the entry *and* the state beside it rather than only the state, so the
/// write cannot be skipped by a caller that already knows the answer.
fn decide_under(
    index: &Index,
    charge: &Charge,
    ceilings: Ceilings,
) -> Result<(Entry, Control), LeaseRefusal> {
    let tenant = charge.tenant.as_str();
    let run = &charge.run;
    let ticket = charge.ticket.as_ref();
    let cost = charge.cost;
    let max_attempts = ceilings.attempts;
    let max_spend = ceilings.spend;
    // A run this ledger has never charged starts from `Control::fresh`, which is
    // the run's own starting state rather than a stand-in for a failed lookup: the
    // ceilings below then judge a first attempt exactly as they judge a later one.
    let mut next = Control::fresh(tenant);
    if let Some(control) = index.runs.get(run) {
        next = control.clone();
    }

    // A run that already belongs to someone else is refused before any counter
    // moves. This is the same cross-tenant rule the step store enforces, and for
    // the same reason: charging another tenant's run under this one's name is how
    // a shared ledger becomes a cross-tenant write.
    if next.tenant != tenant {
        let refusal = Err(LeaseRefusal::ForeignTenant {
            owner: next.tenant,
            asked: tenant.to_owned(),
        });
        lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "decide_under: returning an error to the caller");
        return refusal;
    }

    let next_attempts = next.attempts.saturating_add(1);
    let next_spend = next.spend.saturating_add(cost);
    if next_attempts > max_attempts || next_spend > max_spend {
        let refusal = Err(LeaseRefusal::BudgetSpent {
            attempts: next_attempts,
            max_attempts,
        });
        lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "decide_under: returning an error to the caller");
        return refusal;
    }

    let mut applied = None;
    if let Some(stamp) = ticket {
        // A ticket minted at an epoch other than the run's current one was minted
        // against a repair state this run is not in. Applying it would replay an
        // authority the run has already moved past, so it is refused; a newer
        // ticket is refused by the same check, so "stale" names the disagreement
        // rather than the direction.
        if stamp.epoch != next.epoch {
            let refusal = Err(LeaseRefusal::StaleEpoch {
                current: next.epoch,
                offered: stamp.epoch,
            });
            lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "decide_under: returning an error to the caller");
            return refusal;
        }
        // The same ticket delivered twice is one repair applied once. The identity
        // is checked against the applied set on this run, so a different ticket
        // for the same run is unaffected.
        if next.applied.contains(&stamp.identity) {
            let refusal = Err(LeaseRefusal::AlreadyApplied);
            lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "decide_under: returning an error to the caller");
            return refusal;
        }
        // Applying a valid repair advances the epoch, which is what makes the
        // *next* delivery of this same ticket stale as well as duplicate: two
        // independent checks, either of which would catch a replay.
        next.epoch = next.epoch.saturating_add(1);
        next.applied.push(stamp.identity.clone());
        applied = Some(stamp.identity.clone());
    }
    next.attempts = next_attempts;
    next.spend = next_spend;

    let entry = Entry {
        run: *run,
        tenant: tenant.to_owned(),
        attempts: next.attempts,
        spend: next.spend,
        epoch: next.epoch,
        applied,
    };
    Ok((entry, next))
}

/// The write and the fold, in one ordered step, under the index the owner holds.
fn write_entry(file: &mut File, index: &mut Index, entry: &Entry) -> Result<(), StoreError> {
    let previous = index.tail;
    let (framed, head) = frame(entry, &previous)?;
    let staged = u64::saturating_from(framed.len());
    let next = index
        .committed
        .checked_add(staged)
        .ok_or(StoreError::Limit {
            kind: StoreLimitKind::StoreBytes,
            requested: u64::MAX,
            limit: super::MAX_STORE_BYTES,
        })?;
    // The length fence, before a byte moves, under the ledger's thread where no
    // other charge can overtake it: the length this handle indexed must be the
    // length on the disk.
    let on_disk = file.metadata().map_err(StoreError::storage)?.len();
    if on_disk != index.committed {
        let refusal = Err(StoreError::Corrupt { at: 0 });
        lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "write_entry: returning an error to the caller");
        return refusal;
    }
    file.write_all(&framed).map_err(StoreError::storage)?;
    file.sync_all().map_err(StoreError::storage)?;
    index.committed = next;
    index.tail = head;
    fold(index, entry);
    Ok(())
}

/// Fold one entry's control state into `index`, creating the run's slot if this
/// is the first entry for it.
///
/// One function because the live write and the replay must agree exactly: a state
/// reached by committing an entry and the same state reached by reading it back
/// have to be the same value, or a ledger's answer would depend on whether the
/// reader crashed. The three counters are assigned rather than accumulated
/// because an entry is the *absolute* state, and the applied-ticket list is
/// extended rather than replaced for the same reason: a ticket applied by an
/// earlier entry is still applied.
fn fold(index: &mut Index, entry: &Entry) {
    let run = index.runs.entry(entry.run).or_insert_with(|| Control {
        tenant: entry.tenant.clone(),
        attempts: 0,
        spend: 0,
        epoch: 0,
        applied: Vec::new(),
    });
    run.tenant.clone_from(&entry.tenant);
    run.attempts = entry.attempts;
    run.spend = entry.spend;
    run.epoch = entry.epoch;
    if let Some(ref ticket) = entry.applied
        && !run.applied.contains(ticket)
    {
        run.applied.push(ticket.clone());
    }
}

/// The frame for one ledger entry: length, payload, head, through the shared
/// grammar.
fn frame(entry: &Entry, previous: &Digest) -> Result<(Vec<u8>, Digest), StoreError> {
    frame::frame_record(
        entry,
        previous,
        MAX_LEDGER_RECORD_BYTES,
        |entry| {
            to_bytes::<WireError>(entry)
                .map(|bytes| bytes.as_ref().to_vec())
                .map_err(|cause| StoreError::Encoding { cause })
        },
        |entry, previous, _| entry.head_from(previous),
        |len| StoreError::Limit {
            kind: StoreLimitKind::RecordBytes,
            requested: u64::saturating_from(len),
            limit: u64::saturating_from(MAX_LEDGER_RECORD_BYTES),
        },
    )
}

/// One complete, decoded frame read from the ledger.
struct Framed {
    /// The entry the frame's payload decoded to.
    entry: Entry,
    /// The chain head the frame recorded for itself.
    head: [u8; HEAD_BYTES],
    /// The payload length the frame declared.
    declared: usize,
}

/// Read the frame at ordinal `at`, or `None` where the file stops holding a
/// whole frame.
///
/// A partial length prefix, or a payload or head cut short, is an append that
/// never finished: it was never anyone's answer, so the scan stops and the
/// caller trims to `committed`. A complete prefix that names a frame this
/// ledger never writes, or a payload that does not decode, cannot be an
/// interrupted append and is refused, and so is a cut whose bytes hold an entry
/// this ledger acknowledged under another length (#262).
fn next_frame(file: &mut File, cursor: &Cursor<'_>) -> Result<Option<Framed>, StoreError> {
    let at = cursor.at;
    let corrupt = || StoreError::Corrupt { at };
    let ceiling = MAX_LEDGER_RECORD_BYTES;
    let entry_head = |previous: &Digest, payload: &[u8]| {
        let aligned = frame::decodable(payload);
        from_bytes::<Entry, WireError>(aligned.as_slice())
            .ok()
            .map(|entry| entry.head_from(previous))
    };
    let Some(raw) = frame::read_raw(
        file,
        cursor,
        ceiling,
        StoreError::storage,
        corrupt,
        entry_head,
    )?
    else {
        return Ok(None);
    };
    let entry = from_bytes::<Entry, WireError>(&raw.payload).map_err(|error| {
        lgwks_std::trace::debug!(?error, at, "next_frame: the payload did not decode");
        StoreError::Corrupt { at }
    })?;
    Ok(Some(Framed {
        entry,
        declared: raw.payload.len(),
        head: raw.head,
    }))
}

/// Read every committed entry, returning the index it implies.
fn replay(file: &mut File) -> Result<Index, StoreError> {
    let total = file.metadata().map_err(StoreError::storage)?.len();
    file.seek(SeekFrom::Start(0)).map_err(StoreError::storage)?;
    let mut header = [0u8; LEDGER_MAGIC.len()];
    if !read_full(file, &mut header)? || header != *LEDGER_MAGIC {
        let refusal = Err(StoreError::NotAStore);
        lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "replay: returning an error to the caller");
        return refusal;
    }
    let mut index = Index {
        runs: HashMap::new(),
        committed: u64::saturating_from(LEDGER_MAGIC.len()),
        tail: genesis_head(),
    };
    let mut previous = genesis_head();
    let mut at = 0u64;
    while let Some(Framed {
        entry,
        head,
        declared,
    }) = next_frame(file, &Cursor::new(at, index.committed, &previous))?
    {
        if entry.head_from(&previous) != Digest::from_bytes(head) {
            let refusal = Err(StoreError::Corrupt { at });
            lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "replay: returning an error to the caller");
            return refusal;
        }
        let frame_len = frame::framed_len(declared);
        if index
            .committed
            .checked_add(frame_len)
            .is_none_or(|end| end > total)
        {
            break;
        }
        index.committed = index.committed.saturating_add(frame_len);
        index.tail = entry.head_from(&previous);
        previous = index.tail;
        at = at.saturating_add(1);
        if at > MAX_LEDGER_RECORDS_PER_RUN {
            let refusal = Err(StoreError::Limit {
                kind: StoreLimitKind::Records,
                requested: at,
                limit: MAX_LEDGER_RECORDS_PER_RUN,
            });
            lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "replay: returning an error to the caller");
            return refusal;
        }
        fold(&mut index, &entry);
    }
    Ok(index)
}

/// Read `buf` in full, reporting whether all of it arrived.
fn read_full(reader: &mut impl Read, buf: &mut [u8]) -> Result<bool, StoreError> {
    let mut filled = 0usize;
    while filled < buf.len() {
        match reader.read(&mut buf[filled..]) {
            Ok(0) => return Ok(false),
            Ok(read) => filled = filled.saturating_add(read),
            Err(ref error) if error.kind() == std::io::ErrorKind::Interrupted => {}
            Err(cause) => {
                let refusal = Err(StoreError::storage(cause));
                lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "read_full: returning an error to the caller");
                return refusal;
            }
        }
    }
    Ok(true)
}

/// Why charging a run failed: a device refusal, or a logical one.
#[derive(Debug)]
pub enum CommitError {
    /// The ledger's device refused the write. Nothing was charged.
    Store(StoreError),
    /// The charge was refused on its merits. Nothing was charged and no epoch was
    /// minted, so the run's authority and budget are exactly where they were.
    Refused(LeaseRefusal),
}

impl fmt::Display for CommitError {
    /// The refusal, whether it came from the device or from the decision.
    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
        match *self {
            Self::Store(ref cause) => write!(formatter, "{cause}"),
            Self::Refused(ref cause) => write!(formatter, "{cause}"),
        }
    }
}

impl std::error::Error for CommitError {
    /// The device's own error, when the refusal came from the device.
    fn source(&self) -> Option<&(dyn std::error::Error + 'static)> {
        match *self {
            Self::Store(ref cause) => Some(cause),
            Self::Refused(_) => None,
        }
    }
}

impl fmt::Debug for RunLedger {
    /// The path and the indexed run count, never the counters.
    ///
    /// The values are not rendered: a report or a log line must not carry a run's
    /// budget into a reader that asked where the ledger is.
    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
        formatter
            .debug_struct("RunLedger")
            .field("path", &self.inner.path)
            .field("runs", &owner::lock(&self.inner.index).runs.len())
            .finish()
    }
}

/// An entry whose length prefix lies is refused, and an append that was cut is
/// repaired (#262), on the same grammar the run store reads.
#[cfg(test)]
mod tests {
    use super::*;
    use crate::journal::frame::probe::{Scratch, declared_at, frame_starts, with_prefix};

    type TestResult = Result<(), Box<dyn std::error::Error>>;

    /// The ledger file of `count` charges, written to `path`, and its bytes.
    fn written(path: &Path, count: u64) -> Result<Vec<u8>, Box<dyn std::error::Error>> {
        let run = RunId::from_hex(&format!("b{}", "0".repeat(31)))?;
        let mut bytes = LEDGER_MAGIC.to_vec();
        let mut previous = genesis_head();
        for charge in 1..=count {
            let entry = Entry {
                run,
                tenant: "acme".to_owned(),
                attempts: charge,
                spend: charge.saturating_mul(3),
                epoch: charge,
                applied: Some(vec![u8::try_from(charge)?; 24]),
            };
            let (framed, head) = frame(&entry, &previous)?;
            bytes.extend_from_slice(&framed);
            previous = head;
        }
        std::fs::write(path, &bytes)?;
        Ok(bytes)
    }

    /// An acknowledged final entry whose prefix grows is refused as `Corrupt` at
    /// its index, and the file is byte-identical afterwards.
    #[test]
    fn a_lengthened_acknowledged_final_entry_is_refused_not_trimmed() -> TestResult {
        let scratch = Scratch::new("ledger-lengthened")?;
        let bytes = written(scratch.path(), 3)?;
        let last = frame_starts(&bytes, LEDGER_MAGIC.len())?[2];
        let declared = declared_at(&bytes, last);
        for extra in (1u32..=40).chain([100, 1024]) {
            let lied = with_prefix(&bytes, last, declared + extra);
            std::fs::write(scratch.path(), &lied)?;
            assert!(
                matches!(
                    RunLedger::open(scratch.path()),
                    Err(StoreError::Corrupt { at: 2 })
                ),
                "L+{extra}: an acknowledged entry must be refused as corrupt at 2"
            );
            assert_eq!(
                std::fs::read(scratch.path())?,
                lied,
                "L+{extra}: bytes move"
            );
        }
        Ok(())
    }

    /// A lying prefix over an entry that is itself damaged still refuses: the entry
    /// behind it authenticates, from the head stored before it, at a payload offset
    /// the archive was not written at.
    #[test]
    fn a_damaged_lengthened_entry_with_an_acknowledged_one_behind_it_is_refused() -> TestResult {
        let scratch = Scratch::new("ledger-damaged")?;
        let bytes = written(scratch.path(), 3)?;
        let middle = frame_starts(&bytes, LEDGER_MAGIC.len())?[1];
        let mut lied = with_prefix(&bytes, middle, u32::try_from(bytes.len() - middle)? + 9);
        if let Some(byte) = lied.get_mut(middle + 10) {
            *byte ^= 0x55;
        }
        std::fs::write(scratch.path(), &lied)?;
        assert!(matches!(
            RunLedger::open(scratch.path()),
            Err(StoreError::Corrupt { at: 1 })
        ));
        assert_eq!(std::fs::read(scratch.path())?, lied, "refused bytes move");
        Ok(())
    }

    /// The control: an append cut inside the final entry is repaired to exactly the
    /// acknowledged prefix.
    #[test]
    fn an_append_cut_inside_the_final_entry_is_repaired() -> TestResult {
        let scratch = Scratch::new("ledger-cut")?;
        let bytes = written(scratch.path(), 3)?;
        let last = frame_starts(&bytes, LEDGER_MAGIC.len())?[2];
        let whole = bytes.len() - last;
        for cut in [1, 3, 4, 5, whole >> 1, whole - 33, whole - 32, whole - 1] {
            std::fs::write(
                scratch.path(),
                bytes.get(..last + cut).ok_or("past the end")?,
            )?;
            drop(RunLedger::open(scratch.path())?);
            assert_eq!(
                std::fs::metadata(scratch.path())?.len(),
                u64::try_from(last)?,
                "cut {cut}: repaired to the acknowledged prefix"
            );
        }
        Ok(())
    }
}