#[path = "alert_render.rs"]
pub mod alert_render;
use std::path::Path;
use crate::context::AppContext;
use crate::protocol::Response;
pub fn finalize_response(
response: &mut Response,
ctx: &AppContext,
session_id: &str,
attach_command: &str,
) {
finalize_response_with_bg_completions(response, ctx, session_id, attach_command, true);
}
pub fn finalize_response_with_bg_completions(
response: &mut Response,
ctx: &AppContext,
session_id: &str,
attach_command: &str,
allow_bg_completions: bool,
) {
if allow_bg_completions {
attach_bg_completions(response, ctx, session_id, attach_command);
}
let plane_live = publish_fleet_status(response, ctx, session_id);
if response.data.get("text").is_none()
&& !alert_render::is_excluded_finalization_command(attach_command)
{
attach_status_bar_after_publish(response, ctx, plane_live);
}
}
pub fn finalize_response_for_dispatch_root(
response: &mut Response,
ctx: &AppContext,
alerts: &mut alert_render::AlertEngine,
session_id: &str,
dispatch_root: &Path,
attach_command: &str,
allow_bg_completions: bool,
) {
if allow_bg_completions {
attach_bg_completions(response, ctx, session_id, attach_command);
}
let _ = publish_fleet_status(response, ctx, session_id);
attach_alert_block(response, alerts, session_id, dispatch_root, attach_command);
}
fn attach_alert_block(
response: &mut Response,
alerts: &mut alert_render::AlertEngine,
session_id: &str,
dispatch_root: &Path,
command: &str,
) {
let Some(text) = response
.data
.as_object_mut()
.and_then(|data| data.get_mut("text"))
.and_then(|value| value.as_str())
.map(str::to_string)
else {
return;
};
if text.contains("<system-reminder>") {
return;
}
let Some(alert) = alerts.finalize(session_id, dispatch_root, command) else {
return;
};
let joined = if text.is_empty() {
alert.text
} else {
format!("{text}\n\n{}", alert.text)
};
if let Some(data) = response.data.as_object_mut() {
data.insert("text".to_string(), serde_json::Value::String(joined));
}
}
pub enum DispatchOutcome {
Immediate(Response),
Deferred(PendingResponse),
}
pub type PendingResponsePoll = Box<dyn FnMut(&AppContext) -> Option<Response>>;
pub type PendingResponseShutdown = Box<dyn FnMut(&AppContext) -> Response>;
pub struct PendingResponse {
pub request_id: String,
pub session_id: String,
pub attach_command: String,
pub poll: PendingResponsePoll,
pub on_shutdown: Option<PendingResponseShutdown>,
}
pub struct ResolvedPending {
pub response: Response,
pub session_id: String,
pub attach_command: String,
}
#[derive(Default)]
pub struct PendingResponses {
entries: Vec<PendingResponse>,
}
impl PendingResponses {
pub fn register(&mut self, pending: PendingResponse) {
self.entries
.retain(|entry| entry.request_id != pending.request_id);
self.entries.push(pending);
}
pub fn poll_ready(&mut self, ctx: &AppContext) -> Vec<ResolvedPending> {
let mut ready = Vec::new();
let mut waiting = Vec::with_capacity(self.entries.len());
for mut pending in self.entries.drain(..) {
if let Some(response) = (pending.poll)(ctx) {
ready.push(ResolvedPending {
response,
session_id: pending.session_id,
attach_command: pending.attach_command,
});
} else {
waiting.push(pending);
}
}
self.entries = waiting;
ready
}
pub fn is_empty(&self) -> bool {
self.entries.is_empty()
}
pub fn drain_on_shutdown(&mut self) {
self.entries.clear();
}
pub fn drain_on_shutdown_with(&mut self, ctx: &AppContext) -> Vec<ResolvedPending> {
self.entries
.drain(..)
.filter_map(|mut pending| {
let response = (pending.on_shutdown.as_mut()?)(ctx);
Some(ResolvedPending {
response,
session_id: pending.session_id,
attach_command: pending.attach_command,
})
})
.collect()
}
}
pub fn attach_bg_completions(
response: &mut Response,
ctx: &AppContext,
session_id: &str,
command: &str,
) {
if matches!(
command,
"configure"
| "bash_abort_inflight"
| "bash_status"
| "bash_write"
| "bash_promote"
| "bash_wait_detach"
| "bash_regex_match"
| "bash_drain_completions"
| "bash_notify"
| "bash_unnotify"
| "bash_ack_completions"
) {
return;
}
if !ctx
.bash_background()
.has_completions_for_session(Some(session_id))
{
return;
}
let completions = ctx
.bash_background()
.drain_completions_for_session(Some(session_id));
if completions.is_empty() {
return;
}
let value = serde_json::json!(completions);
match response.data.as_object_mut() {
Some(data) => {
data.insert("bg_completions".to_string(), value);
}
None => {
response.data = serde_json::json!({ "bg_completions": value });
}
}
}
fn aft_status_segment(counts: &crate::context::StatusBarCounts) -> String {
let stale_mark = if counts.tier2_stale { "~" } else { "" };
format!(
"E{} W{} | {}D{} U{} C{} | T{}",
counts.errors,
counts.warnings,
stale_mark,
counts.dead_code,
counts.unused_exports,
counts.duplicates,
counts.todos
)
}
fn holder_owns_status_bar(plane_live: bool, harness: Option<&crate::harness::Harness>) -> bool {
plane_live && matches!(harness, Some(crate::harness::Harness::Opencode))
}
fn publish_fleet_status(
response: &mut Response,
ctx: &AppContext,
session_id: &str,
) -> Option<bool> {
if response
.data
.as_object_mut()
.and_then(|data| data.remove("_aft_suppress_status_bar"))
.is_some()
{
return None;
}
let local_counts = ctx.status_bar_counts();
let harness = ctx.harness_opt();
let plane_live = ctx.fleet_status_client().is_some_and(|client| {
let config = ctx.config();
let Some(project_root) = config.project_root.as_deref() else {
return false;
};
let harness_label = harness
.as_ref()
.map(crate::harness::Harness::wire_label)
.unwrap_or_else(|| "unknown".to_string());
let aft_text = local_counts
.as_ref()
.map(aft_status_segment)
.unwrap_or_default();
client.publish(project_root, &harness_label, session_id, &aft_text)
});
Some(plane_live)
}
pub fn attach_status_bar(
response: &mut Response,
ctx: &AppContext,
session_id: &str,
command: &str,
) {
if alert_render::is_excluded_finalization_command(command) {
return;
}
let plane_live = publish_fleet_status(response, ctx, session_id);
attach_status_bar_after_publish(response, ctx, plane_live);
}
fn attach_status_bar_after_publish(
response: &mut Response,
ctx: &AppContext,
plane_live: Option<bool>,
) {
let Some(plane_live) = plane_live else {
return;
};
let harness = ctx.harness_opt();
if holder_owns_status_bar(plane_live, harness.as_ref()) {
return;
}
let Some(counts) = ctx.status_bar_counts() else {
return;
};
if !ctx.should_emit_status_bar(&counts) {
return;
}
let value = serde_json::json!({
"errors": counts.errors,
"warnings": counts.warnings,
"dead_code": counts.dead_code,
"unused_exports": counts.unused_exports,
"duplicates": counts.duplicates,
"todos": counts.todos,
"tier2_stale": counts.tier2_stale,
});
match response.data.as_object_mut() {
Some(data) => {
data.insert("status_bar".to_string(), value);
}
None => {
response.data = serde_json::json!({ "status_bar": value });
}
}
}
#[cfg(test)]
mod tests {
use std::path::PathBuf;
use super::{
aft_status_segment, finalize_response_with_bg_completions, holder_owns_status_bar,
PendingResponse, PendingResponses,
};
use crate::config::Config;
use crate::context::{AppContext, StatusBarCounts};
use crate::fleet_status::FleetStatusClient;
use crate::harness::Harness;
use crate::parser::TreeSitterProvider;
use crate::protocol::Response;
#[test]
fn live_holder_retires_only_opencode_response_bars() {
assert_eq!(
(
holder_owns_status_bar(true, Some(&Harness::Opencode)),
holder_owns_status_bar(true, Some(&Harness::Runner)),
),
(true, false)
);
assert!(!holder_owns_status_bar(true, Some(&Harness::Pi)));
assert!(!holder_owns_status_bar(false, Some(&Harness::Opencode)));
}
#[test]
fn pre_discovery_publish_does_not_trip_the_holder_ownership_gate() {
let ctx = AppContext::new(
Box::new(TreeSitterProvider::new()),
Config {
project_root: Some(PathBuf::from("/tmp/project")),
..Config::default()
},
);
ctx.set_harness(Harness::Opencode);
ctx.update_status_bar_tier2(Some(21), Some(12), Some(13), Some(14), false);
let (client, mut wire_rx) = FleetStatusClient::dial_channel(1);
ctx.install_fleet_status_client(Some(client));
let mut response = Response::success("status", serde_json::json!({}));
finalize_response_with_bg_completions(&mut response, &ctx, "session-1", "echo", false);
assert_eq!(response.data["status_bar"]["dead_code"], 21);
let publish = wire_rx.try_recv().expect("single discovery publish");
assert_eq!(publish.body()["text"], "E0 W0 | D21 U12 C13 | T14");
assert!(
wire_rx.try_recv().is_err(),
"response published more than once"
);
publish.complete_unavailable();
}
#[test]
fn solo_bar_bytes_remain_the_existing_golden() {
let counts = StatusBarCounts {
errors: 2,
warnings: 5,
dead_code: 331,
unused_exports: 221,
duplicates: 1159,
todos: 8,
tier2_stale: false,
};
assert_eq!(
format!("[AFT {}]", aft_status_segment(&counts)),
"[AFT E2 W5 | D331 U221 C1159 | T8]"
);
}
#[test]
fn shutdown_delivery_emits_terminal_before_removing_entry() {
let ctx = AppContext::new(Box::new(TreeSitterProvider::new()), Config::default());
let mut pending = PendingResponses::default();
pending.register(PendingResponse {
request_id: "inspect-shutdown".to_string(),
session_id: String::new(),
attach_command: String::new(),
poll: Box::new(|_| None),
on_shutdown: Some(Box::new(|_| {
Response::error("inspect-shutdown", "daemon_shutdown", "shutdown")
})),
});
let resolved = pending.drain_on_shutdown_with(&ctx);
assert_eq!(resolved.len(), 1);
assert_eq!(resolved[0].response.id, "inspect-shutdown");
assert!(pending.is_empty());
}
#[test]
fn solo_bar_stale_marker_bytes_remain_the_existing_golden() {
let counts = StatusBarCounts {
dead_code: 10,
tier2_stale: true,
..StatusBarCounts::default()
};
assert_eq!(
format!("[AFT {}]", aft_status_segment(&counts)),
"[AFT E0 W0 | ~D10 U0 C0 | T0]"
);
}
}