Skip to main content

aion/query/
service.rs

1//! Query dispatch service for live workflow processes.
2
3use std::sync::Arc;
4use std::time::Duration;
5
6use aion_core::{Payload, WorkflowId};
7use tokio::sync::oneshot;
8use tokio::time;
9
10use crate::engine_seam::{
11    EngineHandle, EngineSeamError, WorkflowMailboxMessage, WorkflowProcessHandle, WorkflowResidency,
12};
13
14/// Result sent by workflow query handlers over a query reply channel.
15pub type QueryResult = Result<Payload, QueryError>;
16
17/// Result returned by [`QueryService::query`].
18pub type QueryServiceResult = Result<Payload, QueryError>;
19
20/// Typed failures surfaced by live workflow query dispatch.
21#[derive(thiserror::Error, Debug, Clone, PartialEq, Eq)]
22pub enum QueryError {
23    /// The resident workflow has no registered handler for the requested query name.
24    #[error("unknown query {0}")]
25    UnknownQuery(String),
26
27    /// No query reply arrived before the engine-configured timeout elapsed.
28    #[error("query reply timed out")]
29    Timeout,
30
31    /// The workflow cannot answer a live query because it is not currently running.
32    #[error("workflow {0} is not running")]
33    NotRunning(WorkflowId),
34
35    /// The engine does not know the requested workflow.
36    #[error("workflow {0} is unknown")]
37    Unknown(WorkflowId),
38
39    /// The workflow query reply channel closed before a handler response was sent.
40    #[error("query reply channel closed before a handler response was sent")]
41    ReplyDropped,
42
43    /// The workflow's query handler ran and reported an application-level failure.
44    #[error("query handler failed: {message}")]
45    HandlerFailed {
46        /// Failure reason reported by the workflow's query handler.
47        message: String,
48    },
49
50    /// The caller's query arguments payload is not a well-formed JSON document.
51    ///
52    /// Query arguments cross into workflow code as the JSON text the caller
53    /// supplied. A payload that is not valid UTF-8 JSON can never reach a
54    /// handler, so it is refused at the delivery boundary and attributed to
55    /// the caller — never forwarded to become a misleading handler failure.
56    #[error("query arguments are invalid: {reason}")]
57    InvalidArguments {
58        /// Why the arguments payload could not be carried to the handler.
59        reason: String,
60    },
61
62    /// The engine seam failed while resolving or delivering the query.
63    #[error("query engine seam failed: {0}")]
64    Engine(#[from] EngineSeamError),
65}
66
67/// Non-recording, non-disruptive live workflow query dispatcher.
68///
69/// `QueryService` depends only on the engine seam's residency and mailbox-delivery operations. It
70/// has no durable history dependency and no persistence method, so query dispatch is structurally a
71/// read-only interaction. The delivered [`WorkflowMailboxMessage::Query`] is a distinct message
72/// kind carrying a one-shot reply channel; AE/workflow processes answer it at deterministic yield
73/// points from registered read-only handlers so in-progress workflow steps are not preempted or
74/// mutated.
75#[derive(Debug)]
76pub struct QueryService<H: ?Sized> {
77    engine: Arc<H>,
78    query_timeout: Duration,
79}
80
81impl<H> QueryService<H>
82where
83    H: EngineHandle + ?Sized,
84{
85    /// Creates a query service with an engine-configured timeout.
86    #[must_use]
87    pub fn new(engine: Arc<H>, query_timeout: Duration) -> Self {
88        Self {
89            engine,
90            query_timeout,
91        }
92    }
93
94    /// Dispatches a read-only query to a resident workflow and returns the handler reply payload.
95    ///
96    /// # Errors
97    ///
98    /// Returns [`QueryError::Unknown`] for unknown workflows, [`QueryError::NotRunning`] for
99    /// terminal or non-resident workflows, [`QueryError::Timeout`] when no handler reply arrives
100    /// before the configured timeout, [`QueryError::UnknownQuery`] when the workflow replies that
101    /// no handler exists, and [`QueryError::Engine`] for seam failures.
102    pub async fn query(
103        &self,
104        workflow_id: &WorkflowId,
105        name: impl Into<String>,
106        args: Payload,
107    ) -> QueryServiceResult {
108        let process = match self.engine.resolve_workflow(workflow_id)? {
109            WorkflowResidency::Resident(process) => process,
110            WorkflowResidency::NonResident | WorkflowResidency::Terminal => {
111                return Err(QueryError::NotRunning(workflow_id.clone()));
112            }
113            WorkflowResidency::Unknown => return Err(QueryError::Unknown(workflow_id.clone())),
114        };
115        self.query_process(process, name, args).await
116    }
117
118    /// Dispatches a read-only query to an already-resolved workflow process.
119    ///
120    /// Run-exact variant of [`Self::query`] for callers that resolved the
121    /// target handle themselves (the engine seam resolves `(workflow, run)`
122    /// before delegation, so re-resolving by workflow id here would race
123    /// continue-as-new and multi-run histories).
124    ///
125    /// # Errors
126    ///
127    /// Returns [`QueryError::Timeout`] when no handler reply arrives before the configured
128    /// timeout, [`QueryError::UnknownQuery`] when the workflow replies that no handler exists,
129    /// [`QueryError::HandlerFailed`] when the handler ran and reported failure,
130    /// [`QueryError::ReplyDropped`] when the workflow ended before answering, and
131    /// [`QueryError::Engine`] for seam failures.
132    pub async fn query_process(
133        &self,
134        process: WorkflowProcessHandle,
135        name: impl Into<String>,
136        args: Payload,
137    ) -> QueryServiceResult {
138        let (reply_to, reply_from) = oneshot::channel();
139        self.engine.deliver_workflow_message(
140            process,
141            WorkflowMailboxMessage::Query {
142                name: name.into(),
143                payload: args,
144                reply_to,
145            },
146        )?;
147
148        match time::timeout(self.query_timeout, reply_from).await {
149            Ok(Ok(reply)) => reply,
150            Ok(Err(_)) => Err(QueryError::ReplyDropped),
151            Err(_) => Err(QueryError::Timeout),
152        }
153    }
154
155    /// Returns the engine-configured timeout used for query replies.
156    #[must_use]
157    pub const fn query_timeout(&self) -> Duration {
158        self.query_timeout
159    }
160}
161
162#[cfg(test)]
163mod tests {
164    use std::collections::HashMap;
165    use std::sync::{Arc, Mutex, MutexGuard};
166    use std::time::Duration;
167
168    use aion_core::{ContentType, Event, Payload, TimerId, WorkflowId};
169    use aion_store::{InMemoryStore, ReadableEventStore};
170
171    use super::{QueryError, QueryService};
172    use crate::Pid;
173    use crate::engine_seam::{
174        ChildWorkflowSpawnRequest, ChildWorkflowSpawnResult, EngineHandle, EngineSeamError,
175        TimerWheelEntry, WorkflowMailboxMessage, WorkflowProcessHandle, WorkflowResidency,
176    };
177
178    const QUERY_TIMEOUT: Duration = Duration::from_millis(10);
179
180    #[derive(Clone)]
181    enum QueryBehavior {
182        Reply(Payload),
183        Fail(String),
184        HoldSender,
185    }
186
187    #[derive(Default)]
188    struct FakeQueryWorkflow {
189        handlers: HashMap<String, QueryBehavior>,
190        query_count: usize,
191        last_payload: Option<Payload>,
192    }
193
194    #[derive(Default)]
195    struct FakeQueryEngineState {
196        residency: HashMap<WorkflowId, WorkflowResidency>,
197        workflows: HashMap<WorkflowProcessHandle, FakeQueryWorkflow>,
198        held_replies: Vec<crate::engine_seam::QueryReplySender>,
199    }
200
201    #[derive(Default)]
202    struct FakeQueryEngine {
203        state: Mutex<FakeQueryEngineState>,
204    }
205
206    impl FakeQueryEngine {
207        fn set_resident_workflow(
208            &self,
209            workflow_id: WorkflowId,
210            process: WorkflowProcessHandle,
211            workflow: FakeQueryWorkflow,
212        ) -> Result<(), EngineSeamError> {
213            let mut state = self.state()?;
214            state
215                .residency
216                .insert(workflow_id, WorkflowResidency::Resident(process));
217            state.workflows.insert(process, workflow);
218            Ok(())
219        }
220
221        fn set_residency(
222            &self,
223            workflow_id: WorkflowId,
224            residency: WorkflowResidency,
225        ) -> Result<(), EngineSeamError> {
226            self.state()?.residency.insert(workflow_id, residency);
227            Ok(())
228        }
229
230        fn query_count(&self, process: WorkflowProcessHandle) -> Result<usize, EngineSeamError> {
231            Ok(self
232                .state()?
233                .workflows
234                .get(&process)
235                .map_or(0, |workflow| workflow.query_count))
236        }
237
238        fn last_payload(
239            &self,
240            process: WorkflowProcessHandle,
241        ) -> Result<Option<Payload>, EngineSeamError> {
242            Ok(self
243                .state()?
244                .workflows
245                .get(&process)
246                .and_then(|workflow| workflow.last_payload.clone()))
247        }
248
249        fn state(&self) -> Result<MutexGuard<'_, FakeQueryEngineState>, EngineSeamError> {
250            self.state.lock().map_err(|_| EngineSeamError::Delivery {
251                reason: "fake query engine state lock was poisoned".to_owned(),
252            })
253        }
254    }
255
256    impl EngineHandle for FakeQueryEngine {
257        fn resolve_workflow(
258            &self,
259            workflow_id: &WorkflowId,
260        ) -> Result<WorkflowResidency, EngineSeamError> {
261            Ok(self
262                .state()?
263                .residency
264                .get(workflow_id)
265                .copied()
266                .unwrap_or(WorkflowResidency::Unknown))
267        }
268
269        fn deliver_workflow_message(
270            &self,
271            process: WorkflowProcessHandle,
272            message: WorkflowMailboxMessage,
273        ) -> Result<(), EngineSeamError> {
274            match message {
275                WorkflowMailboxMessage::Query {
276                    name,
277                    payload,
278                    reply_to,
279                } => {
280                    let mut state = self.state()?;
281                    let behavior = {
282                        let workflow = state.workflows.get_mut(&process).ok_or_else(|| {
283                            EngineSeamError::Delivery {
284                                reason: "query target process was not registered".to_owned(),
285                            }
286                        })?;
287                        workflow.last_payload = Some(payload);
288                        workflow.query_count += 1;
289                        workflow.handlers.get(&name).cloned()
290                    };
291
292                    match behavior {
293                        Some(QueryBehavior::Reply(payload)) => {
294                            if reply_to.send(Ok(payload)).is_err() {
295                                return Err(EngineSeamError::Delivery {
296                                    reason: "query caller dropped reply receiver".to_owned(),
297                                });
298                            }
299                        }
300                        Some(QueryBehavior::Fail(message)) => {
301                            if reply_to
302                                .send(Err(QueryError::HandlerFailed { message }))
303                                .is_err()
304                            {
305                                return Err(EngineSeamError::Delivery {
306                                    reason: "query caller dropped reply receiver".to_owned(),
307                                });
308                            }
309                        }
310                        None => {
311                            if reply_to.send(Err(QueryError::UnknownQuery(name))).is_err() {
312                                return Err(EngineSeamError::Delivery {
313                                    reason: "query caller dropped reply receiver".to_owned(),
314                                });
315                            }
316                        }
317                        Some(QueryBehavior::HoldSender) => state.held_replies.push(reply_to),
318                    }
319                    Ok(())
320                }
321                _ => Err(EngineSeamError::Delivery {
322                    reason: "fake query engine only accepts query messages".to_owned(),
323                }),
324            }
325        }
326
327        fn spawn_child_workflow(
328            &self,
329            request: ChildWorkflowSpawnRequest,
330        ) -> Result<ChildWorkflowSpawnResult, EngineSeamError> {
331            Err(EngineSeamError::ChildSpawn {
332                reason: format!(
333                    "fake query engine does not spawn child workflow {}",
334                    request.workflow_type
335                ),
336            })
337        }
338
339        fn terminate_linked_child_workflow(
340            &self,
341            parent_workflow_id: &WorkflowId,
342            child_process: WorkflowProcessHandle,
343            correlation: u64,
344        ) -> Result<(), EngineSeamError> {
345            Err(EngineSeamError::ChildTermination {
346                reason: format!(
347                    "fake query engine does not terminate child workflow process {} for parent {parent_workflow_id} with correlation {correlation}",
348                    child_process.pid()
349                ),
350            })
351        }
352
353        fn terminate_linked_activity(
354            &self,
355            parent_workflow_id: &WorkflowId,
356            activity_process: Pid,
357            correlation: u64,
358        ) -> Result<(), EngineSeamError> {
359            Err(EngineSeamError::ChildTermination {
360                reason: format!(
361                    "fake query engine does not terminate activity process {activity_process} for parent {parent_workflow_id} with correlation {correlation}"
362                ),
363            })
364        }
365
366        fn arm_timer(&self, entry: TimerWheelEntry) -> Result<(), EngineSeamError> {
367            Err(EngineSeamError::TimerWheel {
368                reason: format!("fake query engine does not arm timer {}", entry.timer_id),
369            })
370        }
371
372        fn disarm_timer(
373            &self,
374            process: WorkflowProcessHandle,
375            timer_id: &TimerId,
376        ) -> Result<(), EngineSeamError> {
377            Err(EngineSeamError::TimerWheel {
378                reason: format!(
379                    "fake query engine does not disarm timer {timer_id} for process {}",
380                    process.pid()
381                ),
382            })
383        }
384
385        fn record_workflow_event(
386            &self,
387            workflow_id: &WorkflowId,
388            event: Event,
389        ) -> Result<crate::engine_seam::RecordOutcome, EngineSeamError> {
390            Err(EngineSeamError::Recorder {
391                reason: format!(
392                    "queries must not record event {} for workflow {workflow_id}",
393                    event.seq()
394                ),
395            })
396        }
397    }
398
399    fn payload(label: &str) -> Payload {
400        Payload::new(
401            ContentType::Json,
402            format!("{{\"label\":\"{label}\"}}").into_bytes(),
403        )
404    }
405
406    fn known_workflow(reply: Payload) -> FakeQueryWorkflow {
407        let mut handlers = HashMap::new();
408        handlers.insert("state".to_owned(), QueryBehavior::Reply(reply));
409        FakeQueryWorkflow {
410            handlers,
411            query_count: 0,
412            last_payload: None,
413        }
414    }
415
416    #[tokio::test]
417    async fn query_returns_registered_handler_reply() -> Result<(), Box<dyn std::error::Error>> {
418        let engine = Arc::new(FakeQueryEngine::default());
419        let workflow_id = WorkflowId::new_v4();
420        let process = WorkflowProcessHandle::new(7);
421        let reply = payload("answer");
422        engine.set_resident_workflow(
423            workflow_id.clone(),
424            process,
425            known_workflow(reply.clone()),
426        )?;
427        let service = QueryService::new(Arc::clone(&engine), QUERY_TIMEOUT);
428
429        let returned = service
430            .query(&workflow_id, "state", payload("args"))
431            .await?;
432
433        assert_eq!(returned, reply);
434        assert_eq!(engine.query_count(process)?, 1);
435        assert_eq!(engine.last_payload(process)?, Some(payload("args")));
436        Ok(())
437    }
438
439    #[tokio::test]
440    async fn query_does_not_record_events() -> Result<(), Box<dyn std::error::Error>> {
441        let store = InMemoryStore::default();
442        let engine = Arc::new(FakeQueryEngine::default());
443        let workflow_id = WorkflowId::new_v4();
444        let process = WorkflowProcessHandle::new(8);
445        engine.set_resident_workflow(
446            workflow_id.clone(),
447            process,
448            known_workflow(payload("visible-state")),
449        )?;
450        let service = QueryService::new(engine, QUERY_TIMEOUT);
451
452        let reply = service
453            .query(&workflow_id, "state", payload("args"))
454            .await?;
455        assert_eq!(reply, payload("visible-state"));
456
457        let history = store.read_history(&workflow_id).await?;
458        assert!(history.is_empty());
459        Ok(())
460    }
461
462    #[tokio::test]
463    async fn unknown_query_returns_typed_error_and_workflow_remains_live()
464    -> Result<(), Box<dyn std::error::Error>> {
465        let engine = Arc::new(FakeQueryEngine::default());
466        let workflow_id = WorkflowId::new_v4();
467        let process = WorkflowProcessHandle::new(9);
468        engine.set_resident_workflow(
469            workflow_id.clone(),
470            process,
471            known_workflow(payload("known")),
472        )?;
473        let service = QueryService::new(Arc::clone(&engine), QUERY_TIMEOUT);
474
475        let result = service
476            .query(&workflow_id, "missing", payload("args"))
477            .await;
478
479        assert_eq!(result, Err(QueryError::UnknownQuery("missing".to_owned())));
480        assert_eq!(
481            engine.resolve_workflow(&workflow_id)?,
482            WorkflowResidency::Resident(process)
483        );
484        assert_eq!(engine.query_count(process)?, 1);
485        Ok(())
486    }
487
488    #[tokio::test]
489    async fn non_replying_workflow_times_out() -> Result<(), Box<dyn std::error::Error>> {
490        let engine = Arc::new(FakeQueryEngine::default());
491        let workflow_id = WorkflowId::new_v4();
492        let process = WorkflowProcessHandle::new(10);
493        let mut handlers = HashMap::new();
494        handlers.insert("slow".to_owned(), QueryBehavior::HoldSender);
495        engine.set_resident_workflow(
496            workflow_id.clone(),
497            process,
498            FakeQueryWorkflow {
499                handlers,
500                query_count: 0,
501                last_payload: None,
502            },
503        )?;
504        let service = QueryService::new(engine, QUERY_TIMEOUT);
505
506        let result = service.query(&workflow_id, "slow", payload("args")).await;
507
508        assert_eq!(result, Err(QueryError::Timeout));
509        Ok(())
510    }
511
512    #[tokio::test]
513    async fn terminal_and_non_resident_workflows_are_not_running()
514    -> Result<(), Box<dyn std::error::Error>> {
515        let engine = Arc::new(FakeQueryEngine::default());
516        let terminal_id = WorkflowId::new_v4();
517        let non_resident_id = WorkflowId::new_v4();
518        engine.set_residency(terminal_id.clone(), WorkflowResidency::Terminal)?;
519        engine.set_residency(non_resident_id.clone(), WorkflowResidency::NonResident)?;
520        let service = QueryService::new(engine, QUERY_TIMEOUT);
521
522        let terminal_result = service.query(&terminal_id, "state", payload("args")).await;
523        let non_resident_result = service
524            .query(&non_resident_id, "state", payload("args"))
525            .await;
526
527        assert_eq!(terminal_result, Err(QueryError::NotRunning(terminal_id)));
528        assert_eq!(
529            non_resident_result,
530            Err(QueryError::NotRunning(non_resident_id))
531        );
532        Ok(())
533    }
534
535    #[tokio::test]
536    async fn query_process_dispatches_to_the_resolved_process_without_resolving()
537    -> Result<(), Box<dyn std::error::Error>> {
538        let engine = Arc::new(FakeQueryEngine::default());
539        // The workflow id is deliberately never registered for residency:
540        // query_process must not resolve, only deliver to the given process.
541        let workflow_id = WorkflowId::new_v4();
542        let process = WorkflowProcessHandle::new(11);
543        let reply = payload("run-exact");
544        engine.set_resident_workflow(
545            workflow_id.clone(),
546            process,
547            known_workflow(reply.clone()),
548        )?;
549        engine.set_residency(workflow_id, WorkflowResidency::Unknown)?;
550        let service = QueryService::new(Arc::clone(&engine), QUERY_TIMEOUT);
551
552        let returned = service
553            .query_process(process, "state", payload("args"))
554            .await?;
555
556        assert_eq!(returned, reply);
557        assert_eq!(engine.query_count(process)?, 1);
558        Ok(())
559    }
560
561    #[tokio::test]
562    async fn handler_failure_propagates_as_typed_handler_failed()
563    -> Result<(), Box<dyn std::error::Error>> {
564        let engine = Arc::new(FakeQueryEngine::default());
565        let workflow_id = WorkflowId::new_v4();
566        let process = WorkflowProcessHandle::new(12);
567        let mut handlers = HashMap::new();
568        handlers.insert(
569            "state".to_owned(),
570            QueryBehavior::Fail("handler raised".to_owned()),
571        );
572        engine.set_resident_workflow(
573            workflow_id.clone(),
574            process,
575            FakeQueryWorkflow {
576                handlers,
577                query_count: 0,
578                last_payload: None,
579            },
580        )?;
581        let service = QueryService::new(Arc::clone(&engine), QUERY_TIMEOUT);
582
583        let resolved = service.query(&workflow_id, "state", payload("args")).await;
584        let run_exact = service
585            .query_process(process, "state", payload("args"))
586            .await;
587
588        let expected = Err(QueryError::HandlerFailed {
589            message: "handler raised".to_owned(),
590        });
591        assert_eq!(resolved, expected);
592        assert_eq!(run_exact, expected);
593        Ok(())
594    }
595
596    #[tokio::test]
597    async fn unknown_workflow_returns_typed_unknown_error() -> Result<(), Box<dyn std::error::Error>>
598    {
599        let engine = Arc::new(FakeQueryEngine::default());
600        let workflow_id = WorkflowId::new_v4();
601        let service = QueryService::new(engine, QUERY_TIMEOUT);
602
603        let result = service.query(&workflow_id, "state", payload("args")).await;
604
605        assert_eq!(result, Err(QueryError::Unknown(workflow_id)));
606        Ok(())
607    }
608}