mod state_machine;
mod store;
mod types;
pub use types::{ConnectionHealth, ConnectionRecord, ConnectionState, PeerSnapshot};
use types::{ConnectionStore, ReplacementHandoffOwner};
use crate::heartbeat::state_machine::peer_health;
use state_machine::{default_health_for_state, merge_upsert_reason, merge_upsert_state};
use std::collections::HashSet;
use std::sync::Arc;
use store::{
add_index, add_indexes, build_peer_snapshot, canonical_peer_id, collect_peer_ids,
migrate_peer_state, remove_index, resolve_peer_id, unix_ms_now,
};
#[cfg(any(feature = "iroh-carrier-core", test))]
use tokio::sync::OwnedRwLockWriteGuard;
use tokio::sync::RwLock;
#[derive(Default)]
pub struct ConnectionManager {
store: Arc<RwLock<ConnectionStore>>,
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub(crate) struct IncomingTransportGenerationFence {
pub incoming_transport_stable_id: u64,
pub physical_transport_stable_id: Option<u64>,
pub managed_transport_stable_id: Option<u64>,
pub manager_record_exists: bool,
pub owned: bool,
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub(crate) enum TerminalReopenAuthority {
None,
FreshSessionCapability,
CurrentDesiredUserDevice,
}
#[cfg(any(feature = "iroh-carrier-core", test))]
pub(crate) struct TransportReplacementCommit {
record: ConnectionRecord,
_store: OwnedRwLockWriteGuard<ConnectionStore>,
}
#[cfg(any(feature = "iroh-carrier-core", test))]
impl TransportReplacementCommit {
pub(crate) fn finish(self) -> ConnectionRecord {
self.record
}
}
fn clear_generation_bound_carrier_route(record: &mut ConnectionRecord) -> bool {
let active_is_carrier =
crate::transport_label::is_generation_bound_iroh_carrier(&record.active_transport);
let parallel_is_carrier = record
.parallel_transport
.as_deref()
.is_some_and(crate::transport_label::is_generation_bound_iroh_carrier);
if active_is_carrier {
record.active_transport = crate::transport_label::IROH.to_string();
record.parallel_transport = None;
} else if parallel_is_carrier {
record.parallel_transport = None;
}
active_is_carrier || parallel_is_carrier
}
#[cfg(any(feature = "iroh-carrier-core", test))]
fn replacement_source_proves_current_carrier(
record: &ConnectionRecord,
transport_source: Option<&str>,
) -> bool {
crate::transport_label::is_generation_bound_iroh_carrier(&record.active_transport)
&& transport_source.and_then(|source| source.strip_suffix("-upgrade"))
== Some(record.active_transport.as_str())
}
impl ConnectionManager {
pub fn new() -> Self {
Self::default()
}
#[cfg(test)]
pub(crate) async fn incoming_transport_generation_is_owned(
&self,
connection_id: &str,
physical_transport_stable_id: Option<u64>,
incoming_transport_stable_id: u64,
) -> bool {
self.incoming_transport_generation_fence(
connection_id,
physical_transport_stable_id,
incoming_transport_stable_id,
)
.await
.owned
}
pub(crate) async fn incoming_transport_generation_fence(
&self,
connection_id: &str,
physical_transport_stable_id: Option<u64>,
incoming_transport_stable_id: u64,
) -> IncomingTransportGenerationFence {
let store = self.store.read().await;
let record = store.by_id.get(connection_id);
let owned = match record {
None => physical_transport_stable_id == Some(incoming_transport_stable_id),
Some(record)
if matches!(
record.state,
ConnectionState::Closing | ConnectionState::Closed | ConnectionState::Failed
) =>
{
false
}
Some(record) => record.transport_stable_id == Some(incoming_transport_stable_id),
};
IncomingTransportGenerationFence {
incoming_transport_stable_id,
physical_transport_stable_id,
managed_transport_stable_id: record.and_then(|record| record.transport_stable_id),
manager_record_exists: record.is_some(),
owned,
}
}
#[cfg(any(test, target_arch = "wasm32"))]
pub(crate) async fn incoming_authenticated_reconnect_candidate(
&self,
connection_id: &str,
physical_transport_stable_id: Option<u64>,
incoming_transport_stable_id: u64,
peer_requested_exclusion: bool,
) -> bool {
if !peer_requested_exclusion
|| physical_transport_stable_id != Some(incoming_transport_stable_id)
{
return false;
}
self.store
.read()
.await
.by_id
.get(connection_id)
.is_some_and(|record| {
record.state == ConnectionState::Closed
&& crate::lifecycle_reason::LifecycleReasonCode::from_text(
record.status_reason.as_deref(),
) == Some(crate::lifecycle_reason::LifecycleReasonCode::ManualDisconnect)
})
}
pub(crate) async fn with_current_transport_generation<T>(
&self,
connection_id: &str,
endpoint_id: &str,
transport_stable_id: u64,
apply: impl FnOnce(&ConnectionRecord) -> T,
) -> Option<T> {
let store = self.store.read().await;
let record = store.by_id.get(connection_id)?;
if matches!(
record.state,
ConnectionState::Closing | ConnectionState::Closed | ConnectionState::Failed
) || record.transport_stable_id != Some(transport_stable_id)
|| record.endpoint_id.as_deref().or(record.node_id.as_deref()) != Some(endpoint_id)
{
return None;
}
Some(apply(record))
}
pub async fn upsert_pending(
&self,
connection_id: String,
node_id: Option<String>,
device_id_hint: Option<String>,
endpoint_id: Option<String>,
) -> ConnectionRecord {
self.upsert(
connection_id,
node_id,
device_id_hint,
endpoint_id,
ConnectionState::Pending,
None,
)
.await
}
pub(crate) async fn upsert_pending_for_physical_observation(
&self,
connection_id: String,
node_id: Option<String>,
device_id_hint: Option<String>,
endpoint_id: Option<String>,
) -> Option<ConnectionRecord> {
self.upsert_when(
connection_id,
node_id,
device_id_hint,
endpoint_id,
ConnectionState::Pending,
None,
|existing, history| {
existing.is_some()
|| !history.is_some_and(|history| {
crate::lifecycle_reason::is_terminal(
history.last_disconnect_reason.as_deref(),
)
})
},
)
.await
}
pub(crate) async fn materialize_manual_disconnect_tombstone(
&self,
connection_id: String,
node_id: Option<String>,
device_id_hint: Option<String>,
endpoint_id: Option<String>,
) -> Option<ConnectionRecord> {
self.upsert_when(
connection_id,
node_id,
device_id_hint,
endpoint_id,
ConnectionState::Closed,
Some(crate::lifecycle_reason::REASON_MANUAL_DISCONNECT.to_string()),
|existing, history| {
existing.is_none()
&& history.is_some_and(|history| {
crate::lifecycle_reason::LifecycleReasonCode::from_text(
history.last_disconnect_reason.as_deref(),
) == Some(crate::lifecycle_reason::LifecycleReasonCode::ManualDisconnect)
})
},
)
.await
}
pub async fn set_connecting(&self, connection_id: &str) -> Option<ConnectionRecord> {
self.transition_state(connection_id, ConnectionState::Connecting, None)
.await
}
pub(crate) async fn begin_automatic_connect(
&self,
connection_id: String,
node_id: Option<String>,
device_id_hint: Option<String>,
endpoint_id: Option<String>,
allowed: impl FnOnce(bool) -> bool,
) -> Option<ConnectionRecord> {
self.upsert_when(
connection_id,
node_id,
device_id_hint,
endpoint_id,
ConnectionState::Connecting,
None,
|existing, history| {
use crate::lifecycle_reason::LifecycleReasonCode;
let reason = existing
.map(|record| record.status_reason.as_deref())
.unwrap_or_else(|| {
history.and_then(|entry| entry.last_disconnect_reason.as_deref())
});
let requires_assignment = matches!(
LifecycleReasonCode::from_text(reason),
Some(
LifecycleReasonCode::SessionTokenRevoked
| LifecycleReasonCode::SessionAdmissionRejected
| LifecycleReasonCode::DesiredPeerWithdrawn
)
);
(!crate::lifecycle_reason::is_terminal(reason) || requires_assignment)
&& allowed(requires_assignment)
},
)
.await
}
pub async fn set_connected(
&self,
connection_id: &str,
endpoint_id: Option<String>,
) -> Option<ConnectionRecord> {
self.set_connected_with_transport(connection_id, endpoint_id, None, None)
.await
}
pub async fn set_connected_with_transport(
&self,
connection_id: &str,
endpoint_id: Option<String>,
transport_stable_id: Option<u64>,
transport_source: Option<String>,
) -> Option<ConnectionRecord> {
self.set_connected_with_transport_when(
connection_id,
endpoint_id,
transport_stable_id,
transport_source,
|_| {},
)
.await
}
#[cfg(test)]
pub(crate) async fn commit_current_transport(
&self,
connection_id: &str,
endpoint_id: Option<String>,
transport_stable_id: u64,
transport_source: Option<String>,
before_publish: impl FnOnce(Option<u64>),
) -> Option<ConnectionRecord> {
self.set_connected_with_transport_when(
connection_id,
endpoint_id,
Some(transport_stable_id),
transport_source,
before_publish,
)
.await
}
#[cfg(test)]
pub(crate) async fn commit_admitted_transport(
&self,
connection_id: &str,
endpoint_id: Option<String>,
authoritative_device_id: Option<String>,
transport_stable_id: u64,
transport_source: Option<String>,
authorize_before_publish: impl FnOnce(bool, Option<u64>) -> bool,
normalize_scopes: impl FnOnce(&mut HashSet<String>),
) -> Option<ConnectionRecord> {
self.set_connected_with_transport_when_allowing_authenticated_peer_reconnect(
connection_id,
endpoint_id,
authoritative_device_id.clone(),
authoritative_device_id,
Some(transport_stable_id),
transport_source,
authorize_before_publish,
normalize_scopes,
TerminalReopenAuthority::FreshSessionCapability,
)
.await
}
#[cfg(test)]
pub(crate) async fn commit_authenticated_peer_reconnect_transport(
&self,
connection_id: &str,
endpoint_id: Option<String>,
device_id_hint: Option<String>,
transport_stable_id: u64,
transport_source: Option<String>,
authorize_before_publish: impl FnOnce(bool, Option<u64>) -> bool,
) -> Option<ConnectionRecord> {
self.set_connected_with_transport_when_allowing_authenticated_peer_reconnect(
connection_id,
endpoint_id,
device_id_hint.clone(),
device_id_hint,
Some(transport_stable_id),
transport_source,
authorize_before_publish,
|_| {},
TerminalReopenAuthority::None,
)
.await
}
async fn set_connected_with_transport_when(
&self,
connection_id: &str,
endpoint_id: Option<String>,
transport_stable_id: Option<u64>,
transport_source: Option<String>,
before_publish: impl FnOnce(Option<u64>),
) -> Option<ConnectionRecord> {
self.set_connected_with_transport_when_allowing_authenticated_peer_reconnect(
connection_id,
endpoint_id,
None,
None,
transport_stable_id,
transport_source,
|requires_authenticated_peer_reconnect, previous| {
if requires_authenticated_peer_reconnect {
return false;
}
before_publish(previous);
true
},
|_| {},
TerminalReopenAuthority::None,
)
.await
}
pub(crate) async fn set_connected_with_transport_when_allowing_authenticated_peer_reconnect(
&self,
connection_id: &str,
endpoint_id: Option<String>,
device_id_hint: Option<String>,
authoritative_device_id: Option<String>,
transport_stable_id: Option<u64>,
transport_source: Option<String>,
authorize_before_publish: impl FnOnce(bool, Option<u64>) -> bool,
normalize_scopes: impl FnOnce(&mut HashSet<String>),
terminal_reopen_authority: TerminalReopenAuthority,
) -> Option<ConnectionRecord> {
let mut store = self.store.write().await;
let terminal_history_reason = store
.history_by_connection_id
.get(connection_id)
.and_then(|history| history.last_disconnect_reason.as_deref());
let missing = !store.by_id.contains_key(connection_id);
let current_terminal_reason = store.by_id.get(connection_id).and_then(|current| {
current
.status_reason
.as_deref()
.or(current.last_disconnect_reason.as_deref())
});
let terminal_reason_code = crate::lifecycle_reason::LifecycleReasonCode::from_text(
current_terminal_reason.or(terminal_history_reason),
);
let missing_manual_tombstone = missing
&& terminal_reason_code
== Some(crate::lifecycle_reason::LifecycleReasonCode::ManualDisconnect);
let current_manual_tombstone = store.by_id.get(connection_id).is_some_and(|current| {
current.state == ConnectionState::Closed
&& crate::lifecycle_reason::LifecycleReasonCode::from_text(
current.status_reason.as_deref(),
) == Some(crate::lifecycle_reason::LifecycleReasonCode::ManualDisconnect)
});
let requires_authenticated_peer_reconnect =
missing_manual_tombstone || current_manual_tombstone;
let allows_fresh_session_reauthorization = terminal_reopen_authority
== TerminalReopenAuthority::FreshSessionCapability
&& terminal_reason_code.is_some_and(|reason| {
matches!(
reason,
crate::lifecycle_reason::LifecycleReasonCode::SessionTokenRevoked
| crate::lifecycle_reason::LifecycleReasonCode::SessionAdmissionRejected
| crate::lifecycle_reason::LifecycleReasonCode::IncomingTransportLost
| crate::lifecycle_reason::LifecycleReasonCode::DesiredPeerWithdrawn
) || reason.is_transient()
});
let allows_current_desired_user_device = terminal_reopen_authority
== TerminalReopenAuthority::CurrentDesiredUserDevice
&& terminal_reason_code.is_some_and(|reason| {
reason == crate::lifecycle_reason::LifecycleReasonCode::DesiredPeerWithdrawn
|| reason.is_transient()
});
let terminal = store.by_id.get(connection_id).is_some_and(|current| {
matches!(
current.state,
ConnectionState::Closing | ConnectionState::Closed | ConnectionState::Failed
)
});
if missing
&& !missing_manual_tombstone
&& !allows_fresh_session_reauthorization
&& !allows_current_desired_user_device
{
return None;
}
if terminal
&& !requires_authenticated_peer_reconnect
&& !allows_fresh_session_reauthorization
&& !allows_current_desired_user_device
{
return None;
}
let previous_transport_stable_id = store
.by_id
.get(connection_id)
.and_then(|current| current.transport_stable_id);
if !authorize_before_publish(
requires_authenticated_peer_reconnect,
previous_transport_stable_id,
) {
return None;
}
if missing {
let mut record = ConnectionRecord::new(
connection_id.to_string(),
endpoint_id.clone(),
device_id_hint.clone(),
endpoint_id.clone(),
);
if let Some(history) = store.history_by_connection_id.get(connection_id).cloned() {
record.transport_generation = history.transport_generation;
record.route_generation = history.route_generation;
record.transition_count = history.transition_count;
record.connecting_transition_count = history.connecting_transition_count;
record.replacement_count = history.replacement_count;
record.retire_count = history.retire_count;
record.last_disconnect_reason = history.last_disconnect_reason;
record.last_reconnect_reason = history.last_reconnect_reason;
}
add_indexes(&mut store, &record);
store.by_id.insert(connection_id.to_string(), record);
}
let current = store.by_id.get(connection_id)?;
let previous_peer_id = canonical_peer_id(current);
normalize_scopes(
store
.scopes_by_peer
.entry(previous_peer_id.clone())
.or_default(),
);
if store
.scopes_by_peer
.get(&previous_peer_id)
.is_some_and(HashSet::is_empty)
{
store.scopes_by_peer.remove(&previous_peer_id);
}
let previous_endpoint = store
.by_id
.get(connection_id)
.and_then(|record| record.endpoint_id.clone());
let new_endpoint = endpoint_id.clone();
if endpoint_id.is_some() {
remove_index(
&mut store.by_endpoint_id,
previous_endpoint.as_deref(),
connection_id,
);
add_index(
&mut store.by_endpoint_id,
new_endpoint.as_deref(),
connection_id,
);
}
if authoritative_device_id.is_some() {
let existing = store.by_id.get(connection_id)?;
let existing_device_id = existing.device_id.clone();
let existing_device_hint = existing.device_id_hint.clone();
remove_index(
&mut store.by_device_id,
existing_device_id.as_deref(),
connection_id,
);
remove_index(
&mut store.by_device_hint_id,
existing_device_hint.as_deref(),
connection_id,
);
}
let previous = store.by_id.get(connection_id)?.clone();
let previous_peer_id = canonical_peer_id(&previous);
let record = store.by_id.get_mut(connection_id)?;
if endpoint_id.is_some() {
record.endpoint_id = endpoint_id;
}
if let Some(device_id) = authoritative_device_id.as_deref() {
let prior_hint = record
.device_id_hint
.clone()
.or_else(|| record.device_id.clone());
record.device_id = Some(device_id.to_string());
record.device_id_hint = if prior_hint.as_deref() == Some(device_id) {
None
} else {
prior_hint
};
}
let transport_changed = transport_stable_id
.map(|stable_id| record.transport_stable_id != Some(stable_id))
.unwrap_or(false);
let state_changed = record.state != ConnectionState::Connected;
if state_changed {
record.transition_count = record.transition_count.saturating_add(1);
}
if transport_changed {
record.transport_generation = record.transport_generation.saturating_add(1);
record.route_generation = 0;
clear_generation_bound_carrier_route(record);
record.transport_stable_id = transport_stable_id;
record.transport_source = transport_source;
let now = unix_ms_now();
record.last_transport_change_at_ms = now;
record.last_route_change_at_ms = now;
record.replacement_count = record.replacement_count.saturating_add(1);
} else {
if transport_stable_id.is_some() {
record.transport_stable_id = transport_stable_id;
}
if transport_source.is_some() {
record.transport_source = transport_source;
}
}
record.state = ConnectionState::Connected;
record.status_reason = None;
record.last_disconnect_reason = None;
record.last_reconnect_reason = record.transport_source.clone();
let now = unix_ms_now();
if state_changed {
record.last_state_change_at_ms = now;
}
record.updated_at_ms = now;
let updated = record.clone();
if authoritative_device_id.is_some() {
add_index(
&mut store.by_device_id,
updated.device_id.as_deref(),
connection_id,
);
add_index(
&mut store.by_device_hint_id,
updated.device_id_hint.as_deref(),
connection_id,
);
}
let next_peer_id = canonical_peer_id(&updated);
if updated.transport_stable_id.is_some() {
store
.replacement_handoff_by_connection_id
.remove(connection_id);
}
if transport_changed {
store
.health_by_peer
.insert(next_peer_id.clone(), ConnectionHealth::Unknown);
} else {
store
.health_by_peer
.entry(next_peer_id.clone())
.or_insert_with(|| default_health_for_state(&updated.state));
}
migrate_peer_state(&mut store, &previous_peer_id, &next_peer_id);
Some(updated)
}
pub async fn mark_transport_replaced(
&self,
connection_id: &str,
transport_stable_id: Option<u64>,
transport_source: Option<String>,
status_reason: Option<String>,
) -> Option<ConnectionRecord> {
let mut store = self.store.write().await;
let (updated, previous_transport) = {
let record = store.by_id.get_mut(connection_id)?;
let previous_transport =
record
.transport_stable_id
.map(|transport_stable_id| ReplacementHandoffOwner {
transport_stable_id,
transport_generation: record.transport_generation,
route_generation: record.route_generation,
});
let now = unix_ms_now();
let transport_changed = transport_stable_id
.map(|stable_id| record.transport_stable_id != Some(stable_id))
.unwrap_or(false);
if transport_changed {
record.transport_generation = record.transport_generation.saturating_add(1);
record.route_generation = 0;
clear_generation_bound_carrier_route(record);
record.replacement_count = record.replacement_count.saturating_add(1);
} else if transport_stable_id.is_none() && clear_generation_bound_carrier_route(record)
{
record.route_generation = record.route_generation.saturating_add(1);
}
record.transport_stable_id = transport_stable_id;
record.transport_source = transport_source;
record.last_transport_change_at_ms = now;
record.last_route_change_at_ms = now;
record.status_reason = status_reason;
record.last_reconnect_reason = record.transport_source.clone();
record.updated_at_ms = now;
(record.clone(), previous_transport)
};
if updated.transport_stable_id.is_none() {
if let Some(owner) = previous_transport {
store
.replacement_handoff_by_connection_id
.insert(connection_id.to_string(), owner);
}
} else {
store
.replacement_handoff_by_connection_id
.remove(connection_id);
}
let peer_id = canonical_peer_id(&updated);
store
.health_by_peer
.insert(peer_id, ConnectionHealth::Unknown);
Some(updated)
}
pub async fn mark_transport_replaced_if_current(
&self,
connection_id: &str,
expected_transport_stable_id: u64,
transport_source: Option<String>,
status_reason: Option<String>,
) -> Option<ConnectionRecord> {
let mut store = self.store.write().await;
let (updated, owner) = {
let record = store.by_id.get_mut(connection_id)?;
if record.transport_stable_id != Some(expected_transport_stable_id) {
return None;
}
let owner = ReplacementHandoffOwner {
transport_stable_id: expected_transport_stable_id,
transport_generation: record.transport_generation,
route_generation: record.route_generation,
};
let now = unix_ms_now();
if clear_generation_bound_carrier_route(record) {
record.route_generation = record.route_generation.saturating_add(1);
}
record.transport_stable_id = None;
record.transport_source = transport_source;
record.last_transport_change_at_ms = now;
record.last_route_change_at_ms = now;
record.status_reason = status_reason;
record.last_reconnect_reason = record.transport_source.clone();
record.updated_at_ms = now;
(record.clone(), owner)
};
store
.replacement_handoff_by_connection_id
.insert(connection_id.to_string(), owner);
let peer_id = canonical_peer_id(&updated);
store
.health_by_peer
.insert(peer_id, ConnectionHealth::Stale);
Some(updated)
}
pub async fn mark_transport_replaced_by_if_current(
&self,
connection_id: &str,
expected_transport_stable_id: u64,
replacement_transport_stable_id: u64,
transport_source: Option<String>,
) -> Option<ConnectionRecord> {
let mut store = self.store.write().await;
let updated = {
let record = store.by_id.get_mut(connection_id)?;
if record.transport_stable_id != Some(expected_transport_stable_id) {
return None;
}
let now = unix_ms_now();
if replacement_transport_stable_id != expected_transport_stable_id {
record.transport_generation = record.transport_generation.saturating_add(1);
record.route_generation = 0;
clear_generation_bound_carrier_route(record);
record.replacement_count = record.replacement_count.saturating_add(1);
}
if record.state != ConnectionState::Connected {
record.transition_count = record.transition_count.saturating_add(1);
record.last_state_change_at_ms = now;
}
record.state = ConnectionState::Connected;
record.transport_stable_id = Some(replacement_transport_stable_id);
record.transport_source = transport_source;
record.status_reason = None;
record.last_disconnect_reason = None;
record.last_reconnect_reason = record.transport_source.clone();
record.last_transport_change_at_ms = now;
record.last_route_change_at_ms = now;
record.updated_at_ms = now;
record.clone()
};
store
.replacement_handoff_by_connection_id
.remove(connection_id);
let peer_id = canonical_peer_id(&updated);
store
.health_by_peer
.insert(peer_id, ConnectionHealth::Unknown);
Some(updated)
}
#[cfg(test)]
pub(crate) async fn install_transport_replacement_if_generation_current(
&self,
connection_id: &str,
expected_transport_stable_id: u64,
expected_transport_generation: u64,
_expected_route_generation: u64,
replacement_transport_stable_id: u64,
transport_source: Option<String>,
) -> Option<ConnectionRecord> {
self.begin_transport_replacement_commit_if_generation_current(
connection_id,
expected_transport_stable_id,
expected_transport_generation,
_expected_route_generation,
replacement_transport_stable_id,
transport_source,
|_| {},
)
.await
.map(TransportReplacementCommit::finish)
}
#[cfg(any(feature = "iroh-carrier-core", test))]
pub(crate) async fn begin_transport_replacement_commit_if_generation_current(
&self,
connection_id: &str,
expected_transport_stable_id: u64,
expected_transport_generation: u64,
_expected_route_generation: u64,
replacement_transport_stable_id: u64,
transport_source: Option<String>,
before_publish: impl FnOnce(Option<u64>),
) -> Option<TransportReplacementCommit> {
let mut store = self.store.clone().write_owned().await;
let updated = {
let record = store.by_id.get_mut(connection_id)?;
if !matches!(
record.state,
ConnectionState::Connecting | ConnectionState::Connected
) {
return None;
}
let already_installed = record.transport_stable_id
== Some(replacement_transport_stable_id)
&& record.transport_generation == expected_transport_generation.saturating_add(1);
if already_installed {
before_publish(record.transport_stable_id);
return Some(TransportReplacementCommit {
record: record.clone(),
_store: store,
});
}
let base_is_current = record.transport_stable_id == Some(expected_transport_stable_id)
&& record.transport_generation == expected_transport_generation;
let base_is_in_handoff = record.transport_stable_id.is_none()
&& record.transport_generation == expected_transport_generation
&& record.status_reason.as_deref()
== Some(crate::lifecycle_reason::REASON_REPLACEMENT_IN_PROGRESS);
if !base_is_current && !base_is_in_handoff {
return None;
}
before_publish(record.transport_stable_id);
let now = unix_ms_now();
record.transport_generation = record.transport_generation.saturating_add(1);
record.route_generation = 0;
if !replacement_source_proves_current_carrier(record, transport_source.as_deref()) {
clear_generation_bound_carrier_route(record);
}
record.replacement_count = record.replacement_count.saturating_add(1);
if record.state != ConnectionState::Connected {
record.transition_count = record.transition_count.saturating_add(1);
record.last_state_change_at_ms = now;
}
record.state = ConnectionState::Connected;
record.transport_stable_id = Some(replacement_transport_stable_id);
record.transport_source = transport_source;
record.status_reason = None;
record.last_disconnect_reason = None;
record.last_reconnect_reason = record.transport_source.clone();
record.last_transport_change_at_ms = now;
record.last_route_change_at_ms = now;
record.updated_at_ms = now;
record.clone()
};
store
.replacement_handoff_by_connection_id
.remove(connection_id);
let peer_id = canonical_peer_id(&updated);
store
.health_by_peer
.insert(peer_id, ConnectionHealth::Unknown);
Some(TransportReplacementCommit {
record: updated,
_store: store,
})
}
pub async fn set_closed_if_current(
&self,
connection_id: &str,
expected_transport_stable_id: u64,
reason: Option<String>,
) -> Option<ConnectionRecord> {
let mut store = self.store.write().await;
let updated = {
let record = store.by_id.get_mut(connection_id)?;
if record.transport_stable_id != Some(expected_transport_stable_id) {
return None;
}
let previous_state = record.state.clone();
let now = unix_ms_now();
record.state = ConnectionState::Closed;
record.transport_stable_id = None;
record.status_reason = reason;
record.updated_at_ms = now;
record.last_transport_change_at_ms = now;
record.last_route_change_at_ms = now;
if previous_state != ConnectionState::Closed {
record.last_state_change_at_ms = now;
record.transition_count = record.transition_count.saturating_add(1);
}
if !matches!(
previous_state,
ConnectionState::Closing | ConnectionState::Closed | ConnectionState::Failed
) {
record.retire_count = record.retire_count.saturating_add(1);
}
record.last_disconnect_reason = record.status_reason.clone();
record.clone()
};
store
.replacement_handoff_by_connection_id
.remove(connection_id);
let peer_id = canonical_peer_id(&updated);
store
.health_by_peer
.insert(peer_id, ConnectionHealth::Stale);
Some(updated)
}
pub(crate) async fn set_closed_if_current_or_replacement_handoff(
&self,
connection_id: &str,
expected_transport_stable_id: u64,
reason: Option<String>,
) -> Option<ConnectionRecord> {
let mut store = self.store.write().await;
let handoff_matches = {
let record = store.by_id.get(connection_id)?;
store
.replacement_handoff_by_connection_id
.get(connection_id)
.is_some_and(|owner| {
record.transport_stable_id.is_none()
&& owner.transport_stable_id == expected_transport_stable_id
&& owner.transport_generation == record.transport_generation
&& owner.route_generation == record.route_generation
&& record.status_reason.as_deref()
== Some(crate::lifecycle_reason::REASON_REPLACEMENT_IN_PROGRESS)
})
};
let updated = {
let record = store.by_id.get_mut(connection_id)?;
if record.transport_stable_id != Some(expected_transport_stable_id) && !handoff_matches
{
return None;
}
let previous_state = record.state.clone();
let now = unix_ms_now();
record.state = ConnectionState::Closed;
record.transport_stable_id = None;
record.status_reason = reason;
record.updated_at_ms = now;
record.last_transport_change_at_ms = now;
record.last_route_change_at_ms = now;
if previous_state != ConnectionState::Closed {
record.last_state_change_at_ms = now;
record.transition_count = record.transition_count.saturating_add(1);
}
if !matches!(
previous_state,
ConnectionState::Closing | ConnectionState::Closed | ConnectionState::Failed
) {
record.retire_count = record.retire_count.saturating_add(1);
}
record.last_disconnect_reason = record.status_reason.clone();
record.clone()
};
store
.replacement_handoff_by_connection_id
.remove(connection_id);
let peer_id = canonical_peer_id(&updated);
store
.health_by_peer
.insert(peer_id, ConnectionHealth::Stale);
Some(updated)
}
pub async fn set_failed_if_current(
&self,
connection_id: &str,
expected_transport_stable_id: u64,
reason: Option<String>,
) -> Option<ConnectionRecord> {
let mut store = self.store.write().await;
let updated = {
let record = store.by_id.get_mut(connection_id)?;
if record.transport_stable_id != Some(expected_transport_stable_id) {
return None;
}
let previous_state = record.state.clone();
let now = unix_ms_now();
record.state = ConnectionState::Failed;
record.transport_stable_id = None;
record.status_reason = reason;
record.updated_at_ms = now;
record.last_transport_change_at_ms = now;
record.last_route_change_at_ms = now;
if previous_state != ConnectionState::Failed {
record.last_state_change_at_ms = now;
record.transition_count = record.transition_count.saturating_add(1);
}
if !matches!(
previous_state,
ConnectionState::Closing | ConnectionState::Closed | ConnectionState::Failed
) {
record.retire_count = record.retire_count.saturating_add(1);
}
record.last_disconnect_reason = record.status_reason.clone();
record.clone()
};
store
.replacement_handoff_by_connection_id
.remove(connection_id);
let peer_id = canonical_peer_id(&updated);
store
.health_by_peer
.insert(peer_id, ConnectionHealth::Stale);
Some(updated)
}
pub(crate) async fn set_failed_if_dial_attempt_current(
&self,
connection_id: &str,
expected_transport_generation: u64,
expected_transport_stable_id: Option<u64>,
reason: Option<String>,
) -> Option<ConnectionRecord> {
let mut store = self.store.write().await;
let updated = {
let record = store.by_id.get_mut(connection_id)?;
if !matches!(
record.state,
ConnectionState::Pending | ConnectionState::Connecting
) || record.transport_generation != expected_transport_generation
|| record.transport_stable_id != expected_transport_stable_id
{
return None;
}
let now = unix_ms_now();
record.state = ConnectionState::Failed;
record.transport_stable_id = None;
record.status_reason = reason;
record.updated_at_ms = now;
record.last_state_change_at_ms = now;
record.last_transport_change_at_ms = now;
record.last_route_change_at_ms = now;
record.transition_count = record.transition_count.saturating_add(1);
record.retire_count = record.retire_count.saturating_add(1);
record.last_disconnect_reason = record.status_reason.clone();
record.clone()
};
store
.replacement_handoff_by_connection_id
.remove(connection_id);
let peer_id = canonical_peer_id(&updated);
store
.health_by_peer
.insert(peer_id, ConnectionHealth::Stale);
Some(updated)
}
pub async fn current_transport_matches(
&self,
connection_id: &str,
transport_stable_id: Option<u64>,
) -> bool {
let Some(record) = self.store.read().await.by_id.get(connection_id).cloned() else {
return false;
};
match transport_stable_id {
Some(stable_id) => record.transport_stable_id == Some(stable_id),
None => record.transport_stable_id.is_none(),
}
}
#[cfg(not(target_arch = "wasm32"))]
pub(crate) async fn replacement_handoff_matches(
&self,
connection_id: &str,
transport_stable_id: u64,
) -> bool {
self.store
.read()
.await
.replacement_handoff_by_connection_id
.get(connection_id)
.map(|owner| owner.transport_stable_id == transport_stable_id)
.unwrap_or(false)
}
pub async fn set_device_id(
&self,
connection_id: &str,
device_id: String,
) -> Option<ConnectionRecord> {
let mut store = self.store.write().await;
let existing = store.by_id.get(connection_id)?.clone();
let previous_peer_id = canonical_peer_id(&existing);
remove_index(
&mut store.by_device_id,
existing.device_id.as_deref(),
connection_id,
);
remove_index(
&mut store.by_device_hint_id,
existing.device_id_hint.as_deref(),
connection_id,
);
let mut updated = existing;
let prior_hint = updated
.device_id_hint
.clone()
.or_else(|| updated.device_id.clone());
updated.device_id = Some(device_id.clone());
if prior_hint.as_deref() == Some(device_id.as_str()) {
updated.device_id_hint = None;
} else {
updated.device_id_hint = prior_hint;
}
updated.updated_at_ms = unix_ms_now();
add_indexes(&mut store, &updated);
let next_peer_id = canonical_peer_id(&updated);
migrate_peer_state(&mut store, &previous_peer_id, &next_peer_id);
store
.by_id
.insert(connection_id.to_string(), updated.clone());
Some(updated)
}
pub async fn set_closing(
&self,
connection_id: &str,
reason: Option<String>,
) -> Option<ConnectionRecord> {
self.transition_state(connection_id, ConnectionState::Closing, reason)
.await
}
pub async fn set_closed(
&self,
connection_id: &str,
reason: Option<String>,
) -> Option<ConnectionRecord> {
self.transition_state(connection_id, ConnectionState::Closed, reason)
.await
}
pub async fn set_failed(
&self,
connection_id: &str,
reason: Option<String>,
) -> Option<ConnectionRecord> {
self.transition_state(connection_id, ConnectionState::Failed, reason)
.await
}
pub async fn remove(&self, connection_id: &str) -> Option<ConnectionRecord> {
let mut store = self.store.write().await;
let removed = store.by_id.remove(connection_id)?;
store
.replacement_handoff_by_connection_id
.remove(connection_id);
store.history_by_connection_id.insert(
connection_id.to_string(),
types::ConnectionHistory::from(&removed),
);
let peer_id = canonical_peer_id(&removed);
remove_index(
&mut store.by_node_id,
removed.node_id.as_deref(),
connection_id,
);
remove_index(
&mut store.by_device_id,
removed.device_id.as_deref(),
connection_id,
);
remove_index(
&mut store.by_device_hint_id,
removed.device_id_hint.as_deref(),
connection_id,
);
remove_index(
&mut store.by_endpoint_id,
removed.endpoint_id.as_deref(),
connection_id,
);
let peer_has_records = store
.by_id
.values()
.any(|r| canonical_peer_id(r) == peer_id);
if !peer_has_records {
store.health_by_peer.remove(&peer_id);
store.scopes_by_peer.remove(&peer_id);
}
Some(removed)
}
pub async fn get_by_connection_id(&self, connection_id: &str) -> Option<ConnectionRecord> {
self.store.read().await.by_id.get(connection_id).cloned()
}
#[cfg(test)]
pub async fn set_last_transport_change_at_ms_for_test(
&self,
connection_id: &str,
timestamp_ms: i64,
) -> Option<ConnectionRecord> {
let mut store = self.store.write().await;
let record = store.by_id.get_mut(connection_id)?;
record.last_transport_change_at_ms = timestamp_ms;
record.updated_at_ms = timestamp_ms.max(record.updated_at_ms);
Some(record.clone())
}
pub async fn get_by_node_id(&self, node_id: &str) -> Vec<ConnectionRecord> {
self.lookup_indexed(node_id, IndexType::Node).await
}
pub async fn get_by_device_id(&self, device_id: &str) -> Vec<ConnectionRecord> {
self.lookup_indexed(device_id, IndexType::Device).await
}
pub async fn get_by_endpoint_id(&self, endpoint_id: &str) -> Vec<ConnectionRecord> {
self.lookup_indexed(endpoint_id, IndexType::Endpoint).await
}
pub async fn list_all(&self) -> Vec<ConnectionRecord> {
self.store.read().await.by_id.values().cloned().collect()
}
pub async fn list_active(&self) -> Vec<ConnectionRecord> {
self.store
.read()
.await
.by_id
.values()
.filter(|record| {
matches!(
record.state,
ConnectionState::Pending
| ConnectionState::Connecting
| ConnectionState::Connected
)
})
.cloned()
.collect()
}
pub async fn add_scope(&self, id: &str, scope: &str) -> Vec<String> {
let mut store = self.store.write().await;
let Some(peer_id) = resolve_peer_id(&store, id) else {
return Vec::new();
};
let scopes = store.scopes_by_peer.entry(peer_id).or_default();
scopes.insert(scope.to_string());
scopes.iter().cloned().collect()
}
pub async fn release_scope(&self, id: &str, scope: Option<&str>) -> Vec<String> {
let mut store = self.store.write().await;
let Some(peer_id) = resolve_peer_id(&store, id) else {
return Vec::new();
};
let Some(scopes) = store.scopes_by_peer.get_mut(&peer_id) else {
return Vec::new();
};
match scope {
Some(scope) => {
scopes.remove(scope);
}
None => {
scopes.clear();
}
}
if scopes.is_empty() {
store.scopes_by_peer.remove(&peer_id);
return Vec::new();
}
scopes.iter().cloned().collect()
}
pub async fn get_scopes(&self, id: &str) -> Vec<String> {
let store = self.store.read().await;
let Some(peer_id) = resolve_peer_id(&store, id) else {
return Vec::new();
};
store
.scopes_by_peer
.get(&peer_id)
.map(|scopes| scopes.iter().cloned().collect())
.unwrap_or_default()
}
pub async fn are_same_peer(&self, left: &str, right: &str) -> bool {
let store = self.store.read().await;
let left_peer = resolve_peer_id(&store, left);
let right_peer = resolve_peer_id(&store, right);
left_peer.is_some() && left_peer == right_peer
}
#[cfg(test)]
pub(crate) async fn report_transport_status(
&self,
connection_id: &str,
active_transport: String,
parallel_transport: Option<String>,
) -> Option<PeerSnapshot> {
let mut store = self.store.write().await;
let record = store.by_id.get_mut(connection_id)?;
let transport_changed = record.active_transport != active_transport
|| record.parallel_transport != parallel_transport;
record.active_transport = active_transport;
record.parallel_transport = parallel_transport;
let now = unix_ms_now();
record.updated_at_ms = now;
if transport_changed {
record.route_generation = record.route_generation.saturating_add(1);
record.last_route_change_at_ms = now;
}
let peer_id = canonical_peer_id(record);
build_peer_snapshot(&store, &peer_id)
}
pub async fn report_transport_status_if_current(
&self,
connection_id: &str,
expected_transport_stable_id: u64,
expected_transport_generation: u64,
expected_route_generation: u64,
active_transport: String,
parallel_transport: Option<String>,
) -> Option<PeerSnapshot> {
let mut store = self.store.write().await;
let record = store.by_id.get_mut(connection_id)?;
if record.transport_stable_id != Some(expected_transport_stable_id)
|| record.transport_generation != expected_transport_generation
|| record.route_generation != expected_route_generation
{
return None;
}
let transport_changed = record.active_transport != active_transport
|| record.parallel_transport != parallel_transport;
record.active_transport = active_transport;
record.parallel_transport = parallel_transport;
let now = unix_ms_now();
record.updated_at_ms = now;
if transport_changed {
record.route_generation = record.route_generation.saturating_add(1);
record.last_route_change_at_ms = now;
}
let peer_id = canonical_peer_id(record);
build_peer_snapshot(&store, &peer_id)
}
pub async fn set_health(&self, id: &str, health: ConnectionHealth) -> Option<PeerSnapshot> {
let mut store = self.store.write().await;
let peer_id = resolve_peer_id(&store, id)?;
store.health_by_peer.insert(peer_id.clone(), health);
build_peer_snapshot(&store, &peer_id)
}
pub async fn set_health_if_current(
&self,
connection_id: &str,
expected_transport_stable_id: u64,
health: ConnectionHealth,
) -> Option<PeerSnapshot> {
let mut store = self.store.write().await;
let peer_id = {
let record = store.by_id.get_mut(connection_id)?;
if record.transport_stable_id != Some(expected_transport_stable_id) {
return None;
}
if matches!(health, ConnectionHealth::Healthy)
&& record.status_reason.as_deref()
== Some(crate::lifecycle_reason::REASON_REPLACEMENT_IN_PROGRESS)
{
record.status_reason = None;
record.updated_at_ms = unix_ms_now();
}
canonical_peer_id(record)
};
store.health_by_peer.insert(peer_id.clone(), health);
build_peer_snapshot(&store, &peer_id)
}
pub async fn set_settled_if_current(
&self,
connection_id: &str,
expected_transport_stable_id: u64,
expected_transport_generation: u64,
expected_route_generation: u64,
settled: bool,
) -> Option<PeerSnapshot> {
let mut store = self.store.write().await;
let record = store.by_id.get_mut(connection_id)?;
if record.transport_stable_id != Some(expected_transport_stable_id)
|| record.transport_generation != expected_transport_generation
|| record.route_generation != expected_route_generation
{
return None;
}
if settled
&& record.status_reason.as_deref()
== Some(crate::lifecycle_reason::REASON_REPLACEMENT_IN_PROGRESS)
{
record.status_reason = None;
record.updated_at_ms = unix_ms_now();
}
let peer_id = canonical_peer_id(record);
store.health_by_peer.insert(
peer_id.clone(),
if settled {
ConnectionHealth::Healthy
} else {
ConnectionHealth::Unknown
},
);
build_peer_snapshot(&store, &peer_id)
}
pub async fn clear_status_reason(
&self,
connection_id: &str,
expected_reason: Option<&str>,
) -> Option<PeerSnapshot> {
let mut store = self.store.write().await;
let record = store.by_id.get_mut(connection_id)?;
if expected_reason
.map(|reason| record.status_reason.as_deref() == Some(reason))
.unwrap_or(true)
{
record.status_reason = None;
record.updated_at_ms = unix_ms_now();
}
let peer_id = canonical_peer_id(record);
build_peer_snapshot(&store, &peer_id)
}
pub async fn report_transport_health(
&self,
id: &str,
transport: &str,
health: ConnectionHealth,
) -> Option<PeerSnapshot> {
let mut store = self.store.write().await;
let peer_id = resolve_peer_id(&store, id)?;
store
.transport_health_by_peer
.entry(peer_id.clone())
.or_default()
.insert(transport.to_string(), health);
let aggregated = peer_health(store.transport_health_by_peer.get(&peer_id).unwrap());
store.health_by_peer.insert(peer_id.clone(), aggregated);
build_peer_snapshot(&store, &peer_id)
}
pub async fn clear_transport_health(&self, id: &str) {
let mut store = self.store.write().await;
if let Some(peer_id) = resolve_peer_id(&store, id) {
store.transport_health_by_peer.remove(&peer_id);
}
}
pub async fn peer_snapshot(&self, id: &str) -> Option<PeerSnapshot> {
let store = self.store.read().await;
let peer_id = resolve_peer_id(&store, id)?;
build_peer_snapshot(&store, &peer_id)
}
pub async fn list_peer_snapshots(&self) -> Vec<PeerSnapshot> {
let store = self.store.read().await;
let peer_ids = collect_peer_ids(&store);
peer_ids
.iter()
.filter_map(|peer_id| build_peer_snapshot(&store, peer_id))
.collect()
}
pub async fn best_connection_for_peer(&self, id: &str) -> Option<ConnectionRecord> {
let store = self.store.read().await;
let peer_id = resolve_peer_id(&store, id)?;
let snapshot = build_peer_snapshot(&store, &peer_id);
let active_transport_stable_id = snapshot
.as_ref()
.and_then(|value| value.active_transport_stable_id);
let active_transport_generation = snapshot
.as_ref()
.map(|value| value.active_transport_generation)
.unwrap_or_default();
let active_node_id = snapshot.as_ref().and_then(|value| value.node_id.as_deref());
store
.by_id
.values()
.filter(|record| canonical_peer_id(record) == peer_id)
.filter(|record| matches!(record.state, ConnectionState::Connected))
.cloned()
.max_by(|left, right| {
let left_active = left.transport_stable_id == active_transport_stable_id
&& left.transport_generation == active_transport_generation;
let right_active = right.transport_stable_id == active_transport_stable_id
&& right.transport_generation == active_transport_generation;
let left_node_match = active_node_id == left.node_id.as_deref();
let right_node_match = active_node_id == right.node_id.as_deref();
left_active
.cmp(&right_active)
.then_with(|| left_node_match.cmp(&right_node_match))
.then_with(|| left.transport_generation.cmp(&right.transport_generation))
.then_with(|| left.updated_at_ms.cmp(&right.updated_at_ms))
})
}
async fn upsert(
&self,
connection_id: String,
node_id: Option<String>,
device_id_hint: Option<String>,
endpoint_id: Option<String>,
state: ConnectionState,
status_reason: Option<String>,
) -> ConnectionRecord {
self.upsert_when(
connection_id,
node_id,
device_id_hint,
endpoint_id,
state,
status_reason,
|_, _| true,
)
.await
.expect("unconditional connection upsert")
}
async fn upsert_when(
&self,
connection_id: String,
node_id: Option<String>,
device_id_hint: Option<String>,
endpoint_id: Option<String>,
state: ConnectionState,
status_reason: Option<String>,
allowed: impl FnOnce(Option<&ConnectionRecord>, Option<&types::ConnectionHistory>) -> bool,
) -> Option<ConnectionRecord> {
let mut store = self.store.write().await;
if !allowed(
store.by_id.get(&connection_id),
store.history_by_connection_id.get(&connection_id),
) {
return None;
}
if let Some(mut existing) = store.by_id.get(&connection_id).cloned() {
let previous_peer_id = canonical_peer_id(&existing);
remove_index(
&mut store.by_node_id,
existing.node_id.as_deref(),
&connection_id,
);
remove_index(
&mut store.by_device_id,
existing.device_id.as_deref(),
&connection_id,
);
remove_index(
&mut store.by_device_hint_id,
existing.device_id_hint.as_deref(),
&connection_id,
);
remove_index(
&mut store.by_endpoint_id,
existing.endpoint_id.as_deref(),
&connection_id,
);
if node_id.is_some() {
existing.node_id = node_id;
}
if let Some(device_id_hint) = device_id_hint {
if existing.device_id.is_none()
|| existing.device_id.as_deref() == Some(device_id_hint.as_str())
{
existing.device_id_hint = Some(device_id_hint);
}
}
if endpoint_id.is_some() {
existing.endpoint_id = endpoint_id;
}
let merged_state = merge_upsert_state(&existing.state, &state);
let merged_reason =
merge_upsert_reason(&existing, &merged_state, &state, status_reason);
if existing.state != merged_state {
existing.last_state_change_at_ms = unix_ms_now();
}
existing.state = merged_state;
existing.status_reason = merged_reason;
existing.updated_at_ms = unix_ms_now();
add_indexes(&mut store, &existing);
let next_peer_id = canonical_peer_id(&existing);
store
.health_by_peer
.entry(next_peer_id.clone())
.or_insert_with(|| default_health_for_state(&existing.state));
migrate_peer_state(&mut store, &previous_peer_id, &next_peer_id);
store.by_id.insert(connection_id, existing.clone());
return Some(existing);
}
let mut record =
ConnectionRecord::new(connection_id.clone(), node_id, device_id_hint, endpoint_id);
if let Some(history) = store.history_by_connection_id.get(&connection_id).cloned() {
record.transport_generation = history.transport_generation;
record.route_generation = history.route_generation;
record.transition_count = history.transition_count;
record.connecting_transition_count = history.connecting_transition_count;
record.replacement_count = history.replacement_count;
record.retire_count = history.retire_count;
record.last_disconnect_reason = history.last_disconnect_reason;
record.last_reconnect_reason = history.last_reconnect_reason;
}
if record.state != state {
record.last_state_change_at_ms = unix_ms_now();
}
record.state = state;
record.status_reason = status_reason;
add_indexes(&mut store, &record);
store
.health_by_peer
.entry(canonical_peer_id(&record))
.or_insert_with(|| default_health_for_state(&record.state));
store.by_id.insert(connection_id, record.clone());
Some(record)
}
async fn transition_state(
&self,
connection_id: &str,
next_state: ConnectionState,
status_reason: Option<String>,
) -> Option<ConnectionRecord> {
let mut store = self.store.write().await;
let mut record = store.by_id.get(connection_id)?.clone();
let peer_id = canonical_peer_id(&record);
let previous_state = record.state.clone();
record.state = next_state.clone();
record.status_reason = status_reason;
let now = unix_ms_now();
record.updated_at_ms = now;
if previous_state != next_state {
record.last_state_change_at_ms = now;
record.transition_count = record.transition_count.saturating_add(1);
if matches!(next_state, ConnectionState::Connecting) {
record.connecting_transition_count =
record.connecting_transition_count.saturating_add(1);
}
}
if matches!(
next_state,
ConnectionState::Closing | ConnectionState::Closed | ConnectionState::Failed
) {
store
.replacement_handoff_by_connection_id
.remove(connection_id);
if !matches!(
previous_state,
ConnectionState::Closing | ConnectionState::Closed | ConnectionState::Failed
) {
record.retire_count = record.retire_count.saturating_add(1);
}
record.last_disconnect_reason = record.status_reason.clone();
} else if matches!(next_state, ConnectionState::Connected) {
record.last_reconnect_reason = record.status_reason.clone();
}
match &next_state {
ConnectionState::Connected => {
store
.health_by_peer
.entry(peer_id)
.or_insert(ConnectionHealth::Unknown);
}
ConnectionState::Pending | ConnectionState::Connecting => {
store
.health_by_peer
.insert(peer_id, ConnectionHealth::Unknown);
}
ConnectionState::Closing | ConnectionState::Closed | ConnectionState::Failed => {
store
.health_by_peer
.insert(peer_id, ConnectionHealth::Stale);
}
}
store
.by_id
.insert(connection_id.to_string(), record.clone());
Some(record)
}
async fn lookup_indexed(&self, key: &str, index_type: IndexType) -> Vec<ConnectionRecord> {
let store = self.store.read().await;
let ids: Vec<String> = match index_type {
IndexType::Node => store
.by_node_id
.get(key)
.map(|values| values.iter().cloned().collect())
.unwrap_or_default(),
IndexType::Device => {
let mut ids = HashSet::new();
if let Some(values) = store.by_device_id.get(key) {
ids.extend(values.iter().cloned());
}
if let Some(values) = store.by_device_hint_id.get(key) {
ids.extend(values.iter().cloned());
}
ids.into_iter().collect()
}
IndexType::Endpoint => store
.by_endpoint_id
.get(key)
.map(|values| values.iter().cloned().collect())
.unwrap_or_default(),
};
ids.into_iter()
.filter_map(|id| store.by_id.get(&id).cloned())
.collect()
}
}
enum IndexType {
Node,
Device,
Endpoint,
}
#[cfg(test)]
mod tests {
use super::*;
#[tokio::test]
async fn incoming_ingress_bootstrap_and_final_fences_reject_latched_handoffs() {
let manager = Arc::new(ConnectionManager::new());
let connection_id = "browser-ingress-replacement";
assert!(
manager
.incoming_transport_generation_is_owned(connection_id, Some(10), 10)
.await,
"before a logical record exists, the physical registry remains the bootstrap owner",
);
let (bootstrap_snapshot_tx, bootstrap_snapshot_rx) = tokio::sync::oneshot::channel();
let (release_bootstrap_fence_tx, release_bootstrap_fence_rx) =
tokio::sync::oneshot::channel();
let bootstrap_manager = Arc::clone(&manager);
let bootstrap = tokio::spawn(async move {
let physical_transport_stable_id = Some(10);
bootstrap_snapshot_tx
.send(())
.expect("test observes bootstrap physical snapshot");
release_bootstrap_fence_rx
.await
.expect("test releases bootstrap final fence");
bootstrap_manager
.incoming_transport_generation_is_owned(
connection_id,
physical_transport_stable_id,
10,
)
.await
});
bootstrap_snapshot_rx
.await
.expect("bootstrap physical snapshot captured before logical handoff");
manager
.upsert_pending(
connection_id.to_string(),
Some("node-remote".to_string()),
Some("device-remote".to_string()),
Some("endpoint-remote".to_string()),
)
.await;
release_bootstrap_fence_tx
.send(())
.expect("release bootstrap final manager fence");
assert!(
!bootstrap.await.expect("bootstrap ingress task"),
"an ACK that observed the physical bootstrap generation before the logical handoff must still fail its final manager fence",
);
assert!(
!manager
.incoming_transport_generation_is_owned(connection_id, Some(10), 10)
.await,
"an existing record without a stable ID is a handoff and must fail closed",
);
manager
.set_connected_with_transport(
connection_id,
Some("endpoint-remote".to_string()),
Some(10),
Some("browser-iroh".to_string()),
)
.await
.expect("connected record");
let (physical_snapshot_tx, physical_snapshot_rx) = tokio::sync::oneshot::channel();
let (release_final_fence_tx, release_final_fence_rx) = tokio::sync::oneshot::channel();
let ingress_manager = Arc::clone(&manager);
let ingress = tokio::spawn(async move {
let physical_transport_stable_id = Some(11);
physical_snapshot_tx
.send(())
.expect("test observes physical snapshot");
release_final_fence_rx
.await
.expect("test releases final manager fence");
ingress_manager
.incoming_transport_generation_is_owned(
connection_id,
physical_transport_stable_id,
10,
)
.await
});
physical_snapshot_rx
.await
.expect("physical snapshot captured before handoff");
manager
.mark_transport_replaced_if_current(
connection_id,
10,
Some("replacement-handoff".to_string()),
Some(crate::lifecycle_reason::REASON_REPLACEMENT_IN_PROGRESS.to_string()),
)
.await
.expect("incumbent enters handoff");
release_final_fence_tx
.send(())
.expect("release ingress final fence");
assert!(
!ingress.await.expect("ingress task"),
"the final manager fence must reject the incumbent after a latched physical winner appears",
);
assert!(
!manager
.incoming_transport_generation_is_owned(connection_id, Some(11), 11)
.await,
"the candidate also remains blocked until the lifecycle owner commits it",
);
manager
.set_connected_with_transport(
connection_id,
Some("endpoint-remote".to_string()),
Some(11),
Some("replacement-committed".to_string()),
)
.await
.expect("replacement commit");
assert!(
manager
.incoming_transport_generation_is_owned(connection_id, Some(11), 11)
.await,
"the exact committed replacement is admitted",
);
}
#[tokio::test]
async fn indexes_stay_consistent_across_state_transitions() {
let manager = ConnectionManager::new();
let record = manager
.upsert_pending(
"conn-1".to_string(),
Some("node-A".to_string()),
Some("device-A".to_string()),
Some("endpoint-A".to_string()),
)
.await;
assert_eq!(record.state, ConnectionState::Pending);
assert_eq!(manager.get_by_node_id("node-A").await.len(), 1);
assert_eq!(manager.get_by_device_id("device-A").await.len(), 1);
assert_eq!(manager.get_by_endpoint_id("endpoint-A").await.len(), 1);
manager.set_connecting("conn-1").await;
manager
.set_connected("conn-1", Some("endpoint-B".to_string()))
.await;
let updated = manager
.get_by_connection_id("conn-1")
.await
.expect("updated record");
assert_eq!(updated.state, ConnectionState::Connected);
assert_eq!(updated.endpoint_id.as_deref(), Some("endpoint-B"));
assert_eq!(manager.get_by_endpoint_id("endpoint-A").await.len(), 0);
assert_eq!(manager.get_by_endpoint_id("endpoint-B").await.len(), 1);
manager
.set_closed("conn-1", Some("test-close".to_string()))
.await;
let closed = manager
.get_by_connection_id("conn-1")
.await
.expect("closed record");
assert_eq!(closed.state, ConnectionState::Closed);
assert_eq!(manager.list_active().await.len(), 0);
manager.remove("conn-1").await;
assert!(manager.get_by_connection_id("conn-1").await.is_none());
assert_eq!(manager.get_by_node_id("node-A").await.len(), 0);
assert_eq!(manager.get_by_device_id("device-A").await.len(), 0);
assert_eq!(manager.get_by_endpoint_id("endpoint-B").await.len(), 0);
}
#[tokio::test]
async fn set_device_id_reindexes_without_changing_state() {
let manager = ConnectionManager::new();
manager
.upsert_pending(
"conn-2".to_string(),
Some("node-A".to_string()),
Some("device-A".to_string()),
Some("endpoint-A".to_string()),
)
.await;
manager
.set_connected("conn-2", Some("endpoint-A".to_string()))
.await;
let updated = manager
.set_device_id("conn-2", "device-B".to_string())
.await
.expect("updated record");
assert_eq!(updated.state, ConnectionState::Connected);
assert_eq!(updated.device_id.as_deref(), Some("device-B"));
assert_eq!(updated.device_id_hint.as_deref(), Some("device-A"));
assert_eq!(manager.get_by_device_id("device-A").await.len(), 1);
assert_eq!(manager.get_by_device_id("device-B").await.len(), 1);
assert!(manager.are_same_peer("device-A", "device-B").await);
}
#[tokio::test]
async fn test_pending_to_connected_to_closed_lifecycle() {
let manager = ConnectionManager::new();
let record = manager
.upsert_pending(
"lc-1".to_string(),
Some("node-X".to_string()),
Some("dev-X".to_string()),
None,
)
.await;
assert_eq!(record.state, ConnectionState::Pending);
assert_eq!(manager.list_active().await.len(), 1);
let connecting = manager.set_connecting("lc-1").await.unwrap();
assert_eq!(connecting.state, ConnectionState::Connecting);
assert_eq!(manager.list_active().await.len(), 1);
let connected = manager
.set_connected("lc-1", Some("ep-X".to_string()))
.await
.unwrap();
assert_eq!(connected.state, ConnectionState::Connected);
assert_eq!(connected.endpoint_id.as_deref(), Some("ep-X"));
assert_eq!(manager.list_active().await.len(), 1);
let closing = manager
.set_closing("lc-1", Some("graceful".into()))
.await
.unwrap();
assert_eq!(closing.state, ConnectionState::Closing);
assert_eq!(manager.list_active().await.len(), 0);
let closed = manager
.set_closed("lc-1", Some("done".into()))
.await
.unwrap();
assert_eq!(closed.state, ConnectionState::Closed);
assert_eq!(manager.list_active().await.len(), 0);
assert!(manager.get_by_connection_id("lc-1").await.is_some());
assert_eq!(manager.get_by_node_id("node-X").await.len(), 1);
}
#[tokio::test]
async fn test_upsert_preserves_existing_fields() {
let manager = ConnectionManager::new();
manager
.upsert_pending(
"up-1".to_string(),
Some("node-1".to_string()),
Some("dev-1".to_string()),
Some("ep-1".to_string()),
)
.await;
let updated = manager
.upsert(
"up-1".to_string(),
None,
None,
None,
ConnectionState::Connected,
None,
)
.await;
assert_eq!(updated.state, ConnectionState::Connected);
assert_eq!(updated.node_id.as_deref(), Some("node-1"));
assert_eq!(updated.device_id, None);
assert_eq!(updated.device_id_hint.as_deref(), Some("dev-1"));
assert_eq!(updated.endpoint_id.as_deref(), Some("ep-1"));
}
#[tokio::test]
async fn test_set_failed_retains_indexes() {
let manager = ConnectionManager::new();
manager
.upsert_pending(
"fail-1".to_string(),
Some("node-F".to_string()),
Some("dev-F".to_string()),
None,
)
.await;
manager.set_failed("fail-1", Some("timeout".into())).await;
let failed = manager
.get_by_connection_id("fail-1")
.await
.expect("record exists");
assert_eq!(failed.state, ConnectionState::Failed);
assert_eq!(failed.status_reason.as_deref(), Some("timeout"));
assert_eq!(manager.get_by_node_id("node-F").await.len(), 1);
assert_eq!(manager.get_by_device_id("dev-F").await.len(), 1);
assert_eq!(manager.list_active().await.len(), 0);
}
#[tokio::test]
async fn test_list_active_excludes_terminal_states() {
let manager = ConnectionManager::new();
manager.upsert_pending("a-1".into(), None, None, None).await;
manager.upsert_pending("a-2".into(), None, None, None).await;
manager.upsert_pending("a-3".into(), None, None, None).await;
manager.set_connecting("a-1").await;
manager.set_connected("a-2", Some("ep".into())).await;
manager.set_failed("a-3", None).await;
let active = manager.list_active().await;
assert_eq!(active.len(), 2);
let ids: Vec<_> = active.iter().map(|r| r.connection_id.as_str()).collect();
assert!(ids.contains(&"a-1"));
assert!(ids.contains(&"a-2"));
assert!(!ids.contains(&"a-3"));
}
#[tokio::test]
async fn test_concurrent_connections_per_node() {
let manager = ConnectionManager::new();
manager
.upsert_pending(
"multi-1".into(),
Some("shared-node".into()),
Some("dev-A".into()),
None,
)
.await;
manager
.upsert_pending(
"multi-2".into(),
Some("shared-node".into()),
Some("dev-A".into()),
None,
)
.await;
let by_node = manager.get_by_node_id("shared-node").await;
assert_eq!(by_node.len(), 2);
let by_device = manager.get_by_device_id("dev-A").await;
assert_eq!(by_device.len(), 2);
manager.remove("multi-1").await;
assert_eq!(manager.get_by_node_id("shared-node").await.len(), 1);
assert_eq!(manager.get_by_device_id("dev-A").await.len(), 1);
}
#[tokio::test]
async fn test_connected_with_transport_tracks_generation_changes() {
let manager = ConnectionManager::new();
manager
.upsert_pending(
"conn-1".into(),
Some("node-1".into()),
Some("device-1".into()),
Some("endpoint-1".into()),
)
.await;
let first = manager
.set_connected_with_transport(
"conn-1",
Some("endpoint-1".into()),
Some(10),
Some("incoming".into()),
)
.await
.expect("record exists");
assert_eq!(first.transport_generation, 1);
assert_eq!(first.transport_stable_id, Some(10));
assert_eq!(first.transport_source.as_deref(), Some("incoming"));
let unchanged = manager
.set_connected_with_transport(
"conn-1",
Some("endpoint-1".into()),
Some(10),
Some("incoming-refresh".into()),
)
.await
.expect("record exists");
assert_eq!(unchanged.transport_generation, 1);
assert_eq!(unchanged.transport_stable_id, Some(10));
assert_eq!(
unchanged.transport_source.as_deref(),
Some("incoming-refresh")
);
manager
.report_transport_status(
"conn-1",
"webrtc".to_string(),
Some("iroh-quic".to_string()),
)
.await
.expect("route report");
assert_eq!(
manager
.get_by_connection_id("conn-1")
.await
.expect("route record")
.route_generation,
1,
);
let replaced = manager
.set_connected_with_transport(
"conn-1",
Some("endpoint-1".into()),
Some(11),
Some("replacement".into()),
)
.await
.expect("record exists");
assert_eq!(replaced.transport_generation, 2);
assert_eq!(replaced.route_generation, 0);
assert_eq!(replaced.transport_stable_id, Some(11));
assert_eq!(replaced.transport_source.as_deref(), Some("replacement"));
assert!(manager.current_transport_matches("conn-1", Some(11)).await);
assert!(!manager.current_transport_matches("conn-1", Some(10)).await);
}
#[tokio::test]
async fn generation_fenced_replacement_accepts_close_handoff_and_rejects_stale_attempts() {
let manager = ConnectionManager::new();
manager
.upsert_pending(
"conn-replacement".into(),
Some("node-replacement".into()),
Some("device-replacement".into()),
Some("endpoint-replacement".into()),
)
.await;
let base = manager
.set_connected_with_transport(
"conn-replacement",
Some("endpoint-replacement".into()),
Some(20),
Some("iroh-relay".into()),
)
.await
.expect("base transport");
manager
.mark_transport_replaced_if_current(
"conn-replacement",
20,
Some("incoming-replacement-churn".into()),
Some(crate::lifecycle_reason::REASON_REPLACEMENT_IN_PROGRESS.into()),
)
.await
.expect("base enters the coordinated handoff window");
let replacement = manager
.install_transport_replacement_if_generation_current(
"conn-replacement",
20,
base.transport_generation,
base.route_generation,
21,
Some("ble-upgrade".into()),
)
.await
.expect("exact replacement may cross the close handoff");
assert_eq!(replacement.transport_stable_id, Some(21));
assert_eq!(
replacement.transport_generation,
base.transport_generation + 1
);
assert_eq!(replacement.route_generation, 0);
let idempotent = manager
.install_transport_replacement_if_generation_current(
"conn-replacement",
20,
base.transport_generation,
base.route_generation,
21,
Some("ble-upgrade".into()),
)
.await
.expect("same replacement callback is idempotent");
assert_eq!(
idempotent.transport_generation,
replacement.transport_generation
);
assert!(
manager
.install_transport_replacement_if_generation_current(
"conn-replacement",
20,
base.transport_generation,
base.route_generation,
22,
Some("stale-ble-upgrade".into()),
)
.await
.is_none(),
"a delayed replacement cannot overwrite the winning generation",
);
}
#[tokio::test]
async fn replacement_commit_fences_logical_retirement_until_physical_install_finishes() {
let manager = Arc::new(ConnectionManager::new());
manager
.upsert_pending(
"conn-atomic-replacement".into(),
Some("node-atomic-replacement".into()),
Some("device-atomic-replacement".into()),
Some("endpoint-atomic-replacement".into()),
)
.await;
let base = manager
.set_connected_with_transport(
"conn-atomic-replacement",
Some("endpoint-atomic-replacement".into()),
Some(40),
Some("iroh-relay".into()),
)
.await
.expect("base transport");
let commit = manager
.begin_transport_replacement_commit_if_generation_current(
"conn-atomic-replacement",
40,
base.transport_generation,
base.route_generation,
41,
Some("webrtc-upgrade".into()),
|_| {},
)
.await
.expect("logical replacement commit");
let manager_for_retirement = manager.clone();
let retirement = tokio::spawn(async move {
manager_for_retirement
.remove("conn-atomic-replacement")
.await
});
tokio::task::yield_now().await;
assert!(
!retirement.is_finished(),
"logical retirement must wait while physical insertion owns the commit guard",
);
let committed = commit.finish();
assert_eq!(committed.transport_stable_id, Some(41));
let removed = retirement
.await
.expect("retirement task")
.expect("retired committed record");
assert_eq!(removed.transport_stable_id, Some(41));
assert!(
manager
.get_by_connection_id("conn-atomic-replacement")
.await
.is_none(),
"retirement proceeds only after the replacement can exist physically",
);
}
#[tokio::test]
async fn physical_replacement_survives_stale_carrier_observations() {
let manager = ConnectionManager::new();
manager
.upsert_pending(
"conn-route-refresh".into(),
Some("node-route-refresh".into()),
Some("device-route-refresh".into()),
Some("endpoint-route-refresh".into()),
)
.await;
let base = manager
.set_connected_with_transport(
"conn-route-refresh",
Some("endpoint-route-refresh".into()),
Some(30),
Some("iroh-relay".into()),
)
.await
.expect("base transport");
let observed = manager
.report_transport_status_if_current(
"conn-route-refresh",
30,
base.transport_generation,
base.route_generation,
"iroh-direct".into(),
None,
)
.await
.expect("route observation");
assert!(observed.active_route_generation > base.route_generation);
let replacement = manager
.install_transport_replacement_if_generation_current(
"conn-route-refresh",
30,
base.transport_generation,
base.route_generation,
31,
Some("webrtc-upgrade".into()),
)
.await
.expect("same physical generation remains replaceable");
assert_eq!(replacement.transport_stable_id, Some(31));
assert_eq!(
replacement.transport_generation,
base.transport_generation + 1
);
assert_eq!(replacement.route_generation, 0);
}
#[tokio::test]
async fn proven_carrier_replacement_preserves_its_route_during_atomic_commit() {
let manager = ConnectionManager::new();
manager
.upsert_pending(
"conn-carrier-commit".into(),
Some("node-carrier-commit".into()),
Some("device-carrier-commit".into()),
Some("endpoint-carrier-commit".into()),
)
.await;
let base = manager
.set_connected_with_transport(
"conn-carrier-commit",
Some("endpoint-carrier-commit".into()),
Some(60),
Some("iroh-relay".into()),
)
.await
.expect("base transport");
manager
.report_transport_status(
"conn-carrier-commit",
"moq".into(),
Some("iroh-relay".into()),
)
.await
.expect("proven carrier observation");
let replacement = manager
.install_transport_replacement_if_generation_current(
"conn-carrier-commit",
60,
base.transport_generation,
base.route_generation,
61,
Some("moq-upgrade".into()),
)
.await
.expect("proven carrier replacement");
assert_eq!(replacement.transport_stable_id, Some(61));
assert_eq!(replacement.active_transport, "moq");
assert_eq!(
replacement.parallel_transport.as_deref(),
Some("iroh-relay")
);
}
#[tokio::test]
async fn report_transport_status_updates_peer_snapshot_without_replacing_transport() {
let manager = ConnectionManager::new();
manager
.upsert_pending(
"conn-status".into(),
Some("node-status".into()),
Some("device-status".into()),
Some("endpoint-status".into()),
)
.await;
manager
.set_connected_with_transport(
"conn-status",
Some("endpoint-status".into()),
Some(42),
Some("incoming".into()),
)
.await;
let before = manager
.get_by_connection_id("conn-status")
.await
.expect("record exists before status report")
.last_transport_change_at_ms;
tokio::time::sleep(std::time::Duration::from_millis(1)).await;
let snapshot = manager
.report_transport_status(
"conn-status",
"webrtc".to_string(),
Some("iroh".to_string()),
)
.await
.expect("peer snapshot");
assert_eq!(snapshot.active_transport, "webrtc");
assert_eq!(snapshot.parallel_transport.as_deref(), Some("iroh"));
assert_eq!(snapshot.active_transport_stable_id, Some(42));
assert_eq!(snapshot.active_transport_generation, 1);
assert_eq!(snapshot.active_route_generation, 1);
let record = manager
.get_by_connection_id("conn-status")
.await
.expect("record exists after status report");
assert_eq!(record.transport_stable_id, Some(42));
assert_eq!(record.replacement_count, 1);
let after = manager
.get_by_connection_id("conn-status")
.await
.expect("record exists after status report");
assert_eq!(after.last_transport_change_at_ms, before);
assert!(after.last_route_change_at_ms >= before);
}
#[tokio::test]
async fn route_migration_is_idempotent_and_does_not_churn_the_connection_lifecycle() {
let manager = ConnectionManager::new();
manager
.upsert_pending(
"conn-path".into(),
Some("node-path".into()),
Some("device-path".into()),
Some("endpoint-path".into()),
)
.await;
manager
.set_connected_with_transport(
"conn-path",
Some("endpoint-path".into()),
Some(7),
Some("incoming".into()),
)
.await;
let baseline = manager
.get_by_connection_id("conn-path")
.await
.expect("connected baseline");
let routes = [
("iroh-relay", None),
("iroh-lan", Some("iroh-relay")),
("iroh-quic", Some("iroh-relay")),
("webrtc", Some("iroh-quic")),
("moq", Some("iroh-quic")),
("iroh-quic", Some("moq")),
];
for (index, (active, parallel)) in routes.into_iter().enumerate() {
let snapshot = manager
.report_transport_status(
"conn-path",
active.to_string(),
parallel.map(str::to_string),
)
.await
.expect("route snapshot");
assert_eq!(snapshot.status, ConnectionState::Connected);
assert_eq!(snapshot.active_transport, active);
assert_eq!(snapshot.parallel_transport.as_deref(), parallel);
assert_eq!(snapshot.active_transport_stable_id, Some(7));
assert_eq!(snapshot.transition_count, baseline.transition_count);
assert_eq!(snapshot.replacement_count, baseline.replacement_count);
assert_eq!(
snapshot.active_transport_generation, baseline.transport_generation,
"route changes must not impersonate physical transport replacement",
);
assert_eq!(
snapshot.active_route_generation,
baseline.route_generation + index as u64 + 1,
);
}
let before_duplicate = manager
.get_by_connection_id("conn-path")
.await
.expect("record before duplicate route report");
let duplicate = manager
.report_transport_status(
"conn-path",
"iroh-quic".to_string(),
Some("moq".to_string()),
)
.await
.expect("duplicate route snapshot");
assert_eq!(
duplicate.active_route_generation, before_duplicate.route_generation,
"duplicate route reports must be idempotent",
);
let record = manager
.get_by_connection_id("conn-path")
.await
.expect("record exists after path migration");
assert_eq!(record.state, ConnectionState::Connected);
assert_eq!(record.transport_stable_id, Some(7));
assert_eq!(
record.last_transport_change_at_ms,
baseline.last_transport_change_at_ms,
);
assert_eq!(record.transition_count, baseline.transition_count);
assert_eq!(record.replacement_count, baseline.replacement_count);
}
#[tokio::test]
async fn idempotent_route_observation_does_not_advance_lifecycle_time() {
let manager = ConnectionManager::new();
manager
.upsert_pending(
"conn-observation".into(),
Some("node-observation".into()),
Some("device-observation".into()),
Some("endpoint-observation".into()),
)
.await;
manager
.set_connected_with_transport(
"conn-observation",
Some("endpoint-observation".into()),
Some(9),
Some("test".into()),
)
.await;
let first = manager
.report_transport_status("conn-observation", "iroh-lan".into(), None)
.await
.expect("initial route observation");
tokio::time::sleep(std::time::Duration::from_millis(2)).await;
let repeated = manager
.report_transport_status("conn-observation", "iroh-lan".into(), None)
.await
.expect("repeated route observation");
assert_eq!(
repeated.active_route_generation,
first.active_route_generation
);
assert_eq!(
repeated.last_lifecycle_transition_at_ms,
first.last_lifecycle_transition_at_ms,
);
assert!(repeated.last_seen_at_ms >= first.last_seen_at_ms);
}
#[tokio::test]
async fn test_mark_transport_replaced_advances_generation_without_state_reset() {
let manager = ConnectionManager::new();
manager
.upsert_pending(
"conn-2".into(),
Some("node-2".into()),
Some("device-2".into()),
Some("endpoint-2".into()),
)
.await;
manager
.set_connected_with_transport(
"conn-2",
Some("endpoint-2".into()),
Some(20),
Some("outgoing".into()),
)
.await;
let replaced = manager
.mark_transport_replaced(
"conn-2",
Some(21),
Some("incoming-replacement".into()),
Some("replacement-in-progress".into()),
)
.await
.expect("record exists");
assert_eq!(replaced.transport_generation, 2);
assert_eq!(replaced.transport_stable_id, Some(21));
assert_eq!(
replaced.transport_source.as_deref(),
Some("incoming-replacement")
);
assert_eq!(
replaced.status_reason.as_deref(),
Some("replacement-in-progress")
);
assert_eq!(replaced.state, ConnectionState::Connected);
let peer = manager
.peer_snapshot("device-2")
.await
.expect("peer snapshot");
assert_eq!(peer.health, ConnectionHealth::Unknown);
}
#[tokio::test]
async fn physical_replacement_expires_generation_bound_carrier_route() {
let manager = ConnectionManager::new();
manager
.upsert_pending(
"conn-carrier-replacement".into(),
Some("node-carrier-replacement".into()),
Some("device-carrier-replacement".into()),
Some("endpoint-carrier-replacement".into()),
)
.await;
manager
.set_connected_with_transport(
"conn-carrier-replacement",
Some("endpoint-carrier-replacement".into()),
Some(40),
Some("outgoing".into()),
)
.await;
manager
.report_transport_status(
"conn-carrier-replacement",
"webrtc".into(),
Some("iroh-relay".into()),
)
.await
.expect("carrier route observation");
let replaced = manager
.mark_transport_replaced(
"conn-carrier-replacement",
Some(41),
Some("incoming-replacement".into()),
None,
)
.await
.expect("replacement record");
assert_eq!(replaced.transport_generation, 2);
assert_eq!(replaced.route_generation, 0);
assert_eq!(replaced.active_transport, "iroh");
assert_eq!(replaced.parallel_transport, None);
}
#[tokio::test]
async fn clear_status_reason_removes_matching_transient_reason() {
let manager = ConnectionManager::new();
manager
.upsert_pending(
"conn-clear-reason".into(),
Some("node-clear-reason".into()),
Some("device-clear-reason".into()),
Some("endpoint-clear-reason".into()),
)
.await;
manager
.set_connected_with_transport(
"conn-clear-reason",
Some("endpoint-clear-reason".into()),
Some(30),
Some("outgoing".into()),
)
.await;
manager
.mark_transport_replaced(
"conn-clear-reason",
Some(31),
Some("wasm-accept-closed".into()),
Some("replacement-in-progress".into()),
)
.await
.expect("replacement record");
manager
.clear_status_reason("conn-clear-reason", Some("replacement-in-progress"))
.await
.expect("peer snapshot");
let record = manager
.get_by_connection_id("conn-clear-reason")
.await
.expect("record exists");
assert_eq!(record.status_reason.as_deref(), None);
manager
.mark_transport_replaced(
"conn-clear-reason",
Some(32),
Some("wasm-accept-closed".into()),
Some("replacement-in-progress".into()),
)
.await
.expect("replacement record");
manager
.clear_status_reason("conn-clear-reason", Some("different-reason"))
.await
.expect("peer snapshot");
let record = manager
.get_by_connection_id("conn-clear-reason")
.await
.expect("record exists");
assert_eq!(
record.status_reason.as_deref(),
Some("replacement-in-progress")
);
}
#[tokio::test]
async fn mark_transport_replaced_without_new_binding_preserves_generation() {
let manager = ConnectionManager::new();
manager
.upsert_pending(
"conn-2b".into(),
Some("node-2b".into()),
Some("device-2b".into()),
Some("endpoint-2b".into()),
)
.await;
manager
.set_connected_with_transport(
"conn-2b",
Some("endpoint-2b".into()),
Some(30),
Some("outgoing".into()),
)
.await
.expect("connected record");
let replaced = manager
.mark_transport_replaced(
"conn-2b",
None,
Some("outgoing-closed".into()),
Some("replacement-in-progress".into()),
)
.await
.expect("record exists");
assert_eq!(replaced.transport_generation, 1);
assert_eq!(replaced.transport_stable_id, None);
assert_eq!(
replaced.status_reason.as_deref(),
Some("replacement-in-progress")
);
assert_eq!(replaced.replacement_count, 1);
assert_eq!(replaced.state, ConnectionState::Connected);
}
#[tokio::test]
async fn stale_transport_loss_cannot_clear_a_replacement_generation() {
let manager = ConnectionManager::new();
manager
.upsert_pending(
"conn-generation-fence".into(),
Some("node-generation-fence".into()),
Some("device-generation-fence".into()),
Some("endpoint-generation-fence".into()),
)
.await;
manager
.set_connected_with_transport(
"conn-generation-fence",
Some("endpoint-generation-fence".into()),
Some(41),
Some("initial".into()),
)
.await
.expect("initial transport");
manager
.set_connected_with_transport(
"conn-generation-fence",
Some("endpoint-generation-fence".into()),
Some(42),
Some("replacement".into()),
)
.await
.expect("replacement transport");
assert!(manager
.mark_transport_replaced_if_current(
"conn-generation-fence",
41,
Some("delayed-old-close".into()),
Some("replacement-in-progress".into()),
)
.await
.is_none());
assert!(manager
.set_closed_if_current(
"conn-generation-fence",
41,
Some("delayed-terminal-close".into()),
)
.await
.is_none());
assert!(manager
.report_transport_status_if_current(
"conn-generation-fence",
41,
1,
0,
"webrtc".into(),
Some("iroh-relay".into()),
)
.await
.is_none());
assert!(manager
.set_settled_if_current("conn-generation-fence", 41, 1, 0, true)
.await
.is_none());
let current = manager
.get_by_connection_id("conn-generation-fence")
.await
.expect("current record");
assert_eq!(current.transport_stable_id, Some(42));
assert_eq!(current.transport_source.as_deref(), Some("replacement"));
assert_eq!(current.status_reason, None);
assert_eq!(current.active_transport, "iroh");
assert_eq!(
manager
.peer_snapshot("conn-generation-fence")
.await
.expect("replacement snapshot")
.health,
ConnectionHealth::Unknown,
);
let cleared = manager
.mark_transport_replaced_if_current(
"conn-generation-fence",
42,
Some("current-close".into()),
Some("replacement-in-progress".into()),
)
.await
.expect("current generation cleared");
assert_eq!(cleared.transport_stable_id, None);
assert_eq!(cleared.transport_source.as_deref(), Some("current-close"));
}
#[tokio::test]
async fn stale_dial_failure_cannot_terminalize_an_admitted_successor() {
let manager = ConnectionManager::new();
let connection_id = "conn-stale-dial-failure";
manager
.upsert_pending(
connection_id.into(),
Some("node-stale-dial-failure".into()),
Some("device-stale-dial-failure".into()),
Some("endpoint-stale-dial-failure".into()),
)
.await;
let attempt = manager
.set_connecting(connection_id)
.await
.expect("connecting attempt");
manager
.set_connected_with_transport(
connection_id,
Some("endpoint-stale-dial-failure".into()),
Some(77),
Some("accepted-inbound".into()),
)
.await
.expect("successor transport");
assert!(
manager
.set_failed_if_dial_attempt_current(
connection_id,
attempt.transport_generation,
attempt.transport_stable_id,
Some("stale timeout".into()),
)
.await
.is_none(),
"an older dial attempt must not fail the accepted successor",
);
assert!(
manager
.commit_admitted_transport(
connection_id,
Some("endpoint-stale-dial-failure".into()),
Some("device-stale-dial-failure".into()),
77,
Some("incoming-admission".into()),
|_, _| true,
|_| {},
)
.await
.is_some(),
"the current authenticated successor must remain publishable",
);
let current_attempt = manager
.set_connecting(connection_id)
.await
.expect("current connecting attempt");
assert!(
manager
.set_failed_if_dial_attempt_current(
connection_id,
current_attempt.transport_generation,
current_attempt.transport_stable_id,
Some("current timeout".into()),
)
.await
.is_some(),
"a timeout still owns and fails its unchanged connecting attempt",
);
}
#[tokio::test]
async fn authenticated_inbound_generation_reopens_concurrent_dial_timeout() {
let manager = ConnectionManager::new();
let connection_id = "conn-inbound-after-dial-timeout";
manager
.upsert_pending(
connection_id.into(),
Some("endpoint-inbound-after-timeout".into()),
Some("device-inbound-after-timeout".into()),
Some("endpoint-inbound-after-timeout".into()),
)
.await;
let attempt = manager
.set_connecting(connection_id)
.await
.expect("outgoing attempt");
manager
.set_failed_if_dial_attempt_current(
connection_id,
attempt.transport_generation,
attempt.transport_stable_id,
Some(format!(
"{}:ensure_connected_addr-timeout",
crate::lifecycle_reason::REASON_AUTO_CONNECT_FAILURE
)),
)
.await
.expect("outgoing attempt times out before inbound admission commits");
let admitted = manager
.commit_admitted_transport(
connection_id,
Some("endpoint-inbound-after-timeout".into()),
Some("device-inbound-after-timeout".into()),
5,
Some("incoming-admission".into()),
|requires_authenticated_peer_reconnect, previous| {
assert!(!requires_authenticated_peer_reconnect);
assert_eq!(previous, None);
true
},
|_| {},
)
.await
.expect("fresh authenticated inbound generation must win the dial-timeout race");
assert_eq!(admitted.state, ConnectionState::Connected);
assert_eq!(admitted.transport_stable_id, Some(5));
assert_eq!(admitted.status_reason, None);
}
#[tokio::test]
async fn delayed_physical_commit_cannot_resurrect_terminal_session() {
let manager = ConnectionManager::new();
let connection_id = "conn-terminal-physical-commit";
manager
.upsert_pending(
connection_id.into(),
Some("node-terminal".into()),
Some("device-terminal".into()),
Some("endpoint-terminal".into()),
)
.await;
manager
.commit_current_transport(
connection_id,
Some("endpoint-terminal".into()),
61,
Some("initial".into()),
|_| {},
)
.await
.expect("initial transport");
manager
.set_closed_if_current(connection_id, 61, Some("manual disconnect".into()))
.await
.expect("terminal close");
let stale_security_rebind = std::sync::Arc::new(std::sync::atomic::AtomicBool::new(false));
let stale_security_rebind_for_commit = stale_security_rebind.clone();
assert!(manager
.commit_current_transport(
connection_id,
Some("endpoint-terminal".into()),
61,
Some("delayed-admission".into()),
move |_| {
stale_security_rebind_for_commit
.store(true, std::sync::atomic::Ordering::SeqCst);
},
)
.await
.is_none());
assert!(!stale_security_rebind.load(std::sync::atomic::Ordering::SeqCst));
let terminal = manager
.get_by_connection_id(connection_id)
.await
.expect("terminal record remains");
assert_eq!(terminal.state, ConnectionState::Closed);
assert_eq!(terminal.transport_stable_id, None);
assert_eq!(
terminal.last_disconnect_reason.as_deref(),
Some("manual disconnect")
);
manager
.upsert_pending(
connection_id.into(),
Some("node-terminal".into()),
Some("device-terminal".into()),
Some("endpoint-terminal".into()),
)
.await;
let reconnected = manager
.commit_current_transport(
connection_id,
Some("endpoint-terminal".into()),
62,
Some("explicit-reconnect".into()),
|_| {},
)
.await
.expect("explicit reconnect transport");
assert_eq!(reconnected.state, ConnectionState::Connected);
assert_eq!(reconnected.transport_stable_id, Some(62));
}
#[tokio::test]
async fn rejected_admission_commit_preserves_terminal_record_and_reconnect_intent() {
let manager = ConnectionManager::new();
let connection_id = "conn-rejected-admission-commit";
manager
.upsert_pending(
connection_id.into(),
Some("node-rejected-admission".into()),
Some("device-rejected-admission".into()),
Some("endpoint-rejected-admission".into()),
)
.await;
manager
.commit_current_transport(
connection_id,
Some("endpoint-rejected-admission".into()),
71,
Some("initial".into()),
|_| {},
)
.await
.expect("initial transport");
manager.add_scope(connection_id, "drive-grant:old").await;
manager
.set_closed_if_current(
connection_id,
71,
Some(crate::lifecycle_reason::REASON_MANUAL_DISCONNECT.into()),
)
.await
.expect("manual terminal");
let exclusion_consumed = std::sync::Arc::new(std::sync::atomic::AtomicBool::new(false));
let exclusion_consumed_for_commit = exclusion_consumed.clone();
let scopes_normalized = std::sync::Arc::new(std::sync::atomic::AtomicBool::new(false));
let scopes_normalized_for_commit = scopes_normalized.clone();
assert!(manager
.commit_admitted_transport(
connection_id,
Some("endpoint-rejected-admission".into()),
Some("device-rejected-admission".into()),
72,
Some("rejected-admission".into()),
move |requires_authenticated_peer_reconnect, _| {
assert!(requires_authenticated_peer_reconnect);
let _ = &exclusion_consumed_for_commit;
false
},
move |_| {
scopes_normalized_for_commit.store(true, std::sync::atomic::Ordering::SeqCst);
},
)
.await
.is_none());
assert!(!exclusion_consumed.load(std::sync::atomic::Ordering::SeqCst));
assert!(!scopes_normalized.load(std::sync::atomic::Ordering::SeqCst));
let terminal = manager
.get_by_connection_id(connection_id)
.await
.expect("terminal record remains");
assert_eq!(terminal.state, ConnectionState::Closed);
assert_eq!(terminal.transport_stable_id, None);
assert_eq!(
manager.get_scopes(connection_id).await,
vec!["drive-grant:old".to_string()]
);
}
#[tokio::test]
async fn admitted_scope_class_is_replaced_before_connected_snapshot_is_visible() {
let manager = ConnectionManager::new();
let connection_id = "conn-atomic-scope-class";
manager
.upsert_pending(
connection_id.into(),
Some("node-atomic-scope".into()),
Some("device-hint-atomic-scope".into()),
Some("endpoint-atomic-scope".into()),
)
.await;
manager.add_scope(connection_id, "drive-grant:old").await;
let connected = manager
.commit_admitted_transport(
connection_id,
Some("endpoint-atomic-scope".into()),
Some("device-atomic-scope".into()),
73,
Some("admitted-user-device".into()),
|requires_authenticated_peer_reconnect, _| !requires_authenticated_peer_reconnect,
|scopes| {
scopes.clear();
scopes.insert("user-device".to_string());
},
)
.await
.expect("admitted transport");
assert_eq!(connected.state, ConnectionState::Connected);
assert_eq!(connected.device_id.as_deref(), Some("device-atomic-scope"));
assert_eq!(
manager
.peer_snapshot(connection_id)
.await
.expect("settled snapshot")
.scopes,
vec!["user-device".to_string()]
);
}
#[tokio::test]
async fn authenticated_peer_reconnect_reopens_only_peer_requested_manual_terminal() {
let manager = ConnectionManager::new();
let connection_id = "conn-authenticated-peer-reconnect";
manager
.upsert_pending(
connection_id.into(),
Some("node-peer-reconnect".into()),
Some("device-peer-reconnect".into()),
Some("endpoint-peer-reconnect".into()),
)
.await;
manager
.commit_current_transport(
connection_id,
Some("endpoint-peer-reconnect".into()),
81,
Some("initial".into()),
|_| {},
)
.await
.expect("initial transport");
manager
.set_closed_if_current(
connection_id,
81,
Some(
"ApplicationClosed(ApplicationClose { error_code: 1, reason: b\"manual disconnect\" })"
.into(),
),
)
.await
.expect("peer-requested close");
let peer_exclusion_consumed =
std::sync::Arc::new(std::sync::atomic::AtomicBool::new(false));
let peer_exclusion_consumed_for_commit = peer_exclusion_consumed.clone();
let security_committed = std::sync::Arc::new(std::sync::atomic::AtomicBool::new(false));
let security_committed_for_commit = security_committed.clone();
let reopened = manager
.commit_authenticated_peer_reconnect_transport(
connection_id,
Some("endpoint-peer-reconnect".into()),
Some("device-peer-reconnect".into()),
82,
Some("authenticated-peer-reconnect".into()),
move |requires_authenticated_peer_reconnect, _| {
assert!(requires_authenticated_peer_reconnect);
security_committed_for_commit.store(true, std::sync::atomic::Ordering::SeqCst);
peer_exclusion_consumed_for_commit
.store(true, std::sync::atomic::Ordering::SeqCst);
true
},
)
.await
.expect("authenticated peer reconnect");
assert!(peer_exclusion_consumed.load(std::sync::atomic::Ordering::SeqCst));
assert!(security_committed.load(std::sync::atomic::Ordering::SeqCst));
assert_eq!(reopened.state, ConnectionState::Connected);
assert_eq!(reopened.transport_stable_id, Some(82));
assert_eq!(reopened.device_id.as_deref(), Some("device-peer-reconnect"));
manager
.set_closed_if_current(
connection_id,
82,
Some(crate::lifecycle_reason::REASON_MANUAL_DISCONNECT.into()),
)
.await
.expect("local manual close");
let rejected_security_commit =
std::sync::Arc::new(std::sync::atomic::AtomicBool::new(false));
let rejected_security_commit_for_commit = rejected_security_commit.clone();
assert!(manager
.commit_authenticated_peer_reconnect_transport(
connection_id,
Some("endpoint-peer-reconnect".into()),
Some("device-peer-reconnect".into()),
83,
Some("unauthorized-peer-reconnect".into()),
move |requires_authenticated_peer_reconnect, _| {
assert!(requires_authenticated_peer_reconnect);
let _ = &rejected_security_commit_for_commit;
false
},
)
.await
.is_none());
assert!(!rejected_security_commit.load(std::sync::atomic::Ordering::SeqCst));
assert_eq!(
manager
.get_by_connection_id(connection_id)
.await
.expect("local terminal remains")
.state,
ConnectionState::Closed,
);
manager
.set_failed(
connection_id,
Some(crate::lifecycle_reason::REASON_SESSION_TOKEN_REVOKED.into()),
)
.await
.expect("revoked terminal");
let revoked_reopen_checked = std::sync::Arc::new(std::sync::atomic::AtomicBool::new(false));
let revoked_reopen_checked_for_commit = revoked_reopen_checked.clone();
assert!(manager
.commit_authenticated_peer_reconnect_transport(
connection_id,
Some("endpoint-peer-reconnect".into()),
Some("device-peer-reconnect".into()),
84,
Some("revoked-peer-reconnect".into()),
move |_, _| {
revoked_reopen_checked_for_commit
.store(true, std::sync::atomic::Ordering::SeqCst);
true
},
)
.await
.is_none());
assert!(!revoked_reopen_checked.load(std::sync::atomic::Ordering::SeqCst));
assert_eq!(
manager
.get_by_connection_id(connection_id)
.await
.expect("revoked terminal remains")
.state,
ConnectionState::Failed,
);
}
#[tokio::test]
async fn retired_terminal_history_blocks_physical_recreation_until_authenticated_peer_resume() {
let manager = ConnectionManager::new();
let connection_id = "conn-retired-manual-tombstone";
manager
.upsert_pending(
connection_id.into(),
Some("endpoint-retired-manual".into()),
Some("device-retired-manual".into()),
Some("endpoint-retired-manual".into()),
)
.await;
manager
.commit_current_transport(
connection_id,
Some("endpoint-retired-manual".into()),
91,
Some("initial".into()),
|_| {},
)
.await
.expect("initial transport");
manager
.set_closed_if_current(
connection_id,
91,
Some(crate::lifecycle_reason::REASON_MANUAL_DISCONNECT.into()),
)
.await
.expect("manual terminal");
manager.remove(connection_id).await.expect("retired row");
assert!(manager
.upsert_pending_for_physical_observation(
connection_id.into(),
Some("endpoint-retired-manual".into()),
Some("device-retired-manual".into()),
Some("endpoint-retired-manual".into()),
)
.await
.is_none());
assert!(manager.get_by_connection_id(connection_id).await.is_none());
manager
.materialize_manual_disconnect_tombstone(
connection_id.into(),
Some("endpoint-retired-manual".into()),
Some("device-retired-manual".into()),
Some("endpoint-retired-manual".into()),
)
.await
.expect("manual tombstone");
assert!(
manager
.incoming_authenticated_reconnect_candidate(connection_id, Some(92), 92, true)
.await
);
assert!(
!manager
.incoming_authenticated_reconnect_candidate(connection_id, Some(92), 92, false)
.await
);
assert!(
!manager
.incoming_authenticated_reconnect_candidate(connection_id, Some(93), 92, true)
.await
);
let reopened = manager
.commit_authenticated_peer_reconnect_transport(
connection_id,
Some("endpoint-retired-manual".into()),
Some("device-retired-manual".into()),
92,
Some("authenticated-peer-reconnect".into()),
|requires_authenticated_peer_reconnect, _| requires_authenticated_peer_reconnect,
)
.await
.expect("authenticated peer reconnect recreates the retired row");
assert_eq!(reopened.state, ConnectionState::Connected);
assert_eq!(reopened.transport_stable_id, Some(92));
assert_eq!(reopened.device_id.as_deref(), Some("device-retired-manual"));
assert!(reopened.device_id_hint.is_none());
manager
.set_failed(
connection_id,
Some(crate::lifecycle_reason::REASON_SESSION_TOKEN_REVOKED.into()),
)
.await
.expect("revoked terminal");
manager
.remove(connection_id)
.await
.expect("retired revoked row");
assert!(manager
.materialize_manual_disconnect_tombstone(
connection_id.into(),
Some("endpoint-retired-manual".into()),
Some("device-retired-manual".into()),
Some("endpoint-retired-manual".into()),
)
.await
.is_none());
assert!(manager.get_by_connection_id(connection_id).await.is_none());
let revoked_resume_checked = std::sync::Arc::new(std::sync::atomic::AtomicBool::new(false));
let revoked_resume_checked_for_commit = revoked_resume_checked.clone();
assert!(manager
.commit_authenticated_peer_reconnect_transport(
connection_id,
Some("endpoint-retired-manual".into()),
Some("device-retired-manual".into()),
93,
Some("revoked-peer-reconnect".into()),
move |_, _| {
revoked_resume_checked_for_commit
.store(true, std::sync::atomic::Ordering::SeqCst);
true
},
)
.await
.is_none());
assert!(!revoked_resume_checked.load(std::sync::atomic::Ordering::SeqCst));
assert!(manager.get_by_connection_id(connection_id).await.is_none());
}
#[tokio::test]
async fn fresh_session_capability_reopens_only_revoked_session_history_after_security_commit() {
let manager = ConnectionManager::new();
let connection_id = "conn-fresh-session-reauthorization";
manager
.upsert_pending(
connection_id.into(),
Some("endpoint-fresh-session".into()),
None,
Some("endpoint-fresh-session".into()),
)
.await;
manager
.commit_current_transport(
connection_id,
Some("endpoint-fresh-session".into()),
101,
Some("initial-session".into()),
|_| {},
)
.await
.expect("initial session");
manager
.add_scope(connection_id, "drive-grant:revoked")
.await;
manager
.set_failed(
connection_id,
Some(crate::lifecycle_reason::REASON_SESSION_TOKEN_REVOKED.into()),
)
.await
.expect("revoked session");
manager
.remove(connection_id)
.await
.expect("retired session");
assert!(manager
.upsert_pending_for_physical_observation(
connection_id.into(),
Some("endpoint-fresh-session".into()),
None,
Some("endpoint-fresh-session".into()),
)
.await
.is_none());
let rejected_scope_normalization =
std::sync::Arc::new(std::sync::atomic::AtomicBool::new(false));
let rejected_scope_normalization_for_commit = rejected_scope_normalization.clone();
assert!(manager
.commit_admitted_transport(
connection_id,
Some("endpoint-fresh-session".into()),
None,
102,
Some("rejected-fresh-session".into()),
|requires_authenticated_peer_reconnect, previous| {
assert!(!requires_authenticated_peer_reconnect);
assert_eq!(previous, None);
false
},
move |_| {
rejected_scope_normalization_for_commit
.store(true, std::sync::atomic::Ordering::SeqCst);
},
)
.await
.is_none());
assert!(!rejected_scope_normalization.load(std::sync::atomic::Ordering::SeqCst));
assert!(manager.get_by_connection_id(connection_id).await.is_none());
let security_committed = std::sync::Arc::new(std::sync::atomic::AtomicBool::new(false));
let security_committed_for_commit = security_committed.clone();
let reopened = manager
.commit_admitted_transport(
connection_id,
Some("endpoint-fresh-session".into()),
None,
103,
Some("fresh-session-capability".into()),
move |requires_authenticated_peer_reconnect, previous| {
assert!(!requires_authenticated_peer_reconnect);
assert_eq!(previous, None);
security_committed_for_commit.store(true, std::sync::atomic::Ordering::SeqCst);
true
},
|scopes| {
scopes.clear();
scopes.insert("drive-grant:fresh".to_string());
},
)
.await
.expect("fresh session capability reopens revoked history");
assert!(security_committed.load(std::sync::atomic::Ordering::SeqCst));
assert_eq!(reopened.state, ConnectionState::Connected);
assert_eq!(reopened.transport_stable_id, Some(103));
assert_eq!(
manager.get_scopes(connection_id).await,
vec!["drive-grant:fresh".to_string()]
);
manager
.set_closed_if_current(
connection_id,
103,
Some(crate::lifecycle_reason::REASON_MANUAL_DISCONNECT.into()),
)
.await
.expect("manual terminal");
manager
.remove(connection_id)
.await
.expect("retired manual session");
let manual_authorization_checked =
std::sync::Arc::new(std::sync::atomic::AtomicBool::new(false));
let manual_authorization_checked_for_commit = manual_authorization_checked.clone();
assert!(manager
.commit_admitted_transport(
connection_id,
Some("endpoint-fresh-session".into()),
None,
104,
Some("manual-close-must-remain-terminal".into()),
move |requires_authenticated_peer_reconnect, _| {
manual_authorization_checked_for_commit
.store(true, std::sync::atomic::Ordering::SeqCst);
assert!(requires_authenticated_peer_reconnect);
false
},
|_| {},
)
.await
.is_none());
assert!(manual_authorization_checked.load(std::sync::atomic::Ordering::SeqCst));
assert!(manager.get_by_connection_id(connection_id).await.is_none());
}
#[tokio::test]
async fn fresh_session_capability_reopens_transient_retirement_but_observation_cannot_publish_connected(
) {
let manager = ConnectionManager::new();
let connection_id = "conn-transient-session-recovery";
manager
.upsert_pending(
connection_id.into(),
Some("endpoint-transient-session".into()),
None,
Some("endpoint-transient-session".into()),
)
.await;
manager
.commit_current_transport(
connection_id,
Some("endpoint-transient-session".into()),
201,
Some("initial-session".into()),
|_| {},
)
.await
.expect("initial session");
manager
.set_closed_if_current(
connection_id,
201,
Some(crate::lifecycle_reason::REASON_STALE_ACTIVE_CONNECTION_RECONNECT.into()),
)
.await
.expect("transient retirement");
manager
.remove(connection_id)
.await
.expect("retired transient row");
let observed = manager
.upsert_pending_for_physical_observation(
connection_id.into(),
Some("endpoint-transient-session".into()),
None,
Some("endpoint-transient-session".into()),
)
.await
.expect("transient history permits a non-routable observation");
assert_eq!(observed.state, ConnectionState::Pending);
assert_eq!(observed.transport_stable_id, None);
manager
.set_closed(
connection_id,
Some(crate::lifecycle_reason::REASON_STALE_ACTIVE_CONNECTION_RECONNECT.into()),
)
.await
.expect("second transient retirement");
manager
.remove(connection_id)
.await
.expect("second retired transient row");
let rejected_scope_normalization =
std::sync::Arc::new(std::sync::atomic::AtomicBool::new(false));
let rejected_scope_normalization_for_commit = rejected_scope_normalization.clone();
assert!(manager
.commit_admitted_transport(
connection_id,
Some("endpoint-transient-session".into()),
Some("device-transient-session".into()),
202,
Some("rejected-session-after-transient-retirement".into()),
|_, _| false,
move |_| {
rejected_scope_normalization_for_commit
.store(true, std::sync::atomic::Ordering::SeqCst);
},
)
.await
.is_none());
assert!(!rejected_scope_normalization.load(std::sync::atomic::Ordering::SeqCst));
assert!(manager.get_by_connection_id(connection_id).await.is_none());
let security_committed = std::sync::Arc::new(std::sync::atomic::AtomicBool::new(false));
let security_committed_for_commit = security_committed.clone();
let reopened = manager
.commit_admitted_transport(
connection_id,
Some("endpoint-transient-session".into()),
Some("device-transient-session".into()),
202,
Some("fresh-session-after-transient-retirement".into()),
move |requires_authenticated_peer_reconnect, previous| {
assert!(!requires_authenticated_peer_reconnect);
assert_eq!(previous, None);
security_committed_for_commit.store(true, std::sync::atomic::Ordering::SeqCst);
true
},
|_| {},
)
.await
.expect("fresh capability reopens transient retirement");
assert!(security_committed.load(std::sync::atomic::Ordering::SeqCst));
assert_eq!(reopened.state, ConnectionState::Connected);
assert_eq!(reopened.transport_stable_id, Some(202));
assert_eq!(
reopened.device_id.as_deref(),
Some("device-transient-session")
);
}
#[tokio::test]
async fn withdrawn_session_requires_fresh_capability_before_reopening() {
for remove_record in [false, true] {
let manager = ConnectionManager::new();
let id = "withdrawn-room-peer";
manager
.upsert_pending(id.into(), Some("endpoint".into()), None, None)
.await;
manager
.commit_current_transport(id, Some("endpoint".into()), 1, None, |_| {})
.await
.expect("initial transport");
manager
.set_closed_if_current(
id,
1,
Some(crate::lifecycle_reason::REASON_DESIRED_PEER_WITHDRAWN.into()),
)
.await
.expect("withdrawal");
if remove_record {
manager.remove(id).await.expect("history retained");
}
assert!(manager
.commit_current_transport(id, Some("endpoint".into()), 2, None, |_| panic!(
"observation cannot authorize withdrawal"
))
.await
.is_none());
assert!(manager
.commit_admitted_transport(
id,
Some("endpoint".into()),
None,
2,
None,
|manual, _| {
assert!(!manual);
false
},
|_| panic!("rejected capability cannot mutate scope")
)
.await
.is_none());
let admitted = manager
.commit_admitted_transport(
id,
Some("endpoint".into()),
None,
2,
None,
|manual, _| {
assert!(!manual);
true
},
|scopes| {
scopes.clear();
scopes.insert("room:rejoined".into());
},
)
.await
.expect("fresh capability authorizes a new room session");
assert_eq!(admitted.state, ConnectionState::Connected);
assert_eq!(admitted.transport_stable_id, Some(2));
assert!(
admitted.device_id.is_none(),
"room admission is not durable device identity"
);
assert_eq!(
manager.get_scopes(id).await,
vec!["room:rejoined".to_string()]
);
}
}
#[tokio::test]
async fn automatic_session_reauthorization_requires_current_private_assignment() {
use crate::lifecycle_reason::*;
use crate::session_token::{now_unix_ms, SessionTokenRegistry};
let peer = iroh::SecretKey::generate().public().to_string();
let stranger = iroh::SecretKey::generate().public().to_string();
let scope = "v2:room:automatic-reauthorization";
for reason in [
REASON_SESSION_TOKEN_REVOKED,
REASON_SESSION_ADMISSION_REJECTED,
REASON_DESIRED_PEER_WITHDRAWN,
REASON_MANUAL_DISCONNECT,
] {
for remove_record in [false, true] {
let manager = ConnectionManager::new();
let registry = SessionTokenRegistry::new();
let id = "closed-room-peer";
manager
.upsert_pending(id.into(), Some(peer.clone()), None, None)
.await;
manager.set_closed(id, Some(reason.into())).await.unwrap();
if remove_record {
manager.remove(id).await.unwrap();
}
registry.require_scope_peer_admission(scope).unwrap();
for (revision, expires, peers) in [
(1, now_unix_ms() + 60_000, vec![]),
(2, 0, vec![peer.clone()]),
(3, now_unix_ms() + 60_000, vec![stranger.clone()]),
] {
registry
.update_scope_peer_admission(scope, revision, expires, peers)
.unwrap();
assert!(
manager
.begin_automatic_connect(
id.into(),
Some(peer.clone()),
None,
Some(peer.clone()),
|requires| {
assert!(requires);
registry.is_peer_assigned_to_scope(scope, &peer)
}
)
.await
.is_none(),
"reason={reason}, removed={remove_record}"
);
}
registry
.update_scope_peer_admission(
scope,
4,
now_unix_ms() + 60_000,
vec![peer.clone()],
)
.unwrap();
let attempt = manager
.begin_automatic_connect(
id.into(),
Some(peer.clone()),
None,
Some(peer.clone()),
|requires| {
assert!(requires);
registry.is_peer_assigned_to_scope(scope, &peer)
},
)
.await;
if reason == REASON_MANUAL_DISCONNECT {
assert!(
attempt.is_none(),
"a room lease cannot override manual intent"
);
} else {
let attempt = attempt.expect("current assignment permits admission attempt");
assert_eq!(attempt.state, ConnectionState::Connecting);
assert_eq!(attempt.transport_stable_id, None);
assert!(
!registry.is_connection_admitted(id),
"dial is not admission"
);
}
}
}
}
#[tokio::test]
async fn current_desired_user_device_commit_closes_with_the_terminal_race() {
use crate::lifecycle_reason::*;
for (reason, allowed) in [
(REASON_DESIRED_PEER_WITHDRAWN, true),
(REASON_NETWORK_CHANGE_RECONNECT, true),
(REASON_SESSION_TOKEN_REVOKED, false),
(REASON_SESSION_ADMISSION_REJECTED, false),
(REASON_MANUAL_DISCONNECT, false),
] {
for remove_record in [false, true] {
let manager = ConnectionManager::new();
let id = "desired-user-device-terminal-race";
manager
.upsert_pending(
id.into(),
Some("node-desired".into()),
Some("device-desired".into()),
Some("node-desired".into()),
)
.await;
manager.set_closed(id, Some(reason.into())).await.unwrap();
if remove_record {
manager.remove(id).await.unwrap();
}
let committed = manager
.set_connected_with_transport_when_allowing_authenticated_peer_reconnect(
id,
Some("node-desired".into()),
Some("device-desired".into()),
None,
Some(2),
Some("iroh-relay".into()),
|manual_reconnect, _| {
assert!(!manual_reconnect || reason == REASON_MANUAL_DISCONNECT);
!manual_reconnect
},
|_| {},
TerminalReopenAuthority::CurrentDesiredUserDevice,
)
.await;
assert_eq!(
committed.is_some(),
allowed,
"reason={reason}, removed={remove_record}"
);
}
}
}
#[tokio::test]
async fn current_desired_user_device_reopens_after_automatic_begin() {
let manager = ConnectionManager::new();
let id = "desired-user-device-automatic-reopen";
manager
.upsert_pending(
id.into(),
Some("node-desired".into()),
Some("device-desired".into()),
Some("node-desired".into()),
)
.await;
manager
.set_closed(
id,
Some(crate::lifecycle_reason::REASON_DESIRED_PEER_WITHDRAWN.into()),
)
.await
.expect("withdraw desired peer");
let connecting = manager
.begin_automatic_connect(
id.into(),
Some("node-desired".into()),
Some("device-desired".into()),
Some("node-desired".into()),
|requires_assignment| requires_assignment,
)
.await
.expect("current private assignment begins reconnect");
assert_eq!(connecting.state, ConnectionState::Connecting);
let reopened = manager
.set_connected_with_transport_when_allowing_authenticated_peer_reconnect(
id,
Some("node-desired".into()),
Some("device-desired".into()),
None,
Some(2),
Some("wasm-accepted".into()),
|requires_authenticated_peer_reconnect, previous| {
assert!(!requires_authenticated_peer_reconnect);
assert_eq!(previous, None);
true
},
|_| {},
TerminalReopenAuthority::CurrentDesiredUserDevice,
)
.await
.expect("current desired assignment commits accepted generation");
assert_eq!(reopened.state, ConnectionState::Connected);
assert_eq!(reopened.transport_stable_id, Some(2));
assert_eq!(reopened.status_reason, None);
}
#[tokio::test]
async fn automatic_connect_begin_cannot_reopen_terminal_peer_intent() {
let manager = ConnectionManager::new();
let connection_id = "conn-terminal-automatic-begin";
manager
.upsert_pending(
connection_id.into(),
Some("node-terminal-auto".into()),
Some("device-terminal-auto".into()),
Some("endpoint-terminal-auto".into()),
)
.await;
manager
.commit_current_transport(
connection_id,
Some("endpoint-terminal-auto".into()),
71,
Some("initial".into()),
|_| {},
)
.await
.expect("initial transport");
manager
.set_closed_if_current(
connection_id,
71,
Some(crate::lifecycle_reason::REASON_MANUAL_DISCONNECT.into()),
)
.await
.expect("terminal close");
assert!(manager
.begin_automatic_connect(
connection_id.into(),
Some("node-terminal-auto".into()),
Some("device-terminal-auto".into()),
Some("endpoint-terminal-auto".into()),
|_| true,
)
.await
.is_none());
let terminal = manager
.get_by_connection_id(connection_id)
.await
.expect("terminal record remains");
assert_eq!(terminal.state, ConnectionState::Closed);
assert_eq!(terminal.transport_stable_id, None);
assert_eq!(
terminal.status_reason.as_deref(),
Some(crate::lifecycle_reason::REASON_MANUAL_DISCONNECT)
);
assert!(manager.list_active().await.is_empty());
manager
.upsert_pending(
connection_id.into(),
Some("node-terminal-auto".into()),
Some("device-terminal-auto".into()),
Some("endpoint-terminal-auto".into()),
)
.await;
manager
.commit_current_transport(
connection_id,
Some("endpoint-terminal-auto".into()),
72,
Some("explicit-reconnect".into()),
|_| {},
)
.await
.expect("explicit reconnect transport");
manager
.set_closed_if_current(
connection_id,
72,
Some(crate::lifecycle_reason::REASON_NETWORK_CHANGE_RECONNECT.into()),
)
.await
.expect("recoverable close");
let retry = manager
.begin_automatic_connect(
connection_id.into(),
Some("node-terminal-auto".into()),
Some("device-terminal-auto".into()),
Some("endpoint-terminal-auto".into()),
|_| true,
)
.await
.expect("recoverable automatic begin");
assert_eq!(retry.state, ConnectionState::Connecting);
}
#[tokio::test]
async fn route_observation_requires_the_complete_current_generation_tuple() {
let manager = ConnectionManager::new();
manager
.upsert_pending(
"conn-route-fence".into(),
Some("node-route-fence".into()),
Some("device-route-fence".into()),
Some("endpoint-route-fence".into()),
)
.await;
let connected = manager
.set_connected_with_transport(
"conn-route-fence",
Some("endpoint-route-fence".into()),
Some(51),
Some("initial".into()),
)
.await
.expect("current generation");
assert!(manager
.report_transport_status_if_current(
"conn-route-fence",
51,
connected.transport_generation.saturating_add(1),
connected.route_generation,
"webrtc".into(),
Some("iroh-relay".into()),
)
.await
.is_none());
assert!(manager
.report_transport_status_if_current(
"conn-route-fence",
51,
connected.transport_generation,
connected.route_generation.saturating_add(1),
"webrtc".into(),
Some("iroh-relay".into()),
)
.await
.is_none());
let updated = manager
.report_transport_status_if_current(
"conn-route-fence",
51,
connected.transport_generation,
connected.route_generation,
"webrtc".into(),
Some("iroh-relay".into()),
)
.await
.expect("current tuple accepted");
assert_eq!(updated.active_transport, "webrtc");
assert!(updated.active_route_generation > connected.route_generation);
assert!(manager
.report_transport_status_if_current(
"conn-route-fence",
51,
connected.transport_generation,
connected.route_generation,
"iroh-relay".into(),
None,
)
.await
.is_none());
let record = manager
.get_by_connection_id("conn-route-fence")
.await
.expect("record remains current");
assert_eq!(record.active_transport, "webrtc");
}
#[tokio::test]
async fn upsert_pending_does_not_downgrade_connected_replacement_record() {
let manager = ConnectionManager::new();
manager
.upsert_pending(
"conn-stable".into(),
Some("node-1".into()),
Some("device-1".into()),
Some("endpoint-1".into()),
)
.await;
manager
.set_connected_with_transport(
"conn-stable",
Some("endpoint-1".into()),
Some(101),
Some("outgoing".into()),
)
.await
.expect("connected record");
manager
.mark_transport_replaced(
"conn-stable",
Some(102),
Some("incoming-replacement".into()),
Some("replacement-in-progress".into()),
)
.await
.expect("replacement record");
let updated = manager
.upsert_pending(
"conn-stable".into(),
Some("node-1".into()),
Some("device-1".into()),
Some("endpoint-1".into()),
)
.await;
assert_eq!(updated.state, ConnectionState::Connected);
assert_eq!(
updated.status_reason.as_deref(),
Some("replacement-in-progress")
);
assert_eq!(updated.transport_generation, 2);
assert_eq!(updated.transport_stable_id, Some(102));
}
#[tokio::test]
async fn peer_snapshot_uses_live_connected_record_for_transport_metadata() {
let manager = ConnectionManager::new();
manager
.upsert_pending(
"conn-live".into(),
Some("node-live".into()),
Some("device-1".into()),
Some("endpoint-live".into()),
)
.await;
manager
.set_connected_with_transport(
"conn-live",
Some("endpoint-live".into()),
Some(101),
Some("incoming-live".into()),
)
.await
.expect("live connected record");
manager
.report_transport_status("conn-live", "webrtc".into(), Some("iroh-relay".into()))
.await
.expect("live peer snapshot");
manager
.upsert_pending(
"conn-pending-replacement".into(),
Some("node-live".into()),
Some("device-1".into()),
Some("endpoint-replacement".into()),
)
.await;
manager
.set_connecting("conn-pending-replacement")
.await
.expect("pending replacement record");
let snapshot = manager
.peer_snapshot("device-1")
.await
.expect("peer snapshot");
assert_eq!(snapshot.status, ConnectionState::Connected);
assert_eq!(
snapshot.connection_ids.first().map(String::as_str),
Some("conn-live")
);
assert_eq!(snapshot.active_transport_stable_id, Some(101));
assert_eq!(snapshot.active_transport_generation, 1);
assert_eq!(snapshot.active_route_generation, 1);
assert_eq!(snapshot.active_transport, "webrtc");
assert_eq!(snapshot.parallel_transport.as_deref(), Some("iroh-relay"));
}
#[tokio::test]
async fn transition_counters_increment_without_idle_churn() {
let manager = ConnectionManager::new();
manager
.upsert_pending(
"conn-counters".into(),
Some("node-1".into()),
Some("device-1".into()),
Some("endpoint-1".into()),
)
.await;
manager.set_connecting("conn-counters").await;
manager
.set_connected_with_transport(
"conn-counters",
Some("endpoint-1".into()),
Some(201),
Some("initial".into()),
)
.await
.expect("connected record");
manager
.mark_transport_replaced(
"conn-counters",
Some(202),
Some("replacement".into()),
Some("replacement-in-progress".into()),
)
.await
.expect("replacement record");
let record = manager
.get_by_connection_id("conn-counters")
.await
.expect("record");
assert_eq!(record.transition_count, 2);
assert_eq!(record.connecting_transition_count, 1);
assert_eq!(record.replacement_count, 2);
assert_eq!(record.last_reconnect_reason.as_deref(), Some("replacement"));
let peer = manager.peer_snapshot("device-1").await.expect("peer");
assert_eq!(peer.transition_count, 2);
assert_eq!(peer.connecting_transition_count, 1);
assert_eq!(peer.replacement_count, 2);
assert_eq!(peer.retire_count, 0);
}
#[tokio::test]
async fn remove_and_recreate_preserves_monotonic_generation_history() {
let manager = ConnectionManager::new();
manager
.upsert_pending(
"conn-history".into(),
Some("node-1".into()),
Some("device-1".into()),
Some("endpoint-1".into()),
)
.await;
manager.set_connecting("conn-history").await;
manager
.set_connected_with_transport(
"conn-history",
Some("endpoint-1".into()),
Some(301),
Some("initial".into()),
)
.await
.expect("connected record");
manager
.mark_transport_replaced(
"conn-history",
Some(302),
Some("replacement".into()),
Some("replacement-in-progress".into()),
)
.await
.expect("replacement record");
manager
.remove("conn-history")
.await
.expect("removed record");
let recreated = manager
.upsert_pending(
"conn-history".into(),
Some("node-1".into()),
Some("device-1".into()),
Some("endpoint-1".into()),
)
.await;
assert_eq!(recreated.transport_generation, 2);
assert_eq!(recreated.transition_count, 2);
assert_eq!(recreated.connecting_transition_count, 1);
assert_eq!(recreated.replacement_count, 2);
let reconnect = manager
.set_connected_with_transport(
"conn-history",
Some("endpoint-1".into()),
Some(303),
Some("recreated".into()),
)
.await
.expect("reconnected record");
assert_eq!(reconnect.transport_generation, 3);
assert_eq!(reconnect.transition_count, 3);
assert_eq!(reconnect.replacement_count, 3);
}
#[tokio::test]
async fn peer_snapshot_uses_connection_id_until_device_hint_exists() {
let manager = ConnectionManager::new();
manager
.upsert_pending("conn-provisional".into(), Some("node-1".into()), None, None)
.await;
let provisional = manager
.peer_snapshot("conn-provisional")
.await
.expect("provisional snapshot");
assert_eq!(provisional.peer_id, "conn-provisional");
assert_eq!(provisional.device_id, None);
assert_eq!(provisional.device_id_hint, None);
assert_eq!(provisional.node_id.as_deref(), Some("node-1"));
manager
.upsert_pending(
"conn-provisional".into(),
Some("node-1".into()),
Some("device-hint-1".into()),
None,
)
.await;
let hinted = manager
.peer_snapshot("device-hint-1")
.await
.expect("hinted snapshot");
assert_eq!(hinted.peer_id, "device-hint-1");
assert_eq!(hinted.device_id, None);
assert_eq!(hinted.device_id_hint.as_deref(), Some("device-hint-1"));
assert_eq!(hinted.node_id.as_deref(), Some("node-1"));
assert!(
manager
.are_same_peer("device-hint-1", "conn-provisional")
.await
);
assert!(manager.are_same_peer("device-hint-1", "node-1").await);
}
#[tokio::test]
async fn authoritative_device_binding_rekeys_peer_and_preserves_state() {
let manager = ConnectionManager::new();
manager
.upsert_pending(
"conn-authoritative".into(),
Some("node-2".into()),
Some("device-hint-2".into()),
None,
)
.await;
manager.set_connecting("conn-authoritative").await;
manager.add_scope("device-hint-2", "persistent").await;
manager
.set_health("device-hint-2", ConnectionHealth::Suspect)
.await;
let updated = manager
.set_device_id("conn-authoritative", "device-2".into())
.await
.expect("authoritative record");
assert_eq!(updated.device_id.as_deref(), Some("device-2"));
assert_eq!(updated.device_id_hint.as_deref(), Some("device-hint-2"));
let authoritative = manager
.peer_snapshot("device-2")
.await
.expect("authoritative snapshot");
assert_eq!(authoritative.peer_id, "device-2");
assert_eq!(authoritative.device_id.as_deref(), Some("device-2"));
assert_eq!(
authoritative.device_id_hint.as_deref(),
Some("device-hint-2")
);
assert_eq!(authoritative.status, ConnectionState::Connecting);
assert_eq!(authoritative.health, ConnectionHealth::Suspect);
assert_eq!(authoritative.scopes, vec!["persistent".to_string()]);
assert!(
manager
.are_same_peer("device-2", "conn-authoritative")
.await
);
assert!(manager.are_same_peer("device-2", "node-2").await);
assert!(manager.are_same_peer("device-hint-2", "device-2").await);
let hinted = manager
.peer_snapshot("device-hint-2")
.await
.expect("hinted snapshot");
assert_eq!(hinted.peer_id, "device-2");
assert_eq!(hinted.device_id.as_deref(), Some("device-2"));
assert_eq!(hinted.device_id_hint.as_deref(), Some("device-hint-2"));
}
#[tokio::test]
async fn multiple_transports_collapse_into_one_device_peer() {
let manager = ConnectionManager::new();
manager
.upsert_pending(
"conn-a".into(),
Some("node-a".into()),
Some("device-3".into()),
None,
)
.await;
manager
.upsert_pending(
"conn-b".into(),
Some("node-b".into()),
Some("device-3".into()),
None,
)
.await;
manager.set_connecting("conn-a").await;
manager
.set_connected_with_transport("conn-b", None, Some(30), Some("incoming".into()))
.await;
let snapshots = manager.list_peer_snapshots().await;
assert_eq!(snapshots.len(), 1);
assert_eq!(snapshots[0].peer_id, "device-3");
assert_eq!(snapshots[0].device_id, None);
assert_eq!(snapshots[0].device_id_hint.as_deref(), Some("device-3"));
assert_eq!(snapshots[0].status, ConnectionState::Connected);
assert_eq!(snapshots[0].connection_ids.len(), 2);
}
#[tokio::test]
async fn best_connection_for_peer_prefers_latest_active_transport() {
let manager = ConnectionManager::new();
manager
.upsert_pending(
"conn-old".into(),
Some("node-1".into()),
Some("device-1".into()),
Some("endpoint-1".into()),
)
.await;
manager
.set_device_id("conn-old", "device-1".into())
.await
.expect("old device binding");
manager
.set_connected_with_transport(
"conn-old",
Some("endpoint-1".into()),
Some(10),
Some("incoming".into()),
)
.await
.expect("old record");
manager
.upsert_pending(
"conn-new".into(),
Some("node-1".into()),
Some("device-1".into()),
Some("endpoint-1".into()),
)
.await;
manager
.set_device_id("conn-new", "device-1".into())
.await
.expect("new device binding");
manager
.set_connected_with_transport(
"conn-new",
Some("endpoint-1".into()),
Some(11),
Some("replacement".into()),
)
.await
.expect("new record");
let best = manager
.best_connection_for_peer("device-1")
.await
.expect("best record");
assert_eq!(best.connection_id, "conn-new");
assert_eq!(best.transport_stable_id, Some(11));
}
#[tokio::test]
async fn best_connection_for_peer_returns_none_for_unknown_peer() {
let manager = ConnectionManager::new();
assert!(
manager
.best_connection_for_peer("nonexistent-peer")
.await
.is_none(),
"unknown peer should return None"
);
}
#[tokio::test]
async fn peer_snapshot_surfaces_failed_over_older_connected_same_device() {
let manager = ConnectionManager::new();
manager
.upsert_pending(
"conn-stale".into(),
Some("node-peer".into()),
Some("device-peer".into()),
Some("node-peer".into()),
)
.await;
manager.set_connecting("conn-stale").await;
manager
.set_connected("conn-stale", Some("node-peer".into()))
.await;
manager
.set_device_id("conn-stale", "device-peer-uuid".into())
.await
.expect("set device");
manager
.upsert_pending(
"conn-dial".into(),
Some("node-peer".into()),
Some("device-peer".into()),
Some("node-peer".into()),
)
.await;
manager.set_connecting("conn-dial").await;
manager
.set_device_id("conn-dial", "device-peer-uuid".into())
.await
.expect("set device on dial row");
manager
.set_failed("conn-dial", Some("ensure_connected_addr-timeout".into()))
.await
.expect("mark dial failed");
let snapshot = manager
.peer_snapshot("device-peer-uuid")
.await
.expect("merged peer");
assert_eq!(snapshot.status, ConnectionState::Failed);
}
}