use std::sync::Arc;
use async_trait::async_trait;
use serde_json::{Value, json};
use crate::core::{RunId, Seq, StoreError};
use crate::journal::{Record, RecordKind};
use super::delivery::Projection;
use super::{PushAuthentication, PushConfig, PushStore};
pub const OPERATOR_PREFIX: &str = "operator:";
#[must_use]
pub fn is_operator_id(id: &str) -> bool {
id.starts_with(OPERATOR_PREFIX)
}
#[derive(Debug, Clone)]
pub struct Destination {
pub name: String,
pub url: String,
pub authentication: Option<PushAuthentication>,
}
impl Destination {
#[must_use]
pub fn new(name: impl Into<String>, url: impl Into<String>) -> Self {
Self {
name: name.into(),
url: url.into(),
authentication: None,
}
}
#[must_use]
pub fn authenticated(
mut self,
scheme: impl Into<String>,
credentials: crate::core::Secret,
) -> Self {
self.authentication = Some(PushAuthentication {
scheme: scheme.into(),
credentials,
});
self
}
#[must_use]
pub fn registration_id(&self) -> String {
format!("{OPERATOR_PREFIX}{}", self.name)
}
fn config_for(&self, run: RunId) -> PushConfig {
PushConfig {
id: self.registration_id(),
task: run,
url: self.url.clone(),
token: None,
authentication: self.authentication.clone(),
}
}
}
#[derive(Debug, Clone)]
pub struct Outbox {
destinations: Vec<Destination>,
store: Arc<dyn PushStore>,
}
impl Outbox {
#[must_use]
pub fn new(store: Arc<dyn PushStore>, destinations: Vec<Destination>) -> Self {
let mut names = std::collections::BTreeSet::new();
for destination in &destinations {
assert!(
!destination.name.trim().is_empty(),
"an outbox destination needs a name: it is half of the stored \
registration id"
);
assert!(
names.insert(destination.name.clone()),
"outbox destination '{}' is configured twice — the two would share \
one cursor, so one's acknowledgement would discard the other's \
backlog",
destination.name
);
}
Self {
destinations,
store,
}
}
#[must_use]
pub fn destinations(&self) -> &[Destination] {
&self.destinations
}
#[cfg(feature = "keyring")]
#[must_use]
pub fn sealed(
mut self,
keys: Arc<dyn crate::keyring::KeyRing>,
tenant: crate::core::TenantId,
) -> Self {
self.store = crate::keyring::SealedPush::wrap(Arc::clone(&self.store), keys, tenant);
self
}
pub async fn open(&self, run: RunId) -> Result<(), StoreError> {
for destination in &self.destinations {
let from: Seq = 1;
self.store.put(&destination.config_for(run), from).await?;
}
Ok(())
}
}
#[derive(Debug, Clone)]
pub struct RunCompleted {
source: String,
event_type: String,
}
impl RunCompleted {
#[must_use]
pub fn new(source: impl Into<String>) -> Self {
Self {
source: source.into(),
event_type: "io.agentplane.run.completed".to_owned(),
}
}
#[must_use]
pub fn event_type(mut self, event_type: impl Into<String>) -> Self {
self.event_type = event_type.into();
self
}
}
#[async_trait]
impl Projection for RunCompleted {
async fn payloads(&self, record: &Record) -> Result<Vec<Value>, StoreError> {
let RecordKind::RunSealed {
outcome,
chain_head,
} = record.kind()
else {
return Ok(Vec::new());
};
Ok(vec![json!({
"specversion": "1.0",
"type": self.event_type,
"source": self.source,
"id": record.body.run.to_string(),
"datacontenttype": "application/json",
"data": {
"run": record.body.run.to_string(),
"case": record.body.case.map(|case| case.to_string()),
"outcome": outcome,
"chain_head": chain_head.to_hex(),
},
})])
}
fn namespace(&self) -> super::PushNamespace {
super::PushNamespace::Operator
}
}