swarm-controller: answer 503 on a disconnected queue, not 500 or a hang
put_matrix_account writes the credential to bao before checking that the queue it must notify is actually connected — only that a queue is configured, via state.status.as_ref(). While the queue is Pending or Disconnected, client.flush().await hangs (async-nats does not process commands during the initial connect retry), so the request hangs until nginx times out and the credential is already stored. Add the ensure_connected check every sibling queue route already makes (wanted.rs, term_stream.rs, agent_state_stream.rs), placed immediately before store::connect() so a request that cannot be delivered never reaches the store. Same shape, lower impact, in webhook::announce_knowledge_change: it awaits publish/flush inline in the webhook handler, so a disconnected queue can run past forgejo's short delivery timeout. Add the same ensure_connected guard, warn and return. set_agent_state and get_hive_wanted map every writer error to 500, including swarm_queue_client::Error::NotConnected surfaced through WantedWriter::view/set. Add wanted_error_status, which downcasts the anyhow::Error back to the concrete type and maps NotConnected to 503 (retryable) while leaving every other failure at 500. Update both routes' OpenAPI descriptions to say so. Closes #4688
This commit is contained in:
parent
0bfe354b6d
commit
6a87b25488
3 changed files with 327 additions and 14 deletions
|
|
@ -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:?}"
|
||||
);
|
||||
}
|
||||
}
|
||||
|
|
|
|||
Loading…
Reference in a new issue