use std::num::NonZeroU64;
use std::sync::Arc;
use aion_client::Client;
use aion_core::{RunId, WorkflowId, status_from_events};
use serde::Serialize;
use tokio::runtime::Runtime;
use super::client::phase_from_status;
use super::error::WorkflowCallError;
use super::id::WorkflowIdentity;
use super::observe::WorkflowObservation;
use super::outcome::WorkflowInspection;
pub struct WorkflowConversation {
pub(super) runtime: Arc<Runtime>,
pub(super) client: Client,
pub(super) identity: WorkflowIdentity,
}
impl std::fmt::Debug for WorkflowConversation {
fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
formatter
.debug_struct("WorkflowConversation")
.field("identity", &self.identity)
.finish_non_exhaustive()
}
}
impl WorkflowConversation {
#[must_use]
pub const fn identity(&self) -> WorkflowIdentity {
self.identity
}
pub fn contribute<T>(&self, name: &str, value: &T) -> Result<(), WorkflowCallError>
where
T: Serialize + ?Sized,
{
let workflow_id = self.workflow_id();
let run_id = self.run_id();
self.runtime
.block_on(
self.client
.signal_typed(&workflow_id, Some(&run_id), name, value),
)
.map_err(|error| WorkflowCallError::from_client(&error))
}
#[must_use]
pub fn observe_from(&self, resume_from: NonZeroU64) -> WorkflowObservation {
let stream = self
.client
.subscribe_workflow_from(&self.workflow_id(), resume_from);
WorkflowObservation::new(Arc::clone(&self.runtime), stream)
}
pub fn inspect(&self) -> Result<WorkflowInspection, WorkflowCallError> {
let workflow_id = self.workflow_id();
let run_id = self.run_id();
let description = self
.runtime
.block_on(self.client.describe(&workflow_id, Some(&run_id)))
.map_err(|error| WorkflowCallError::from_client(&error))?;
Ok(WorkflowInspection {
phase: phase_from_status(status_from_events(&description.history)),
recorded_events: description.history.len() as u64,
})
}
fn workflow_id(&self) -> WorkflowId {
WorkflowId::new(self.identity.conversation.as_uuid())
}
fn run_id(&self) -> RunId {
RunId::new(self.identity.run.as_uuid())
}
}