diff --git a/crates/runtime/src/provider/edge_connection.rs b/crates/runtime/src/provider/edge_connection.rs index 644c3c3041..a37d7ffc71 100644 --- a/crates/runtime/src/provider/edge_connection.rs +++ b/crates/runtime/src/provider/edge_connection.rs @@ -435,6 +435,7 @@ mod tests { _user_id: &str, _edge_agent_id: &str, _edge_id_header: &str, + _registration_claim_id: Option<&str>, ) -> Result<(), astra_services::multi_agent::HeartbeatError> { Ok(()) } diff --git a/crates/runtime/src/server/edge/edge_callback_handlers.rs b/crates/runtime/src/server/edge/edge_callback_handlers.rs index 038a9e744f..a28d1eb412 100644 --- a/crates/runtime/src/server/edge/edge_callback_handlers.rs +++ b/crates/runtime/src/server/edge/edge_callback_handlers.rs @@ -1639,7 +1639,7 @@ pub(crate) async fn post_agents_edge_heartbeat_handler( state .execution .edge_registry_service - .heartbeat(&user.user_id, &body.edge_agent_id, &edge_id) + .heartbeat(&user.user_id, &body.edge_agent_id, &edge_id, None) .await .map_err(|e| match e { astra_services::multi_agent::HeartbeatError::Superseded => error_response( diff --git a/crates/runtime/src/server/edge/edge_ws_handler.rs b/crates/runtime/src/server/edge/edge_ws_handler.rs index e9a8ecd466..b5e846fddd 100644 --- a/crates/runtime/src/server/edge/edge_ws_handler.rs +++ b/crates/runtime/src/server/edge/edge_ws_handler.rs @@ -14,6 +14,8 @@ use axum::response::IntoResponse; use futures_util::StreamExt; use futures_util::stream::{SplitSink, SplitStream}; use std::collections::HashMap; +use std::future::Future; +use std::pin::Pin; use std::sync::Arc; use std::sync::atomic::{AtomicUsize, Ordering}; use std::time::Duration; @@ -23,6 +25,52 @@ use tokio::sync::mpsc; const MAX_EDGE_WS_CONNECTIONS: usize = 1024; const EDGE_REGISTRY_UNREGISTER_ATTEMPTS: usize = 3; const EDGE_HEARTBEAT_STORAGE_FAILURE_BUDGET: usize = 3; +const EDGE_REGISTRY_RELEASE_RETRY_BASE: Duration = Duration::from_secs(1); +const EDGE_REGISTRY_RELEASE_RETRY_MAX: Duration = Duration::from_secs(30); +const EDGE_REGISTRY_RELEASE_ATTEMPT_TIMEOUT: Duration = Duration::from_secs(5); + +type EdgeRegistryReleaseReconciliation = Pin + Send>>; + +fn reconcile_edge_registry_release( + registry: Arc, + lease: astra_services::multi_agent::EdgeRegistrationLease, +) -> EdgeRegistryReleaseReconciliation { + Box::pin(async move { + let mut retry_delay = EDGE_REGISTRY_RELEASE_RETRY_BASE; + loop { + match tokio::time::timeout( + EDGE_REGISTRY_RELEASE_ATTEMPT_TIMEOUT, + registry.release_registration(&lease), + ) + .await + { + Ok(Ok(settled)) => return settled, + Ok(Err(error)) => { + tracing::warn!( + target: "astra_runtime::edge_ws", + user_id = %lease.current.user_id, + edge_agent_id = %lease.current.edge_agent_id, + %error, + "edge WebSocket: durable registration publication remains unresolved" + ); + } + Err(_) => { + tracing::warn!( + target: "astra_runtime::edge_ws", + user_id = %lease.current.user_id, + edge_agent_id = %lease.current.edge_agent_id, + timeout_ms = EDGE_REGISTRY_RELEASE_ATTEMPT_TIMEOUT.as_millis(), + "edge WebSocket: durable registration publication attempt timed out" + ); + } + } + tokio::time::sleep(retry_delay).await; + retry_delay = retry_delay + .saturating_mul(2) + .min(EDGE_REGISTRY_RELEASE_RETRY_MAX); + } + }) +} /// Global counter of active edge WebSocket connections. static EDGE_WS_CONNECTION_COUNT: AtomicUsize = AtomicUsize::new(0); @@ -628,46 +676,15 @@ async fn handle_edge_connection( workspace_id.clone(), pool_tx, ); - match edge_registry - .release_registration(®istration_lease) - .await - { - Ok(false) if registration_lease.claim_id.is_some() => { - // A definite claim mismatch means another pod already owns the - // durable generation. Fail closed instead of publishing a local - // connection that cross-pod routing cannot consistently target. - state.edge_connection_pool.unregister_generation( - &user_id, - &edge_agent_id, - pool_generation, - ); - forward_task.abort(); - drop(reconnect_guard); - state - .edge_connection_pool - .gc_reconnect_lock(&user_id, &edge_agent_id); - let _ = send_edge_msg( - &ws_sink, - EdgeServerMessage::AuthError { - message: "edge registry registration claim lost".into(), - }, - ) - .await; - return; - } - Err(error) => { - // The claim has a DB-side expiry, so an outcome-unknown release - // delays another cross-pod reconnect but cannot fence this healthy - // connection forever. - tracing::error!( - target: "astra_runtime::edge_ws", - user_id = %user_id, - edge_agent_id = %edge_agent_id, - %error, - "edge WebSocket: failed to release durable registration claim" - ); - } - _ => {} + // Independently poll publication and its deadline even while a message or + // heartbeat handler awaits database resources. JoinSet aborts on owner drop; + // normal teardown also joins cancellation before durable unregister. + let mut registration_release = tokio::task::JoinSet::new(); + if registration_lease.claim_id.is_some() { + registration_release.spawn(reconcile_edge_registry_release( + edge_registry.clone(), + registration_lease.clone(), + )); } drop(reconnect_guard); // Release the per-key reconnect lock entry now that the reconnect is done. @@ -1002,7 +1019,12 @@ async fn handle_edge_connection( // edge_id so a stale connection cannot refresh a row that a // newer connection has already claimed. match edge_registry - .heartbeat(&user_id, &edge_agent_id, &edge_id_for_registry) + .heartbeat( + &user_id, + &edge_agent_id, + &edge_id_for_registry, + registration_lease.claim_id.as_deref(), + ) .await { Ok(()) => consecutive_heartbeat_storage_failures = 0, @@ -1042,6 +1064,51 @@ async fn handle_edge_connection( } } } + released = registration_release.join_next(), if !registration_release.is_empty() => { + let released = match released { + Some(Ok(released)) => released, + Some(Err(error)) => { + tracing::error!( + target: "astra_runtime::edge_ws", + user_id = %user_id, + edge_agent_id = %edge_agent_id, + %error, + "edge WebSocket: registration publication task failed" + ); + let _ = send_edge_msg( + &ws_sink_write, + EdgeServerMessage::Closing { + reason: "edge registry publication unavailable".into(), + }, + ).await; + break; + } + None => unreachable!("publication task exists while join is enabled"), + }; + if released { + tracing::info!( + target: "astra_runtime::edge_ws", + user_id = %user_id, + edge_agent_id = %edge_agent_id, + "edge WebSocket: reconciled durable registration publication" + ); + } else { + tracing::info!( + target: "astra_runtime::edge_ws", + user_id = %user_id, + edge_agent_id = %edge_agent_id, + "edge WebSocket: durable registration was superseded during publication reconciliation" + ); + let _ = send_edge_msg( + &ws_sink_write, + EdgeServerMessage::Closing { + reason: "edge registry registration claim lost".into(), + }, + ) + .await; + break; + } + } } } }; @@ -1049,6 +1116,9 @@ async fn handle_edge_connection( read_loop.await; // ── Cleanup ────────────────────────────────────────────────────── + // Aborting alone only schedules cancellation. Joining guarantees the + // publication future has dropped its transaction/connection first. + registration_release.shutdown().await; forward_task.abort(); // Drop cancel sender so the dispatch task can break its loop cleanly. drop(dispatch_cancel_tx); @@ -1575,6 +1645,7 @@ mod tests { _user_id: &str, _edge_agent_id: &str, _edge_id_header: &str, + _registration_claim_id: Option<&str>, ) -> Result<(), astra_services::multi_agent::HeartbeatError> { unreachable!("heartbeat is not used by unregister retry tests") } diff --git a/crates/runtime/src/server/runtime_tool_executor.rs b/crates/runtime/src/server/runtime_tool_executor.rs index 097612c4b2..e1bece2f18 100644 --- a/crates/runtime/src/server/runtime_tool_executor.rs +++ b/crates/runtime/src/server/runtime_tool_executor.rs @@ -8213,6 +8213,7 @@ esac _user_id: &str, _edge_agent_id: &str, _edge_id_header: &str, + _registration_claim_id: Option<&str>, ) -> Result<(), astra_services::multi_agent::HeartbeatError> { Err(astra_services::multi_agent::HeartbeatError::StorageFailure( "MCP tools must not heartbeat edge registry".to_string(), @@ -8382,6 +8383,7 @@ esac _user_id: &str, _edge_agent_id: &str, _edge_id_header: &str, + _registration_claim_id: Option<&str>, ) -> Result<(), astra_services::multi_agent::HeartbeatError> { Ok(()) } diff --git a/crates/runtime/src/server/tool_transport/tests.rs b/crates/runtime/src/server/tool_transport/tests.rs index f45e5ea67a..f9f8bca730 100644 --- a/crates/runtime/src/server/tool_transport/tests.rs +++ b/crates/runtime/src/server/tool_transport/tests.rs @@ -846,6 +846,7 @@ impl astra_services::multi_agent::EdgeRegistryService for StaticEdgeRegistry { _user_id: &str, _edge_agent_id: &str, _edge_id_header: &str, + _registration_claim_id: Option<&str>, ) -> Result<(), astra_services::multi_agent::HeartbeatError> { Ok(()) } diff --git a/crates/runtime/tests/edge_ws_e2e.rs b/crates/runtime/tests/edge_ws_e2e.rs index 45f3a9a8b4..13e0a0c2f1 100644 --- a/crates/runtime/tests/edge_ws_e2e.rs +++ b/crates/runtime/tests/edge_ws_e2e.rs @@ -1165,25 +1165,65 @@ struct BlockingLeaseEdgeRegistry { claim_release_started: Notify, claim_release_gate: Notify, rollback_count: AtomicUsize, - release_succeeds: bool, + release_attempts: AtomicUsize, + release_outcomes: std::sync::Mutex>>, + block_first_release: bool, + block_every_release: bool, + active_releases: AtomicUsize, + database_connection: tokio::sync::Semaphore, + live_pool: Option, + heartbeat_completed: Notify, + heartbeat_started: Notify, + unregister_completed: Notify, +} + +struct ReleaseDropSentinel<'a>(&'a AtomicUsize); + +impl Drop for ReleaseDropSentinel<'_> { + fn drop(&mut self) { + self.0.fetch_sub(1, Ordering::SeqCst); + } } impl BlockingLeaseEdgeRegistry { fn new() -> Self { + Self::with_release_outcomes([Ok(true)], true) + } + + fn with_claim_loss() -> Self { + Self::with_release_outcomes([Ok(false)], true) + } + + fn with_release_recovery() -> Self { + Self::with_release_outcomes( + [ + Err("simulated release outcome unknown".to_string()), + Ok(true), + ], + false, + ) + } + + fn with_release_outcomes( + outcomes: impl IntoIterator>, + block_first_release: bool, + ) -> Self { Self { registration_started: Notify::new(), release_registration: Notify::new(), claim_release_started: Notify::new(), claim_release_gate: Notify::new(), rollback_count: AtomicUsize::new(0), - release_succeeds: true, - } - } - - fn with_claim_loss() -> Self { - Self { - release_succeeds: false, - ..Self::new() + release_attempts: AtomicUsize::new(0), + release_outcomes: std::sync::Mutex::new(outcomes.into_iter().collect()), + block_first_release, + block_every_release: false, + active_releases: AtomicUsize::new(0), + database_connection: tokio::sync::Semaphore::new(1), + live_pool: None, + heartbeat_completed: Notify::new(), + heartbeat_started: Notify::new(), + unregister_completed: Notify::new(), } } } @@ -1246,9 +1286,23 @@ impl astra_services::multi_agent::EdgeRegistryService for BlockingLeaseEdgeRegis &self, _lease: &astra_services::multi_agent::EdgeRegistrationLease, ) -> Result { + let attempt = self.release_attempts.fetch_add(1, Ordering::SeqCst); + let _connection = self.database_connection.acquire().await.unwrap(); + let _transaction = match &self.live_pool { + Some(pool) => Some(pool.begin().await.expect("publication transaction")), + None => None, + }; + self.active_releases.fetch_add(1, Ordering::SeqCst); + let _sentinel = ReleaseDropSentinel(&self.active_releases); self.claim_release_started.notify_one(); - self.claim_release_gate.notified().await; - Ok(self.release_succeeds) + if self.block_every_release || (self.block_first_release && attempt == 0) { + self.claim_release_gate.notified().await; + } + self.release_outcomes + .lock() + .unwrap() + .pop_front() + .unwrap_or(Ok(true)) } async fn heartbeat( @@ -1256,7 +1310,16 @@ impl astra_services::multi_agent::EdgeRegistryService for BlockingLeaseEdgeRegis _user_id: &str, _edge_agent_id: &str, _edge_id_header: &str, + _registration_claim_id: Option<&str>, ) -> Result<(), astra_services::multi_agent::HeartbeatError> { + if self.block_every_release { + // A one-connection pool makes heartbeat wait for the pending + // publication attempt to time out and drop its resource. + assert_eq!(self.active_releases.load(Ordering::SeqCst), 1); + self.heartbeat_started.notify_one(); + let _connection = self.database_connection.acquire().await.unwrap(); + self.heartbeat_completed.notify_one(); + } Ok(()) } @@ -1281,6 +1344,19 @@ impl astra_services::multi_agent::EdgeRegistryService for BlockingLeaseEdgeRegis _edge_agent_id: &str, _edge_id_header: &str, ) -> Result { + assert_eq!( + self.active_releases.load(Ordering::SeqCst), + 0, + "release must be dropped before durable unregister starts" + ); + let _connection = self.database_connection.acquire().await.unwrap(); + if let Some(pool) = &self.live_pool { + sqlx::query("SELECT 1") + .execute(pool) + .await + .expect("unregister can acquire the released database connection"); + } + self.unregister_completed.notify_one(); Ok(true) } } @@ -1363,14 +1439,49 @@ async fn edge_ws_close_during_registration_rolls_back_without_pool_commit() { } #[tokio::test] -async fn edge_ws_auth_ok_precedes_claim_release_wait() { - let registry = Arc::new(BlockingLeaseEdgeRegistry::new()); +async fn pending_claim_release_does_not_block_websocket_messages_or_disconnect() { + pending_release_disconnect(false, None).await; +} + +#[tokio::test] +async fn pending_claim_release_deadline_runs_while_heartbeat_waits_for_its_resource() { + pending_release_disconnect(true, None).await; +} + +#[tokio::test] +#[ignore = "requires live MatrixOne: ASTRA_TEST_DB_IT=1 and MATRIXONE_* settings"] +async fn pending_claim_release_disconnect_returns_live_database_connection() { + assert_eq!(std::env::var("ASTRA_TEST_DB_IT").as_deref(), Ok("1")); + let settings = astra_core::MatrixOneSettings::from_env(); + let pool = sqlx::mysql::MySqlPoolOptions::new() + .max_connections(1) + .acquire_timeout(std::time::Duration::from_secs(2)) + .connect_with( + sqlx::mysql::MySqlConnectOptions::new() + .host(&settings.host) + .port(settings.port) + .username(&settings.user) + .password(&settings.password) + .database(&settings.database), + ) + .await + .unwrap(); + pending_release_disconnect(false, Some(pool.clone())).await; + pool.close().await; +} + +async fn pending_release_disconnect(wait_for_heartbeat: bool, live_pool: Option) { + let mut registry = BlockingLeaseEdgeRegistry::new(); + registry.block_every_release = wait_for_heartbeat; + registry.live_pool = live_pool; + let registry = Arc::new(registry); let state = AppState::new( ServiceInfo::new("edge-auth-order-test", "0.0.0-test", ""), Arc::new(StubHealthChecker), ) .with_auth_service(Arc::new(StubAuthService)) .with_edge_registry_service(registry.clone()); + let observed_state = state.clone(); let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); let addr = listener.local_addr().unwrap(); let server = tokio::spawn(async move { @@ -1421,8 +1532,67 @@ async fn edge_ws_auth_ok_precedes_claim_release_wait() { let first: serde_json::Value = serde_json::from_str(&first).unwrap(); assert_eq!(first["type"], "edge_auth_ok"); - registry.claim_release_gate.notify_one(); - ws.close(None).await.ok(); + // Keep release_registration permanently pending. The application data + // plane must remain responsive while durable publication is reconciling. + ws.send(Message::Text( + json!({ "type": "edge_ping" }).to_string().into(), + )) + .await + .expect("send ping while release is pending"); + let pong = tokio::time::timeout(std::time::Duration::from_secs(2), ws.next()) + .await + .expect("pong timeout while release is pending") + .expect("pong frame while release is pending") + .expect("valid pong while release is pending"); + let pong: serde_json::Value = serde_json::from_str(&pong.into_text().unwrap()).unwrap(); + assert_eq!(pong["type"], "edge_pong"); + + if wait_for_heartbeat { + // Establish the socket with real time, then deterministically start a + // fresh pending publication attempt just before heartbeat is due. + tokio::time::pause(); + tokio::time::advance(std::time::Duration::from_secs(5)).await; + tokio::task::yield_now().await; + tokio::time::advance(std::time::Duration::from_secs(24)).await; + tokio::time::timeout( + std::time::Duration::from_millis(100), + registry.claim_release_started.notified(), + ) + .await + .expect("retry acquired the connection before heartbeat"); + tokio::time::advance(std::time::Duration::from_secs(1)).await; + tokio::time::timeout( + std::time::Duration::from_millis(100), + registry.heartbeat_started.notified(), + ) + .await + .expect("heartbeat started while publication holds the connection"); + tokio::time::advance(std::time::Duration::from_secs(5)).await; + tokio::time::timeout( + std::time::Duration::from_secs(1), + registry.heartbeat_completed.notified(), + ) + .await + .expect("heartbeat must acquire the resource released by publication timeout"); + tokio::time::resume(); + } + ws.close(None).await.expect("close websocket"); + tokio::time::timeout( + std::time::Duration::from_secs(2), + registry.unregister_completed.notified(), + ) + .await + .expect("durable unregister must complete after release cancellation"); + tokio::time::timeout(std::time::Duration::from_secs(2), async { + while observed_state + .edge_connection_pool + .has_connected_edge("test-user-1") + { + tokio::task::yield_now().await; + } + }) + .await + .expect("pending release future must be cancelled during disconnect cleanup"); server.abort(); } @@ -1485,7 +1655,7 @@ async fn claim_loss_after_pool_commit_removes_the_unpublished_connection() { serde_json::from_str(&second.into_text().unwrap()).unwrap(); assert!(matches!( second, - astra_server_types::edge_ws_protocol::EdgeServerMessage::AuthError { .. } + astra_server_types::edge_ws_protocol::EdgeServerMessage::Closing { .. } )); assert!( !observed_state @@ -1517,6 +1687,7 @@ impl astra_services::multi_agent::EdgeRegistryService for FailingEdgeRegistry { _user_id: &str, _edge_agent_id: &str, _edge_id_header: &str, + _registration_claim_id: Option<&str>, ) -> Result<(), astra_services::multi_agent::HeartbeatError> { Ok(()) } @@ -1594,6 +1765,91 @@ async fn edge_ws_rejects_connection_when_db_registration_fails() { server.abort(); } +#[tokio::test(flavor = "current_thread")] +async fn release_outcome_unknown_is_reconciled_while_the_connection_is_alive() { + let registry = Arc::new(BlockingLeaseEdgeRegistry::with_release_recovery()); + let state = AppState::new( + ServiceInfo::new("edge-release-recovery-test", "0.0.0-test", ""), + Arc::new(StubHealthChecker), + ) + .with_auth_service(Arc::new(StubAuthService)) + .with_edge_registry_service(registry.clone()); + let observed_state = state.clone(); + let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); + let addr = listener.local_addr().unwrap(); + let server = tokio::spawn(async move { + axum::serve(listener, astra_runtime::build_app(state)) + .await + .unwrap() + }); + + let (mut ws, _) = connect_async(ws_request(addr, "test-edge-token")) + .await + .expect("WS connect"); + ws.send(Message::Text( + json!({ + "type": "edge_auth", + "edge_agent_id": "edge-release-recovery", + "interaction_api_major": astra_server_types::AGENT_INTERACTION_API_MAJOR, + "hostname": "host", + "workspace_dir": "/workspace", + "capabilities": edge_capabilities("edge-release-recovery") + }) + .to_string() + .into(), + )) + .await + .unwrap(); + registry.registration_started.notified().await; + registry.release_registration.notify_one(); + + let auth = ws.next().await.unwrap().unwrap(); + let auth: astra_server_types::edge_ws_protocol::EdgeServerMessage = + serde_json::from_str(&auth.into_text().unwrap()).unwrap(); + assert!(matches!( + auth, + astra_server_types::edge_ws_protocol::EdgeServerMessage::AuthOk { .. } + )); + + ws.send(Message::Text( + json!({ "type": "edge_ping" }).to_string().into(), + )) + .await + .expect("send readiness ping"); + let pong = ws + .next() + .await + .expect("readiness pong frame") + .expect("pong"); + let pong: astra_server_types::edge_ws_protocol::EdgeServerMessage = + serde_json::from_str(&pong.into_text().unwrap()).unwrap(); + assert!(matches!( + pong, + astra_server_types::edge_ws_protocol::EdgeServerMessage::Pong { .. } + )); + + tokio::time::pause(); + tokio::time::advance(std::time::Duration::from_secs(2)).await; + tokio::time::timeout(std::time::Duration::from_secs(1), async { + while registry.release_attempts.load(Ordering::SeqCst) < 2 { + tokio::task::yield_now().await; + } + }) + .await + .expect("durable publication retry"); + tokio::time::resume(); + + assert_eq!(registry.release_attempts.load(Ordering::SeqCst), 2); + assert!( + observed_state + .edge_connection_pool + .has_connected_edge("test-user-1"), + "a transient release failure must converge without discarding the healthy socket" + ); + ws.close(None).await.ok(); + server.abort(); +} + // ── B3: heartbeat tick updates DB registry ──────────────────────────── #[derive(Default)] @@ -1634,6 +1890,7 @@ impl astra_services::multi_agent::EdgeRegistryService for RecordingEdgeRegistry user_id: &str, edge_agent_id: &str, edge_id_header: &str, + _registration_claim_id: Option<&str>, ) -> Result<(), astra_services::multi_agent::HeartbeatError> { self.heartbeats.lock().unwrap().push(( user_id.to_string(), diff --git a/crates/services/src/multi_agent/edge_registry.rs b/crates/services/src/multi_agent/edge_registry.rs index ddb9b6d020..0f00bb5276 100644 --- a/crates/services/src/multi_agent/edge_registry.rs +++ b/crates/services/src/multi_agent/edge_registry.rs @@ -10,6 +10,9 @@ use serde::{Deserialize, Serialize}; use super::metrics::SharedMultiAgentMetrics; use crate::db_row::RowExt as EdgeRegistryDbRow; +const CURRENT_READ_MAX_ATTEMPTS: u32 = 6; +const CURRENT_READ_BASE_BACKOFF_MS: u64 = 25; + // MatrixOne exposes JSON columns with SQL type JSON, while RowExt intentionally // decodes this optional payload as text before serde_json validation. Wrapping // the parsed value in a one-element JSON array lets JSON_UNQUOTE produce @@ -26,6 +29,58 @@ const EDGE_AGENT_RECORD_COLUMNS: &str = "registry_id, user_id, edge_agent_id, ed CAST(registered_at AS CHAR) AS registered_at, \ CAST(last_heartbeat_at AS CHAR) AS last_heartbeat_at"; +// One generation-scoped cleanup statement covers both possible ownership +// positions during a reconnect: +// - `edge_id = target`: the target is the current/predecessor generation. Its +// private metadata is scrubbed and it becomes inactive. A successor claim is +// preserved while state is 0/1 so setup can still complete. +// - `registration_previous_edge_id = target`: the successor is finalized but +// unpublished. Clearing only the predecessor marker records that rollback +// must not resurrect the disconnected target. +// Assignment order is intentional because MySQL-compatible engines evaluate +// single-table UPDATE assignments from left to right. +const DEACTIVATE_EDGE_GENERATION_SQL: &str = "UPDATE edge_agent_registry \ + SET registration_claim_id = CASE \ + WHEN (registration_state IN (0, 1) AND edge_id = ?) \ + OR (registration_state = 2 AND registration_previous_edge_id = ?) \ + THEN registration_claim_id ELSE NULL END, \ + registration_claim_expires_at = CASE \ + WHEN (registration_state IN (0, 1) AND edge_id = ?) \ + OR (registration_state = 2 AND registration_previous_edge_id = ?) \ + THEN registration_claim_expires_at ELSE NULL END, \ + hostname = CASE WHEN edge_id = ? THEN NULL ELSE hostname END, \ + worktree_path = CASE WHEN edge_id = ? THEN NULL ELSE worktree_path END, \ + capabilities_json = CASE WHEN edge_id = ? THEN NULL ELSE capabilities_json END, \ + workspace_id = CASE WHEN edge_id = ? THEN NULL ELSE workspace_id END, \ + registration_state = CASE \ + WHEN registration_state = 2 AND registration_previous_edge_id = ? \ + THEN registration_state ELSE 0 END, \ + registration_previous_edge_id = NULL \ + WHERE user_id = ? AND edge_agent_id = ? \ + AND (edge_id = ? \ + OR (registration_state = 2 AND registration_previous_edge_id = ?))"; + +fn deactivate_edge_generation_query<'q>( + user_id: &'q str, + edge_agent_id: &'q str, + edge_id: &'q str, +) -> sqlx::query::Query<'q, sqlx::MySql, sqlx::mysql::MySqlArguments> { + sqlx::query(DEACTIVATE_EDGE_GENERATION_SQL) + .bind(edge_id) + .bind(edge_id) + .bind(edge_id) + .bind(edge_id) + .bind(edge_id) + .bind(edge_id) + .bind(edge_id) + .bind(edge_id) + .bind(edge_id) + .bind(user_id) + .bind(edge_agent_id) + .bind(edge_id) + .bind(edge_id) +} + fn serialize_edge_capabilities( capabilities: Option<&serde_json::Value>, ) -> Result, String> { @@ -65,6 +120,218 @@ pub struct EdgeRegistrationLease { pub claim_id: Option, } +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +enum RegistrationTransition { + Finalize, + Release, + Rollback, +} + +impl RegistrationTransition { + fn operation(self) -> &'static str { + match self { + Self::Finalize => "finalize", + Self::Release => "release", + Self::Rollback => "rollback", + } + } +} + +#[derive(Clone, Debug, PartialEq, Eq)] +struct RegistrationState { + edge_id: String, + claim_id: Option, + state: i8, + previous_edge_id: Option, +} + +#[derive(Clone, Debug, PartialEq, Eq)] +struct RegistrationGenerationState { + edge_id: String, + claim_id: Option, + previous_edge_id: Option, + state: i8, +} + +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +enum RegistrationTransitionDecision { + AlreadyApplied, + Apply, + Superseded, +} + +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +enum RegistrationAttemptDecision { + Retry, + Transition(RegistrationTransitionDecision), + OutcomeUnknown, +} + +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +enum GenerationMutation { + Heartbeat, + Unregister, +} + +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +enum GenerationMutationOutcome { + Applied, + AlreadyApplied, + Superseded, +} + +fn heartbeat_generation_is_owned( + row: &RegistrationGenerationState, + edge_id: &str, + claim_id: Option<&str>, +) -> bool { + (row.state == 1 && row.edge_id == edge_id) + || (row.state == 2 + && (row.previous_edge_id.as_deref() == Some(edge_id) + || (row.edge_id == edge_id + && claim_id.is_some() + && row.claim_id.as_deref() == claim_id))) +} + +fn unregister_generation_is_owned(row: &RegistrationGenerationState, edge_id: &str) -> bool { + row.edge_id == edge_id || (row.state == 2 && row.previous_edge_id.as_deref() == Some(edge_id)) +} + +fn unregister_generation_is_already_applied( + row: &RegistrationGenerationState, + edge_id: &str, +) -> bool { + row.edge_id == edge_id && row.state == 0 +} + +fn heartbeat_generation_query<'q>( + user_id: &'q str, + edge_agent_id: &'q str, + edge_id: &'q str, + claim_id: Option<&'q str>, +) -> sqlx::query::Query<'q, sqlx::MySql, sqlx::mysql::MySqlArguments> { + sqlx::query( + "UPDATE edge_agent_registry \ + SET last_heartbeat_at = NOW(6), \ + registration_claim_expires_at = CASE \ + WHEN registration_state = 2 AND edge_id = ? \ + AND registration_claim_id = ? \ + THEN DATE_ADD(NOW(6), INTERVAL 120 SECOND) \ + ELSE registration_claim_expires_at END \ + WHERE user_id = ? AND edge_agent_id = ? \ + AND ((registration_state = 1 AND edge_id = ?) \ + OR (registration_state = 2 AND registration_previous_edge_id = ?) \ + OR (registration_state = 2 AND edge_id = ? \ + AND registration_claim_id = ?))", + ) + .bind(edge_id) + .bind(claim_id) + .bind(user_id) + .bind(edge_agent_id) + .bind(edge_id) + .bind(edge_id) + .bind(edge_id) + .bind(claim_id) +} + +fn registration_predecessor_is_live(row: &RegistrationState, predecessor_edge_id: &str) -> bool { + (row.state == 1 && row.edge_id == predecessor_edge_id) + || (row.state == 2 && row.previous_edge_id.as_deref() == Some(predecessor_edge_id)) +} + +fn registration_transition_decision( + transition: RegistrationTransition, + lease: &EdgeRegistrationLease, + claim_id: &str, + row: Option<&RegistrationState>, +) -> RegistrationTransitionDecision { + match transition { + RegistrationTransition::Finalize => match row { + Some(row) + if row.edge_id == lease.current.edge_id + && row.claim_id.as_deref() == Some(claim_id) + && row.state == 2 => + { + RegistrationTransitionDecision::AlreadyApplied + } + Some(row) if row.claim_id.as_deref() == Some(claim_id) => { + RegistrationTransitionDecision::Apply + } + Some(row) if row.claim_id.as_deref() != Some(claim_id) => { + RegistrationTransitionDecision::Superseded + } + None | Some(_) => RegistrationTransitionDecision::Superseded, + }, + RegistrationTransition::Release => match row { + Some(row) + if row.edge_id == lease.current.edge_id + && row.claim_id.is_none() + && row.state == 1 => + { + RegistrationTransitionDecision::AlreadyApplied + } + Some(row) if row.claim_id.as_deref() == Some(claim_id) => { + RegistrationTransitionDecision::Apply + } + Some(row) if row.claim_id.as_deref() != Some(claim_id) => { + RegistrationTransitionDecision::Superseded + } + None | Some(_) => RegistrationTransitionDecision::Superseded, + }, + RegistrationTransition::Rollback => match (&lease.previous, row) { + (_, Some(row)) + if row.edge_id == lease.current.edge_id + && row.claim_id.is_none() + && row.state == 0 => + { + RegistrationTransitionDecision::AlreadyApplied + } + (Some(previous), Some(row)) + if row.edge_id == previous.edge_id && row.claim_id.is_none() && row.state == 1 => + { + RegistrationTransitionDecision::AlreadyApplied + } + (_, Some(row)) if row.claim_id.as_deref() == Some(claim_id) => { + RegistrationTransitionDecision::Apply + } + (_, Some(_)) | (_, None) => RegistrationTransitionDecision::Superseded, + }, + } +} + +fn registration_before_mutation_decision( + transition: RegistrationTransition, + lease: &EdgeRegistrationLease, + claim_id: &str, + row: Option<&RegistrationState>, + attempt: u32, + max_attempts: u32, +) -> RegistrationAttemptDecision { + if row.is_none() { + if attempt + 1 < max_attempts { + return RegistrationAttemptDecision::Retry; + } + return RegistrationAttemptDecision::OutcomeUnknown; + } + RegistrationAttemptDecision::Transition(registration_transition_decision( + transition, lease, claim_id, row, + )) +} + +async fn transition_error_with_rollback( + transaction: sqlx::Transaction<'_, sqlx::MySql>, + operation: &str, + registry_id: &str, + error: impl std::fmt::Display, +) -> String { + match transaction.rollback().await { + Ok(()) => format!("edge_registry {operation} registration {registry_id}: {error}"), + Err(rollback_error) => format!( + "edge_registry {operation} registration {registry_id}: {error}; transaction rollback failed: {rollback_error}" + ), + } +} + /// Structured error for `EdgeRegistryService::heartbeat`. /// /// Callers must treat these two variants differently: @@ -138,7 +405,10 @@ pub trait EdgeRegistryService: Send + Sync { } /// Undo a claimed generation only if it still owns the registry row. - /// Returns false when a newer generation has already taken over. + /// Durable claiming backends return true when rollback is applied or was + /// already applied, false only after verifying another generation, and an + /// error when ownership cannot be established. Non-claiming backends use + /// their ordinary generation-scoped cleanup result. async fn rollback_registration(&self, lease: &EdgeRegistrationLease) -> Result { match &lease.previous { Some(previous) => { @@ -162,23 +432,33 @@ pub trait EdgeRegistryService: Send + Sync { } /// Release the cross-pod setup claim after the connection is published. - /// Backends without durable claim support have nothing to release. + /// Durable claiming backends return true when release is applied or was + /// already applied, false only after verifying another generation, and an + /// error when ownership cannot be established. Backends without durable + /// claim support have nothing to release. async fn release_registration(&self, _lease: &EdgeRegistrationLease) -> Result { Ok(false) } /// Finalize the durable generation while retaining the cross-pod claim. - /// The default registration path is already final, so non-claiming - /// backends have nothing else to do. + /// Durable claiming backends return true when finalization is applied or + /// was already applied, false only after verifying another generation, and + /// an error when ownership cannot be established. The default registration + /// path is already final, so non-claiming backends have nothing else to do. async fn finalize_registration(&self, _lease: &EdgeRegistrationLease) -> Result { Ok(true) } + /// Refresh this exact connection generation. While a finalized durable + /// claim is awaiting release, `registration_claim_id` also fences and + /// renews that claim; ordinary published/non-durable registrations pass + /// `None`. async fn heartbeat( &self, user_id: &str, edge_agent_id: &str, edge_id_header: &str, + registration_claim_id: Option<&str>, ) -> Result<(), HeartbeatError>; /// Find the most-recently-active registry record for a given edge_agent_id, @@ -202,7 +482,11 @@ pub trait EdgeRegistryService: Send + Sync { async fn list_by_user(&self, user_id: &str) -> Result, String>; /// Remove only the exact connection incarnation registered by this socket. - /// Returns false when a newer connection already replaced it. + /// Durable backends retain an inactive, non-routable owner row so a later + /// retry has authoritative idempotence evidence; they return false only + /// after verifying that a newer connection replaced it. Non-durable + /// backends may return false because there is no persistent generation to + /// deactivate. async fn unregister_generation( &self, user_id: &str, @@ -268,6 +552,499 @@ impl DatabaseEdgeRegistryService { self } + async fn load_registration_state_for_update( + transaction: &mut sqlx::Transaction<'_, sqlx::MySql>, + user_id: &str, + registry_id: &str, + ) -> Result, sqlx::Error> { + let row: Option<(String, Option, i8, Option)> = sqlx::query_as( + "SELECT edge_id, registration_claim_id, registration_state, \ + registration_previous_edge_id \ + FROM edge_agent_registry WHERE user_id = ? AND registry_id = ? FOR UPDATE", + ) + .bind(user_id) + .bind(registry_id) + .fetch_optional(&mut **transaction) + .await?; + Ok(row.map( + |(edge_id, claim_id, state, previous_edge_id)| RegistrationState { + edge_id, + claim_id, + state, + previous_edge_id, + }, + )) + } + + /// Establish a write-write conflict boundary before trusting a registry + /// observation. MatrixOne optimistic transactions do not acquire a current + /// row lock for `SELECT ... FOR UPDATE`, while a no-op UPDATE participates + /// in commit-time conflict detection. + async fn establish_registry_current_read( + transaction: &mut sqlx::Transaction<'_, sqlx::MySql>, + user_id: &str, + registry_id: &str, + ) -> Result<(), sqlx::Error> { + sqlx::query( + "UPDATE edge_agent_registry SET last_heartbeat_at = last_heartbeat_at \ + WHERE user_id = ? AND registry_id = ?", + ) + .bind(user_id) + .bind(registry_id) + .execute(&mut **transaction) + .await?; + Ok(()) + } + + async fn establish_generation_current_read( + transaction: &mut sqlx::Transaction<'_, sqlx::MySql>, + user_id: &str, + edge_agent_id: &str, + ) -> Result<(), sqlx::Error> { + sqlx::query( + "UPDATE edge_agent_registry SET last_heartbeat_at = last_heartbeat_at \ + WHERE user_id = ? AND edge_agent_id = ?", + ) + .bind(user_id) + .bind(edge_agent_id) + .execute(&mut **transaction) + .await?; + Ok(()) + } + + async fn load_generation_state( + transaction: &mut sqlx::Transaction<'_, sqlx::MySql>, + user_id: &str, + edge_agent_id: &str, + ) -> Result, sqlx::Error> { + let row: Option<(String, Option, Option, i8)> = sqlx::query_as( + "SELECT edge_id, registration_claim_id, registration_previous_edge_id, registration_state \ + FROM edge_agent_registry WHERE user_id = ? AND edge_agent_id = ?", + ) + .bind(user_id) + .bind(edge_agent_id) + .fetch_optional(&mut **transaction) + .await?; + Ok(row.map( + |(edge_id, claim_id, previous_edge_id, state)| RegistrationGenerationState { + edge_id, + claim_id, + previous_edge_id, + state, + }, + )) + } + + async fn settle_generation_mutation_after_miss( + &self, + user_id: &str, + edge_agent_id: &str, + edge_id: &str, + registration_claim_id: Option<&str>, + mutation: GenerationMutation, + ) -> Result { + let operation = match mutation { + GenerationMutation::Heartbeat => "heartbeat", + GenerationMutation::Unregister => "unregister", + }; + for attempt in 0..CURRENT_READ_MAX_ATTEMPTS { + if attempt > 0 { + if let Some(ref metrics) = self.metrics { + metrics.registry_retry_total.fetch_add(1, Ordering::Relaxed); + } + tokio::time::sleep(std::time::Duration::from_millis( + CURRENT_READ_BASE_BACKOFF_MS * (1 << (attempt - 1)), + )) + .await; + } + + let mut transaction = self.pool.begin().await.map_err(|error| { + format!("edge_registry {operation} verification begin (attempt {attempt}): {error}") + })?; + if let Err(error) = + Self::establish_generation_current_read(&mut transaction, user_id, edge_agent_id) + .await + { + return Err(transition_error_with_rollback( + transaction, + operation, + edge_agent_id, + format!("current-read barrier failed on attempt {attempt}: {error}"), + ) + .await); + } + let row = + match Self::load_generation_state(&mut transaction, user_id, edge_agent_id).await { + Ok(row) => row, + Err(error) => { + return Err(transition_error_with_rollback( + transaction, + operation, + edge_agent_id, + format!("ownership lookup failed on attempt {attempt}: {error}"), + ) + .await); + } + }; + + let Some(row) = row else { + if attempt + 1 < CURRENT_READ_MAX_ATTEMPTS { + transaction.rollback().await.map_err(|error| { + format!( + "edge_registry {operation} retry rollback (attempt {attempt}): {error}" + ) + })?; + continue; + } + transaction.rollback().await.map_err(|error| { + format!( + "edge_registry {operation} absent-state rollback (attempt {attempt}): {error}" + ) + })?; + return Err(format!( + "edge_registry generation {operation} remained absent after {CURRENT_READ_MAX_ATTEMPTS} current-read attempts" + )); + }; + if mutation == GenerationMutation::Unregister + && unregister_generation_is_already_applied(&row, edge_id) + { + transaction.commit().await.map_err(|error| { + format!("edge_registry {operation} idempotent verification commit: {error}") + })?; + return Ok(GenerationMutationOutcome::AlreadyApplied); + } + let owned = match mutation { + GenerationMutation::Heartbeat => { + heartbeat_generation_is_owned(&row, edge_id, registration_claim_id) + } + GenerationMutation::Unregister => unregister_generation_is_owned(&row, edge_id), + }; + if !owned { + transaction.commit().await.map_err(|error| { + format!("edge_registry {operation} supersession verification commit: {error}") + })?; + return Ok(GenerationMutationOutcome::Superseded); + } + + let affected = match mutation { + GenerationMutation::Heartbeat => { + heartbeat_generation_query( + user_id, + edge_agent_id, + edge_id, + registration_claim_id, + ) + .execute(&mut *transaction) + .await + } + GenerationMutation::Unregister => { + deactivate_edge_generation_query(user_id, edge_agent_id, edge_id) + .execute(&mut *transaction) + .await + } + } + .map_err(|error| { + format!("edge_registry {operation} retry mutation (attempt {attempt}): {error}") + })? + .rows_affected(); + if affected == 0 { + transaction.rollback().await.map_err(|error| { + format!("edge_registry {operation} retry rollback (attempt {attempt}): {error}") + })?; + continue; + } + transaction.commit().await.map_err(|error| { + format!( + "edge_registry {operation} verification commit (attempt {attempt}): {error}" + ) + })?; + return Ok(GenerationMutationOutcome::Applied); + } + + Err(format!( + "edge_registry generation {operation} remained ambiguous after {CURRENT_READ_MAX_ATTEMPTS} current-read attempts" + )) + } + + async fn settle_registration_transition( + &self, + lease: &EdgeRegistrationLease, + transition: RegistrationTransition, + ) -> Result { + let operation = transition.operation(); + let registry_id = lease.current.registry_id.as_str(); + let claim_id = lease.claim_id.as_deref().ok_or_else(|| { + format!("edge_registry {operation} registration {registry_id} has no durable claim") + })?; + let current_capabilities = + serialize_edge_capabilities(lease.current.capabilities.as_ref())?; + let previous_capabilities = match &lease.previous { + Some(previous) => serialize_edge_capabilities(previous.capabilities.as_ref())?, + None => None, + }; + + for attempt in 0..CURRENT_READ_MAX_ATTEMPTS { + if attempt > 0 { + if let Some(ref metrics) = self.metrics { + metrics.registry_retry_total.fetch_add(1, Ordering::Relaxed); + } + tokio::time::sleep(std::time::Duration::from_millis( + CURRENT_READ_BASE_BACKOFF_MS * (1 << (attempt - 1)), + )) + .await; + } + + let mut transaction = self.pool.begin().await.map_err(|error| { + format!( + "edge_registry {operation} registration {registry_id} begin (attempt {attempt}): {error}" + ) + })?; + if let Err(error) = Self::establish_registry_current_read( + &mut transaction, + &lease.current.user_id, + registry_id, + ) + .await + { + return Err(transition_error_with_rollback( + transaction, + operation, + registry_id, + format!("current-read barrier failed on attempt {attempt}: {error}"), + ) + .await); + } + let before = match Self::load_registration_state_for_update( + &mut transaction, + &lease.current.user_id, + registry_id, + ) + .await + { + Ok(row) => row, + Err(error) => { + return Err(transition_error_with_rollback( + transaction, + operation, + registry_id, + format!("state lookup failed on attempt {attempt}: {error}"), + ) + .await); + } + }; + + let before_decision = registration_before_mutation_decision( + transition, + lease, + claim_id, + before.as_ref(), + attempt, + CURRENT_READ_MAX_ATTEMPTS, + ); + match before_decision { + RegistrationAttemptDecision::Retry => { + transaction.rollback().await.map_err(|error| { + format!( + "edge_registry {operation} registration {registry_id} retry rollback (attempt {attempt}): {error}" + ) + })?; + continue; + } + RegistrationAttemptDecision::OutcomeUnknown => { + transaction.rollback().await.map_err(|error| { + format!( + "edge_registry {operation} registration {registry_id} absent-state rollback: {error}" + ) + })?; + return Err(format!( + "edge_registry {operation} registration {registry_id} remained absent after {CURRENT_READ_MAX_ATTEMPTS} current-read attempts" + )); + } + RegistrationAttemptDecision::Transition( + RegistrationTransitionDecision::AlreadyApplied, + ) => { + transaction.commit().await.map_err(|error| { + format!( + "edge_registry {operation} registration {registry_id} verification commit: {error}" + ) + })?; + return Ok(true); + } + RegistrationAttemptDecision::Transition( + RegistrationTransitionDecision::Superseded, + ) => { + // Commit the no-op write barrier before reporting a newer + // owner. On an optimistic snapshot, a stale predecessor + // observation conflicts here instead of becoming a false + // supersession result. + transaction.commit().await.map_err(|error| { + format!( + "edge_registry {operation} registration {registry_id} superseded verification commit: {error}" + ) + })?; + return Ok(false); + } + RegistrationAttemptDecision::Transition(RegistrationTransitionDecision::Apply) => {} + } + + let owned_before = before + .as_ref() + .expect("an applicable registration transition requires an observed owner row"); + let live_previous = lease.previous.as_ref().filter(|previous| { + registration_predecessor_is_live(owned_before, &previous.edge_id) + }); + + let execution = match transition { + RegistrationTransition::Finalize => { + sqlx::query( + "UPDATE edge_agent_registry \ + SET edge_id = ?, hostname = ?, worktree_path = ?, capabilities_json = ?, \ + workspace_id = ?, last_heartbeat_at = NOW(6), registration_state = 2, \ + registration_previous_edge_id = ? \ + WHERE user_id = ? AND registry_id = ? AND registration_claim_id = ?", + ) + .bind(&lease.current.edge_id) + .bind(&lease.current.hostname) + .bind(&lease.current.worktree_path) + .bind(¤t_capabilities) + .bind(&lease.current.workspace_id) + .bind(live_previous.map(|previous| previous.edge_id.as_str())) + .bind(&lease.current.user_id) + .bind(registry_id) + .bind(claim_id) + .execute(&mut *transaction) + .await + } + RegistrationTransition::Release => { + sqlx::query( + "UPDATE edge_agent_registry \ + SET registration_claim_id = NULL, registration_claim_expires_at = NULL, \ + registration_state = 1, registration_previous_edge_id = NULL \ + WHERE user_id = ? AND registry_id = ? AND edge_id = ? \ + AND registration_claim_id = ? AND registration_state = 2", + ) + .bind(&lease.current.user_id) + .bind(registry_id) + .bind(&lease.current.edge_id) + .bind(claim_id) + .execute(&mut *transaction) + .await + } + RegistrationTransition::Rollback => match live_previous { + Some(previous) => { + sqlx::query( + "UPDATE edge_agent_registry \ + SET edge_id = ?, hostname = ?, worktree_path = ?, capabilities_json = ?, \ + workspace_id = ?, last_heartbeat_at = NOW(6), \ + registration_claim_id = NULL, registration_claim_expires_at = NULL, \ + registration_state = 1, registration_previous_edge_id = NULL \ + WHERE user_id = ? AND registry_id = ? AND registration_claim_id = ?", + ) + .bind(&previous.edge_id) + .bind(&previous.hostname) + .bind(&previous.worktree_path) + .bind(&previous_capabilities) + .bind(&previous.workspace_id) + .bind(&lease.current.user_id) + .bind(registry_id) + .bind(claim_id) + .execute(&mut *transaction) + .await + } + None => { + // A first registration, or a replacement whose + // predecessor disconnected during setup, rolls back to + // a skeletal inactive owner. It is excluded from + // routing but remains durable settlement evidence. + sqlx::query( + "UPDATE edge_agent_registry \ + SET edge_id = ?, hostname = NULL, worktree_path = NULL, \ + capabilities_json = NULL, workspace_id = NULL, \ + last_heartbeat_at = NOW(6), \ + registration_claim_id = NULL, registration_claim_expires_at = NULL, \ + registration_state = 0, registration_previous_edge_id = NULL \ + WHERE user_id = ? AND registry_id = ? AND registration_claim_id = ?", + ) + .bind(&lease.current.edge_id) + .bind(&lease.current.user_id) + .bind(registry_id) + .bind(claim_id) + .execute(&mut *transaction) + .await + } + }, + }; + if let Err(error) = execution { + return Err(transition_error_with_rollback( + transaction, + operation, + registry_id, + format!("mutation failed on attempt {attempt}: {error}"), + ) + .await); + } + + let after = match Self::load_registration_state_for_update( + &mut transaction, + &lease.current.user_id, + registry_id, + ) + .await + { + Ok(row) => row, + Err(error) => { + return Err(transition_error_with_rollback( + transaction, + operation, + registry_id, + format!("state verification failed on attempt {attempt}: {error}"), + ) + .await); + } + }; + if after.is_none() { + return Err(transition_error_with_rollback( + transaction, + operation, + registry_id, + format!("mutation produced an unexpected absent state on attempt {attempt}"), + ) + .await); + } + let after_decision = + registration_transition_decision(transition, lease, claim_id, after.as_ref()); + match after_decision { + RegistrationTransitionDecision::AlreadyApplied => { + transaction.commit().await.map_err(|error| { + format!( + "edge_registry {operation} registration {registry_id} commit (attempt {attempt}): {error}" + ) + })?; + return Ok(true); + } + RegistrationTransitionDecision::Superseded => { + transaction.commit().await.map_err(|error| { + format!( + "edge_registry {operation} registration {registry_id} superseded verification commit (attempt {attempt}): {error}" + ) + })?; + return Ok(false); + } + RegistrationTransitionDecision::Apply => { + transaction.rollback().await.map_err(|error| { + format!( + "edge_registry {operation} registration {registry_id} retry rollback (attempt {attempt}): {error}" + ) + })?; + } + } + } + + Err(format!( + "edge_registry {operation} registration {registry_id} remained owned but did not reach its target state after {CURRENT_READ_MAX_ATTEMPTS} attempts" + )) + } + #[allow(clippy::too_many_arguments)] async fn claim_registration( &self, @@ -279,7 +1056,7 @@ impl DatabaseEdgeRegistryService { capabilities: Option, workspace_id: Option<&str>, ) -> Result { - let cap_json = serialize_edge_capabilities(capabilities.as_ref())?; + serialize_edge_capabilities(capabilities.as_ref())?; const MAX_RETRIES: u32 = 5; let claim_id = uuid::Uuid::new_v4().to_string(); @@ -322,10 +1099,22 @@ impl DatabaseEdgeRegistryService { // unchanged until finalize_registration(), so the published // predecessor remains heartbeatable and routable while setup is // pending. + // Taking an expired finalized claim abandons its unpublished + // owner and predecessor atomically. Keep the old edge_id for + // fenced cleanup, but use state 0 so that cleanup preserves the + // new setup claim instead of clearing it as the owner's claim. let updated = sqlx::query( "UPDATE edge_agent_registry \ SET registration_claim_id = ?, \ - registration_claim_expires_at = DATE_ADD(NOW(6), INTERVAL 120 SECOND) \ + registration_claim_expires_at = DATE_ADD(NOW(6), INTERVAL 120 SECOND), \ + hostname = CASE WHEN registration_state = 1 THEN hostname ELSE NULL END, \ + worktree_path = CASE WHEN registration_state = 1 THEN worktree_path ELSE NULL END, \ + capabilities_json = CASE WHEN registration_state = 1 THEN capabilities_json ELSE NULL END, \ + workspace_id = CASE WHEN registration_state = 1 THEN workspace_id ELSE NULL END, \ + registration_previous_edge_id = CASE WHEN registration_state = 2 \ + THEN NULL ELSE registration_previous_edge_id END, \ + registration_state = CASE WHEN registration_state = 2 \ + THEN 0 ELSE registration_state END \ WHERE user_id = ? AND registry_id = ? AND edge_id = ? \ AND (registration_claim_id IS NULL \ OR registration_claim_expires_at < NOW(6))", @@ -348,9 +1137,10 @@ impl DatabaseEdgeRegistryService { let now = chrono::Utc::now() .format("%Y-%m-%d %H:%M:%S%.6f") .to_string(); - // State 1 is the only published state. State 0 is a never- - // published insert and state 2 is a finalized generation whose - // owner crashed before releasing its claim; neither is safe to + // State 1 is the only published state. State 0 is an inactive + // owner (either never published or disconnected while its + // successor holds the claim), and state 2 is a finalized + // generation whose claim is not released; neither is safe to // resurrect as a rollback target. let published_previous = (registration_state == 1).then_some(previous.clone()); let current = EdgeAgentRecord { @@ -384,20 +1174,15 @@ impl DatabaseEdgeRegistryService { let registry_id = uuid::Uuid::new_v4().to_string(); let inserted = sqlx::query( "INSERT INTO edge_agent_registry \ - (registry_id, user_id, edge_agent_id, edge_id, hostname, worktree_path, \ - capabilities_json, workspace_id, registered_at, last_heartbeat_at, \ + (registry_id, user_id, edge_agent_id, edge_id, registered_at, last_heartbeat_at, \ registration_claim_id, registration_claim_expires_at, registration_state) \ - VALUES (?, ?, ?, ?, ?, ?, ?, ?, NOW(6), NOW(6), ?, \ + VALUES (?, ?, ?, ?, NOW(6), NOW(6), ?, \ DATE_ADD(NOW(6), INTERVAL 120 SECOND), 0)", ) .bind(®istry_id) .bind(user_id) .bind(edge_agent_id) .bind(edge_id_header) - .bind(hostname) - .bind(worktree_path) - .bind(&cap_json) - .bind(workspace_id) .bind(&claim_id) .execute(&mut *transaction) .await; @@ -726,35 +1511,46 @@ impl EdgeRegistryService for DatabaseEdgeRegistryService { user_id: &str, edge_agent_id: &str, edge_id_header: &str, + registration_claim_id: Option<&str>, ) -> Result<(), HeartbeatError> { // Guard on edge_id so a stale connection cannot refresh (or resurrect) // the row after a newer connection has replaced it. register_or_update // already set edge_id to the current connection's value, so we only // touch last_heartbeat_at and never rewrite edge_id here. If a newer - // connection has taken over (edge_id differs), this matches 0 rows and - // the stale connection's heartbeat correctly returns Superseded. - let n = sqlx::query( - "UPDATE edge_agent_registry SET last_heartbeat_at = NOW(6) \ - WHERE user_id = ? AND edge_agent_id = ? \ - AND ((registration_state = 1 AND edge_id = ?) \ - OR (registration_state = 2 AND registration_previous_edge_id = ?))", + // connection has taken over (edge_id differs), this matches 0 rows. A + // finalized current generation also remains healthy if releasing its + // durable claim had an outcome-unknown storage failure. + let n = heartbeat_generation_query( + user_id, + edge_agent_id, + edge_id_header, + registration_claim_id, ) - .bind(user_id) - .bind(edge_agent_id) - .bind(edge_id_header) - .bind(edge_id_header) .execute(&self.pool) .await .map_err(|e| HeartbeatError::StorageFailure(format!("edge heartbeat: {e}")))? .rows_affected(); if n == 0 { - // The row is gone or belongs to a newer connection — this connection - // has been superseded and must not keep the DB entry alive. tracing::warn!( edge_id = %edge_id_header, - "edge_registry: heartbeat matched no row (unregistered or superseded by newer connection)" + "edge_registry: heartbeat matched no row; verifying durable generation ownership" ); - return Err(HeartbeatError::Superseded); + return match self + .settle_generation_mutation_after_miss( + user_id, + edge_agent_id, + edge_id_header, + registration_claim_id, + GenerationMutation::Heartbeat, + ) + .await + .map_err(HeartbeatError::StorageFailure)? + { + GenerationMutationOutcome::Applied | GenerationMutationOutcome::AlreadyApplied => { + Ok(()) + } + GenerationMutationOutcome::Superseded => Err(HeartbeatError::Superseded), + }; } Ok(()) } @@ -765,17 +1561,28 @@ impl EdgeRegistryService for DatabaseEdgeRegistryService { edge_agent_id: &str, edge_id_header: &str, ) -> Result { - let deleted = sqlx::query( - "DELETE FROM edge_agent_registry \ - WHERE user_id = ? AND edge_agent_id = ? AND edge_id = ?", - ) - .bind(user_id) - .bind(edge_agent_id) - .bind(edge_id_header) - .execute(&self.pool) - .await - .map_err(|e| format!("edge_registry unregister: {e}"))?; - Ok(deleted.rows_affected() > 0) + let deactivated = deactivate_edge_generation_query(user_id, edge_agent_id, edge_id_header) + .execute(&self.pool) + .await + .map_err(|e| format!("edge_registry unregister: {e}"))?; + if deactivated.rows_affected() > 0 { + return Ok(true); + } + match self + .settle_generation_mutation_after_miss( + user_id, + edge_agent_id, + edge_id_header, + None, + GenerationMutation::Unregister, + ) + .await? + { + GenerationMutationOutcome::Applied | GenerationMutationOutcome::AlreadyApplied => { + Ok(true) + } + GenerationMutationOutcome::Superseded => Ok(false), + } } async fn rollback_registration(&self, lease: &EdgeRegistrationLease) -> Result { @@ -788,104 +1595,27 @@ impl EdgeRegistryService for DatabaseEdgeRegistryService { ) .await; }; - let Some(previous) = &lease.previous else { - let deleted = sqlx::query( - "DELETE FROM edge_agent_registry \ - WHERE user_id = ? AND registry_id = ? AND registration_claim_id = ?", - ) - .bind(&lease.current.user_id) - .bind(&lease.current.registry_id) - .bind(claim_id) - .execute(&self.pool) + debug_assert!(!claim_id.is_empty()); + self.settle_registration_transition(lease, RegistrationTransition::Rollback) .await - .map_err(|e| format!("edge_registry rollback inserted registration: {e}"))?; - return Ok(deleted.rows_affected() > 0); - }; - let capabilities_json = previous - .capabilities - .as_ref() - .map(serde_json::to_string) - .transpose() - .map_err(|e| format!("edge_registry rollback capabilities json: {e}"))?; - let restored = sqlx::query( - "UPDATE edge_agent_registry \ - SET edge_id = ?, hostname = ?, worktree_path = ?, capabilities_json = ?, \ - workspace_id = ?, last_heartbeat_at = NOW(6), \ - registration_claim_id = NULL, registration_claim_expires_at = NULL, \ - registration_state = 1, registration_previous_edge_id = NULL \ - WHERE user_id = ? AND registry_id = ? AND registration_claim_id = ?", - ) - .bind(&previous.edge_id) - .bind(&previous.hostname) - .bind(&previous.worktree_path) - .bind(&capabilities_json) - .bind(&previous.workspace_id) - .bind(&lease.current.user_id) - .bind(&lease.current.registry_id) - .bind(claim_id) - .execute(&self.pool) - .await - .map_err(|e| format!("edge_registry rollback registration: {e}"))?; - Ok(restored.rows_affected() > 0) } async fn finalize_registration(&self, lease: &EdgeRegistrationLease) -> Result { let Some(claim_id) = lease.claim_id.as_deref() else { return Ok(true); }; - let capabilities_json = lease - .current - .capabilities - .as_ref() - .map(serde_json::to_string) - .transpose() - .map_err(|e| format!("edge_registry finalize capabilities json: {e}"))?; - let finalized = sqlx::query( - "UPDATE edge_agent_registry \ - SET edge_id = ?, hostname = ?, worktree_path = ?, capabilities_json = ?, \ - workspace_id = ?, last_heartbeat_at = NOW(6), registration_state = 2, \ - registration_previous_edge_id = ? \ - WHERE user_id = ? AND registry_id = ? AND registration_claim_id = ?", - ) - .bind(&lease.current.edge_id) - .bind(&lease.current.hostname) - .bind(&lease.current.worktree_path) - .bind(&capabilities_json) - .bind(&lease.current.workspace_id) - .bind( - lease - .previous - .as_ref() - .map(|previous| previous.edge_id.as_str()), - ) - .bind(&lease.current.user_id) - .bind(&lease.current.registry_id) - .bind(claim_id) - .execute(&self.pool) - .await - .map_err(|e| format!("edge_registry finalize registration: {e}"))?; - Ok(finalized.rows_affected() > 0) + debug_assert!(!claim_id.is_empty()); + self.settle_registration_transition(lease, RegistrationTransition::Finalize) + .await } async fn release_registration(&self, lease: &EdgeRegistrationLease) -> Result { let Some(claim_id) = lease.claim_id.as_deref() else { return Ok(false); }; - let released = sqlx::query( - "UPDATE edge_agent_registry \ - SET registration_claim_id = NULL, registration_claim_expires_at = NULL, \ - registration_state = 1, registration_previous_edge_id = NULL \ - WHERE user_id = ? AND registry_id = ? AND edge_id = ? \ - AND registration_claim_id = ? AND registration_state = 2", - ) - .bind(&lease.current.user_id) - .bind(&lease.current.registry_id) - .bind(&lease.current.edge_id) - .bind(claim_id) - .execute(&self.pool) - .await - .map_err(|e| format!("edge_registry release registration claim: {e}"))?; - Ok(released.rows_affected() > 0) + debug_assert!(!claim_id.is_empty()); + self.settle_registration_transition(lease, RegistrationTransition::Release) + .await } async fn current_edge_id( @@ -1026,6 +1756,7 @@ impl EdgeRegistryService for UnconfiguredEdgeRegistryService { _user_id: &str, _edge_agent_id: &str, _edge_id_header: &str, + _registration_claim_id: Option<&str>, ) -> Result<(), HeartbeatError> { Ok(()) } @@ -1049,6 +1780,265 @@ impl EdgeRegistryService for UnconfiguredEdgeRegistryService { mod tests { use super::*; + fn registration_lease(previous_edge_id: Option<&str>) -> EdgeRegistrationLease { + let record = |edge_id: &str| EdgeAgentRecord { + registry_id: "registry-1".to_string(), + user_id: "user-1".to_string(), + edge_agent_id: "agent-1".to_string(), + edge_id: edge_id.to_string(), + hostname: None, + worktree_path: None, + capabilities: None, + workspace_id: Some("workspace-1".to_string()), + registered_at: "2026-09-02 00:00:00.000000".to_string(), + last_heartbeat_at: "2026-09-02 00:00:00.000000".to_string(), + }; + EdgeRegistrationLease { + current: record("edge-new"), + previous: previous_edge_id.map(record), + claim_id: Some("claim-1".to_string()), + } + } + + fn registration_state(edge_id: &str, claim_id: Option<&str>, state: i8) -> RegistrationState { + RegistrationState { + edge_id: edge_id.to_string(), + claim_id: claim_id.map(ToString::to_string), + state, + previous_edge_id: None, + } + } + + fn registration_state_with_previous( + edge_id: &str, + claim_id: Option<&str>, + state: i8, + previous_edge_id: Option<&str>, + ) -> RegistrationState { + RegistrationState { + previous_edge_id: previous_edge_id.map(ToString::to_string), + ..registration_state(edge_id, claim_id, state) + } + } + + #[test] + fn release_retries_an_owned_finalized_claim_and_accepts_idempotent_success() { + let lease = registration_lease(Some("edge-old")); + let before_finalize = registration_state("edge-old", Some("claim-1"), 1); + assert_eq!( + registration_transition_decision( + RegistrationTransition::Release, + &lease, + "claim-1", + Some(&before_finalize), + ), + RegistrationTransitionDecision::Apply, + "the same claim may still expose its pre-finalize row through another DB session" + ); + + let owned = registration_state("edge-new", Some("claim-1"), 2); + assert_eq!( + registration_transition_decision( + RegistrationTransition::Release, + &lease, + "claim-1", + Some(&owned), + ), + RegistrationTransitionDecision::Apply + ); + + let released = registration_state("edge-new", None, 1); + assert_eq!( + registration_transition_decision( + RegistrationTransition::Release, + &lease, + "claim-1", + Some(&released), + ), + RegistrationTransitionDecision::AlreadyApplied + ); + } + + #[test] + fn registration_transition_only_reports_superseded_for_a_different_claim() { + let lease = registration_lease(Some("edge-old")); + let successor = registration_state("edge-successor", Some("claim-2"), 1); + for transition in [ + RegistrationTransition::Finalize, + RegistrationTransition::Release, + RegistrationTransition::Rollback, + ] { + assert_eq!( + registration_transition_decision(transition, &lease, "claim-1", Some(&successor),), + RegistrationTransitionDecision::Superseded + ); + } + } + + #[test] + fn rollback_accepts_both_restored_and_inactive_idempotent_states() { + let replacement = registration_lease(Some("edge-old")); + let restored = registration_state("edge-old", None, 1); + assert_eq!( + registration_transition_decision( + RegistrationTransition::Rollback, + &replacement, + "claim-1", + Some(&restored), + ), + RegistrationTransitionDecision::AlreadyApplied + ); + + let first_registration = registration_lease(None); + let inactive = registration_state("edge-new", None, 0); + assert_eq!( + registration_transition_decision( + RegistrationTransition::Rollback, + &first_registration, + "claim-1", + Some(&inactive), + ), + RegistrationTransitionDecision::AlreadyApplied + ); + + let replacement_without_predecessor = registration_state("edge-new", None, 0); + assert_eq!( + registration_transition_decision( + RegistrationTransition::Rollback, + &replacement, + "claim-1", + Some(&replacement_without_predecessor), + ), + RegistrationTransitionDecision::AlreadyApplied + ); + } + + #[test] + fn predecessor_liveness_tracks_disconnects_before_and_after_finalize() { + let pending_live = registration_state("edge-old", Some("claim-1"), 1); + assert!(registration_predecessor_is_live(&pending_live, "edge-old")); + + let pending_disconnected = registration_state("edge-old", Some("claim-1"), 0); + assert!(!registration_predecessor_is_live( + &pending_disconnected, + "edge-old" + )); + + let finalized_live = + registration_state_with_previous("edge-new", Some("claim-1"), 2, Some("edge-old")); + assert!(registration_predecessor_is_live( + &finalized_live, + "edge-old" + )); + + let finalized_disconnected = + registration_state_with_previous("edge-new", Some("claim-1"), 2, None); + assert!(!registration_predecessor_is_live( + &finalized_disconnected, + "edge-old" + )); + } + + #[test] + fn first_registration_rollback_retries_a_stale_empty_read_before_deactivating_owned_claim() { + let lease = registration_lease(None); + let owned = registration_state("edge-new", Some("claim-1"), 0); + + assert_eq!( + registration_before_mutation_decision( + RegistrationTransition::Rollback, + &lease, + "claim-1", + None, + 0, + CURRENT_READ_MAX_ATTEMPTS, + ), + RegistrationAttemptDecision::Retry + ); + assert_eq!( + registration_before_mutation_decision( + RegistrationTransition::Rollback, + &lease, + "claim-1", + Some(&owned), + 1, + CURRENT_READ_MAX_ATTEMPTS, + ), + RegistrationAttemptDecision::Transition(RegistrationTransitionDecision::Apply), + "the retried current read must recover the committed pending claim" + ); + assert_eq!( + registration_before_mutation_decision( + RegistrationTransition::Rollback, + &lease, + "claim-1", + None, + CURRENT_READ_MAX_ATTEMPTS - 1, + CURRENT_READ_MAX_ATTEMPTS, + ), + RegistrationAttemptDecision::OutcomeUnknown, + "absence cannot prove that a committed owner row was removed" + ); + let inactive = registration_state("edge-new", None, 0); + assert_eq!( + registration_transition_decision( + RegistrationTransition::Rollback, + &lease, + "claim-1", + Some(&inactive), + ), + RegistrationTransitionDecision::AlreadyApplied, + "the retained inactive owner row is authoritative rollback evidence" + ); + } + + #[test] + fn generation_ownership_covers_release_unknown_and_rejects_a_successor() { + let release_unknown = RegistrationGenerationState { + edge_id: "edge-new".to_string(), + claim_id: Some("claim-1".to_string()), + previous_edge_id: Some("edge-old".to_string()), + state: 2, + }; + assert!(heartbeat_generation_is_owned( + &release_unknown, + "edge-new", + Some("claim-1") + )); + assert!(!heartbeat_generation_is_owned( + &release_unknown, + "edge-new", + Some("claim-2") + )); + assert!(heartbeat_generation_is_owned( + &release_unknown, + "edge-old", + None + )); + assert!(unregister_generation_is_owned(&release_unknown, "edge-new")); + assert!(unregister_generation_is_owned(&release_unknown, "edge-old")); + + let inactive = RegistrationGenerationState { + edge_id: "edge-new".to_string(), + claim_id: None, + previous_edge_id: None, + state: 0, + }; + assert!(unregister_generation_is_already_applied( + &inactive, "edge-new" + )); + assert!(!heartbeat_generation_is_owned(&inactive, "edge-new", None)); + + let successor = RegistrationGenerationState { + edge_id: "edge-successor".to_string(), + claim_id: None, + previous_edge_id: None, + state: 1, + }; + assert!(!heartbeat_generation_is_owned(&successor, "edge-new", None)); + assert!(!unregister_generation_is_owned(&successor, "edge-new")); + } + #[test] fn edge_agent_record_projection_extracts_matrixone_json_as_unbounded_text() { assert!( diff --git a/crates/services/tests/edge_dispatch_db_it.rs b/crates/services/tests/edge_dispatch_db_it.rs index 89a0634096..2132761178 100644 --- a/crates/services/tests/edge_dispatch_db_it.rs +++ b/crates/services/tests/edge_dispatch_db_it.rs @@ -10,7 +10,8 @@ use std::sync::Arc; use astra_services::multi_agent::{ DatabaseEdgeDispatchService, DatabaseEdgeRegistryService, EdgeDispatchAdmission, - EdgeDispatchAdmissionError, EdgeDispatchIdentity, EdgeDispatchService, EdgeRegistryService, + EdgeDispatchAdmissionError, EdgeDispatchIdentity, EdgeDispatchService, EdgeRegistrationLease, + EdgeRegistryService, }; use sqlx::Row; use uuid::Uuid; @@ -40,6 +41,84 @@ fn require_env() { ); } +type RegistryPrivacyState = ( + String, + i8, + Option, + Option, + Option, + Option, + Option, +); + +async fn registry_privacy_state( + pool: &sqlx::Pool, + user_id: &str, + edge_agent_id: &str, +) -> RegistryPrivacyState { + sqlx::query_as( + "SELECT edge_id, registration_state, registration_claim_id, \ + hostname, worktree_path, capabilities_json, workspace_id \ + FROM edge_agent_registry WHERE user_id = ? AND edge_agent_id = ?", + ) + .bind(user_id) + .bind(edge_agent_id) + .fetch_one(pool) + .await + .expect("read edge registry privacy state") +} + +async fn publish_registry_predecessor( + service: &DatabaseEdgeRegistryService, + user_id: &str, + edge_agent_id: &str, +) -> EdgeRegistrationLease { + let lease = service + .register_or_update_with_lease( + user_id, + edge_agent_id, + "edge-old", + Some("old-private-host"), + Some("/old/private/worktree"), + Some(serde_json::json!({"generation": "old"})), + Some("workspace-old"), + ) + .await + .expect("claim predecessor"); + assert!( + service + .finalize_registration(&lease) + .await + .expect("finalize predecessor") + ); + assert!( + service + .release_registration(&lease) + .await + .expect("publish predecessor") + ); + lease +} + +async fn claim_registry_successor( + service: &DatabaseEdgeRegistryService, + user_id: &str, + edge_agent_id: &str, +) -> EdgeRegistrationLease { + service + .register_or_update_with_lease( + user_id, + edge_agent_id, + "edge-new", + Some("new-host"), + Some("/new/worktree"), + Some(serde_json::json!({"generation": "new"})), + Some("workspace-new"), + ) + .await + .expect("claim successor") +} + // ═══════════════════════════════════════════════════════════════════════ // DatabaseEdgeDispatchService // ═══════════════════════════════════════════════════════════════════════ @@ -661,6 +740,110 @@ async fn edge_registry_register_list_unregister() { // List returns empty let list = svc.list_by_user(&user_id).await.expect("list_by_user"); assert_eq!(list.len(), 0); + + let inactive: (String, i8, Option) = sqlx::query_as( + "SELECT edge_id, registration_state, registration_claim_id \ + FROM edge_agent_registry WHERE user_id = ? AND edge_agent_id = ?", + ) + .bind(&user_id) + .bind(&edge_agent_id) + .fetch_one(&pool) + .await + .expect("read inactive registry owner"); + assert_eq!(inactive, (replacement_edge_id.clone(), 0, None)); + assert!( + svc.unregister_generation(&user_id, &edge_agent_id, &replacement_edge_id) + .await + .expect("repeat unregister current generation"), + "the retained inactive owner row makes repeated cleanup authoritative" + ); +} + +#[tokio::test] +#[ignore = "requires live DB: run with ASTRA_TEST_DB_IT=1"] +async fn edge_registry_first_registration_rollback_retains_an_inactive_owner() { + require_env(); + let pool = common::setup_pool().await.get().clone(); + let svc = DatabaseEdgeRegistryService::new(pool.clone()); + let user_id = format!("user_{}", unique_suffix()); + let edge_agent_id = format!("agent_{}", unique_suffix()); + let edge_id = format!("edge_{}", unique_suffix()); + + let lease = svc + .register_or_update_with_lease( + &user_id, + &edge_agent_id, + &edge_id, + Some("private-host"), + Some("/private/worktree"), + Some(serde_json::json!({"private": "capability"})), + Some("private-workspace"), + ) + .await + .expect("claim first registration"); + let pending = registry_privacy_state(&pool, &user_id, &edge_agent_id).await; + assert_eq!(pending.0, edge_id.clone()); + assert_eq!(pending.1, 0); + assert_eq!(pending.2.as_deref(), lease.claim_id.as_deref()); + assert_eq!( + (pending.3, pending.4, pending.5, pending.6), + (None, None, None, None), + "an unpublished first registration must not persist private metadata" + ); + assert!(svc.rollback_registration(&lease).await.unwrap()); + assert!( + svc.rollback_registration(&lease).await.unwrap(), + "the inactive owner is durable idempotence evidence" + ); + + let inactive = registry_privacy_state(&pool, &user_id, &edge_agent_id).await; + assert_eq!(inactive, (edge_id.clone(), 0, None, None, None, None, None)); + assert!(svc.list_by_user(&user_id).await.unwrap().is_empty()); + + let successor = svc + .register_or_update_with_lease( + &user_id, + &edge_agent_id, + "edge-successor", + Some("successor-host"), + Some("/successor/worktree"), + Some(serde_json::json!({"generation": "successor"})), + Some("successor-workspace"), + ) + .await + .expect("reuse inactive owner for a later registration"); + assert!(successor.previous.is_none()); + let pending_successor = registry_privacy_state(&pool, &user_id, &edge_agent_id).await; + assert_eq!(pending_successor.0, edge_id); + assert_eq!(pending_successor.1, 0); + assert_eq!( + pending_successor.2.as_deref(), + successor.claim_id.as_deref() + ); + assert_eq!( + ( + pending_successor.3, + pending_successor.4, + pending_successor.5, + pending_successor.6, + ), + (None, None, None, None), + "reusing an inactive owner must not persist successor metadata before finalize" + ); + assert!(svc.finalize_registration(&successor).await.unwrap()); + assert!(svc.release_registration(&successor).await.unwrap()); + let published = svc.list_by_user(&user_id).await.unwrap(); + assert_eq!(published.len(), 1); + assert_eq!(published[0].edge_id, "edge-successor"); + assert_eq!(published[0].hostname.as_deref(), Some("successor-host")); + assert_eq!( + published[0].worktree_path.as_deref(), + Some("/successor/worktree") + ); + assert_eq!( + published[0].workspace_id.as_deref(), + Some("successor-workspace") + ); } #[tokio::test] @@ -675,11 +858,11 @@ async fn edge_registry_heartbeat_unregistered() { let edge_id_header = format!("edge_{}", unique_suffix()); let result = svc - .heartbeat(&user_id, &edge_agent_id, &edge_id_header) + .heartbeat(&user_id, &edge_agent_id, &edge_id_header, None) .await; assert!(matches!( result.unwrap_err(), - astra_services::HeartbeatError::Superseded + astra_services::HeartbeatError::StorageFailure(_) )); } @@ -719,7 +902,7 @@ async fn edge_registry_concurrent_register() { async fn edge_registry_registration_lease_restores_the_exact_predecessor() { require_env(); let pool = common::setup_pool().await.get().clone(); - let svc = DatabaseEdgeRegistryService::new(pool); + let svc = DatabaseEdgeRegistryService::new(pool.clone()); let user_id = format!("user_{}", unique_suffix()); let edge_agent_id = format!("agent_{}", unique_suffix()); @@ -780,6 +963,12 @@ async fn edge_registry_registration_lease_restores_the_exact_predecessor() { .await .expect("rollback replacement") ); + assert!( + svc.rollback_registration(&replacement) + .await + .expect("repeat rollback replacement"), + "an already restored predecessor is an idempotent rollback success" + ); let restored = svc .list_by_user(&user_id) @@ -793,6 +982,16 @@ async fn edge_registry_registration_lease_restores_the_exact_predecessor() { assert_eq!(restored.worktree_path.as_deref(), Some("/old/worktree")); assert_eq!(restored.capabilities, Some(old_capabilities)); assert_eq!(restored.workspace_id.as_deref(), Some("workspace-old")); + let persisted: (String, i8, Option) = sqlx::query_as( + "SELECT edge_id, registration_state, registration_claim_id \ + FROM edge_agent_registry WHERE user_id = ? AND edge_agent_id = ?", + ) + .bind(&user_id) + .bind(&edge_agent_id) + .fetch_one(&pool) + .await + .expect("read persisted rollback state"); + assert_eq!(persisted, (old_edge_id, 1, None)); } #[tokio::test] @@ -800,7 +999,7 @@ async fn edge_registry_registration_lease_restores_the_exact_predecessor() { async fn edge_registry_two_phase_registration_keeps_pending_metadata_unroutable() { require_env(); let pool = common::setup_pool().await.get().clone(); - let svc = DatabaseEdgeRegistryService::new(pool); + let svc = DatabaseEdgeRegistryService::new(pool.clone()); let user_id = format!("user_{}", unique_suffix()); let edge_agent_id = format!("agent_{}", unique_suffix()); @@ -838,24 +1037,154 @@ async fn edge_registry_two_phase_registration_keeps_pending_metadata_unroutable( still_published[0].workspace_id.as_deref(), Some("workspace-old") ); - svc.heartbeat(&user_id, &edge_agent_id, "edge-old") + svc.heartbeat(&user_id, &edge_agent_id, "edge-old", None) .await .expect("published predecessor remains healthy during claim"); assert!(svc.finalize_registration(&pending).await.unwrap()); + assert!( + svc.finalize_registration(&pending).await.unwrap(), + "finalization must be idempotent while the same claim is retained" + ); assert!( svc.list_by_user(&user_id).await.unwrap().is_empty(), "finalized generation stays unroutable until pool commit releases the claim" ); assert!(svc.release_registration(&pending).await.unwrap()); assert!( - !svc.release_registration(&pending).await.unwrap(), - "a committed claim cannot be released twice" + svc.release_registration(&pending).await.unwrap(), + "an already committed claim is an idempotent release success" ); let current = svc.list_by_user(&user_id).await.unwrap(); assert_eq!(current.len(), 1); assert_eq!(current[0].edge_id, "edge-new"); assert_eq!(current[0].workspace_id.as_deref(), Some("workspace-new")); + let persisted: (String, i8, Option) = sqlx::query_as( + "SELECT edge_id, registration_state, registration_claim_id \ + FROM edge_agent_registry WHERE user_id = ? AND edge_agent_id = ?", + ) + .bind(&user_id) + .bind(&edge_agent_id) + .fetch_one(&pool) + .await + .expect("read persisted release state"); + assert_eq!(persisted, ("edge-new".to_string(), 1, None)); +} + +#[tokio::test] +#[ignore = "requires live DB: run with ASTRA_TEST_DB_IT=1"] +async fn edge_registry_predecessor_disconnect_before_finalize_preserves_successor_claim() { + require_env(); + let predecessor_pool = common::setup_pool().await.get().clone(); + let successor_pool = common::setup_pool().await.get().clone(); + let predecessor_pod = DatabaseEdgeRegistryService::new(predecessor_pool.clone()); + let successor_pod = DatabaseEdgeRegistryService::new(successor_pool); + let user_id = format!("user_{}", unique_suffix()); + let edge_agent_id = format!("agent_{}", unique_suffix()); + + publish_registry_predecessor(&predecessor_pod, &user_id, &edge_agent_id).await; + let successor = claim_registry_successor(&successor_pod, &user_id, &edge_agent_id).await; + assert!( + predecessor_pod + .unregister_generation(&user_id, &edge_agent_id, "edge-old") + .await + .expect("deactivate predecessor without erasing successor claim") + ); + + let pending = registry_privacy_state(&predecessor_pool, &user_id, &edge_agent_id).await; + assert_eq!(pending.0, "edge-old"); + assert_eq!(pending.1, 0); + assert_eq!(pending.2.as_deref(), successor.claim_id.as_deref()); + assert_eq!( + (pending.3, pending.4, pending.5, pending.6), + (None, None, None, None) + ); + + assert!( + successor_pod + .finalize_registration(&successor) + .await + .expect("finalize successor after predecessor disconnect") + ); + assert!( + successor_pod + .release_registration(&successor) + .await + .expect("publish successor after predecessor disconnect") + ); + let published = predecessor_pod.list_by_user(&user_id).await.unwrap(); + assert_eq!(published.len(), 1); + assert_eq!(published[0].edge_id, "edge-new"); + assert_eq!(published[0].hostname.as_deref(), Some("new-host")); +} + +#[tokio::test] +#[ignore = "requires live DB: run with ASTRA_TEST_DB_IT=1"] +async fn edge_registry_predecessor_disconnect_after_finalize_prevents_rollback_resurrection() { + require_env(); + let predecessor_pool = common::setup_pool().await.get().clone(); + let successor_pool = common::setup_pool().await.get().clone(); + let predecessor_pod = DatabaseEdgeRegistryService::new(predecessor_pool.clone()); + let successor_pod = DatabaseEdgeRegistryService::new(successor_pool); + let user_id = format!("user_{}", unique_suffix()); + let edge_agent_id = format!("agent_{}", unique_suffix()); + + publish_registry_predecessor(&predecessor_pod, &user_id, &edge_agent_id).await; + let successor = claim_registry_successor(&successor_pod, &user_id, &edge_agent_id).await; + assert!( + successor_pod + .finalize_registration(&successor) + .await + .expect("finalize successor") + ); + assert!( + predecessor_pod + .unregister_generation(&user_id, &edge_agent_id, "edge-old") + .await + .expect("record predecessor disconnect after finalize") + ); + + let finalized: (String, i8, Option, Option, Option) = sqlx::query_as( + "SELECT edge_id, registration_state, registration_claim_id, \ + registration_previous_edge_id, hostname \ + FROM edge_agent_registry WHERE user_id = ? AND edge_agent_id = ?", + ) + .bind(&user_id) + .bind(&edge_agent_id) + .fetch_one(&predecessor_pool) + .await + .expect("read finalized successor after predecessor disconnect"); + assert_eq!(finalized.0, "edge-new"); + assert_eq!(finalized.1, 2); + assert_eq!(finalized.2.as_deref(), successor.claim_id.as_deref()); + assert_eq!(finalized.3, None); + assert_eq!(finalized.4.as_deref(), Some("new-host")); + + assert!( + successor_pod + .rollback_registration(&successor) + .await + .expect("rollback successor without resurrecting disconnected predecessor") + ); + assert!( + successor_pod + .rollback_registration(&successor) + .await + .expect("repeat rollback remains idempotent") + ); + assert!( + predecessor_pod + .list_by_user(&user_id) + .await + .unwrap() + .is_empty() + ); + + let inactive = registry_privacy_state(&predecessor_pool, &user_id, &edge_agent_id).await; + assert_eq!( + inactive, + ("edge-new".to_string(), 0, None, None, None, None, None) + ); } #[tokio::test] @@ -1071,6 +1400,166 @@ async fn edge_registry_registration_claim_serializes_cross_pod_setup() { ); } +#[tokio::test] +#[ignore = "requires live DB: run with ASTRA_TEST_DB_IT=1"] +async fn edge_registry_finalized_claim_is_renewed_and_fences_a_third_generation() { + expired_finalized_claim_survives_displaced_cleanup(false).await; +} + +#[tokio::test] +#[ignore = "requires live DB: run with ASTRA_TEST_DB_IT=1"] +async fn edge_registry_expired_finalized_claim_publishes_after_displaced_cleanup() { + expired_finalized_claim_survives_displaced_cleanup(true).await; +} + +async fn expired_finalized_claim_survives_displaced_cleanup(publish_third: bool) { + require_env(); + let pool = common::setup_pool().await.get().clone(); + let first_pod = DatabaseEdgeRegistryService::new(pool.clone()); + let second_pod = DatabaseEdgeRegistryService::new(common::setup_pool().await.get().clone()); + let third_pod = DatabaseEdgeRegistryService::new(common::setup_pool().await.get().clone()); + let user_id = format!("user_{}", unique_suffix()); + let edge_agent_id = format!("agent_{}", unique_suffix()); + + let first = first_pod + .register_or_update_with_lease( + &user_id, + &edge_agent_id, + "edge-first", + None, + None, + None, + None, + ) + .await + .expect("claim first generation"); + assert!(first_pod.finalize_registration(&first).await.unwrap()); + assert!(first_pod.release_registration(&first).await.unwrap()); + + let second = second_pod + .register_or_update_with_lease( + &user_id, + &edge_agent_id, + "edge-second", + None, + None, + None, + None, + ) + .await + .expect("claim second generation"); + assert!(second_pod.finalize_registration(&second).await.unwrap()); + let second_claim = second.claim_id.as_deref().expect("durable second claim"); + + sqlx::query( + "UPDATE edge_agent_registry \ + SET registration_claim_expires_at = DATE_SUB(NOW(6), INTERVAL 1 SECOND) \ + WHERE user_id = ? AND edge_agent_id = ?", + ) + .bind(&user_id) + .bind(&edge_agent_id) + .execute(&pool) + .await + .expect("expire second claim for renewal test"); + second_pod + .heartbeat(&user_id, &edge_agent_id, "edge-second", Some(second_claim)) + .await + .expect("the finalized owner renews its exact claim"); + + let renewed: (Option, i8) = sqlx::query_as( + "SELECT registration_claim_id, registration_claim_expires_at > NOW(6) \ + FROM edge_agent_registry WHERE user_id = ? AND edge_agent_id = ?", + ) + .bind(&user_id) + .bind(&edge_agent_id) + .fetch_one(&pool) + .await + .expect("read renewed claim"); + assert_eq!(renewed, (Some(second_claim.to_string()), 1)); + + let blocked_third = third_pod + .register_or_update_with_lease( + &user_id, + &edge_agent_id, + "edge-third-blocked", + None, + None, + None, + None, + ) + .await; + assert!( + blocked_third.is_err(), + "a healthy finalized owner must not lose its renewed claim" + ); + + sqlx::query( + "UPDATE edge_agent_registry \ + SET registration_claim_expires_at = DATE_SUB(NOW(6), INTERVAL 1 SECOND) \ + WHERE user_id = ? AND edge_agent_id = ?", + ) + .bind(&user_id) + .bind(&edge_agent_id) + .execute(&pool) + .await + .expect("expire abandoned second claim"); + let third = third_pod + .register_or_update_with_lease( + &user_id, + &edge_agent_id, + "edge-third", + None, + None, + None, + None, + ) + .await + .expect("third generation takes an expired finalized claim"); + + assert!(matches!( + second_pod + .heartbeat(&user_id, &edge_agent_id, "edge-second", Some(second_claim),) + .await, + Err(astra_services::multi_agent::HeartbeatError::Superseded) + )); + assert!(third.previous.is_none()); + assert!( + second_pod + .unregister_generation(&user_id, &edge_agent_id, "edge-second") + .await + .unwrap() + ); + let state: (i8, Option, Option, i8) = sqlx::query_as( + "SELECT registration_state, registration_claim_id, registration_previous_edge_id, \ + registration_claim_expires_at > NOW(6) \ + FROM edge_agent_registry WHERE user_id = ? AND edge_agent_id = ?", + ) + .bind(&user_id) + .bind(&edge_agent_id) + .fetch_one(&pool) + .await + .unwrap(); + assert_eq!(state, (0, third.claim_id.clone(), None, 1)); + + if publish_third { + assert!(third_pod.finalize_registration(&third).await.unwrap()); + assert!(third_pod.release_registration(&third).await.unwrap()); + // Repeated cleanup must also leave the published successor intact. + assert!( + !second_pod + .unregister_generation(&user_id, &edge_agent_id, "edge-second") + .await + .unwrap() + ); + let published = third_pod.list_by_user(&user_id).await.unwrap(); + assert_eq!(published.len(), 1); + assert_eq!(published[0].edge_id, "edge-third"); + } else { + assert!(third_pod.rollback_registration(&third).await.unwrap()); + assert!(third_pod.list_by_user(&user_id).await.unwrap().is_empty()); + } +} + #[tokio::test] #[ignore = "requires live DB: run with ASTRA_TEST_DB_IT=1"] async fn edge_registry_register_twice_updates() { diff --git a/crates/services/tests/multi_agent_integration.rs b/crates/services/tests/multi_agent_integration.rs index 95a81b24b9..71cf1a163a 100644 --- a/crates/services/tests/multi_agent_integration.rs +++ b/crates/services/tests/multi_agent_integration.rs @@ -59,7 +59,7 @@ async fn edge_registry_register_twice_keeps_registry_id() { assert_eq!(second.hostname.as_deref(), Some("h2")); registry - .heartbeat(&user, &edge_agent, "transport-b") + .heartbeat(&user, &edge_agent, "transport-b", None) .await .expect("heartbeat");