use std::pin::Pin;
use aion_core::{ActivityEvent, ActivityId, RunId, WorkflowId};
use aion_proto::{ProtoActivityId, ProtoRunId, ProtoWorkflowId, TranscriptSubscription};
use futures::Stream;
use crate::{Client, ClientError};
pub type TranscriptStream =
Pin<Box<dyn Stream<Item = Result<TranscriptStreamItem, ClientError>> + Send>>;
#[derive(Clone, Debug, PartialEq)]
pub enum TranscriptStreamItem {
Event(Box<ActivityEvent>),
Lagged {
skipped: u64,
},
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub struct TranscriptTarget {
pub workflow_id: WorkflowId,
pub run_id: RunId,
pub activity_id: ActivityId,
pub attempt: u32,
pub after_seq: Option<u64>,
}
impl TranscriptTarget {
fn subscription(self, namespace: &str) -> TranscriptSubscription {
TranscriptSubscription {
namespace: namespace.to_owned(),
workflow_id: Some(ProtoWorkflowId::from(self.workflow_id)),
run_id: Some(ProtoRunId::from(self.run_id)),
activity_id: Some(ProtoActivityId::from(self.activity_id)),
attempt: self.attempt,
after_seq: self.after_seq,
}
}
}
impl Client {
pub async fn subscribe_transcript(
&self,
target: TranscriptTarget,
) -> Result<TranscriptStream, ClientError> {
crate::transport::transcript_ws::open(&self.config, target.subscription(self.namespace()))
.await
}
}