use serde::{Deserialize, Serialize};
use serde_json::Value;
use tokio::sync::broadcast::error::RecvError;
use tokio::sync::mpsc;
use trusty_common::uds::server::{RpcError, RpcStreamItems};
use crate::transport::api_error::ApiError;
use crate::{ActivityFilter, ActivitySource, AppState};
use super::parse_iso_or_bad_request;
const ACTIVITY_DEFAULT_LIMIT: usize = 50;
const ACTIVITY_MAX_LIMIT: usize = 500;
#[derive(Debug, Default, Deserialize)]
pub struct ActivityParams {
#[serde(default)]
pub limit: Option<usize>,
#[serde(default)]
pub offset: Option<usize>,
#[serde(default)]
pub palace: Option<String>,
#[serde(default)]
pub source: Option<String>,
#[serde(default)]
pub since: Option<String>,
#[serde(default)]
pub until: Option<String>,
}
#[derive(Debug, Serialize)]
pub struct ActivityRow {
pub id: u64,
pub timestamp: chrono::DateTime<chrono::Utc>,
pub source: &'static str,
#[serde(skip_serializing_if = "Option::is_none")]
pub palace_id: Option<String>,
pub event_type: String,
pub payload: Value,
}
pub async fn activity(state: &AppState, params: ActivityParams) -> Result<Value, ApiError> {
let limit = params
.limit
.unwrap_or(ACTIVITY_DEFAULT_LIMIT)
.clamp(1, ACTIVITY_MAX_LIMIT);
let offset = params.offset.unwrap_or(0);
let source = match params.source.as_deref() {
Some(s) => match ActivitySource::parse(s) {
Some(parsed) => Some(parsed),
None => {
return Err(ApiError::bad_request(format!(
"unknown source '{s}'; expected one of http, mcp, hook"
)));
}
},
None => None,
};
let filter = ActivityFilter {
palace_id: params.palace.filter(|s| !s.is_empty()),
source,
since: parse_iso_or_bad_request(params.since.as_deref(), "since")?,
until: parse_iso_or_bad_request(params.until.as_deref(), "until")?,
};
let entries = state
.activity_log
.list(&filter, limit, offset)
.map_err(|e| ApiError::internal(format!("activity list: {e:#}")))?;
let total = state
.activity_log
.count()
.map_err(|e| ApiError::internal(format!("activity count: {e:#}")))?;
let rows: Vec<ActivityRow> = entries
.into_iter()
.map(|e| {
let payload = serde_json::from_str::<Value>(&e.payload)
.unwrap_or_else(|_| Value::String(e.payload.clone()));
ActivityRow {
id: e.id,
timestamp: e.timestamp,
source: e.source.as_str(),
palace_id: e.palace_id,
event_type: e.event_type,
payload,
}
})
.collect();
super::to_value(serde_json::json!({
"entries": rows,
"total": total,
"limit": limit,
"offset": offset,
}))
}
const STREAM_BUFFER: usize = 256;
pub async fn activity_stream(state: &AppState) -> Result<RpcStreamItems, RpcError> {
let mut events = state.events.subscribe();
let (tx, items) = mpsc::channel(STREAM_BUFFER);
tokio::spawn(async move {
loop {
let frame = match events.recv().await {
Ok(event) => match serde_json::to_value(&event) {
Ok(value) => Ok(value),
Err(e) => Err(RpcError::internal(format!("serialize activity event: {e}"))),
},
Err(RecvError::Lagged(n)) => {
tracing::warn!(
target: "trusty_memory::activity_stream",
"an activity_stream reader lagged {n} events"
);
Ok(serde_json::json!({ "type": "lagged", "lagged": n }))
}
Err(RecvError::Closed) => break,
};
let terminal = frame.is_err();
if tx.send(frame).await.is_err() {
break;
}
if terminal {
break;
}
}
});
Ok(items)
}