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
//! How a flow fails: located, classed and never silent.

use std::fmt;
use std::sync::Arc;
use std::time::Duration;

use crate::error::{BotError, RetryClass};
use crate::task::StoreError;

use super::Scope;

// ── FlowError ───────────────────────────────────────────────────────────────

/// Why a flow, or one step of it, did not produce its value.
///
/// Every variant carries `at`, the step path where it happened, so a failure in
/// the fortieth item of a nested fan-out reads `sync/each:page#39/retry` rather
/// than a bare message. An error raised with `?` before a location is known is
/// located at the flow boundary by [`FlowError::located`].
///
/// Retry is decided by the variant, never by the message:
/// [`FlowError::is_retryable`] is the only rule [`retry`](crate::script::retry) reads.
#[derive(Debug)]
#[non_exhaustive]
pub enum FlowError {
    /// The scope was stopped before or during the step.
    Cancelled {
        /// Where the stop was observed.
        at: Arc<str>,
    },
    /// A `within` deadline passed.
    TimedOut {
        /// The `within` block that expired.
        at: Arc<str>,
        /// The deadline it declared.
        after: Duration,
    },
    /// A `retry` spent every attempt on retryable failures.
    Exhausted {
        /// The `retry` block.
        at: Arc<str>,
        /// Attempts made.
        attempts: u32,
        /// The last attempt's failure.
        last: Box<FlowError>,
    },
    /// The run's retry budget was spent: failures are systemic, and another
    /// retry would add load rather than recover. Not retryable.
    Throttled {
        /// The `retry` block that was refused a retry.
        at: Arc<str>,
        /// Attempts made before the refusal.
        attempts: u32,
        /// The last attempt's failure.
        last: Box<FlowError>,
    },
    /// An untrusted payload was refused by the [`admit`](crate::script::admit) step.
    ///
    /// Typed rather than a [`FlowError::Failed`] with the refusal's text, because
    /// the whole point of admitting model output as data is that a caller can match
    /// on *why* without parsing a string — and because the [`Provenance`](crate::proposal::Provenance)
    /// of the exact refused bytes travels with it, so no refusal in a report is
    /// unattributable. Not retryable: a payload refused for its content will be
    /// refused the same way however many times it is re-read.
    Refused {
        /// Where the refusal was located.
        at: Arc<str>,
        /// Why the payload did not become work.
        refusal: Box<crate::proposal::Refusal>,
        /// Where the refused bytes came from.
        provenance: crate::proposal::Provenance,
    },
    /// The run reached a finite typed intervention instead of repairing again.
    ///
    /// Distinct from [`FlowError::Refused`]: the refusal is about *this* payload,
    /// while an intervention is about the run's ledger and is reached by repeated
    /// unchanged failure. Not retryable — another attempt would be exactly the
    /// repair the intervention refused.
    Intervention {
        /// Where the intervention was located.
        at: Arc<str>,
        /// What the run reached.
        intervention: Box<crate::proposal::Intervention>,
    },
    /// A permanent failure: repeating the step gives the same answer.
    Failed {
        /// Where it failed.
        at: Arc<str>,
        /// What the step said.
        reason: String,
    },
    /// A transient failure: the same step may succeed if repeated.
    Transient {
        /// Where it failed.
        at: Arc<str>,
        /// What the step said.
        reason: String,
    },
    /// A bot verb failed. Retryable exactly when its
    /// [`RetryClass`] is [`RetryClass::Safe`].
    Bot {
        /// Where it failed.
        at: Arc<str>,
        /// The verb's typed error, boxed so every flow `Result` stays small.
        source: Box<BotError>,
    },
    /// Steps nested past [`MAX_DEPTH`](crate::script::MAX_DEPTH).
    TooDeep {
        /// The scope that refused to go deeper.
        at: Arc<str>,
        /// The limit.
        limit: u16,
    },
    /// A tenant name did not validate.
    InvalidTenant {
        /// What is wrong with it.
        reason: &'static str,
    },
    /// A task's logical name did not validate.
    InvalidName {
        /// What is wrong with it.
        reason: &'static str,
    },
    /// A request key did not validate.
    ///
    /// A request key is the caller-supplied idempotency identity of a durable
    /// submission ([`Host::submit`](crate::task::Host::submit)), so it is
    /// validated before it is hashed into a run identity and refused here
    /// rather than at the first submission.
    InvalidRequestKey {
        /// What is wrong with it.
        reason: &'static str,
    },
    /// A step reached for authority this run does not have.
    ///
    /// The one flow failure that is not a defect in the step and not a transient
    /// condition: the host is willing, the work is well-formed, and the authority
    /// is missing. It is never retryable on its own, because repeating the same
    /// step against the same authority asks the same question and gets the same
    /// answer — a permanent refusal plus repeated `NotApplied` must reach a finite
    /// refusal rather than an unbounded retry loop.
    ///
    /// The unmet requirements ride with it, as a complete
    /// [`Deficit`](crate::cap::Deficit) rather than the first one, so the run's
    /// report can name every presently knowable need at once and the repair
    /// ticket is written once against the whole shortfall.
    Blocked {
        /// The step that reached for the authority.
        at: Arc<str>,
        /// Every capability this step needs and the run does not have.
        deficit: Box<crate::cap::Deficit>,
    },
    /// A bound computed at run time is outside what the block accepts.
    InvalidBound {
        /// Which bound.
        what: &'static str,
        /// The value given.
        value: u64,
        /// The largest value accepted.
        max: u64,
    },
    /// A resumed run's records were written under a different definition, input,
    /// step count or value schema, so none of them is this build's own work.
    ///
    /// Typed rather than a `Failed` string because the caller has to choose: a
    /// changed definition wants a new run, a changed input wants a decision, a
    /// changed schema wants a migration. Raised before any step body is polled
    /// and before any record is written, so a refusal costs nothing and replays
    /// nothing — and, because it is permanent, an enclosing `retry` will not
    /// spend attempts asking the same store the same question.
    Incompatible {
        /// Where the run would have resumed.
        at: Arc<str>,
        /// Which axis disagreed, and what each side holds.
        drift: crate::task::Drift,
    },
    /// The durable run store refused to read or write a step's record.
    ///
    /// Carried as the store's own typed error rather than as a rendering of it,
    /// because the caller's repair depends on *which* refusal it was: an
    /// unreadable device, a declared ceiling, a foreign tenant and a drift are
    /// four different actions, and a `Failed` string leaves the caller parsing
    /// prose to tell them apart. This is the run store's half of INV-BOT-7 —
    /// a read failure is an error, never absence, and never another error's
    /// answer — so an unreadable store reaches the caller as itself and is
    /// never reported as a definition drift.
    Store {
        /// Where it happened.
        at: Arc<str>,
        /// The store's own refusal.
        source: Box<StoreError>,
    },
}

impl FlowError {
    /// A permanent failure with `reason`. Not retried.
    pub fn failed(reason: impl fmt::Display) -> Self {
        Self::Failed {
            at: Arc::from(""),
            reason: reason.to_string(),
        }
    }

    /// A refusal to resume a run under a definition its records do not match.
    ///
    /// Permanent by construction: the records are what they are, so repeating the
    /// resume asks the same store the same question and gets the same answer.
    #[must_use]
    pub fn incompatible(at: &str, drift: crate::task::Drift) -> Self {
        Self::Incompatible {
            at: Arc::from(at),
            drift,
        }
    }

    /// A transient failure with `reason`. Retried by an enclosing `retry`.
    pub fn transient(reason: impl fmt::Display) -> Self {
        Self::Transient {
            at: Arc::from(""),
            reason: reason.to_string(),
        }
    }

    /// Whether repeating the step that produced this could succeed.
    ///
    /// `TimedOut` and `Transient` are; a bot error is when its retry class is
    /// [`RetryClass::Safe`] (the effect definitely did not happen). A
    /// cancellation, a permanent failure, an exhausted retry, a refused payload,
    /// an intervention, and every validation failure are not: repeating them is a
    /// retry storm.
    #[must_use]
    pub fn is_retryable(&self) -> bool {
        match *self {
            Self::TimedOut { .. } | Self::Transient { .. } => true,
            Self::Bot { ref source, .. } => source.retry_class() == RetryClass::Safe,
            Self::Cancelled { .. }
            | Self::Exhausted { .. }
            | Self::Throttled { .. }
            | Self::Failed { .. }
            | Self::Blocked { .. }
            | Self::Refused { .. }
            | Self::Intervention { .. }
            | Self::TooDeep { .. }
            | Self::InvalidTenant { .. }
            | Self::InvalidName { .. }
            | Self::InvalidRequestKey { .. }
            | Self::InvalidBound { .. }
            | Self::Incompatible { .. }
            | Self::Store { .. } => false,
        }
    }

    /// Whether this is a stop rather than a failure.
    #[must_use]
    pub fn is_cancelled(&self) -> bool {
        matches!(*self, Self::Cancelled { .. })
    }

    /// Where it happened; empty for a failure not yet located.
    #[must_use]
    pub fn at(&self) -> &str {
        match *self {
            Self::Cancelled { ref at }
            | Self::TimedOut { ref at, .. }
            | Self::Exhausted { ref at, .. }
            | Self::Throttled { ref at, .. }
            | Self::Failed { ref at, .. }
            | Self::Transient { ref at, .. }
            | Self::Bot { ref at, .. }
            | Self::Blocked { ref at, .. }
            | Self::Refused { ref at, .. }
            | Self::Intervention { ref at, .. }
            | Self::TooDeep { ref at, .. } => at,
            Self::Incompatible { ref at, .. } => at,
            Self::Store { ref at, .. } => at,
            // Every other variant is unlocated: it is raised before a step exists
            // to name, so there is no path to report.
            _ => "",
        }
    }

    /// Give an unlocated error the location `scope` names.
    ///
    /// An error that already has a location keeps it: the innermost location
    /// is the one that says where the failure happened.
    #[must_use]
    pub fn located(self, scope: &Scope) -> Self {
        self.located_at(scope.shared_path())
    }

    /// Whether this error is a refusal for missing authority, and the whole
    /// shortfall if it is.
    ///
    /// The one question a caller asking "was this run blocked rather than broken?"
    /// needs answered, and it answers with the complete shortfall rather than the
    /// first requirement, so the report that builds a repair ticket from it can be
    /// written once against the whole set.
    #[must_use]
    pub fn deficit(&self) -> Option<&crate::cap::Deficit> {
        match *self {
            Self::Blocked { ref deficit, .. } => Some(deficit),
            _ => None,
        }
    }

    /// [`FlowError::located`] against a path already in hand.
    #[must_use]
    pub(super) fn located_at(mut self, path: &Arc<str>) -> Self {
        match self {
            Self::Cancelled { ref mut at }
            | Self::TimedOut { ref mut at, .. }
            | Self::Exhausted { ref mut at, .. }
            | Self::Throttled { ref mut at, .. }
            | Self::Failed { ref mut at, .. }
            | Self::Transient { ref mut at, .. }
            | Self::Bot { ref mut at, .. }
            | Self::Blocked { ref mut at, .. }
            | Self::Refused { ref mut at, .. }
            | Self::Intervention { ref mut at, .. }
            | Self::TooDeep { ref mut at, .. }
            | Self::Incompatible { ref mut at, .. }
            | Self::Store { ref mut at, .. } => {
                if at.is_empty() {
                    *at = Arc::clone(path);
                }
            }
            // Every other variant carries no location: it is raised before a step
            // exists to name, so there is nothing to fill in.
            Self::InvalidTenant { .. }
            | Self::InvalidName { .. }
            | Self::InvalidRequestKey { .. }
            | Self::InvalidBound { .. } => {}
        }
        self
    }
}

impl fmt::Display for FlowError {
    /// `<path>: <what happened>`. Step-supplied text is rendered with `{:?}`,
    /// which escapes control characters, so a reason cannot forge a log line.
    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
        match *self {
            Self::Cancelled { ref at } => write!(formatter, "{at}: cancelled"),
            Self::TimedOut { ref at, after } => {
                write!(formatter, "{at}: timed out after {after:?}")
            }
            Self::Exhausted {
                ref at,
                attempts,
                ref last,
            } => write!(
                formatter,
                "{at}: gave up after {attempts} attempts; last: {last}"
            ),
            Self::Throttled {
                ref at,
                attempts,
                ref last,
            } => write!(
                formatter,
                "{at}: retry budget spent after {attempts} attempts; last: {last}"
            ),
            Self::Failed { ref at, ref reason } => write!(formatter, "{at}: failed: {reason:?}"),
            Self::Transient { ref at, ref reason } => {
                write!(formatter, "{at}: failed (transient): {reason:?}")
            }
            Self::Bot { ref at, ref source } => write!(formatter, "{at}: {source}"),
            Self::Blocked {
                ref at,
                ref deficit,
            } => {
                write!(formatter, "{at}: blocked: {deficit}")
            }
            Self::Refused {
                ref at,
                ref refusal,
                ref provenance,
            } => write!(formatter, "{at}: {refusal} (from {provenance})"),
            Self::Intervention {
                ref at,
                ref intervention,
            } => write!(formatter, "{at}: {intervention}"),
            Self::TooDeep { ref at, limit } => {
                write!(formatter, "{at}: steps nested deeper than {limit}")
            }
            Self::InvalidTenant { reason } => write!(formatter, "invalid tenant: {reason}"),
            Self::InvalidName { reason } => write!(formatter, "invalid task name: {reason}"),
            Self::InvalidRequestKey { reason } => {
                write!(formatter, "invalid request key: {reason}")
            }
            Self::InvalidBound { what, value, max } => {
                write!(formatter, "{what}: {value} is outside 1..={max}")
            }
            Self::Incompatible { ref at, ref drift } => write!(
                formatter,
                "{at}: refusing to resume under a different definition: {drift}"
            ),
            Self::Store { ref at, ref source } => {
                write!(formatter, "{at}: the run store refused: {source}")
            }
        }
    }
}

impl std::error::Error for FlowError {
    fn source(&self) -> Option<&(dyn std::error::Error + 'static)> {
        match *self {
            Self::Bot { ref source, .. } => Some(&**source),
            Self::Refused { ref refusal, .. } => Some(&**refusal),
            Self::Intervention {
                ref intervention, ..
            } => Some(&**intervention),
            Self::Store { ref source, .. } => Some(&**source),
            Self::Exhausted { ref last, .. } | Self::Throttled { ref last, .. } => Some(&**last),
            Self::Cancelled { .. }
            | Self::TimedOut { .. }
            | Self::Failed { .. }
            | Self::Transient { .. }
            | Self::Blocked { .. }
            | Self::TooDeep { .. }
            | Self::InvalidTenant { .. }
            | Self::InvalidName { .. }
            | Self::Incompatible { .. }
            | Self::InvalidRequestKey { .. }
            | Self::InvalidBound { .. } => None,
        }
    }
}

impl From<BotError> for FlowError {
    /// Lets `?` carry a verb's error; the location is filled at the boundary.
    fn from(source: BotError) -> Self {
        Self::Bot {
            at: Arc::from(""),
            source: Box::new(source),
        }
    }
}

impl From<std::io::Error> for FlowError {
    /// The kinds that describe a moment rather than a fact (a timeout, an
    /// interruption, a dropped connection) are transient; the rest are not.
    fn from(error: std::io::Error) -> Self {
        use std::io::ErrorKind;
        match error.kind() {
            ErrorKind::TimedOut
            | ErrorKind::Interrupted
            | ErrorKind::WouldBlock
            | ErrorKind::ConnectionReset
            | ErrorKind::ConnectionAborted
            | ErrorKind::ConnectionRefused
            | ErrorKind::NotConnected
            | ErrorKind::BrokenPipe
            | ErrorKind::UnexpectedEof => Self::transient(error),
            _ => Self::failed(error),
        }
    }
}

/// Turn a foreign `Result` into a flow result, saying whether its failure is
/// worth repeating.
///
/// In scope inside every flow `script!` writes, so a step reads
/// `parse(text).or_fail()?` or `fetch(url).await.or_retry()?`.
pub trait ResultExt<T> {
    /// A permanent failure on `Err`.
    ///
    /// # Errors
    ///
    /// [`FlowError::Failed`] carrying the error's text.
    fn or_fail(self) -> Result<T, FlowError>;

    /// A transient failure on `Err`: an enclosing `retry` repeats the step.
    ///
    /// # Errors
    ///
    /// [`FlowError::Transient`] carrying the error's text.
    fn or_retry(self) -> Result<T, FlowError>;
}

impl<T, E: fmt::Display> ResultExt<T> for Result<T, E> {
    fn or_fail(self) -> Result<T, FlowError> {
        self.map_err(FlowError::failed)
    }

    fn or_retry(self) -> Result<T, FlowError> {
        self.map_err(FlowError::transient)
    }
}

/// Turn an absent value into a permanent failure that says what was missing.
pub trait OptionExt<T> {
    /// The value, or [`FlowError::Failed`] with `reason`.
    ///
    /// # Errors
    ///
    /// [`FlowError::Failed`] when `None`.
    fn or_fail(self, reason: &str) -> Result<T, FlowError>;
}

impl<T> OptionExt<T> for Option<T> {
    fn or_fail(self, reason: &str) -> Result<T, FlowError> {
        self.ok_or_else(|| FlowError::failed(reason))
    }
}