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
//! One declared logical clock: the authority every deadline in this crate names.
//!
//! Three clocks exist in any async system and they are not the same clock. The
//! engine's timer wheel, `std::time::Instant`, and a test's notion of "now" all
//! advance independently, and a bot that mixes them produces deadlines that
//! disagree with each other by exactly the amount nobody measured. This module
//! makes one of them authoritative and makes the difference *nameable*.
//!
//! # The three clocks, and what each governs
//!
//! | Clock | Governs | Sourced from |
//! |---|---|---|
//! | [`Clock`] (logical) | every deadline this crate evaluates: retry, sleep, budget, readiness, continuation | a caller-advanceable counter |
//! | the engine timer | when a `sleep` future actually resolves | the engine's driver |
//! | [`WallClock`] | the watchdog on things a logical clock cannot stop | `std::time::Instant` |
//!
//! # Why the wall clock is never governed by the logical one
//!
//! A logical clock is the right authority for anything a test can model: a retry
//! that backs off, a sleep, a readiness wait, a continuation deadline. It is the
//! wrong authority for a subprocess that has stopped responding, a blocking
//! callback that will never return, or a store whose `fsync` is stuck — none of
//! those are waiting for time to pass, they have stopped making progress
//! entirely, and a paused logical clock would wait on them forever. That is why
//! [`WallClock`] exists, is always on, and is never reachable from a paused
//! logical clock: [`Clock::wall`] hands out a wall watchdog that keeps running
//! with logical time frozen.
//!
//! The rule this encodes, which the module docs make a contract: **pausing
//! logical time never disables the wall-clock watchdog.** A test that pauses
//! time to examine a three-day outage still cannot hang on a process that
//! stopped responding, because the watchdog it was given reads real elapsed
//! time, not the counter it paused.
//!
//! # Determinism, and what is *not* deterministic
//!
//! Advancing a [`Clock`] is deterministic in the order work becomes eligible,
//! which is what makes a seeded replay meaningful. It is not a claim about
//! anything outside this process:
//!
//! - Poll *order* across worker threads is not deterministic, and a clock does
//!   not make it so. What a clock determinizes is which deadline is *eligible*,
//!   not which task the OS scheduler runs first.
//! - Actual observed external order — the order a remote peer saw two requests
//!   in, the order bytes hit a socket — is never deterministic and is not
//!   claimed here.
//! - Clock skew between hosts is real and unmeasured. A logical clock has no
//!   cross-host notion at all.
//! - Real-time liveness is the [`WallClock`]'s domain and is the only thing that
//!   says a subprocess is still answering.
//!
//! # Duration and remaining-deadline semantics across a restart
//!
//! A [`Clock`] is a **process-local** counter and it is explicitly *not*
//! serializable. A `std::time::Instant` is an opaque monotonic reading with an
//! undefined epoch, and writing it to a durable store as though it were a
//! timestamp produces a value that means nothing on another host, or after the
//! host's clock source changes. What survives a restart is the **remaining
//! duration**, never the instant:
//!
//! ```
//! # use std::time::Duration;
//! # use lgwks_bot::rt::clock::Clock;
//! let clock = Clock::virtual_at(Duration::ZERO);
//! let budget = Duration::from_secs(5);
//!
//! // A run records how much of its budget it has spent, not a wall instant.
//! clock.advance(Duration::from_secs(2))?;
//! let recorded = clock.snapshot();
//!
//! // A restart re-establishes the origin from the recorded elapsed time, and
//! // the remaining budget is the same fact on either host.
//! let resumed = Clock::virtual_at(recorded.elapsed());
//! assert_eq!(resumed.snapshot().remaining_from(budget), Duration::from_secs(3));
//! # Ok::<(), lgwks_bot::rt::clock::ClockError>(())
//! ```
//!
//! [`ClockSnapshot::remaining_from`] gives the remaining duration of a budget at
//! the recorded point, and a restarted clock restored to that point has the same
//! remaining budget. Nothing on this type can be mistaken for a wall-clock
//! timestamp, because it carries no epoch and no host identity.
//!
//! # No-runtime construction
//!
//! A [`Clock`] is built, advanced and read entirely from `std`. It does not
//! require a runtime, and constructing one inside an existing runtime is
//! allowed. Nothing here starts a runtime, promotes local work to remote
//! execution, or nests one runtime per task.
//!
//! [`Clock`]: crate::clock::Clock
//! [`Clock::wall`]: crate::clock::Clock::wall
//! [`WallClock`]: crate::clock::WallClock
//! [`ClockSnapshot::remaining_from`]: crate::clock::ClockSnapshot::remaining_from

use std::fmt;
use std::sync::Arc;
use std::sync::atomic::{AtomicU64, Ordering};
use std::time::{Duration, Instant};

/// Largest duration a logical instant can represent, saturating rather than
/// wrapping.
///
/// One [`Duration`] below the saturation point. A clock that reached it would
/// be reporting a deadline nobody could wait out in a human lifetime, so the
/// value is a documented ceiling rather than a silently wrong number.
const MAX_ELAPSED: Duration = Duration::from_nanos(u64::MAX);

/// Which time source a [`Clock`] reads, and therefore what "advance" means.
///
/// A sum type rather than a boolean, because the two sources have different
/// contracts — one can be driven by a test and the other cannot — and a
/// `bool` field would let a caller construct a combination the module cannot
/// honour.
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
#[non_exhaustive]
pub enum TimeSource {
    /// Time follows the real monotonic clock. `advance` is not available
    /// because there is nothing to advance past.
    Wall,
    /// Time only moves when a caller advances it, so a test controls every
    /// deadline exactly and a three-day outage costs three arithmetic
    /// operations.
    Virtual,
}

/// The one logical clock a caller declares, and the authority for every deadline
/// this crate evaluates.
///
/// Cheap to clone: every clone shares the same counter, so a body can hold one
/// and a test can drive it from outside. This is what makes it a *declared*
/// clock rather than a value copied around.
#[derive(Debug, Clone)]
pub struct Clock {
    /// Shared so [`Clock::advance`] from a test and [`Clock::now`] from a body
    /// observe the same counter, and so a clock handed into a task does not
    /// become a second, independent timeline.
    inner: Arc<Inner>,
}

/// The shared state behind every [`Clock`] clone.
#[derive(Debug)]
struct Inner {
    /// Nanoseconds since the origin.
    ///
    /// Relaxed is correct and is the whole reason this is a clock rather than a
    /// channel: every read and write is of one `u64` that a caller either drives
    /// with [`Clock::advance`] or observes with [`Clock::now`]. There is no
    /// multi-field invariant to publish, so there is nothing to order against.
    /// A reader that observes a stale value observes a value from a
    /// still-consistent instant of a monotonically increasing counter, which is
    /// what an elapsed-time reading is.
    elapsed: AtomicU64,
    /// Whether this clock can be advanced by a caller.
    source: TimeSource,
    /// Where a wall clock was placed on its timeline, read through
    /// [`Inner::origin`].
    ///
    /// `None` for a virtual clock, whose timeline is the counter above and
    /// nothing else. The `Option` rather than a bare `Instant` because
    /// `Instant::now()` at construction would be a real read a virtual-clock
    /// caller never asked for, and a clock that touches the wall clock when it
    /// is supposed not to is precisely the coupling this module forbids.
    placed: Option<Instant>,
}

impl Inner {
    /// The real monotonic instant real-time readings count from.
    ///
    /// A wall clock's construction instant. A virtual clock has none, so a
    /// watchdog asked of it starts now — the moment it was asked for — which is
    /// the only real origin a virtual clock can honestly offer.
    fn origin(&self) -> Instant {
        match self.placed {
            // A wall clock carries the instant it was placed at, so every
            // reading of that clock counts from one stable origin.
            Some(placed) => placed,
            // A virtual clock has no placement, and `Instant::now()` is not a
            // stand-in for one: it is the only real origin a clock that was
            // never placed has, which is why the watchdog it hands out starts
            // when it is asked for. The two arms are different facts about the
            // clock, not one value covering for another.
            None => Instant::now(),
        }
    }
}

/// Why a clock operation was refused.
///
/// One error type rather than a `bool`, so a caller that ignores the refusal has
/// a value whose name says what it lost, not a silent no-op.
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
#[non_exhaustive]
pub enum ClockError {
    /// The operation requires a caller-advanceable clock and this clock follows
    /// real time. Nothing was advanced and the clock did not move.
    NotVirtual,
    /// The requested advance would leave the representable range. The clock is
    /// unchanged; a caller that needs a longer horizon needs a new origin.
    OutOfRange {
        /// The advance that was refused.
        requested: Duration,
        /// The largest advance that would have been accepted.
        ceiling: Duration,
    },
}

impl fmt::Display for ClockError {
    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
        match *self {
            Self::NotVirtual => formatter
                .write_str("the clock follows real time and cannot be advanced by a caller"),
            Self::OutOfRange { requested, .. } => {
                write!(
                    formatter,
                    "advancing the logical clock by {requested:?} leaves the representable range"
                )
            }
        }
    }
}

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

impl Clock {
    /// A clock that follows real monotonic time.
    ///
    /// The default for production use: nothing has to drive it, and a deadline
    /// it reports is a real elapsed duration rather than a test's opinion.
    #[must_use]
    pub fn wall() -> Self {
        Self {
            inner: Arc::new(Inner {
                elapsed: AtomicU64::new(0),
                source: TimeSource::Wall,
                placed: Some(Instant::now()),
            }),
        }
    }

    /// A wall clock re-established at `elapsed`, for a restart.
    ///
    /// The same fact as [`Clock::wall`] with a known starting point: the origin
    /// is the recorded elapsed time the previous process reached, and real time
    /// advances from there. This is the form a durable record replays, and it is
    /// why a restart restores a **duration** and never a stored `Instant` — an
    /// `Instant` is process-local and means nothing here.
    ///
    /// ```
    /// # use std::time::Duration;
    /// # use lgwks_bot::rt::clock::Clock;
    /// // A run that had spent four of its ten seconds restarts here.
    /// let clock = Clock::wall_at(Duration::from_secs(4));
    /// assert!(clock.now() >= Duration::from_secs(4));
    /// // At most six seconds remain, and real time keeps spending them.
    /// assert!(clock.snapshot().remaining_from(Duration::from_secs(10)) <= Duration::from_secs(6));
    /// ```
    #[must_use]
    pub fn wall_at(elapsed: Duration) -> Self {
        Self {
            inner: Arc::new(Inner {
                elapsed: AtomicU64::new(duration_to_nanos(elapsed)),
                source: TimeSource::Wall,
                placed: Some(Instant::now()),
            }),
        }
    }

    /// A caller-advanceable clock whose origin is `elapsed` nanoseconds after
    /// tick zero.
    ///
    /// The origin argument exists for **restart**, and is a duration since the
    /// origin rather than a wall-clock timestamp. A process that replays a
    /// durable record restores the origin from the recorded elapsed time, not
    /// from a stored `Instant`: an `Instant` is process-local, has no epoch, and
    /// means nothing on another host.
    ///
    /// ```
    /// # use std::time::Duration;
    /// # use lgwks_bot::rt::clock::Clock;
    /// let clock = Clock::virtual_at(Duration::from_secs(3));
    /// assert_eq!(clock.now(), Duration::from_secs(3));
    /// ```
    #[must_use]
    pub fn virtual_at(elapsed: Duration) -> Self {
        Self {
            inner: Arc::new(Inner {
                elapsed: AtomicU64::new(duration_to_nanos(elapsed)),
                source: TimeSource::Virtual,
                placed: None,
            }),
        }
    }

    /// Which time source this clock reads.
    #[must_use]
    pub fn source(&self) -> TimeSource {
        self.inner.source
    }

    /// Elapsed time since this clock's origin.
    ///
    /// A [`TimeSource::Virtual`] clock reports the counter a caller drives, so a
    /// test owns every deadline exactly and a three-day outage costs three
    /// arithmetic operations.
    ///
    /// A [`TimeSource::Wall`] clock reports where it was placed plus the
    /// **real** time elapsed since, sampled from `std::time::Instant` on every
    /// call. A fresh wall clock is placed at zero; one restored with
    /// [`Clock::wall_at`] keeps the time its previous process had already spent,
    /// so a restart does not hand a run its whole budget back. See
    /// [`Clock::placed_at`] for the placement alone.
    #[must_use]
    pub fn now(&self) -> Duration {
        match self.inner.source {
            TimeSource::Wall => self
                .placed_at()
                .saturating_add(self.inner.origin().elapsed()),
            TimeSource::Virtual => Duration::from_nanos(self.inner.elapsed.load(Ordering::Relaxed)),
        }
    }

    /// Where this clock was placed on its timeline, for either source.
    ///
    /// The origin a clock was restored to. For a wall clock it is the recorded
    /// elapsed time the restart re-established, so `now` starts from it and
    /// advances with real time; for a virtual clock it is the current counter,
    /// which is what [`Clock::snapshot`] persists.
    #[must_use]
    pub fn placed_at(&self) -> Duration {
        Duration::from_nanos(self.inner.elapsed.load(Ordering::Relaxed))
    }

    /// Move a caller-advanceable clock forward.
    ///
    /// An advance that would carry the clock past [`Clock::elapsed_ceiling`]
    /// **saturates at the ceiling** rather than refusing: a deadline computed from
    /// a wrapped clock fires immediately and is indistinguishable from one that
    /// legitimately expired, so landing on the ceiling is the honest answer to
    /// "advance further". The return value is where the clock actually landed, so
    /// a caller that asked for more and got less can see it.
    ///
    /// A clock already sitting at the ceiling refuses, because there is nothing
    /// left to land on and reporting success would be a lie about movement that
    /// did not happen.
    ///
    /// # Errors
    ///
    /// [`ClockError::NotVirtual`] on a wall clock, and
    /// [`ClockError::OutOfRange`] when the clock is already at its ceiling. In
    /// both cases the clock is unchanged.
    pub fn advance(&self, by: Duration) -> Result<Duration, ClockError> {
        if self.inner.source == TimeSource::Wall {
            let refusal = Err(ClockError::NotVirtual);
            lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "advance: returning an error to the caller");
            return refusal;
        }
        let now = self.placed_at();
        if now >= Self::elapsed_ceiling() {
            let refusal = Err(ClockError::OutOfRange {
                requested: by,
                ceiling: Duration::ZERO,
            });
            lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "advance: returning an error to the caller");
            return refusal;
        }
        // Clamped rather than refused: see the saturation contract above.
        let by = duration_to_nanos(by.min(Self::elapsed_ceiling().saturating_sub(now)));
        // `fetch_add` would wrap; a compare-exchange loop is the only way to
        // saturate atomically, and it re-reads on each attempt so a concurrent
        // advance from another task cannot lose an update.
        let mut current = duration_to_nanos(now);
        loop {
            let next = current.saturating_add(by);
            match self.inner.elapsed.compare_exchange_weak(
                current,
                next,
                Ordering::Relaxed,
                Ordering::Relaxed,
            ) {
                Ok(_) => return Ok(nanos_to_duration(next)),
                Err(observed) => {
                    // Another task advanced past the ceiling between the check
                    // above and this attempt. Refuse rather than move: the
                    // clock is at a horizon the caller was told about.
                    if observed >= duration_to_nanos(Self::elapsed_ceiling()) {
                        let refusal = Err(ClockError::OutOfRange {
                            requested: nanos_to_duration(by),
                            ceiling: Duration::ZERO,
                        });
                        lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "advance: returning an error to the caller");
                        return refusal;
                    }
                    current = observed;
                }
            }
        }
    }

    /// The largest elapsed time a logical instant can represent.
    ///
    /// Reported rather than discovered, because a caller that needs a horizon
    /// past this needs a new origin and should be told that by name.
    #[must_use]
    pub const fn elapsed_ceiling() -> Duration {
        MAX_ELAPSED
    }

    /// A point-in-time reading suitable for persisting across a restart.
    ///
    /// The reading is a *duration since the origin*, never an absolute instant,
    /// and it carries no host identity — so a record written from one host and
    /// read on another is a duration, and cannot be mistaken for a cross-host
    /// timestamp.
    #[must_use]
    pub fn snapshot(&self) -> ClockSnapshot {
        ClockSnapshot {
            elapsed: self.now(),
        }
    }

    /// The independent wall-clock watchdog, sampling the real monotonic origin.
    fn watchdog(&self) -> WallClock {
        WallClock {
            started: self.inner.origin(),
        }
    }

    /// A wall-clock watchdog that keeps running while logical time is frozen.
    ///
    /// The point of this method is that pausing logical time cannot pause a
    /// subprocess watchdog, a blocking callback or a store hang. Those are not
    /// waiting for time; they have stopped making progress, and only real
    /// elapsed time says so.
    ///
    /// ```
    /// # use std::time::Duration;
    /// # use lgwks_bot::rt::clock::Clock;
    /// let clock = Clock::virtual_at(Duration::ZERO);
    /// let watchdog = clock.wall_watchdog();
    /// clock.advance(Duration::from_secs(86_400))?;
    /// // Three logical days have passed and the watchdog has barely moved.
    /// assert!(watchdog.elapsed() < Duration::from_secs(1));
    /// # Ok::<(), lgwks_bot::rt::clock::ClockError>(())
    /// ```
    #[must_use]
    pub fn wall_watchdog(&self) -> WallClock {
        self.watchdog()
    }
}

/// A persistable reading of a [`Clock`]'s elapsed time.
///
/// Deliberately *not* a `std::time::Instant`. An `Instant` is a process-local
/// monotonic reading with an undefined epoch; persisting one as a timestamp
/// produces a value that means nothing on another host. This carries only a
/// duration, and [`Clock::virtual_at`] takes exactly this back.
#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord)]
#[non_exhaustive]
pub struct ClockSnapshot {
    /// Nanoseconds since the clock's origin.
    elapsed: Duration,
}

impl ClockSnapshot {
    /// Elapsed time the snapshot records.
    #[must_use]
    pub const fn elapsed(&self) -> Duration {
        self.elapsed
    }

    /// How much of `budget` was left at the instant this snapshot was taken.
    ///
    /// The restart-safe form: a deadline persisted as a *remaining duration*
    /// means the same thing on another host, where a persisted instant does not.
    /// A budget already spent yields [`Duration::ZERO`] rather than wrapping.
    #[must_use]
    pub fn remaining_from(&self, budget: Duration) -> Duration {
        budget.saturating_sub(self.elapsed)
    }

    /// Whether `budget` had already been spent at this snapshot.
    #[must_use]
    pub fn is_exhausted(&self, budget: Duration) -> bool {
        self.elapsed >= budget
    }
}

/// The independent wall-clock watchdog.
///
/// Always real time, never the logical counter, and therefore never paused by a
/// test. Constructed through [`Clock::wall_watchdog`] so that every watchdog in a
/// program traces back to a declared clock, rather than being a free-floating
/// [`Instant`].
#[derive(Debug, Clone, Copy)]
pub struct WallClock {
    /// The monotonic instant this watchdog started from.
    started: Instant,
}

impl WallClock {
    /// A watchdog whose origin is `started`.
    ///
    /// Crate-visible rather than public because a caller with an [`Instant`] in
    /// hand and no logical clock in scope has not declared a clock, and this
    /// module's whole contract is that every deadline names one. The one public
    /// constructor is [`Clock::wall_watchdog`].
    ///
    /// Gated with its only caller, `rt::time::Deadline`, which does not exist
    /// without the engine. Every other way to obtain a watchdog — including the
    /// observation phase's per-poll deadline — goes through a declared [`Clock`],
    /// which is the point of the visibility restriction above.
    #[cfg(feature = "rt")]
    pub(crate) fn from_instant(started: Instant) -> Self {
        Self { started }
    }

    /// Real elapsed time since this watchdog was created.
    ///
    /// Saturates at zero if the OS reports the start after the sample, which no
    /// monotonic source should do; the direction is chosen so a caller never
    /// computes a negative remaining budget.
    #[must_use]
    pub fn elapsed(&self) -> Duration {
        Instant::now().saturating_duration_since(self.started)
    }
}

/// Nanoseconds in `duration`, saturating at the representable ceiling.
///
/// `u64::MAX` nanoseconds is about 584 years, so the ceiling is never reached by
/// a real budget; it exists so that a caller who somehow asks for more gets the
/// documented ceiling rather than a wrapped instant.
///
/// Split rather than widen-then-narrow: `Duration` reports seconds and
/// sub-second nanoseconds through two infallible projections, so the only
/// arithmetic left is the scaling, and `saturating_mul`/`saturating_add` put
/// the ceiling in that arithmetic instead of in a conversion that could fail
/// and leave a caller holding a default.
fn duration_to_nanos(duration: Duration) -> u64 {
    const NANOS_PER_SEC: u64 = 1_000_000_000;
    duration
        .as_secs()
        .saturating_mul(NANOS_PER_SEC)
        .saturating_add(u64::from(duration.subsec_nanos()))
}

/// `nanos` nanoseconds as a [`Duration`].
///
/// The inverse of [`duration_to_nanos`] with no intermediate arithmetic, so
/// there is no quotient to truncate and no remainder to narrow.
fn nanos_to_duration(nanos: u64) -> Duration {
    Duration::from_nanos(nanos)
}