use super::*;
#[cfg(feature = "handlers")]
use boatramp_core::project::ProjectRef;
#[cfg(feature = "handlers")]
impl HandlerRuntimeInner {
fn stream_connections_for_site(&self, site: &str) -> usize {
let counts = self.stream_ip_counts.lock().unwrap();
let sub_prefix = format!("{site}/");
counts
.iter()
.filter(|((scope, _), _)| scope == site || scope.starts_with(&sub_prefix))
.map(|(_, n)| *n as usize)
.sum()
}
}
#[cfg(feature = "handlers")]
#[derive(Serialize)]
struct ConsumerStat {
scope: String,
topic: String,
backlog: usize,
dead_letters: usize,
}
#[cfg(feature = "handlers")]
#[derive(Serialize, Default)]
struct OperatorStats {
handlers: Vec<metrics::HandlerStat>,
consumers: Vec<ConsumerStat>,
stream_connections: usize,
}
#[cfg(feature = "handlers")]
pub(super) async fn operator_handler_stats(
State(deploy): State<DeployStore>,
Extension(handlers): Extension<Arc<HandlerRuntime>>,
Path(site): Path<String>,
) -> Response {
let Some(inner) = handlers.inner.as_ref() else {
return Json(OperatorStats::default()).into_response();
};
let handler_stats = inner.metrics.snapshot_site(&site);
let mut consumers = Vec::new();
if let Some(messaging) = &inner.messaging {
match collect_consumer_stats(&deploy, messaging.as_ref(), &site).await {
Ok(stats) => consumers = stats,
Err(err) => return deploy_error_response(err),
}
}
Json(OperatorStats {
handlers: handler_stats,
consumers,
stream_connections: inner.stream_connections_for_site(&site),
})
.into_response()
}
#[cfg(feature = "handlers")]
#[derive(Deserialize)]
#[serde(rename_all = "lowercase")]
enum DlqAction {
Purge,
Redrive,
}
#[cfg(feature = "handlers")]
#[derive(Deserialize)]
pub(super) struct DlqRequest {
topic: String,
#[serde(default)]
alias: Option<String>,
action: DlqAction,
}
#[cfg(feature = "handlers")]
#[derive(Serialize)]
struct DlqResponse {
affected: usize,
}
#[cfg(feature = "handlers")]
pub(super) async fn operator_dlq(
Extension(handlers): Extension<Arc<HandlerRuntime>>,
Path(site): Path<String>,
Json(req): Json<DlqRequest>,
) -> Response {
let Some(inner) = handlers.inner.as_ref() else {
return not_found();
};
let Some(messaging) = inner.messaging.as_ref() else {
return (
StatusCode::SERVICE_UNAVAILABLE,
"messaging backend not configured\n",
)
.into_response();
};
let scope = match &req.alias {
Some(alias) => format!("{site}/{alias}"),
None => site.clone(),
};
let namespaced = format!("{scope}/{}", req.topic);
let result = match req.action {
DlqAction::Purge => messaging.purge_dead_letters(&namespaced).await,
DlqAction::Redrive => messaging.redrive_dead_letters(&namespaced).await,
};
match result {
Ok(affected) => Json(DlqResponse { affected }).into_response(),
Err(err) => (
StatusCode::INTERNAL_SERVER_ERROR,
format!("dead-letter operation failed: {err}\n"),
)
.into_response(),
}
}
#[cfg(feature = "handlers")]
async fn collect_consumer_stats(
deploy: &DeployStore,
messaging: &dyn boatramp_core::messaging::Messaging,
site: &str,
) -> Result<Vec<ConsumerStat>, DeployError> {
let mut out = Vec::new();
let Some(site_config) = deploy.get_site_config(ProjectRef::DEFAULT, site).await? else {
return Ok(out);
};
let Some(site_handlers) = site_config.handlers.as_ref().filter(|h| h.enabled) else {
return Ok(out);
};
let mut active: Vec<(String, String)> = Vec::new();
if let Some(id) = deploy.current_id(ProjectRef::DEFAULT, site).await? {
active.push((id, site.to_string()));
}
for alias in &site_handlers.background_aliases {
if let Some(id) = deploy.get_alias(ProjectRef::DEFAULT, site, alias).await? {
active.push((id, format!("{site}/{alias}")));
}
}
for (id, scope) in active {
let Some(manifest) = deploy.get_manifest(&id).await? else {
continue;
};
for consumer in &manifest.config.consumers {
let namespaced = format!("{scope}/{}", consumer.topic);
out.push(ConsumerStat {
scope: scope.clone(),
topic: consumer.topic.clone(),
backlog: messaging.backlog(&namespaced).await.unwrap_or(0),
dead_letters: messaging.dead_letter_count(&namespaced).await.unwrap_or(0),
});
}
}
Ok(out)
}
#[cfg(feature = "handlers")]
#[derive(Deserialize)]
pub(super) struct LogsQuery {
limit: Option<usize>,
after: Option<u64>,
stream: Option<String>,
}
#[cfg(feature = "handlers")]
use boatramp_core::logs::LogsResponse;
#[cfg(feature = "handlers")]
pub(super) async fn operator_logs(
Extension(handlers): Extension<Arc<HandlerRuntime>>,
Path(site): Path<String>,
Query(query): Query<LogsQuery>,
) -> Response {
let Some(inner) = handlers.inner.as_ref() else {
return Json(LogsResponse {
entries: Vec::new(),
dropped: 0,
})
.into_response();
};
let stream = match query.stream.as_deref() {
Some("stdout") => Some(boatramp_handlers::LogStream::Stdout),
Some("stderr") => Some(boatramp_handlers::LogStream::Stderr),
_ => None,
};
let limit = query.limit.unwrap_or(200).min(1000);
let (entries, dropped) = inner
.logs
.tail(&site, limit, query.after.unwrap_or(0), stream);
Json(LogsResponse { entries, dropped }).into_response()
}
#[cfg(feature = "handlers")]
pub(super) async fn operator_logs_stream(
Extension(handlers): Extension<Arc<HandlerRuntime>>,
Path(site): Path<String>,
) -> Response {
use axum::response::sse::{Event, KeepAlive, Sse};
let Some(inner) = handlers.inner.as_ref() else {
return (StatusCode::NOT_FOUND, "handlers disabled\n").into_response();
};
let rx = inner.logs.subscribe();
let stream = futures::stream::unfold(rx, move |mut rx| {
let site = site.clone();
async move {
loop {
match rx.recv().await {
Ok((scope, entry)) if scope == site => {
let data = serde_json::to_string(&entry).unwrap_or_default();
let event = Event::default()
.id(entry.seq.to_string())
.event("log")
.data(data);
return Some((Ok::<_, std::convert::Infallible>(event), rx));
}
Ok(_) => continue,
Err(tokio::sync::broadcast::error::RecvError::Lagged(_)) => continue,
Err(tokio::sync::broadcast::error::RecvError::Closed) => return None,
}
}
}
});
Sse::new(stream)
.keep_alive(KeepAlive::default())
.into_response()
}
#[cfg_attr(not(feature = "handlers"), allow(unused_variables))]
pub(super) async fn prometheus_metrics(
State(deploy): State<DeployStore>,
Extension(handlers): Extension<Arc<HandlerRuntime>>,
Extension(daemon): Extension<Arc<DaemonRuntime>>,
) -> Response {
let mut body = srvmetrics::server_metrics().render_prometheus();
let generation = daemon.generation().unwrap_or_else(|| "none".to_string());
body.push_str(
"# HELP boatramp_daemon_config_info Active dynamic daemon-config generation.\n\
# TYPE boatramp_daemon_config_info gauge\n",
);
body.push_str(&format!(
"boatramp_daemon_config_info{{generation=\"{generation}\"}} 1\n"
));
#[cfg(feature = "handlers")]
if let Some(inner) = handlers.inner.as_ref() {
body.push_str(&inner.metrics.render_prometheus());
if let Some(messaging) = &inner.messaging {
let mut rows = Vec::new();
if let Ok(sites) = deploy.list_sites(ProjectRef::DEFAULT).await {
for site in sites {
if let Ok(stats) =
collect_consumer_stats(&deploy, messaging.as_ref(), &site).await
{
for s in stats {
rows.push(metrics::ConsumerGauge {
site: site.clone(),
scope: s.scope,
topic: s.topic,
backlog: s.backlog,
dead_letters: s.dead_letters,
});
}
}
}
}
body.push_str(&metrics::render_consumer_gauges(&rows));
}
if let Ok(mut usage) = deploy.list_metering(ProjectRef::DEFAULT).await {
if !usage.is_empty() {
usage.sort_by(|a, b| a.function.cmp(&b.function));
body.push_str(
"# HELP boatramp_function_invocations_total Function invocations metered.\n\
# TYPE boatramp_function_invocations_total counter\n",
);
for m in &usage {
let f = metrics::escape_label(&m.function);
body.push_str(&format!(
"boatramp_function_invocations_total{{function=\"{f}\"}} {}\n",
m.invocations
));
}
body.push_str(
"# HELP boatramp_function_failures_total Function invocations that failed to deliver.\n\
# TYPE boatramp_function_failures_total counter\n",
);
for m in &usage {
let f = metrics::escape_label(&m.function);
body.push_str(&format!(
"boatramp_function_failures_total{{function=\"{f}\"}} {}\n",
m.failures
));
}
body.push_str(
"# HELP boatramp_function_duration_ms_total Summed function wall-clock duration, ms.\n\
# TYPE boatramp_function_duration_ms_total counter\n",
);
for m in &usage {
let f = metrics::escape_label(&m.function);
body.push_str(&format!(
"boatramp_function_duration_ms_total{{function=\"{f}\"}} {}\n",
m.duration_ms_total
));
}
}
}
}
(
[(
header::CONTENT_TYPE,
"text/plain; version=0.0.4; charset=utf-8",
)],
body,
)
.into_response()
}