aion/error.rs
1//! Engine error taxonomy.
2
3use crate::schedule::{ScheduleError, ScheduleEvaluatorError};
4use aion_core::{RunId, ScheduleId, WorkflowId};
5use aion_package::{ContentHash, ContractIdentityError, PackageError};
6use aion_store::StoreError;
7
8use crate::durability::DurabilityError;
9
10/// Errors returned by the embedded workflow engine.
11#[derive(thiserror::Error, Debug)]
12pub enum EngineError {
13 /// The builder was asked to construct an engine without an event store.
14 #[error("engine store is required")]
15 MissingStore,
16
17 /// The builder was asked to construct an engine without a visibility store.
18 #[error(
19 "engine visibility store is required; call EngineBuilder::visibility_store() or EngineBuilder::in_memory_visibility()"
20 )]
21 MissingVisibilityStore,
22
23 /// A workflow package failed to load or validate for engine registration.
24 #[error("workflow package load failed: {reason}")]
25 Load {
26 /// Human-readable load failure reason.
27 reason: String,
28 },
29
30 /// A package offered for deployment declares a contract the engine cannot
31 /// enforce: at least one declared schema does not compile into a validator.
32 ///
33 /// Refused at the door rather than admitted, because the alternative is
34 /// silent: every admission boundary answers an uncompilable schema by
35 /// letting the value through unchecked, so a package that reaches the
36 /// catalog with one runs with that part of its declared contract switched
37 /// off and only a log line to say so.
38 #[error(
39 "workflow package `{workflow_type}` declares {count} type(s) the engine cannot compile into a validator, so nothing it declares there could ever be enforced: {detail}"
40 )]
41 UnenforceableContract {
42 /// Logical workflow type of the refused package.
43 workflow_type: String,
44 /// How many declarations could not be compiled.
45 count: usize,
46 /// Each unenforceable declaration, named, with the compiler's reason.
47 detail: String,
48 },
49
50 /// A route or unload targeted a `(workflow type, version)` that is not loaded.
51 #[error(
52 "workflow `{workflow_type}` version `{version}` is not loaded (loaded versions: {loaded})"
53 )]
54 UnknownVersion {
55 /// Logical workflow type requested by the caller.
56 workflow_type: String,
57 /// Content-hash version requested by the caller.
58 version: ContentHash,
59 /// Comma-separated loaded versions of the type, or `none`.
60 loaded: String,
61 },
62
63 /// An unload was refused because something still pins the version.
64 #[error("cannot unload workflow `{workflow_type}` version `{version}`: {pinned_by}")]
65 VersionPinned {
66 /// Logical workflow type targeted by the unload.
67 workflow_type: String,
68 /// Content-hash version targeted by the unload.
69 version: ContentHash,
70 /// What pins the version, naming the concrete holder.
71 pinned_by: PinHolder,
72 },
73
74 /// An unload was refused because the version is route-active for its type.
75 #[error(
76 "cannot unload workflow `{workflow_type}` version `{version}`: it is the route-active version; route another version first"
77 )]
78 RouteActive {
79 /// Logical workflow type targeted by the unload.
80 workflow_type: String,
81 /// Content-hash version targeted by the unload.
82 version: ContentHash,
83 },
84
85 /// An idempotent re-load presented the resident package identity with a
86 /// different manifest. V4 binds beams and the durable execution contract,
87 /// but not every packaging/admin field, so this remains the wrong-deploy
88 /// tripwire: the resident version is retained and the archive is refused.
89 #[error(
90 "workflow `{workflow_type}` version `{version}` is already loaded with a different manifest (resident digest {resident_digest}, incoming digest {incoming_digest}); rebuild the archive so its complete manifest matches the resident version, or change its contract-bound content"
91 )]
92 ManifestMismatch {
93 /// Logical workflow type targeted by the load.
94 workflow_type: String,
95 /// Content-hash version shared by both archives.
96 version: ContentHash,
97 /// Canonical digest of the resident manifest.
98 resident_digest: String,
99 /// Canonical digest of the incoming manifest.
100 incoming_digest: String,
101 },
102
103 /// The builder was given both `event_streaming` and an explicit event-publisher seam.
104 #[error(
105 "conflicting event publisher configuration: EngineBuilder::event_streaming installs the broadcast publisher and cannot be combined with EngineBuilder::event_publisher"
106 )]
107 ConflictingEventPublisher,
108
109 /// Live event streaming setup failed.
110 #[error("event streaming setup failed: {0}")]
111 EventStreaming(#[from] crate::publish::PublishError),
112
113 /// The configured event store returned an error.
114 #[error("store error: {0}")]
115 Store(#[from] StoreError),
116
117 /// The durability recorder or replay path returned an error.
118 #[error("durability error: {0}")]
119 Durability(#[from] DurabilityError),
120
121 /// A `.aion` package operation returned an error.
122 #[error("package error: {0}")]
123 Package(#[from] PackageError),
124
125 /// The selected package identity predates the `.v4` contract commitment.
126 #[error("workflow `{workflow_type}` cannot start: {source}")]
127 ContractIdentity {
128 /// Workflow type selected for the start.
129 workflow_type: String,
130 /// Typed migration refusal from the package identity boundary.
131 #[source]
132 source: ContractIdentityError,
133 },
134
135 /// The package names activities without a durable queue-scoped contract.
136 #[error(
137 "NO_QUEUE_DECLARATION: workflow `{workflow_type}` version `{version}` has unscoped activities {activities}; re-deploy from a checked AWL contract"
138 )]
139 NoQueueDeclaration {
140 /// Workflow type selected for the start.
141 workflow_type: String,
142 /// Exact `.v4` package identity selected for the run.
143 version: ContentHash,
144 /// Stable comma-separated unscoped activity names.
145 activities: String,
146 },
147
148 /// A start's input did not satisfy the declared input schema of the exact
149 /// package identity the start resolved to.
150 ///
151 /// Returned at the start boundary BEFORE any history is appended and
152 /// before any process is spawned, so a refused start leaves no trace: the
153 /// caller sees their own mistake at the moment they made it, with nothing
154 /// to clean up.
155 #[error(
156 "start input for workflow `{workflow_type}` does not satisfy the input type declared by package version `{version}`: {reason}"
157 )]
158 StartInputRefused {
159 /// Workflow type selected for the start.
160 workflow_type: String,
161 /// Exact `.v4` package identity the start resolved to.
162 version: ContentHash,
163 /// What did not match, naming every field that failed.
164 reason: String,
165 },
166
167 /// A signal was refused at the boundary: its name is not declared by the
168 /// target run's package, or its payload did not satisfy the declared
169 /// payload type.
170 ///
171 /// Returned BEFORE anything is recorded and before the arrival can be
172 /// consumed, so the target run's history is unchanged and it stays parked
173 /// on exactly the wait it was parked on. That ordering is the whole point:
174 /// a signal decoded after being consumed destroys a durable run that a
175 /// refusal merely inconveniences.
176 #[error(
177 "signal `{signal_name}` was refused for workflow `{workflow_id}` run `{run_id}` against the contract declared by package version `{version}`: {reason}"
178 )]
179 SignalRefused {
180 /// Workflow execution the signal targeted.
181 workflow_id: WorkflowId,
182 /// Concrete run the signal targeted.
183 run_id: RunId,
184 /// Signal name the caller sent.
185 signal_name: String,
186 /// Exact `.v4` package identity the target run is pinned to.
187 version: ContentHash,
188 /// Why the signal was refused — an undeclared name, or the fields of
189 /// the payload that did not match the declared type.
190 reason: String,
191 },
192
193 /// The embedded runtime returned an error.
194 #[error("runtime error: {reason}")]
195 Runtime {
196 /// Human-readable runtime failure reason.
197 reason: String,
198 },
199
200 /// A Gate-3 BIF required for tracked local fun spawns was not registered.
201 #[error("required Gate-3 BIF `{module}:{function}/{arity}` was missing during runtime startup")]
202 Gate3BifReplacementMissing {
203 /// Native module containing the required function.
204 module: String,
205 /// Required native function.
206 function: String,
207 /// Required native function arity.
208 arity: u8,
209 },
210
211 /// [`crate::Engine::run_startup_recovery`] was called on an engine whose
212 /// build was not deferred — `build()` already ran startup recovery, and
213 /// running it twice would re-dispatch every in-flight activity.
214 #[error(
215 "startup recovery was not deferred: EngineBuilder::build() already ran it \
216 (call defer_startup_recovery() on the builder to take ownership of the steps)"
217 )]
218 StartupRecoveryNotDeferred,
219
220 /// [`crate::Engine::run_startup_recovery`] was called a second time.
221 #[error("startup recovery already ran: run_startup_recovery() is one-shot")]
222 StartupRecoveryAlreadyRan,
223
224 /// [`crate::Engine::run_startup_catchup`] was called before the
225 /// workflow-recovery leg ran — catch-up delivers owed timer fires to
226 /// resident workflows, so residency recovery must precede it.
227 #[error(
228 "startup catch-up was requested before workflow recovery: call \
229 recover_workflows_on_startup() first"
230 )]
231 StartupCatchupBeforeWorkflowRecovery,
232
233 /// The deferred-startup-recovery slot lock was poisoned.
234 #[error("deferred startup recovery slot was poisoned")]
235 StartupRecoverySlotPoisoned,
236
237 /// The runtime-owned cleanup executor's ownership state was poisoned.
238 #[error("process cleanup executor state was poisoned")]
239 CleanupExecutorPoisoned,
240
241 /// The runtime cleanup worker did not stop within the configured bound.
242 #[error("process cleanup executor did not stop within {timeout_millis}ms")]
243 CleanupExecutorShutdownTimedOut {
244 /// Configured shutdown observation bound in milliseconds.
245 timeout_millis: u128,
246 },
247
248 /// The process-exit registry lifecycle lock was poisoned.
249 #[error("process exit registry lifecycle state was poisoned")]
250 ProcessExitRegistryPoisoned,
251
252 /// A process exit record's installation/abort ownership gate was poisoned.
253 #[error("process exit ownership gate for process {process_id} was poisoned")]
254 ProcessExitOwnershipPoisoned {
255 /// Process whose monitor/abort ownership could not be serialized.
256 process_id: u64,
257 },
258
259 /// A process exit record's fan-out state was poisoned.
260 #[error("process exit outcome state for process {process_id} was poisoned")]
261 ProcessExitStatePoisoned {
262 /// Process whose cached exit state could not be accessed.
263 process_id: u64,
264 },
265
266 /// The scheduler's one exit-event subscription was already claimed.
267 #[error("beamr process exit-event subscription is already owned")]
268 ProcessExitSubscriptionUnavailable,
269
270 /// The singleton process-exit drainer could not be spawned.
271 #[error("process exit drainer could not start: {reason}")]
272 ProcessExitDrainerSpawn {
273 /// Operating-system thread creation failure.
274 reason: String,
275 },
276
277 /// The singleton process-exit drainer's ownership lock was poisoned.
278 #[error("process exit drainer state was poisoned")]
279 ProcessExitDrainerPoisoned,
280
281 /// beamr published an exit event without the promised durable outcome.
282 #[error("process {process_id} exit event had no takeable outcome")]
283 ProcessExitOutcomeMissingAfterEvent {
284 /// Process named by the contract-breaking event.
285 process_id: u64,
286 },
287
288 /// beamr disconnected its event publisher while the runtime still owned it.
289 #[error("beamr process exit-event publisher disconnected")]
290 ProcessExitEventStreamDisconnected,
291
292 /// The process-exit drainer did not stop within the configured bound.
293 #[error("process exit drainer did not stop within {timeout_millis}ms")]
294 ProcessExitDrainerShutdownTimedOut {
295 /// Configured shutdown observation bound in milliseconds.
296 timeout_millis: u128,
297 },
298
299 /// The process-exit drainer thread panicked.
300 #[error("process exit drainer terminated unexpectedly")]
301 ProcessExitDrainerPanicked,
302
303 /// The process-exit callback dispatcher's ownership state was poisoned.
304 #[error("process exit callback dispatcher state was poisoned")]
305 ProcessExitCallbackDispatcherPoisoned,
306
307 /// The process-exit callback dispatcher had already stopped.
308 #[error("process exit callback dispatcher is unavailable")]
309 ProcessExitCallbackDispatcherUnavailable,
310
311 /// The process-exit callback dispatcher did not stop within its configured bound.
312 #[error("process exit callback dispatcher did not stop within {timeout_millis}ms")]
313 ProcessExitCallbackDispatcherShutdownTimedOut {
314 /// Configured shutdown observation bound in milliseconds.
315 timeout_millis: u128,
316 },
317
318 /// A retired process generation cannot accept another outcome consumer.
319 #[error("process {process_id} already reached its terminal runtime outcome")]
320 ProcessExitAlreadyTerminal {
321 /// Process generation whose heavyweight exit record was retired.
322 process_id: u64,
323 },
324
325 /// A workflow's activity-delivery synchronization lock was poisoned.
326 #[error("activity delivery lock for process {process_id} was poisoned")]
327 ActivityDeliveryPoisoned {
328 /// Workflow process whose scoped delivery lock was poisoned.
329 process_id: u64,
330 },
331
332 /// The active workflow registry lock was poisoned.
333 #[error("active workflow registry lock was poisoned")]
334 RegistryPoisoned,
335
336 /// A registered run has no `WorkflowStarted` in the history it was
337 /// reconciled against — the registry and the store disagree that it exists.
338 ///
339 /// Raised by registry reconciliation rather than defaulting the projection.
340 /// `status_from_events` returns `Running` for a slice holding no lifecycle
341 /// event, so a run absent from the history it is projected against would
342 /// otherwise be silently cached as RUNNING — a terminal run reported live,
343 /// produced by the reconciliation whose whole job is to stop exactly that.
344 ///
345 /// Not reachable through a normal start: `WorkflowStarted` is recorded
346 /// before the handle is published. It means a genuine invariant breach, so
347 /// it is surfaced rather than absorbed.
348 #[error(
349 "run {run_id} of workflow {workflow_id} is absent from the history it was reconciled against"
350 )]
351 RunNotInHistory {
352 /// Workflow whose history was read.
353 workflow_id: WorkflowId,
354 /// Run that the history does not contain.
355 run_id: RunId,
356 },
357
358 /// The workflow catalog lock was poisoned.
359 #[error("workflow catalog lock was poisoned")]
360 CatalogPoisoned,
361
362 /// A precondition on the target workflow's current state was not met.
363 ///
364 /// Raised by the reopen operation when the target run is not in a reopenable
365 /// state: not terminal, terminal for a non-reopenable reason
366 /// (Completed/`TimedOut`), or already Running. The `reason` names the actual
367 /// status so callers and operators can see why the reopen was rejected. Maps
368 /// to the `INVALID_STATE` wire code (gRPC `FailedPrecondition` / HTTP 409).
369 #[error("invalid workflow state: {reason}")]
370 InvalidState {
371 /// Human-readable precondition-failure reason naming the actual status.
372 reason: String,
373 },
374
375 /// The engine is already shutting down and no new workflow starts are accepted.
376 #[error("engine is shutting down")]
377 ShuttingDown,
378
379 /// No live, durable, or loaded workflow was found for the request.
380 #[error("workflow `{workflow_type}` was not found")]
381 WorkflowNotFound {
382 /// Logical workflow type requested by the caller.
383 workflow_type: String,
384 },
385
386 /// A terminal-writer reservation could not be taken because the workflow
387 /// already has a writer (#117(c)).
388 ///
389 /// The extraordinary cancellation path exists only for a run that can never
390 /// obtain a handle. A workflow that has one — or that another reservation is
391 /// already writing — is not that case, and taking a second writer would
392 /// break the single-writer invariant this refusal protects.
393 #[error("workflow `{workflow_id}` run `{run_id}` cannot take the terminal writer: {holder}")]
394 TerminalWriterUnavailable {
395 /// Workflow whose writer slot is occupied.
396 workflow_id: String,
397 /// Run the refused reservation named.
398 run_id: String,
399 /// What holds the slot, in the operator's terms.
400 holder: String,
401 },
402
403 /// A handle could not be registered because a terminal-writer reservation
404 /// holds this workflow's writer slot (#117(c)).
405 ///
406 /// The mirror of [`Self::TerminalWriterUnavailable`], and transient by
407 /// construction: a reservation lives only across one terminal transition.
408 #[error(
409 "workflow `{workflow_id}` cannot register a handle: run `{run_id}` holds the terminal writer"
410 )]
411 TerminalWriterHeld {
412 /// Workflow whose writer slot is reserved.
413 workflow_id: String,
414 /// Run holding the reservation.
415 run_id: String,
416 },
417
418 /// A terminal event was about to be appended after the engine-task epoch
419 /// had already closed.
420 ///
421 /// Raised at the append boundary itself, which is the only instant at which
422 /// the hazard it guards is real. The engine that owned this run has been
423 /// shut down or released, so this process is no longer that workflow's
424 /// single writer (invariant 3). Appending here risks two writers.
425 ///
426 /// # This is not only the successor case
427 ///
428 /// The obvious reading — a successor engine is already recovering the same
429 /// history — is the *eventual* case, not the whole of it. `Engine::shutdown`
430 /// closes the epoch as its FIRST act and only stops admitting process-exit
431 /// callbacks several steps later, so this error is also raised for runs that
432 /// exit during **this** engine's own graceful teardown, while no successor
433 /// exists yet. Saying "a successor may already be recovering" would tell an
434 /// operator reading the message during a clean shutdown to go looking for a
435 /// second node that is not there.
436 ///
437 /// Deliberately **not** transient: no later attempt re-opens a closed
438 /// epoch. In both cases the run stays `Running` and a startup sweep — the
439 /// successor's, or this node's own on restart — re-installs a monitor,
440 /// which is the mechanism that actually repairs it.
441 #[error(
442 "run `{run_id}` of workflow `{workflow_id}` could not append its terminal event: the \
443 engine-task epoch closed first, so this engine is no longer the run's single writer"
444 )]
445 EngineTaskEpochClosed {
446 /// Workflow whose terminal event was refused.
447 workflow_id: String,
448 /// Run whose terminal event was refused.
449 run_id: String,
450 },
451
452 /// The extraordinary cancellation path was asked for a run whose pinned
453 /// package resolves right now, so the run is recoverable (#117(c)).
454 ///
455 /// Measured at the moment of the request, never cited from an earlier boot's
456 /// verdict: a redeploy between then and now is exactly the remedy that makes
457 /// the ordinary path work again, and the ordinary path must be used when it
458 /// does.
459 #[error(
460 "workflow `{workflow_id}` run `{run_id}` is recoverable: its pinned package version `{version}` resolves, so it must be recovered and cancelled through the ordinary path"
461 )]
462 RunIsRecoverable {
463 /// Workflow the request named.
464 workflow_id: String,
465 /// Run the request named.
466 run_id: String,
467 /// The pinned package version that resolved.
468 version: String,
469 },
470
471 /// A run holds no handle, cannot obtain one, and this engine has no recorded
472 /// reason why (#117(c)).
473 ///
474 /// Distinct from [`Self::WorkflowNotFound`] on purpose: the run EXISTS and
475 /// its history is readable. What is absent is a verdict from this process's
476 /// startup recovery, so the extraordinary cancellation path — which must
477 /// cite that verdict — has nothing to cite.
478 #[error(
479 "workflow `{workflow_id}` run `{run_id}` is not resident and this engine recorded no reason it could not be made resident; it exists but cannot be cancelled here"
480 )]
481 NoResidencyVerdict {
482 /// Workflow the request named.
483 workflow_id: String,
484 /// Run the request named.
485 run_id: String,
486 },
487
488 /// No durable schedule was found for the request.
489 #[error("schedule `{schedule_id}` was not found")]
490 ScheduleNotFound {
491 /// Schedule identifier requested by the caller.
492 schedule_id: ScheduleId,
493 },
494
495 /// Schedule trigger, projection, or evaluator side effect failed.
496 #[error("schedule error: {reason}")]
497 Schedule {
498 /// Human-readable schedule failure reason.
499 reason: String,
500 },
501
502 /// Native implemented function registration failed.
503 #[error("NIF registration failed: {reason}")]
504 NifRegistration {
505 /// Human-readable native implemented function registration failure reason.
506 reason: String,
507 },
508
509 /// Signal routing failed after the target was resolved.
510 #[error("signal router error: {0}")]
511 SignalRouter(#[from] SignalRouterError),
512
513 /// Live workflow query dispatch failed after the target was resolved.
514 #[error("query error: {0}")]
515 Query(#[from] crate::query::QueryError),
516}
517
518/// What pins a workflow version against unload, naming the concrete holder.
519#[derive(Debug, Clone, PartialEq, Eq)]
520pub enum PinHolder {
521 /// A start resolved this version but has not yet registered a handle.
522 InFlightStart,
523 /// A live, non-terminal run executes on this version.
524 LiveRun {
525 /// Pinning workflow id.
526 workflow_id: WorkflowId,
527 /// Pinning run id.
528 run_id: RunId,
529 },
530 /// A recoverable instance in the store is pinned to this version.
531 RecoverableRun {
532 /// Pinning workflow id.
533 workflow_id: WorkflowId,
534 },
535 /// A recorded-but-never-started child is pinned to this version.
536 RecordedChild {
537 /// Child workflow id pinned to the version.
538 child_workflow_id: WorkflowId,
539 /// Parent workflow whose history records the child.
540 recorded_by: WorkflowId,
541 },
542}
543
544impl std::fmt::Display for PinHolder {
545 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
546 match self {
547 Self::InFlightStart => formatter.write_str("an in-flight start is pinned to it"),
548 Self::LiveRun {
549 workflow_id,
550 run_id,
551 } => write!(
552 formatter,
553 "live run `{workflow_id}/{run_id}` is pinned to it"
554 ),
555 Self::RecoverableRun { workflow_id } => {
556 write!(formatter, "recoverable run `{workflow_id}` is pinned to it")
557 }
558 Self::RecordedChild {
559 child_workflow_id,
560 recorded_by,
561 } => write!(
562 formatter,
563 "child `{child_workflow_id}` recorded by `{recorded_by}` is pinned to it and has not started"
564 ),
565 }
566 }
567}
568
569/// Errors surfaced by the signal routing boundary.
570#[derive(thiserror::Error, Debug, Clone, PartialEq, Eq)]
571pub enum SignalRouterError {
572 /// The target workflow is terminal and cannot receive new signals.
573 #[error("workflow {workflow_id}/{run_id} is terminal")]
574 Terminal {
575 /// Target workflow id.
576 workflow_id: WorkflowId,
577 /// Target run id.
578 run_id: RunId,
579 },
580
581 /// The router could not defer a recorded non-resident signal.
582 #[error("signal resume handoff failed: {reason}")]
583 Handoff {
584 /// Human-readable handoff failure reason.
585 reason: String,
586 },
587
588 /// The signal was durably recorded but could not be delivered to the live mailbox.
589 #[error(
590 "signal `{signal_name}` for workflow {workflow_id}/{run_id} could not be delivered to process {process_id}: {reason}"
591 )]
592 DeliveryFailed {
593 /// Target workflow id.
594 workflow_id: WorkflowId,
595 /// Target run id.
596 run_id: RunId,
597 /// Embedded runtime process identifier selected for delivery.
598 process_id: u64,
599 /// Signal name that was recorded and attempted.
600 signal_name: String,
601 /// Human-readable delivery failure reason.
602 reason: String,
603 },
604}
605
606impl From<ScheduleError> for EngineError {
607 fn from(error: ScheduleError) -> Self {
608 Self::Schedule {
609 reason: error.to_string(),
610 }
611 }
612}
613
614impl From<ScheduleEvaluatorError> for EngineError {
615 fn from(error: ScheduleEvaluatorError) -> Self {
616 match error {
617 ScheduleEvaluatorError::ScheduleNotFound { schedule_id } => {
618 Self::ScheduleNotFound { schedule_id }
619 }
620 other => Self::Schedule {
621 reason: other.to_string(),
622 },
623 }
624 }
625}