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,
in_flight: usize,
#[serde(skip_serializing_if = "Option::is_none")]
oldest_pending_ms: Option<u64>,
lag: 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")]
const DLQ_VIEW_VERSION: u32 = 1;
#[cfg(feature = "handlers")]
#[derive(Deserialize)]
#[serde(rename_all = "lowercase")]
enum DlqAction {
Purge,
Redrive,
Discard,
}
#[cfg(feature = "handlers")]
#[derive(Deserialize, Default)]
pub(super) struct DlqFilterWire {
#[serde(default)]
id: Option<String>,
#[serde(default)]
group: Option<String>,
#[serde(default)]
older_than_ms: Option<u64>,
#[serde(default, rename = "match")]
match_last_error: Option<String>,
#[serde(default)]
limit: Option<usize>,
}
#[cfg(feature = "handlers")]
impl DlqFilterWire {
fn into_core(self) -> boatramp_core::messaging::DeadLetterFilter {
boatramp_core::messaging::DeadLetterFilter {
id: self.id,
group: self.group,
older_than_ms: self.older_than_ms,
match_last_error: self.match_last_error,
limit: self.limit,
}
}
fn is_empty(&self) -> bool {
self.id.is_none()
&& self.group.is_none()
&& self.older_than_ms.is_none()
&& self.match_last_error.is_none()
&& self.limit.is_none()
}
}
#[cfg(feature = "handlers")]
#[derive(Deserialize)]
pub(super) struct DlqRequest {
topic: String,
#[serde(default)]
alias: Option<String>,
action: DlqAction,
#[serde(default)]
filter: DlqFilterWire,
#[serde(default)]
dry_run: bool,
}
#[cfg(feature = "handlers")]
#[derive(Serialize)]
struct DlqEntry {
id: String,
group: String,
attempts: u32,
last_error: Option<String>,
signed_context_present: bool,
#[serde(skip_serializing_if = "Option::is_none")]
payload_b64: Option<String>,
}
#[cfg(feature = "handlers")]
impl DlqEntry {
fn from_dead(dl: boatramp_core::messaging::DeadLetter) -> Self {
use base64::Engine as _;
Self {
id: dl.id,
group: dl.group,
attempts: dl.attempts,
last_error: dl.last_error,
signed_context_present: dl.signed_context.is_some(),
payload_b64: dl
.payload
.map(|p| base64::engine::general_purpose::STANDARD.encode(p)),
}
}
}
#[cfg(feature = "handlers")]
#[derive(Serialize)]
struct DlqResponse {
affected: usize,
#[serde(skip_serializing_if = "Vec::is_empty")]
matched: Vec<DlqEntry>,
}
#[cfg(feature = "handlers")]
#[derive(Serialize)]
struct DlqListResponse {
version: u32,
dead_letters: Vec<DlqEntry>,
}
#[cfg(feature = "handlers")]
#[derive(Deserialize)]
pub(super) struct DlqListQuery {
topic: String,
#[serde(default)]
alias: Option<String>,
#[serde(default)]
show: bool,
#[serde(flatten)]
filter: DlqFilterWire,
}
#[cfg(feature = "handlers")]
fn dlq_namespace(site: &str, alias: &Option<String>, topic: &str) -> String {
match alias {
Some(alias) => format!("{site}/{alias}/{topic}"),
None => format!("{site}/{topic}"),
}
}
#[cfg(feature = "handlers")]
fn bus_namespace(project: &ProjectRef<'_>, topic: &str) -> String {
project.qualified(&format!("bus/{topic}"))
}
#[cfg(feature = "handlers")]
fn messaging_or_unavailable(
handlers: &HandlerRuntime,
) -> Result<&std::sync::Arc<dyn boatramp_core::messaging::Messaging>, Response> {
let Some(inner) = handlers.inner.as_ref() else {
return Err(not_found());
};
inner.messaging.as_ref().ok_or_else(|| {
(
StatusCode::SERVICE_UNAVAILABLE,
"messaging backend not configured\n",
)
.into_response()
})
}
#[cfg(feature = "handlers")]
pub(super) async fn operator_dlq_list(
Extension(handlers): Extension<Arc<HandlerRuntime>>,
Path(site): Path<String>,
axum::extract::Query(q): axum::extract::Query<DlqListQuery>,
) -> Response {
let messaging = match messaging_or_unavailable(&handlers) {
Ok(m) => m,
Err(resp) => return resp,
};
let namespaced = dlq_namespace(&site, &q.alias, &q.topic);
dlq_list_core(messaging.as_ref(), &namespaced, q.show, q.filter).await
}
#[cfg(feature = "handlers")]
async fn dlq_list_core(
messaging: &dyn boatramp_core::messaging::Messaging,
namespaced: &str,
show: bool,
filter: DlqFilterWire,
) -> Response {
if show {
let Some(id) = filter.id.clone() else {
return (StatusCode::BAD_REQUEST, "show requires ?id=<id>\n").into_response();
};
let group = filter.group.clone().unwrap_or_default();
return match messaging.show_dead_letter(namespaced, &group, &id).await {
Ok(Some(dl)) => Json(DlqListResponse {
version: DLQ_VIEW_VERSION,
dead_letters: vec![DlqEntry::from_dead(dl)],
})
.into_response(),
Ok(None) => not_found(),
Err(err) => (
StatusCode::INTERNAL_SERVER_ERROR,
format!("dead-letter show failed: {err}\n"),
)
.into_response(),
};
}
match messaging
.list_dead_letters(namespaced, &filter.into_core())
.await
{
Ok(list) => Json(DlqListResponse {
version: DLQ_VIEW_VERSION,
dead_letters: list.into_iter().map(DlqEntry::from_dead).collect(),
})
.into_response(),
Err(err) => (
StatusCode::INTERNAL_SERVER_ERROR,
format!("dead-letter list failed: {err}\n"),
)
.into_response(),
}
}
#[cfg(feature = "handlers")]
#[derive(Deserialize)]
pub(super) struct QueuePeekQuery {
topic: String,
#[serde(default)]
alias: Option<String>,
#[serde(default)]
limit: Option<usize>,
}
#[cfg(feature = "handlers")]
#[derive(Serialize)]
struct QueuePeekEntry {
id: String,
attempts: u32,
leased: bool,
signed_context_present: bool,
payload_b64: String,
}
#[cfg(feature = "handlers")]
#[derive(Serialize)]
struct QueuePeekResponse {
version: u32,
messages: Vec<QueuePeekEntry>,
}
#[cfg(feature = "handlers")]
const QUEUE_PEEK_MAX: usize = 100;
#[cfg(feature = "handlers")]
pub(super) async fn operator_queue_peek(
Extension(handlers): Extension<Arc<HandlerRuntime>>,
Path(site): Path<String>,
axum::extract::Query(q): axum::extract::Query<QueuePeekQuery>,
) -> Response {
let messaging = match messaging_or_unavailable(&handlers) {
Ok(m) => m,
Err(resp) => return resp,
};
let namespaced = dlq_namespace(&site, &q.alias, &q.topic);
queue_peek_core(messaging.as_ref(), &namespaced, q.limit).await
}
#[cfg(feature = "handlers")]
async fn queue_peek_core(
messaging: &dyn boatramp_core::messaging::Messaging,
namespaced: &str,
limit: Option<usize>,
) -> Response {
let limit = limit.unwrap_or(10).min(QUEUE_PEEK_MAX);
match messaging.peek(namespaced, limit).await {
Ok(msgs) => {
use base64::Engine as _;
let messages = msgs
.into_iter()
.map(|m| QueuePeekEntry {
id: m.id,
attempts: m.attempts,
leased: m.leased,
signed_context_present: m.signed_context.is_some(),
payload_b64: base64::engine::general_purpose::STANDARD.encode(m.payload),
})
.collect();
Json(QueuePeekResponse {
version: DLQ_VIEW_VERSION,
messages,
})
.into_response()
}
Err(err) => (
StatusCode::INTERNAL_SERVER_ERROR,
format!("queue peek failed: {err}\n"),
)
.into_response(),
}
}
#[cfg(feature = "handlers")]
#[derive(Deserialize)]
pub(super) struct QueueReplayQuery {
topic: String,
#[serde(default)]
alias: Option<String>,
#[serde(default)]
after: Option<String>,
#[serde(default)]
limit: Option<usize>,
}
#[cfg(feature = "handlers")]
#[derive(Serialize)]
struct QueueReplayResponse {
version: u32,
messages: Vec<QueuePeekEntry>,
#[serde(skip_serializing_if = "Option::is_none")]
next_after: Option<String>,
}
#[cfg(feature = "handlers")]
pub(super) async fn operator_queue_replay(
Extension(handlers): Extension<Arc<HandlerRuntime>>,
Path(site): Path<String>,
axum::extract::Query(q): axum::extract::Query<QueueReplayQuery>,
) -> Response {
let messaging = match messaging_or_unavailable(&handlers) {
Ok(m) => m,
Err(resp) => return resp,
};
let namespaced = dlq_namespace(&site, &q.alias, &q.topic);
queue_replay_core(messaging.as_ref(), &namespaced, q.after.as_deref(), q.limit).await
}
#[cfg(feature = "handlers")]
async fn queue_replay_core(
messaging: &dyn boatramp_core::messaging::Messaging,
namespaced: &str,
after: Option<&str>,
limit: Option<usize>,
) -> Response {
let limit = limit.unwrap_or(10).min(QUEUE_PEEK_MAX);
match messaging.replay(namespaced, after, limit).await {
Ok(msgs) => {
use base64::Engine as _;
let next_after = msgs.last().map(|m| m.id.clone());
let messages = msgs
.into_iter()
.map(|m| QueuePeekEntry {
id: m.id,
attempts: m.attempts,
leased: m.leased,
signed_context_present: m.signed_context.is_some(),
payload_b64: base64::engine::general_purpose::STANDARD.encode(m.payload),
})
.collect();
Json(QueueReplayResponse {
version: DLQ_VIEW_VERSION,
messages,
next_after,
})
.into_response()
}
Err(err) => (
StatusCode::INTERNAL_SERVER_ERROR,
format!("queue replay failed: {err}\n"),
)
.into_response(),
}
}
#[cfg(feature = "handlers")]
#[derive(Deserialize)]
pub(super) struct QueueGroupsQuery {
topic: String,
#[serde(default)]
alias: Option<String>,
}
#[cfg(feature = "handlers")]
#[derive(Serialize)]
struct GroupEntry {
group: String,
hwm: String,
in_flight: usize,
lag: usize,
}
#[cfg(feature = "handlers")]
#[derive(Serialize)]
struct QueueGroupsResponse {
version: u32,
groups: Vec<GroupEntry>,
}
#[cfg(feature = "handlers")]
pub(super) async fn operator_queue_groups(
Extension(handlers): Extension<Arc<HandlerRuntime>>,
Path(site): Path<String>,
axum::extract::Query(q): axum::extract::Query<QueueGroupsQuery>,
) -> Response {
let messaging = match messaging_or_unavailable(&handlers) {
Ok(m) => m,
Err(resp) => return resp,
};
let namespaced = dlq_namespace(&site, &q.alias, &q.topic);
queue_groups_core(messaging.as_ref(), &namespaced).await
}
#[cfg(feature = "handlers")]
async fn queue_groups_core(
messaging: &dyn boatramp_core::messaging::Messaging,
namespaced: &str,
) -> Response {
match messaging.list_groups(namespaced).await {
Ok(groups) => Json(QueueGroupsResponse {
version: DLQ_VIEW_VERSION,
groups: groups
.into_iter()
.map(|g| GroupEntry {
group: g.group,
hwm: g.hwm,
in_flight: g.in_flight,
lag: g.lag,
})
.collect(),
})
.into_response(),
Err(err) => (
StatusCode::INTERNAL_SERVER_ERROR,
format!("group list failed: {err}\n"),
)
.into_response(),
}
}
#[cfg(feature = "handlers")]
#[derive(Deserialize)]
pub(super) struct QueuePauseRequest {
topic: String,
#[serde(default)]
alias: Option<String>,
paused: bool,
}
#[cfg(feature = "handlers")]
pub(super) async fn operator_queue_pause(
Extension(handlers): Extension<Arc<HandlerRuntime>>,
Path(site): Path<String>,
Json(req): Json<QueuePauseRequest>,
) -> Response {
let messaging = match messaging_or_unavailable(&handlers) {
Ok(m) => m,
Err(resp) => return resp,
};
let namespaced = dlq_namespace(&site, &req.alias, &req.topic);
queue_pause_core(messaging.as_ref(), &namespaced, req.paused).await
}
#[cfg(feature = "handlers")]
async fn queue_pause_core(
messaging: &dyn boatramp_core::messaging::Messaging,
namespaced: &str,
paused: bool,
) -> Response {
match messaging.set_paused(namespaced, paused).await {
Ok(()) => Json(serde_json::json!({ "ok": true, "paused": paused })).into_response(),
Err(err) => (
StatusCode::INTERNAL_SERVER_ERROR,
format!("pause/resume failed: {err}\n"),
)
.into_response(),
}
}
#[cfg(feature = "handlers")]
#[derive(Deserialize)]
pub(super) struct QueuePolicyRequest {
topic: String,
#[serde(default)]
alias: Option<String>,
#[serde(default)]
max_depth: Option<usize>,
#[serde(default)]
max_rate_per_sec: Option<u32>,
#[serde(default)]
max_unflushed: Option<usize>,
}
#[cfg(feature = "handlers")]
pub(super) async fn operator_queue_policy(
Extension(handlers): Extension<Arc<HandlerRuntime>>,
Path(site): Path<String>,
Json(req): Json<QueuePolicyRequest>,
) -> Response {
let messaging = match messaging_or_unavailable(&handlers) {
Ok(m) => m,
Err(resp) => return resp,
};
let namespaced = dlq_namespace(&site, &req.alias, &req.topic);
queue_policy_core(messaging.as_ref(), &namespaced, &req).await
}
#[cfg(feature = "handlers")]
async fn queue_policy_core(
messaging: &dyn boatramp_core::messaging::Messaging,
namespaced: &str,
req: &QueuePolicyRequest,
) -> Response {
let policy = boatramp_core::messaging::TopicPolicy {
max_depth: req.max_depth,
max_rate_per_sec: req.max_rate_per_sec,
max_unflushed: req.max_unflushed,
};
match messaging.set_topic_policy(namespaced, policy).await {
Ok(()) => Json(serde_json::json!({
"ok": true,
"max_depth": req.max_depth,
"max_rate_per_sec": req.max_rate_per_sec,
"max_unflushed": req.max_unflushed,
}))
.into_response(),
Err(boatramp_core::messaging::MessagingError::Unsupported(msg)) => {
(StatusCode::NOT_IMPLEMENTED, format!("{msg}\n")).into_response()
}
Err(err) => (
StatusCode::INTERNAL_SERVER_ERROR,
format!("set policy failed: {err}\n"),
)
.into_response(),
}
}
#[cfg(feature = "handlers")]
#[derive(Deserialize)]
#[serde(rename_all = "lowercase")]
enum GroupAction {
Reset,
Delete,
}
#[cfg(feature = "handlers")]
#[derive(Deserialize)]
pub(super) struct QueueGroupRequest {
topic: String,
#[serde(default)]
alias: Option<String>,
group: String,
action: GroupAction,
#[serde(default)]
start: Option<boatramp_core::messaging::StartPosition>,
}
#[cfg(feature = "handlers")]
pub(super) async fn operator_queue_group(
Extension(handlers): Extension<Arc<HandlerRuntime>>,
Path(site): Path<String>,
Json(req): Json<QueueGroupRequest>,
) -> Response {
let messaging = match messaging_or_unavailable(&handlers) {
Ok(m) => m,
Err(resp) => return resp,
};
let namespaced = dlq_namespace(&site, &req.alias, &req.topic);
queue_group_core(
messaging.as_ref(),
&namespaced,
req.group,
req.action,
req.start,
)
.await
}
#[cfg(feature = "handlers")]
async fn queue_group_core(
messaging: &dyn boatramp_core::messaging::Messaging,
namespaced: &str,
group: String,
action: GroupAction,
start: Option<boatramp_core::messaging::StartPosition>,
) -> Response {
let result = match action {
GroupAction::Reset => {
let start = start.unwrap_or(boatramp_core::messaging::StartPosition::Earliest);
messaging.reset_group(namespaced, &group, start).await
}
GroupAction::Delete => messaging.delete_group(namespaced, &group).await,
};
match result {
Ok(()) => Json(serde_json::json!({ "ok": true })).into_response(),
Err(err) => (
StatusCode::INTERNAL_SERVER_ERROR,
format!("group operation failed: {err}\n"),
)
.into_response(),
}
}
#[cfg(feature = "handlers")]
pub(super) async fn operator_dlq(
Extension(handlers): Extension<Arc<HandlerRuntime>>,
Path(site): Path<String>,
Json(req): Json<DlqRequest>,
) -> Response {
let messaging = match messaging_or_unavailable(&handlers) {
Ok(m) => m,
Err(resp) => return resp,
};
let namespaced = dlq_namespace(&site, &req.alias, &req.topic);
dlq_mutate_core(
messaging.as_ref(),
&namespaced,
req.action,
req.filter,
req.dry_run,
)
.await
}
#[cfg(feature = "handlers")]
async fn dlq_mutate_core(
messaging: &dyn boatramp_core::messaging::Messaging,
namespaced: &str,
action: DlqAction,
filter: DlqFilterWire,
dry_run: bool,
) -> Response {
let selective = !filter.is_empty();
if dry_run {
let filter = filter.into_core();
return match messaging.list_dead_letters(namespaced, &filter).await {
Ok(list) => {
let matched: Vec<DlqEntry> = list.into_iter().map(DlqEntry::from_dead).collect();
Json(DlqResponse {
affected: matched.len(),
matched,
})
.into_response()
}
Err(err) => (
StatusCode::INTERNAL_SERVER_ERROR,
format!("dead-letter dry-run failed: {err}\n"),
)
.into_response(),
};
}
let result = match action {
DlqAction::Purge => {
if selective {
messaging
.discard_dead_letters(namespaced, &filter.into_core())
.await
} else {
messaging.purge_dead_letters(namespaced).await
}
}
DlqAction::Discard => {
messaging
.discard_dead_letters(namespaced, &filter.into_core())
.await
}
DlqAction::Redrive => {
if selective {
messaging
.redrive_dead_letters_filtered(namespaced, &filter.into_core())
.await
} else {
messaging.redrive_dead_letters(namespaced).await
}
}
};
match result {
Ok(affected) => Json(DlqResponse {
affected,
matched: Vec::new(),
})
.into_response(),
Err(err) => (
StatusCode::INTERNAL_SERVER_ERROR,
format!("dead-letter operation failed: {err}\n"),
)
.into_response(),
}
}
#[cfg(feature = "handlers")]
#[derive(Deserialize)]
pub(super) struct BusDlqListQuery {
topic: String,
#[serde(default)]
show: bool,
#[serde(flatten)]
filter: DlqFilterWire,
}
#[cfg(feature = "handlers")]
pub(super) async fn operator_bus_dlq_list(
Extension(handlers): Extension<Arc<HandlerRuntime>>,
Extension(project): Extension<crate::project_scope::ProjectContext>,
axum::extract::Query(q): axum::extract::Query<BusDlqListQuery>,
) -> Response {
let messaging = match messaging_or_unavailable(&handlers) {
Ok(m) => m,
Err(resp) => return resp,
};
let namespaced = bus_namespace(&project.as_ref(), &q.topic);
dlq_list_core(messaging.as_ref(), &namespaced, q.show, q.filter).await
}
#[cfg(feature = "handlers")]
#[derive(Deserialize)]
pub(super) struct BusDlqRequest {
topic: String,
action: DlqAction,
#[serde(default)]
filter: DlqFilterWire,
#[serde(default)]
dry_run: bool,
}
#[cfg(feature = "handlers")]
pub(super) async fn operator_bus_dlq(
Extension(handlers): Extension<Arc<HandlerRuntime>>,
Extension(project): Extension<crate::project_scope::ProjectContext>,
Json(req): Json<BusDlqRequest>,
) -> Response {
let messaging = match messaging_or_unavailable(&handlers) {
Ok(m) => m,
Err(resp) => return resp,
};
let namespaced = bus_namespace(&project.as_ref(), &req.topic);
dlq_mutate_core(
messaging.as_ref(),
&namespaced,
req.action,
req.filter,
req.dry_run,
)
.await
}
#[cfg(feature = "handlers")]
#[derive(Deserialize)]
pub(super) struct BusQueuePeekQuery {
topic: String,
#[serde(default)]
limit: Option<usize>,
}
#[cfg(feature = "handlers")]
pub(super) async fn operator_bus_queue_peek(
Extension(handlers): Extension<Arc<HandlerRuntime>>,
Extension(project): Extension<crate::project_scope::ProjectContext>,
axum::extract::Query(q): axum::extract::Query<BusQueuePeekQuery>,
) -> Response {
let messaging = match messaging_or_unavailable(&handlers) {
Ok(m) => m,
Err(resp) => return resp,
};
let namespaced = bus_namespace(&project.as_ref(), &q.topic);
queue_peek_core(messaging.as_ref(), &namespaced, q.limit).await
}
#[cfg(feature = "handlers")]
#[derive(Deserialize)]
pub(super) struct BusQueueReplayQuery {
topic: String,
#[serde(default)]
after: Option<String>,
#[serde(default)]
limit: Option<usize>,
}
#[cfg(feature = "handlers")]
pub(super) async fn operator_bus_queue_replay(
Extension(handlers): Extension<Arc<HandlerRuntime>>,
Extension(project): Extension<crate::project_scope::ProjectContext>,
axum::extract::Query(q): axum::extract::Query<BusQueueReplayQuery>,
) -> Response {
let messaging = match messaging_or_unavailable(&handlers) {
Ok(m) => m,
Err(resp) => return resp,
};
let namespaced = bus_namespace(&project.as_ref(), &q.topic);
queue_replay_core(messaging.as_ref(), &namespaced, q.after.as_deref(), q.limit).await
}
#[cfg(feature = "handlers")]
#[derive(Deserialize)]
pub(super) struct BusQueueGroupsQuery {
topic: String,
}
#[cfg(feature = "handlers")]
pub(super) async fn operator_bus_queue_groups(
Extension(handlers): Extension<Arc<HandlerRuntime>>,
Extension(project): Extension<crate::project_scope::ProjectContext>,
axum::extract::Query(q): axum::extract::Query<BusQueueGroupsQuery>,
) -> Response {
let messaging = match messaging_or_unavailable(&handlers) {
Ok(m) => m,
Err(resp) => return resp,
};
let namespaced = bus_namespace(&project.as_ref(), &q.topic);
queue_groups_core(messaging.as_ref(), &namespaced).await
}
#[cfg(feature = "handlers")]
#[derive(Deserialize)]
pub(super) struct BusQueuePauseRequest {
topic: String,
paused: bool,
}
#[cfg(feature = "handlers")]
pub(super) async fn operator_bus_queue_pause(
Extension(handlers): Extension<Arc<HandlerRuntime>>,
Extension(project): Extension<crate::project_scope::ProjectContext>,
Json(req): Json<BusQueuePauseRequest>,
) -> Response {
let messaging = match messaging_or_unavailable(&handlers) {
Ok(m) => m,
Err(resp) => return resp,
};
let namespaced = bus_namespace(&project.as_ref(), &req.topic);
queue_pause_core(messaging.as_ref(), &namespaced, req.paused).await
}
#[cfg(feature = "handlers")]
#[derive(Deserialize)]
pub(super) struct BusQueuePolicyRequest {
topic: String,
#[serde(default)]
max_depth: Option<usize>,
#[serde(default)]
max_rate_per_sec: Option<u32>,
#[serde(default)]
max_unflushed: Option<usize>,
}
#[cfg(feature = "handlers")]
pub(super) async fn operator_bus_queue_policy(
Extension(handlers): Extension<Arc<HandlerRuntime>>,
Extension(project): Extension<crate::project_scope::ProjectContext>,
Json(req): Json<BusQueuePolicyRequest>,
) -> Response {
let messaging = match messaging_or_unavailable(&handlers) {
Ok(m) => m,
Err(resp) => return resp,
};
let namespaced = bus_namespace(&project.as_ref(), &req.topic);
let site_req = QueuePolicyRequest {
topic: req.topic.clone(),
alias: None,
max_depth: req.max_depth,
max_rate_per_sec: req.max_rate_per_sec,
max_unflushed: req.max_unflushed,
};
queue_policy_core(messaging.as_ref(), &namespaced, &site_req).await
}
#[cfg(feature = "handlers")]
#[derive(Deserialize)]
pub(super) struct BusQueueGroupRequest {
topic: String,
group: String,
action: GroupAction,
#[serde(default)]
start: Option<boatramp_core::messaging::StartPosition>,
}
#[cfg(feature = "handlers")]
pub(super) async fn operator_bus_queue_group(
Extension(handlers): Extension<Arc<HandlerRuntime>>,
Extension(project): Extension<crate::project_scope::ProjectContext>,
Json(req): Json<BusQueueGroupRequest>,
) -> Response {
let messaging = match messaging_or_unavailable(&handlers) {
Ok(m) => m,
Err(resp) => return resp,
};
let namespaced = bus_namespace(&project.as_ref(), &req.topic);
queue_group_core(
messaging.as_ref(),
&namespaced,
req.group,
req.action,
req.start,
)
.await
}
#[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),
in_flight: messaging.in_flight_count(&namespaced).await.unwrap_or(0),
oldest_pending_ms: messaging
.oldest_pending_ms(&namespaced)
.await
.unwrap_or(None),
lag: if consumer.group.is_empty() {
0
} else {
messaging
.group_lag(&namespaced, &consumer.group)
.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_function_logs(
Extension(handlers): Extension<Arc<HandlerRuntime>>,
Extension(project): Extension<crate::project_scope::ProjectContext>,
Path(name): 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 scope = project.as_ref().qualified(&format!("fn/{name}"));
let (entries, dropped) = inner
.logs
.tail(&scope, 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(feature = "handlers")]
pub(super) async fn operator_function_logs_stream(
Extension(handlers): Extension<Arc<HandlerRuntime>>,
Extension(project): Extension<crate::project_scope::ProjectContext>,
Path(name): 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 want = project.as_ref().qualified(&format!("fn/{name}"));
let rx = inner.logs.subscribe();
let stream = futures::stream::unfold(rx, move |mut rx| {
let want = want.clone();
async move {
loop {
match rx.recv().await {
Ok((scope, entry)) if scope == want => {
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()
}