use super::*;
use crate::bots::TelemetryEvent;
use crate::telemetry::{TelemetrySection, TelemetrySink, TelemetrySinkReport};
use crate::wire::{HookData, HookEvent, HookSource};
use mobius::agent::{FRONTEND_DISCONNECTED_REASON, TURN_INTERRUPTED_REASON, TURN_RESTARTED_REASON};
use serde_json::{Value, json};
#[derive(Default)]
pub(super) struct StorageMeasurements {
completed: Mutex<Option<CompletedStorageMeasurement>>,
revision: AtomicU64,
}
struct CompletedStorageMeasurement {
completed_at: std::time::Instant,
revision: u64,
usage: crate::storage_usage::StorageUsage,
}
pub(super) fn cancellation_reason(reason: &str, failed: bool) -> &str {
match reason {
TURN_INTERRUPTED_REASON | TURN_RESTARTED_REASON | FRONTEND_DISCONNECTED_REASON => reason,
_ if failed => "failed",
_ => "aborted",
}
}
async fn session_cancellation(
checkpoints: &dyn CheckpointStore,
event: &HookEvent,
sequence: Option<u64>,
) -> Result<Option<Value>> {
let HookData::SessionTurnFinished {
session_id,
turn_id,
outcome: outcome @ (ExecutionOutcome::Aborted | ExecutionOutcome::Failed),
} = &event.data
else {
return Ok(None);
};
if !matches!(&event.source, HookSource::Session { session_id: source_id } if source_id == session_id)
{
return Ok(None);
}
let Some((sequence, before_sequence)) =
sequence.and_then(|sequence| sequence.checked_add(1).map(|before| (sequence, before)))
else {
return Ok(None);
};
let page = checkpoints
.event_page(
session_id,
EventPageRequest {
before_sequence: Some(before_sequence),
limit: 1,
},
)
.await?;
Ok(page.events.into_iter().next().and_then(|record| {
let EventMsg::TurnAborted(turn) = record.event.msg else {
return None;
};
(record.sequence == sequence && &turn.turn_id == turn_id).then(|| {
json!({
"reason": cancellation_reason(&turn.reason, *outcome == ExecutionOutcome::Failed),
"timestamp_ms": record.recorded_at_ms,
"submission_id": record.event.submission_id,
})
})
}))
}
impl StorageMeasurements {
pub(super) fn invalidate(&self) {
self.revision.fetch_add(1, Ordering::AcqRel);
}
}
impl GatewayHost {
pub(crate) fn invalidate_storage_usage(&self) {
self.storage_reads.invalidate();
}
pub(crate) async fn configure_telemetry(
&self,
expected: u64,
mut sinks: Vec<TelemetrySink>,
preserve_auth: &[String],
) -> std::result::Result<(), Rejection> {
let _mutation = self.begin_mutation().await?;
let state = self.state.lock().await;
let mut live = state
.config
.lock()
.map_err(|_| internal("configuration lock poisoned"))?;
if live.telemetry.revision != expected {
return Err(invalid_config(
"telemetry revision changed; refresh and retry",
));
}
if preserve_auth.len() > 16 {
return Err(invalid_config("too many preserved credentials"));
}
for id in preserve_auth {
let previous = live
.telemetry
.sinks
.iter()
.find(|sink| &sink.id == id)
.ok_or_else(|| invalid_config("cannot preserve credentials for an unknown sink"))?;
let replacement = sinks
.iter_mut()
.find(|sink| &sink.id == id)
.ok_or_else(|| invalid_config("preserved sink is missing"))?;
if replacement.url != previous.url {
return Err(invalid_config(
"cannot preserve credentials for a changed destination",
));
}
if replacement.bearer_env.is_some() || replacement.bearer_file.is_some() {
return Err(invalid_config(
"cannot both preserve and replace credentials",
));
}
replacement.bearer_env.clone_from(&previous.bearer_env);
replacement.bearer_file.clone_from(&previous.bearer_file);
}
let revision = expected
.checked_add(1)
.ok_or_else(|| invalid_config("telemetry revision overflow"))?;
self.telemetry.configure_after(|| {
let mut previous = std::mem::replace(&mut live.telemetry.sinks, sinks);
live.telemetry.revision = revision;
let publication = match crate::publication::Outcome::applied(
live.validate().and_then(|()| state.store.save(&live))
) {
Ok(publication) => publication,
Err(error) => {
live.telemetry.sinks = previous;
live.telemetry.revision = expected;
return Err(error);
}
};
if let Err(error) = state.bots.sync_telemetry_cursors(&live.telemetry.sinks) {
std::mem::swap(&mut live.telemetry.sinks, &mut previous);
live.telemetry.revision = expected;
match state.store.save(&live) {
Ok(()) => return Err(error),
Err(rollback @ Error::PublicationApplied { .. }) => {
return Err(crate::publication::applied_error(Error::Config(format!(
"{error}; telemetry rollback applied but completion failed: {rollback}"
))));
}
Err(rollback) => {
live.telemetry.sinks = previous;
live.telemetry.revision = revision;
let failure = crate::publication::applied_error(Error::Config(format!(
"telemetry configuration remains changed; cursor synchronization failed: {error}; rollback failed: {rollback}"
)));
return Ok((live.telemetry.clone(), crate::publication::Outcome::applied(Err(failure))?));
}
}
}
Ok((live.telemetry.clone(), publication))
}).map_err(internal)
}
pub(crate) async fn schedule_telemetry(&self, id: String) -> Result<()> {
self.telemetry.request_manual(id)?;
self.telemetry.notify.notify_one();
Ok(())
}
pub(crate) async fn telemetry_pending(&self) -> Result<bool> {
let config = self.telemetry.config()?;
let state = self.state.lock().await;
for sink in config
.sinks
.iter()
.filter(|sink| sink.enabled && !sink.events.is_empty())
{
if self.telemetry.can_drain(&sink.id)? && state.bots.telemetry_count(sink)? > 0 {
return Ok(true);
}
}
Ok(false)
}
pub(crate) async fn storage_usage_request(
&self,
) -> std::result::Result<crate::storage_usage::StorageUsage, Rejection> {
let mut completed = tokio::time::timeout(
std::time::Duration::from_secs(2),
self.storage_reads.completed.lock(),
)
.await
.map_err(|_| internal("storage measurement timed out waiting for an active measurement"))?;
let revision = self.storage_reads.revision.load(Ordering::Acquire);
if let Some(previous) = completed.as_ref()
&& previous.revision == revision
&& previous.completed_at.elapsed() < std::time::Duration::from_secs(5)
{
return Ok(previous.usage.clone());
}
let usage = self.storage_usage().await.map_err(internal)?;
*completed = Some(CompletedStorageMeasurement {
completed_at: std::time::Instant::now(),
revision,
usage: usage.clone(),
});
Ok(usage)
}
pub(crate) async fn telemetry_report(&self) -> Result<(u64, Vec<TelemetrySinkReport>)> {
let config = self.telemetry.config()?;
let state = self.state.lock().await;
let mut reports = Vec::new();
for configured in &config.sinks {
let mut sink = configured.clone();
let mut status = self.telemetry.status(&sink.id)?;
if !sink.events.is_empty() {
status.events_pending = state.bots.telemetry_count(&sink)?;
}
let auth = if sink.bearer_env.is_some() {
crate::telemetry::SinkAuth::BearerEnv
} else if sink.bearer_file.is_some() {
crate::telemetry::SinkAuth::BearerFile
} else {
crate::telemetry::SinkAuth::None
};
sink.redact_report()?;
reports.push(TelemetrySinkReport { sink, auth, status });
}
Ok((config.revision, reports))
}
pub(crate) async fn telemetry_snapshot(
&self,
sections: &[TelemetrySection],
connected_clients: usize,
) -> Result<Value> {
let mut snapshot = json!({"machine_name": local_machine_name()?});
if sections.is_empty() {
return Ok(snapshot);
}
if sections.contains(&TelemetrySection::Activity) {
let activity = self
.runtime_activity()
.await
.map_err(|e| Error::Config(e.message))?;
snapshot["activity"] = json!({"idle": activity.idle && connected_clients == 0, "connected_clients": connected_clients, "activity_revision": activity.activity_revision, "next_routine_at": activity.next_routine_at, "active_sessions": activity.active_sessions, "running_routines": activity.running_routines});
}
if sections.contains(&TelemetrySection::Resources) {
snapshot["resources"] = self.telemetry.resources().await;
}
if sections.contains(&TelemetrySection::Storage) {
let mut usage = self.storage_usage().await?;
usage.bound_telemetry_details();
snapshot["storage"] = serde_json::to_value(usage)?;
}
if sections.contains(&TelemetrySection::Usage) {
let state = self.state.lock().await;
let config = state
.config
.lock()
.map_err(|_| Error::Config("configuration lock poisoned".into()))?;
let cutoff = u64::try_from(Utc::now().timestamp().max(0)).unwrap_or_default() / 86_400;
snapshot["usage"] = serde_json::to_value(
config
.profile()
.daily_usage
.into_iter()
.filter(|day| day.unix_day >= cutoff.saturating_sub(30))
.collect::<Vec<_>>(),
)?;
}
if sections.contains(&TelemetrySection::Runs) {
let checkpoints = Arc::clone(&self.state.lock().await.checkpoints);
snapshot["runs"] = serde_json::to_value(
gateway_run_stats(&gateway_session_summaries(&checkpoints).await?)?.completed,
)?;
}
Ok(snapshot)
}
pub(crate) async fn telemetry_events(
&self,
sink: &TelemetrySink,
) -> Result<(Vec<Value>, Option<i64>, u64)> {
let (batch, pending, checkpoints, activities, bots) = {
let state = self.state.lock().await;
let (batch, pending) = state.bots.telemetry_batch(sink)?;
(
batch,
pending,
Arc::clone(&state.checkpoints),
Arc::clone(&state.activities),
Arc::clone(&state.bots),
)
};
let mut cursor = None;
let mut bytes = 0_usize;
let sessions = if batch.is_empty() {
Vec::new()
} else {
session_catalog(&checkpoints, &activities).await?
};
let mut events = Vec::with_capacity(batch.len());
for TelemetryEvent {
rowid,
event,
journal_sequence,
} in batch
{
let mut value = serde_json::to_value(&event)?;
if let crate::wire::HookSource::Session { session_id } = &event.source
&& let Some(session) = sessions.iter().find(|s| &s.session_id == session_id)
{
let bot = bots.find_bot(&session.session_context.owner_id)?;
let mut enrichment =
json!({"title": session.title, "run_count": session.execution_stats.run_count});
if let Some(bot) = bot {
enrichment["bot_name"] = json!(bot.name);
}
if let Some(cancellation) =
session_cancellation(checkpoints.as_ref(), &event, journal_sequence).await?
{
enrichment["cancellation"] = cancellation;
}
if matches!(
event.data,
crate::wire::HookData::SessionTurnFinished {
outcome: ExecutionOutcome::Completed,
..
}
) {
enrichment["final_message"] = json!(
session
.activity
.message
.as_deref()
.unwrap_or_default()
.chars()
.take(1000)
.collect::<String>()
);
}
value["session"] = enrichment;
}
let size = serde_json::to_vec(&value)?.len();
if !events.is_empty() && bytes.saturating_add(size) > 48 * 1024 {
break;
}
bytes = bytes.saturating_add(size);
cursor = Some(rowid);
events.push(value);
}
Ok((events, cursor, pending))
}
pub(crate) async fn advance_telemetry(
&self,
id: &str,
cursor: i64,
revision: u64,
) -> Result<()> {
let state = self.state.lock().await;
if state
.config
.lock()
.map_err(|_| Error::Config("configuration lock poisoned".into()))?
.telemetry
.revision
!= revision
{
return Ok(());
}
state.bots.advance_telemetry(id, cursor)
}
}
impl GatewayHost {
pub(crate) async fn storage_usage(&self) -> Result<crate::storage_usage::StorageUsage> {
let budget = crate::storage_usage::MeasurementBudget::new();
let (root, checkpoints, files, bots, limit_bytes, tls) = {
let state = tokio::time::timeout_at(budget.deadline(), self.state.lock())
.await
.map_err(|_| {
Error::Config("storage measurement timed out waiting for gateway state".into())
})?;
let config = state
.config
.lock()
.map_err(|_| Error::Config("configuration lock poisoned".into()))?;
(
state.store.state_dir().to_path_buf(),
Arc::clone(&state.checkpoints),
state.session_files.clone(),
Arc::clone(&state.bots),
config.runtime.storage_limit_bytes,
config.tls.clone(),
)
};
crate::storage_usage::measure_usage(
root,
checkpoints,
files,
bots,
limit_bytes,
tls,
budget,
)
.await
}
}