use std::{
sync::{Arc, Mutex, PoisonError},
time::Duration,
};
use serde::{Deserialize, Serialize};
use crate::{
effect::{EffectFamily, EffectId, HandlerKey, Outcome},
error::ErrorReport,
wasm_compat::{WasmCompatSend, WasmCompatSync},
};
mod adapter;
pub(crate) mod sse_tail;
pub use adapter::{
AdapterAnalysis, AdapterContext, AdapterEnding, AdapterErrorBoundary, AdapterErrorEnvelope,
AdapterEvent, AdapterObservation, AdapterUsage, AdapterVerdict, ObservationSink,
diagnostic_url_secrets, scrub_diagnostic,
};
pub(crate) use adapter::{AdapterSlot, ObservedError};
#[cfg(test)]
mod tests;
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct Observation {
pub seq: u64,
pub subject: Subject,
pub stage: Stage,
pub emitter: Emitter,
pub action: Action,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub at: Option<Duration>,
}
impl Observation {
pub fn new(subject: Subject, stage: Stage, emitter: Emitter, action: Action) -> Self {
Self {
seq: 0,
subject,
stage,
emitter,
action,
at: None,
}
}
}
#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
pub struct Subject {
#[serde(default, skip_serializing_if = "Option::is_none")]
pub scope: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub order: Option<u64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub effect: Option<EffectId>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub parent: Option<EffectId>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub key: Option<HandlerKey>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub family: Option<EffectFamily>,
}
impl Subject {
pub fn scoped(scope: impl Into<String>) -> Self {
Self {
scope: Some(scope.into()),
..Self::default()
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum Stage {
Gate,
Dispatch,
Handler,
Collect,
Judge,
Runtime,
Host,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct Emitter {
pub name: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub version: Option<String>,
}
impl Emitter {
pub fn named(name: impl Into<String>) -> Self {
Self {
name: name.into(),
version: None,
}
}
pub fn versioned(name: impl Into<String>, version: impl Into<String>) -> Self {
Self {
name: name.into(),
version: Some(version.into()),
}
}
pub fn unknown() -> Self {
Self::named("unknown")
}
pub fn is_unknown(&self) -> bool {
self.name == "unknown"
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct Reason {
pub code: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub detail: Option<String>,
}
impl Reason {
pub fn code(code: impl Into<String>) -> Self {
Self {
code: code.into(),
detail: None,
}
}
pub fn with_detail(code: impl Into<String>, detail: impl Into<String>) -> Self {
Self {
code: code.into(),
detail: Some(detail.into()),
}
}
pub fn from_report(report: &ErrorReport) -> Self {
Self::with_detail(report.kind.code(), report.message.clone())
}
pub fn unknown() -> Self {
Self::code("unknown")
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(tag = "outcome", rename_all = "snake_case")]
pub enum OutcomeSummary {
Ok {
family: EffectFamily,
},
Err {
reason: Reason,
retryable: bool,
},
}
impl OutcomeSummary {
pub fn of(outcome: &Result<Outcome, ErrorReport>) -> Self {
match outcome {
Ok(outcome) => Self::Ok {
family: outcome.family(),
},
Err(report) => Self::Err {
reason: Reason::from_report(report),
retryable: report.retryable,
},
}
}
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
#[serde(tag = "action", rename_all = "snake_case")]
pub enum Action {
Adapter {
observation: AdapterObservation,
},
Held {
reason: Reason,
},
Released,
Denied {
reason: Reason,
},
Issued,
Refused {
reason: Reason,
},
Landed {
outcome: OutcomeSummary,
},
StreamTruncated {
delivered: usize,
tail: Vec<String>,
errors: Vec<Reason>,
},
Replaced {
recorded: OutcomeSummary,
consumed: OutcomeSummary,
},
Cancelled {
reason: Reason,
},
Ended {
ending: Reason,
},
Host {
kind: String,
payload: serde_json::Value,
},
}
pub const LARGEST_PAYLOAD_BYTES: usize = 64 * 1024;
impl Action {
pub fn stream_truncated(delivered: usize, mut tail: Vec<String>, errors: Vec<Reason>) -> Self {
while !tail.is_empty()
&& serde_json::to_vec(&tail).map_or(usize::MAX, |bytes| bytes.len())
> LARGEST_PAYLOAD_BYTES
{
tail.remove(0);
}
Self::StreamTruncated {
delivered,
tail,
errors,
}
}
}
pub trait HostAction: Serialize + serde::de::DeserializeOwned {
const KIND: &'static str;
fn action(&self) -> Result<Action, serde_json::Error> {
Ok(Action::Host {
kind: Self::KIND.to_owned(),
payload: serde_json::to_value(self)?,
})
}
fn from_action(action: &Action) -> Option<Result<Self, serde_json::Error>> {
match action {
Action::Host { kind, payload } if kind == Self::KIND => {
Some(serde_json::from_value(payload.clone()))
}
_ => None,
}
}
}
pub trait Clock: WasmCompatSend + WasmCompatSync {
fn elapsed(&self) -> Duration;
}
pub trait Witness: WasmCompatSend + WasmCompatSync + 'static {
fn observe(&self, observation: Observation);
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct ObservationTrace {
#[serde(default, skip_serializing_if = "Option::is_none")]
pub session: Option<String>,
pub observations: Vec<Observation>,
#[serde(default)]
pub dropped: u64,
#[serde(default)]
pub finalized: bool,
}
impl ObservationTrace {
pub fn is_complete(&self) -> bool {
self.dropped == 0
}
}
pub struct ObservationLog {
inner: Mutex<LogState>,
capacity: usize,
ring: bool,
clock: Option<Arc<dyn Clock + Send + Sync>>,
}
struct LogState {
session: Option<String>,
observations: std::collections::VecDeque<Observation>,
next: u64,
dropped: u64,
finalized: bool,
}
impl std::fmt::Debug for ObservationLog {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
let state = self.lock();
f.debug_struct("ObservationLog")
.field("observations", &state.observations.len())
.field("dropped", &state.dropped)
.field("capacity", &self.capacity)
.finish_non_exhaustive()
}
}
pub const DEFAULT_CAPACITY: usize = 65_536;
impl Default for ObservationLog {
fn default() -> Self {
Self::with_capacity(DEFAULT_CAPACITY)
}
}
impl ObservationLog {
pub fn with_capacity(capacity: usize) -> Self {
Self {
inner: Mutex::new(LogState {
session: None,
observations: std::collections::VecDeque::new(),
next: 0,
dropped: 0,
finalized: false,
}),
capacity,
ring: false,
clock: None,
}
}
pub fn ring(capacity: usize) -> Self {
Self {
ring: true,
..Self::with_capacity(capacity.max(1))
}
}
pub fn drain(&self) -> ObservationTrace {
let mut state = self.lock();
let trace = ObservationTrace {
session: state.session.clone(),
observations: state.observations.drain(..).collect(),
dropped: state.dropped,
finalized: state.finalized,
};
state.dropped = 0;
trace
}
pub fn with_clock(mut self, clock: Arc<dyn Clock + Send + Sync>) -> Self {
self.clock = Some(clock);
self
}
pub fn with_session(self, session: impl Into<String>) -> Self {
self.lock().session = Some(session.into());
self
}
fn lock(&self) -> std::sync::MutexGuard<'_, LogState> {
self.inner.lock().unwrap_or_else(PoisonError::into_inner)
}
pub fn finalize(&self) {
self.lock().finalized = true;
}
pub fn trace(&self) -> ObservationTrace {
let state = self.lock();
ObservationTrace {
session: state.session.clone(),
observations: state.observations.iter().cloned().collect(),
dropped: state.dropped,
finalized: state.finalized,
}
}
pub fn len(&self) -> usize {
self.lock().observations.len()
}
pub fn is_empty(&self) -> bool {
self.lock().observations.is_empty()
}
}
impl Witness for ObservationLog {
fn observe(&self, mut observation: Observation) {
let at = self.clock.as_ref().map(|clock| clock.elapsed());
let mut state = self.lock();
state.finalized = false;
observation.seq = state.next;
state.next += 1;
if state.observations.len() >= self.capacity {
state.dropped += 1;
if !self.ring {
return;
}
state.observations.pop_front();
}
observation.at = at;
state.observations.push_back(observation);
}
}
impl<W: Witness + ?Sized> Witness for Arc<W> {
fn observe(&self, observation: Observation) {
(**self).observe(observation);
}
}
const _: fn() = || {
fn assert_wire<T: Clone + Send + Sync + 'static + Serialize + serde::de::DeserializeOwned>() {}
assert_wire::<Observation>();
assert_wire::<ObservationTrace>();
};