use crate::a2a::handler::notify::*;
use crate::a2a::test_helpers::helpers::placeholder_service_context;
use crate::a2a::types::*;
use crate::brain::agent::service::restart_recovery::test_guard;
use crate::brain::agent::service::session_routes::{ChannelOwnership, register_session_route};
use crate::brain::agent::{PushOrigin, QueuedUserMessage};
use crate::services::SessionService;
use std::sync::{Arc, Mutex};
fn params(session_id: &str, message: &str) -> serde_json::Value {
serde_json::json!({ "session_id": session_id, "message": message })
}
fn outcome_of(resp: &JsonRpcResponse) -> String {
resp.result
.as_ref()
.expect("success response")
.get("outcome")
.and_then(serde_json::Value::as_str)
.expect("outcome field")
.to_string()
}
#[tokio::test]
async fn dead_uuid_is_refused_without_touching_the_route_table() {
let ctx = placeholder_service_context().await;
let dead = uuid::Uuid::new_v4();
let resp =
handle_session_notify(serde_json::json!(1), params(&dead.to_string(), "ping"), ctx).await;
assert!(
resp.error.is_none(),
"dead uuid is a business outcome, not a protocol error: {resp:?}"
);
assert_eq!(outcome_of(&resp), "no_route");
}
#[tokio::test]
#[allow(clippy::await_holding_lock)]
async fn live_uuid_delivers_through_the_claimed_route() {
let _guard = test_guard();
let ctx = placeholder_service_context().await;
let session = SessionService::new(ctx.clone())
.create_session(Some("#23 test session".to_string()))
.await
.expect("session row created");
let sid = session.id;
let captured: Arc<Mutex<Option<QueuedUserMessage>>> = Arc::new(Mutex::new(None));
let sink = captured.clone();
register_session_route(
sid,
Arc::new(move |_id, queued| {
*sink.lock().unwrap() = Some(queued);
}),
);
let resp =
handle_session_notify(serde_json::json!(2), params(&sid.to_string(), "ping"), ctx).await;
assert!(resp.error.is_none(), "{resp:?}");
assert_eq!(outcome_of(&resp), "delivered");
let queued = captured.lock().unwrap().take().expect("message enqueued");
assert_eq!(queued.origin, PushOrigin::SessionNotify);
assert!(queued.context_text.contains(&format!(
"[session-notify from={CLI_SENDER_PREFIX}{DEFAULT_CLI_SENDER_LABEL}]"
)));
}
#[tokio::test]
#[allow(clippy::await_holding_lock)]
async fn sender_override_rides_the_header() {
let _guard = test_guard();
let ctx = placeholder_service_context().await;
let session = SessionService::new(ctx.clone())
.create_session(Some("#23 sender override".to_string()))
.await
.expect("session row created");
let sid = session.id;
let captured: Arc<Mutex<Option<QueuedUserMessage>>> = Arc::new(Mutex::new(None));
let sink = captured.clone();
register_session_route(
sid,
Arc::new(move |_id, queued| {
*sink.lock().unwrap() = Some(queued);
}),
);
let mut p = params(&sid.to_string(), "ping");
p["sender"] = serde_json::json!("oc-deploy");
let resp = handle_session_notify(serde_json::json!(7), p, ctx).await;
assert!(resp.error.is_none(), "{resp:?}");
assert_eq!(outcome_of(&resp), "delivered");
let queued = captured.lock().unwrap().take().expect("message enqueued");
assert!(
queued
.context_text
.contains("[session-notify from=cli:oc-deploy]"),
"override must ride the cli:-prefixed header: {}",
queued.context_text
);
assert!(
queued.display_text.contains("from oc-deploy"),
"display frame names the overridden sender: {}",
queued.display_text
);
}
#[tokio::test]
async fn malformed_params_are_protocol_errors() {
let ctx = placeholder_service_context().await;
let bad_uuid = handle_session_notify(
serde_json::json!(3),
params("not-a-uuid", "ping"),
ctx.clone(),
)
.await;
assert_eq!(
bad_uuid.error.expect("error response").code,
error_codes::INVALID_PARAMS
);
let empty_msg = handle_session_notify(
serde_json::json!(4),
params(&uuid::Uuid::new_v4().to_string(), " "),
ctx.clone(),
)
.await;
assert_eq!(
empty_msg.error.expect("error response").code,
error_codes::INVALID_PARAMS
);
let mut bad_sender = params(&uuid::Uuid::new_v4().to_string(), "ping");
bad_sender["sender"] = serde_json::json!("bad]label");
let bad_sender_resp = handle_session_notify(serde_json::json!(5), bad_sender, ctx).await;
assert_eq!(
bad_sender_resp.error.expect("error response").code,
error_codes::INVALID_PARAMS
);
}
#[tokio::test]
#[allow(clippy::await_holding_lock)]
async fn archived_session_auto_routes_to_its_successor() {
let _guard = test_guard();
let ctx = placeholder_service_context().await;
let svc = SessionService::new(ctx.clone());
let old = svc
.create_session(Some("#23 old session".to_string()))
.await
.expect("old session row");
svc.archive_session(old.id).await.expect("archived");
let successor = svc
.create_session(Some("#23 successor session".to_string()))
.await
.expect("successor session row");
let occupant = successor.id;
crate::brain::agent::service::session_routes::register_channel_owner_probe(
old.id,
std::sync::Arc::new(move || ChannelOwnership::Occupied { occupant }),
);
let captured: Arc<Mutex<Option<QueuedUserMessage>>> = Arc::new(Mutex::new(None));
let sink = captured.clone();
register_session_route(
successor.id,
Arc::new(move |_id, queued| {
*sink.lock().unwrap() = Some(queued);
}),
);
let resp = handle_session_notify(
serde_json::json!(5),
params(&old.id.to_string(), "ping"),
ctx,
)
.await;
assert!(resp.error.is_none(), "{resp:?}");
assert_eq!(outcome_of(&resp), "delivered");
let detail = resp
.result
.unwrap()
.get("detail")
.unwrap()
.as_str()
.unwrap()
.to_string();
assert!(
detail.contains("redirected"),
"detail should name the redirect: {detail}"
);
let queued = captured
.lock()
.unwrap()
.take()
.expect("successor received the redirect");
assert!(
queued
.context_text
.contains(&format!("originally for session {}", old.id)),
"provenance framing must name the archived session: {}",
queued.context_text
);
}