use std::sync::Arc;
use async_trait::async_trait;
use serde_json::json;
use crate::core::{CloudEvent, RunId, Seq, StoreError};
use crate::journal::{Record, RecordKind};
use super::delivery::Projection;
use super::{BodySigning, PushAuthentication, PushConfig, PushMessage, 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>,
pub signing: Option<BodySigning>,
}
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,
signing: 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 signed_with(mut self, secret: &crate::core::Secret) -> Self {
self.signing = Some(BodySigning::new(secret));
self
}
pub fn try_signed_with(
mut self,
secret: &crate::core::Secret,
) -> Result<Self, super::SigningKeyError> {
self.signing = Some(BodySigning::try_new(secret)?);
Ok(self)
}
#[must_use]
pub fn also_signed_with(mut self, secret: &crate::core::Secret) -> Self {
let signing = self
.signing
.take()
.expect("also_signed_with needs signed_with first: there is no primary key");
self.signing = Some(signing.also_with(secret));
self
}
pub fn try_also_signed_with(
mut self,
secret: &crate::core::Secret,
) -> Result<Self, super::SigningKeyError> {
let signing = self
.signing
.take()
.ok_or(super::SigningKeyError::NoPrimary)?;
self.signing = Some(signing.try_also_with(secret)?);
Ok(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
);
assert!(
reqwest::Url::parse(&destination.url).is_ok(),
"outbox destination '{}' has an unparseable URL '{}'",
destination.name,
destination.url
);
}
Self {
destinations,
store,
}
}
#[must_use]
pub fn destinations(&self) -> &[Destination] {
&self.destinations
}
pub(crate) fn store_tenant(&self) -> &str {
self.store.tenant()
}
#[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,
tenant: Option<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(),
tenant: None,
}
}
#[must_use]
pub fn for_tenant(mut self, tenant: impl Into<String>) -> Self {
self.tenant = Some(tenant.into());
self
}
#[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 messages(&self, record: &Record) -> Result<Vec<PushMessage>, StoreError> {
let RecordKind::RunConcluded {
outcome,
reason,
exhaustion,
live_spend,
chain_head,
} = record.kind()
else {
return Ok(Vec::new());
};
let run = record.body.run.to_string();
let mut event = CloudEvent::new(&self.source, &run, &self.event_type)
.map_err(|error| StoreError::Backend(error.to_string()))?;
if let Some(tenant) = &self.tenant {
event = event
.with_extension("tenantid", serde_json::Value::String(tenant.clone()))
.map_err(|error| StoreError::Backend(error.to_string()))?;
}
let mut data = json!({
"run": run,
"case": record.body.case.map(|case| case.to_string()),
"outcome": outcome,
"chain_head": chain_head.to_hex(),
});
if let Some(reason) = reason {
data["reason"] = json!(reason);
}
if let Some(exhaustion) = exhaustion {
data["exhaustion"] = json!(exhaustion);
}
if !live_spend.is_free() {
data["live_spend"] = json!(live_spend);
}
let event = event
.with_subject(&run)
.with_data(data);
Ok(vec![PushMessage::cloudevent(&event)])
}
fn namespace(&self) -> super::PushNamespace {
super::PushNamespace::Operator
}
}