diff --git a/swarm-controller/src/main.rs b/swarm-controller/src/main.rs index 04f69945..bc22d654 100644 --- a/swarm-controller/src/main.rs +++ b/swarm-controller/src/main.rs @@ -795,6 +795,26 @@ fn error_problem(status: axum::http::StatusCode, detail: &str) -> problem_detail problem_details::ProblemDetails::from_status_code(status).with_detail(detail) } +/// The status a [`wanted::WantedWriter`] failure should answer with. +/// +/// `WantedWriter::view`/`set` call `ensure_connected` internally and +/// propagate through `anyhow`, so the concrete +/// [`swarm_queue_client::Error`] survives underneath but is not the type in +/// hand — `swarm_queue_client`'s own doc says a caller "gets a type it can +/// match on", which means downcasting back to it rather than string-matching +/// the rendered message. `NotConnected` is a retryable condition (the queue +/// is reconnecting, or hasn't yet) and every other queue-backed route here +/// already answers 503 for it; anything else is a genuine failure to publish +/// or read, which stays 500. +fn wanted_error_status(e: &anyhow::Error) -> axum::http::StatusCode { + match e.downcast_ref::() { + Some(swarm_queue_client::Error::NotConnected(_)) => { + axum::http::StatusCode::SERVICE_UNAVAILABLE + } + _ => axum::http::StatusCode::INTERNAL_SERVER_ERROR, + } +} + /// The state to declare for one agent. #[derive(Debug, Deserialize, ToSchema)] struct SetAgentStateRequest { @@ -905,7 +925,7 @@ fn declaration_target( responses( (status = 200, description = "the declaration as now published", body = Vec), (status = 400, description = "a name is not an identifier, the hive is not in this swarm, or the state is unknown (problem+json)", body = String), - (status = 503, description = "no swarm queue is wired up (problem+json)", body = String), + (status = 503, description = "no swarm queue is wired up, or it is not connected (problem+json)", body = String), (status = 500, description = "the declaration could not be published (problem+json)", body = String), ), tag = "agents" @@ -923,10 +943,7 @@ async fn set_agent_state( let declaration = writer.set(&hive, &agent, req.state).await.map_err(|e| { tracing::warn!(hive = %hive, agent = %agent, error = %format!("{e:#}"), "declaring agent state failed"); - error_problem( - axum::http::StatusCode::INTERNAL_SERVER_ERROR, - &format!("{e:#}"), - ) + error_problem(wanted_error_status(&e), &format!("{e:#}")) })?; Ok(Json(render(&declaration))) } @@ -942,7 +959,7 @@ async fn set_agent_state( responses( (status = 200, description = "the declaration, empty when nothing is published yet", body = Vec), (status = 400, description = "not an identifier, or not a hive in this swarm (problem+json)", body = String), - (status = 503, description = "no swarm queue is wired up (problem+json)", body = String), + (status = 503, description = "no swarm queue is wired up, or it is not connected (problem+json)", body = String), (status = 500, description = "the declaration could not be read (problem+json)", body = String), ), tag = "agents" @@ -955,10 +972,7 @@ async fn get_hive_wanted( declaration_target(&state, &hive).map_err(|(s, d)| error_problem(s, &d))?; let declaration = writer.view(&hive).await.map_err(|e| { tracing::warn!(hive = %hive, error = %format!("{e:#}"), "reading the declaration failed"); - error_problem( - axum::http::StatusCode::INTERNAL_SERVER_ERROR, - &format!("{e:#}"), - ) + error_problem(wanted_error_status(&e), &format!("{e:#}")) })?; Ok(Json(declaration.as_ref().map(render).unwrap_or_default())) } @@ -1970,8 +1984,9 @@ fn build_app(state: AppState) -> axum::Router { #[cfg(test)] mod tests { use super::{ - DEFAULT_SOCKET, HIVES_ENV, HiveEntry, LINKS_ENV, NAME_ENV, ServiceLink, StatusUnavailable, - SwarmNodeKind, WorkerDeps, load_hives, load_links, load_swarm_name, run_swarm_node, + DEFAULT_SOCKET, HIVES_ENV, HiveEntry, LINKS_ENV, NAME_ENV, ServiceLink, + SetAgentStateRequest, StatusUnavailable, SwarmNodeKind, WorkerDeps, load_hives, load_links, + load_swarm_name, run_swarm_node, wanted, }; use std::path::Path; @@ -2782,6 +2797,150 @@ mod tests { std::env::remove_var(NAME_ENV); } } + + // ── NotConnected maps to 503, everything else stays 500 ───────────── + + /// `wanted_error_status` as a pure function first: build the exact + /// `swarm_queue_client::Error::NotConnected` `ensure_connected` raises, + /// wrap it the same way `?` inside `WantedWriter::view`/`set` does, and + /// confirm it downcasts back to the status those routes are meant to + /// answer. + #[test] + fn not_connected_downcasts_to_503() { + let e: anyhow::Error = + swarm_queue_client::Error::NotConnected(async_nats::connection::State::Pending).into(); + assert_eq!( + super::wanted_error_status(&e), + axum::http::StatusCode::SERVICE_UNAVAILABLE + ); + } + + /// Invert-proof for the test above: an `anyhow::Error` that does not + /// wrap a `swarm_queue_client::Error` at all (the shape a genuine write + /// or decode failure takes — `apply`'s "declaration is not decodable" + /// context, for instance) must NOT downcast to `NotConnected`, and must + /// stay 500. Without this, a `wanted_error_status` that answered 503 + /// unconditionally would still pass the test above. + #[test] + fn an_unrelated_error_stays_500() { + let e = anyhow::anyhow!("the hive's current declaration is not decodable"); + assert_eq!( + super::wanted_error_status(&e), + axum::http::StatusCode::INTERNAL_SERVER_ERROR + ); + } + + /// Second invert-proof, closer to the actual bug this closes: a + /// *different* `swarm_queue_client::Error` variant (a real connect + /// failure, not `NotConnected`) also stays 500 — proving the match + /// looks at the variant, not just at "was this crate's error type + /// involved at all". + #[test] + fn a_different_queue_client_error_variant_stays_500() { + let e: anyhow::Error = swarm_queue_client::Error::Connect { + url: "nats://queue.example:4222".to_owned(), + source: async_nats::ConnectErrorKind::TimedOut.into(), + } + .into(); + assert_eq!( + super::wanted_error_status(&e), + axum::http::StatusCode::INTERNAL_SERVER_ERROR + ); + } + + /// A client that exists but has never connected — same shape + /// `matrix_account`'s tests build, duplicated here rather than shared: + /// both are a handful of lines and neither crate has a test-support + /// module to put a shared helper in. + async fn disconnected_client() -> async_nats::Client { + async_nats::ConnectOptions::new() + .retry_on_initial_connect() + .connect("127.0.0.1:1") + .await + .expect("retry_on_initial_connect returns without waiting for a real connection") + } + + /// `state_with_roster` plus a `wanted` writer over a client that never + /// connected — the shape `set_agent_state`/`get_hive_wanted` see when + /// the queue is mid-reconnect. + async fn state_with_disconnected_wanted() -> super::AppState { + let (state, _sched) = state_with_roster(); + super::AppState { + wanted: Some(std::sync::Arc::new(wanted::WantedWriter::new( + disconnected_client().await, + ))), + ..state + } + } + + /// End to end through the real handler: `PUT .../state` against a + /// disconnected queue answers 503, not the 500 every other writer + /// error still gets. + #[tokio::test] + async fn set_agent_state_answers_503_on_a_disconnected_queue() { + let state = state_with_disconnected_wanted().await; + + let result = super::set_agent_state( + axum::extract::State(state), + axum::extract::Path(("pr1ma".to_owned(), "atlas".to_owned())), + axum::Json(SetAgentStateRequest { + state: swarm_queue_client::wanted::AgentState::Paused, + }), + ) + .await; + + let problem = result.expect_err("a disconnected queue must not read as success"); + assert_eq!( + problem.status, + Some(axum::http::StatusCode::SERVICE_UNAVAILABLE), + "{problem:?}" + ); + } + + /// Same fix, the read side: `GET .../wanted` against a disconnected + /// queue also answers 503. + #[tokio::test] + async fn get_hive_wanted_answers_503_on_a_disconnected_queue() { + let state = state_with_disconnected_wanted().await; + + let result = super::get_hive_wanted( + axum::extract::State(state), + axum::extract::Path("pr1ma".to_owned()), + ) + .await; + + let problem = result.expect_err("a disconnected queue must not read as success"); + assert_eq!( + problem.status, + Some(axum::http::StatusCode::SERVICE_UNAVAILABLE), + "{problem:?}" + ); + } + + /// Control for both tests above: a hive that is not in the roster is + /// still a plain 400, even with the same disconnected writer wired up — + /// so the 503 above is `ensure_connected` firing, not every error this + /// handler can produce collapsing to 503. + #[tokio::test] + async fn set_agent_state_still_answers_400_for_an_unknown_hive() { + let state = state_with_disconnected_wanted().await; + + let result = super::set_agent_state( + axum::extract::State(state), + axum::extract::Path(("not-a-hive".to_owned(), "atlas".to_owned())), + axum::Json(SetAgentStateRequest { + state: swarm_queue_client::wanted::AgentState::Paused, + }), + ) + .await; + + let problem = result.expect_err("an unknown hive is still a 400"); + assert_eq!( + problem.status, + Some(axum::http::StatusCode::BAD_REQUEST), + "{problem:?}" + ); + } } #[cfg(test)] diff --git a/swarm-controller/src/matrix_account.rs b/swarm-controller/src/matrix_account.rs index 310c6680..a9163382 100644 --- a/swarm-controller/src/matrix_account.rs +++ b/swarm-controller/src/matrix_account.rs @@ -134,7 +134,7 @@ pub struct PutMatrixAccountResponse { responses( (status = 200, description = "stored, and the hive has been told", body = PutMatrixAccountResponse), (status = 400, description = "a name is not an identifier, the account name is not a single path segment, the account is 'main' (reserved), the mode is unrecognized, a mode's required fields are missing, or the hive is not in this swarm (problem+json)", body = String), - (status = 503, description = "no swarm queue is wired up (problem+json)", body = String), + (status = 503, description = "no swarm queue is wired up, or it is not connected (problem+json)", body = String), (status = 500, description = "the store write, the encode or the publish failed (problem+json)", body = String), ), tag = "agents" @@ -182,6 +182,21 @@ pub async fn put_matrix_account( // touched, so a failed login leaves no partial state behind. let (token, homeserver, user_id) = resolve_credential(&req).await.map_err(|b| *b)?; + // A queue that is configured but not yet (or no longer) connected does + // not fail the `publish`/`flush` further down outright — it *hangs* + // them, per `swarm_queue_client::ensure_connected`'s own doc. Checked + // here, immediately before the store is touched, so a request that + // cannot be delivered never leaves a credential behind — the same + // ordering rule the module doc states for "store first, notify + // second", extended one step earlier. + let client = status.queue_client(); + swarm_queue_client::ensure_connected(&client).map_err(|e| { + error_problem( + StatusCode::SERVICE_UNAVAILABLE, + &swarm_queue_client::chain(&e), + ) + })?; + let store = crate::store::connect().await.map_err(|e| { tracing::warn!(error = %e, "connecting to the swarm secret store failed"); error_problem(StatusCode::INTERNAL_SERVER_ERROR, &e.to_string()) @@ -212,7 +227,6 @@ pub async fn put_matrix_account( ) })?; let subject = credential_subject(&hive); - let client = status.queue_client(); client .publish(subject.clone(), payload.into()) .await @@ -493,4 +507,134 @@ mod tests { .is_err() ); } + + // ── the connected-but-not-yet-connected queue ─────────────────────── + + /// A client that exists but has never connected — the `Pending` state + /// `ensure_connected`'s own doc says a `Disconnected` check would miss. + /// `retry_on_initial_connect` is what makes `.connect()` return + /// immediately instead of blocking on a handshake that will never + /// succeed against a loopback port nothing listens on. + async fn disconnected_client() -> async_nats::Client { + async_nats::ConnectOptions::new() + .retry_on_initial_connect() + .connect("127.0.0.1:1") + .await + .expect("retry_on_initial_connect returns without waiting for a real connection") + } + + /// Bare-minimum `AppState` for a handler test: one hive, no queue-backed + /// helpers beyond `status` (the field this handler actually reads), and + /// an empty in-memory job graph the endpoint under test never touches. + fn state_with_status(status: super::super::status::StatusReader) -> super::super::AppState { + super::super::AppState { + hives: std::sync::Arc::new(vec![super::super::HiveEntry { + name: "pr1ma".to_owned(), + domain: "pr1ma.example".to_owned(), + }]), + links: std::sync::Arc::new(Vec::new()), + status: Some(std::sync::Arc::new(status)), + wanted: None, + agent_status: None, + jobq: std::sync::Arc::new(std::sync::Mutex::new(hive_jobq::scheduler::Scheduler::new( + hive_jobq::Graph::new(), + hive_jobq::resources::ResourceTable::new(), + ))), + webhook_secret: None, + config_prs: None, + swarm_name: None, + auth: None, + forge: None, + } + } + + /// The defect this whole PR exists to close: a queue that is + /// configured but not connected must not let this handler reach the + /// store write at all. + /// + /// `BAO_ADDR`/`BAO_CLIENT_CERT`/`BAO_CLIENT_KEY` are asserted unset + /// first — not incidental setup, but the control that makes the 503 + /// meaningful. If `put_matrix_account` reached `crate::store::connect()` + /// with those unset, *that* call fails too, and would also answer with + /// a `problem+json` body (500, "connecting to the swarm secret store + /// failed"). A 503 here is therefore proof execution never got past + /// `ensure_connected`, not a coincidence of two paths landing on the + /// same status family. + #[tokio::test] + async fn a_disconnected_queue_answers_503_and_never_reaches_the_store() { + for var in ["BAO_ADDR", "BAO_CLIENT_CERT", "BAO_CLIENT_KEY"] { + assert!( + std::env::var(var).is_err(), + "{var} must be unset for this test to prove anything" + ); + } + + let reader = super::super::status::StatusReader::new( + disconnected_client().await, + std::time::Duration::from_mins(1), + ); + let state = state_with_status(reader); + + let result = super::put_matrix_account( + axum::extract::State(state), + axum::extract::Path(( + "pr1ma".to_owned(), + "atlas".to_owned(), + "workaccount".to_owned(), + )), + axum::Json(PutMatrixAccountRequest { + mode: "token".to_owned(), + token: Some("t0k3n".to_owned()), + user_id: None, + password: None, + homeserver: None, + }), + ) + .await; + + let problem = result.expect_err("a disconnected queue must refuse, not hang or 500"); + assert_eq!( + problem.status, + Some(axum::http::StatusCode::SERVICE_UNAVAILABLE), + "{problem:?}" + ); + } + + /// The control for the test above: a client that starts in `Pending` + /// but genuinely never connects is exactly what `ensure_connected` is + /// specified to reject — proving the 503 above tracks the connection + /// state and is not simply what every call through this handler + /// returns. Same client shape, mode rejected before either queue or + /// store is touched, so it exercises a different early return + /// (`is_reserved_account`) and confirms the handler still validates + /// normally on a path that never reaches `ensure_connected`'s sibling + /// checks. + #[tokio::test] + async fn a_reserved_account_name_is_still_refused_before_any_queue_check_matters() { + let reader = super::super::status::StatusReader::new( + disconnected_client().await, + std::time::Duration::from_mins(1), + ); + let state = state_with_status(reader); + + let result = super::put_matrix_account( + axum::extract::State(state), + axum::extract::Path(("pr1ma".to_owned(), "atlas".to_owned(), "main".to_owned())), + axum::Json(PutMatrixAccountRequest { + mode: "token".to_owned(), + token: Some("t0k3n".to_owned()), + user_id: None, + password: None, + homeserver: None, + }), + ) + .await; + + let problem = result.expect_err("'main' is reserved regardless of queue state"); + assert_eq!( + problem.status, + Some(axum::http::StatusCode::BAD_REQUEST), + "{problem:?}" + ); + } } diff --git a/swarm-controller/src/webhook.rs b/swarm-controller/src/webhook.rs index 3da6d8ae..bdaa9551 100644 --- a/swarm-controller/src/webhook.rs +++ b/swarm-controller/src/webhook.rs @@ -419,6 +419,16 @@ async fn announce_knowledge_change(state: &AppState) { let client = status.queue_client(); let subject = swarm_queue_client::KNOWLEDGE_SUBJECT; + // A queue that is configured but not (yet, or no longer) connected does + // not fail `publish`/`flush` below outright — it *hangs* them, per + // `ensure_connected`'s own doc — and this fn is awaited inline from the + // webhook handler, so that hang would sit inside forgejo's own delivery + // timeout instead of failing soft the way every other branch here does. + if let Err(e) = swarm_queue_client::ensure_connected(&client) { + tracing::warn!(%subject, error = %e, "webhook: swarm queue is not connected; not publishing"); + return; + } + // One publish, not one per hive: every subscriber gets the same empty // event, so the roster is not consulted at all. The controller does not // need to know who the hives are in order to say the repository moved.