Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -51,6 +51,40 @@ pub enum PairingState {
Failed(String),
}

/// A pairing attempt still waiting for its first connection this long has
/// found nothing: its scan was refused or ended. It is started again then,
/// not after the handshake window.
const WAITING_FOR_CONNECTION_STALE_SECS: u64 = 20;

/// A handshake under way that has not moved for this long is cleared and
/// started again.
const HANDSHAKE_STALE_SECS: u64 = 90;

/// How long the pairing loop waits for a state change before looking at its
/// sessions again, so an attempt that went stale without an event is still
/// restarted promptly.
const PAIRING_LOOP_WAKE_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(10);

/// After a scan the radio refused, how long before the attempt is retried:
/// past the coordinator's 6 s gap between scans (`BleCoordinator.MIN_SCAN_GAP_MS`).
#[cfg(any(all(target_os = "android", feature = "jni"), test))]
const DISCOVERY_RETRY_DELAY: std::time::Duration = std::time::Duration::from_secs(7);

/// How many seconds a session may sit in `state` before the pairing loop
/// clears it and starts again; `None` for a state the loop never clears by
/// age (finished, or already failed and retried on the next pass).
pub(crate) fn stale_after_secs(state: &PairingState) -> Option<u64> {
match state {
PairingState::Complete | PairingState::Failed(_) => None,
PairingState::WaitingForConnection => Some(WAITING_FOR_CONNECTION_STALE_SECS),
PairingState::ReadingIdentity
| PairingState::ExchangingChainTips
| PairingState::AwaitingConfirm
| PairingState::ConfirmSent
| PairingState::UpdatingStatus => Some(HANDSHAKE_STALE_SECS),
}
}

/// Active pairing session for a contact
#[derive(Debug, Clone)]
pub struct PairingSession {
Expand Down Expand Up @@ -189,23 +223,68 @@ impl PairingOrchestrator {
// Spawn without awaiting to avoid blocking the JNI caller path
crate::runtime::get_runtime().spawn(async move {
let orchestrator = crate::bluetooth::get_pairing_orchestrator();
if let Err(e) = orchestrator.start_ble_discovery(contact).await {
log::warn!(
"[PairingOrchestrator] start_ble_discovery secondary path failed: {}",
e
);
} else {
log::info!(
"[PairingOrchestrator] start_ble_discovery secondary path issued successfully"
);
}
let started = orchestrator.start_ble_discovery(contact).await;
orchestrator
.after_discovery_start(contact, started, DISCOVERY_RETRY_DELAY)
.await;
});
}

// Return the role - Kotlin will start the appropriate BLE operation
Ok(should_advertise)
}

/// What follows the radio's answer to looking for `contact`: nothing when
/// it started; when it was refused, after `retry_after` (past the radio's
/// gap between scans) the attempt waiting on it is freed for the next pass
/// ([`Self::discovery_refused`]).
#[cfg(any(all(target_os = "android", feature = "jni"), test))]
pub(crate) async fn after_discovery_start(
&self,
contact: [u8; 32],
started: Result<(), String>,
retry_after: std::time::Duration,
) {
match started {
Ok(()) => log::info!(
"[PairingOrchestrator] BLE discovery started for {:02x}{:02x}...",
contact[0],
contact[1]
),
Err(e) => {
log::warn!(
"[PairingOrchestrator] BLE discovery for {:02x}{:02x}... was refused: {e}; \
retrying in {retry_after:?}",
contact[0],
contact[1]
);
tokio::time::sleep(retry_after).await;
self.discovery_refused(contact).await;
}
}
}

/// The radio refused to look for `contact` (a scan refused by the
/// platform's rate limit or the coordinator's gap between scans): an
/// attempt still waiting for its first connection is marked failed and the
/// loop woken, so the next pass starts it again instead of leaving it
/// waiting on a scan that never ran. An attempt already past that is left
/// alone.
#[cfg(any(all(target_os = "android", feature = "jni"), test))]
pub(crate) async fn discovery_refused(&self, contact: [u8; 32]) {
let mut sessions = self.sessions.write().await;
let Some(session) = sessions.get_mut(&contact) else {
return;
};
if session.state != PairingState::WaitingForConnection {
return;
}
session.state = PairingState::Failed("the scan for the peer was refused".to_string());
session.last_activity = Instant::now();
drop(sessions);
self.signal_state_change();
}

/// Handle BLE identity observed event from AndroidBleBridge
///
/// Called when peer's identity is read from GATT characteristic.
Expand Down Expand Up @@ -796,12 +875,6 @@ impl PairingOrchestrator {

log::info!("[PairingOrchestrator] start_pairing_all_unpaired: loop started");

/// Maximum wall-clock interval the pairing loop waits for a state-change
/// notification before re-evaluating sessions regardless. This is transport
/// runtime control only and ensures stale or silently-dropped BLE sessions
/// are recovered even when no explicit disconnect event fires.
const PAIRING_LOOP_WAKE_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(30);

loop {
let state_changed = self.state_change.notified();

Expand Down Expand Up @@ -880,32 +953,17 @@ impl PairingOrchestrator {
let mut device_id = [0u8; 32];
device_id.copy_from_slice(&contact.device_id);

// BLE sessions stale after 90s. This transport timer bounds
// handshake freshness and reconnect retry windows only.
const STALE_SECS: u64 = 90;

// Determine whether to skip or clear this contact's pairing session.
// Determine whether to skip or clear this contact's pairing session:
// a finished one is skipped, an in-progress one while it is fresh
// ([`stale_after_secs`]), a failed one never.
let should_skip = {
let sessions = self.sessions.read().await;
if let Some(session) = sessions.get(&device_id) {
match &session.state {
// Pairing done — skip unconditionally.
PairingState::Complete => true,
// In-progress states: skip while fresh, clear when stale.
PairingState::ReadingIdentity
| PairingState::ExchangingChainTips
| PairingState::AwaitingConfirm
| PairingState::ConfirmSent
| PairingState::UpdatingStatus
| PairingState::WaitingForConnection => {
session.last_activity.elapsed().as_secs() < STALE_SECS
}
// Failed or stale: do not skip — clear below and retry.
PairingState::Failed(_) => false,
}
} else {
false // no session yet — proceed to initiate
}
sessions.get(&device_id).is_some_and(|session| {
session.state == PairingState::Complete
|| stale_after_secs(&session.state).is_some_and(|stale| {
session.last_activity.elapsed().as_secs() < stale
})
})
};

if should_skip {
Expand All @@ -921,8 +979,8 @@ impl PairingOrchestrator {
let mut sessions = self.sessions.write().await;
let should_clear = sessions.get(&device_id).is_some_and(|s| {
matches!(&s.state, PairingState::Failed(_))
|| (!matches!(&s.state, PairingState::Complete)
&& s.last_activity.elapsed().as_secs() >= STALE_SECS)
|| stale_after_secs(&s.state)
.is_some_and(|stale| s.last_activity.elapsed().as_secs() >= stale)
});
if should_clear {
sessions.remove(&device_id);
Expand Down Expand Up @@ -1139,10 +1197,12 @@ impl PairingOrchestrator {
adv_ok,
scan_ok,
);
if !adv_ok && !scan_ok {
return Err(
"Both startBlePairingAdvertise and startBlePairingScan failed".to_string(),
);
// The scanner finds the peer: without its scan the attempt would
// wait on a scan that never ran.
if !scan_ok {
return Err(format!(
"startBlePairingScan was refused (advertise={adv_ok})"
));
}
}

Expand Down Expand Up @@ -1305,6 +1365,127 @@ mod tests {
orchestrator.cancel_pairing(&device_id).await;
}

/// An attempt still waiting for its first connection is started again
/// after 20 s: its scan was refused or ended, and the handshake window of
/// 90 s would leave the user waiting for two minutes. A handshake under
/// way keeps the 90 s; a finished or failed session is never cleared by
/// age.
#[test]
fn an_attempt_waiting_for_its_first_connection_goes_stale_before_a_handshake() {
assert_eq!(
stale_after_secs(&PairingState::WaitingForConnection),
Some(20)
);
for state in [
PairingState::ReadingIdentity,
PairingState::ExchangingChainTips,
PairingState::AwaitingConfirm,
PairingState::ConfirmSent,
PairingState::UpdatingStatus,
] {
assert_eq!(stale_after_secs(&state), Some(90), "{state:?}");
}
assert_eq!(stale_after_secs(&PairingState::Complete), None);
assert_eq!(stale_after_secs(&PairingState::Failed("x".into())), None);
assert_eq!(
PAIRING_LOOP_WAKE_TIMEOUT,
std::time::Duration::from_secs(10)
);
}

/// A refused scan frees the attempt that waited on it: it is marked failed
/// and the loop is woken, so the next pass starts it again. An attempt
/// already reading the peer's identity is left alone, and a contact with
/// no session gets none.
#[tokio::test]
async fn a_refused_scan_frees_the_attempt_waiting_on_it() {
let orchestrator = PairingOrchestrator::new();
let waiting = [0x51u8; 32];
let reading = [0x52u8; 32];
{
let mut sessions = orchestrator.sessions.write().await;
for (id, state) in [
(waiting, PairingState::WaitingForConnection),
(reading, PairingState::ReadingIdentity),
] {
sessions.insert(
id,
PairingSession {
contact_device_id: id,
state,
ble_address: None,
peer_genesis_hash: None,
peer_chain_tip: None,
last_activity: Instant::now(),
},
);
}
}
let woken = orchestrator.state_change.notified();
tokio::pin!(woken);
woken.as_mut().enable();

orchestrator.discovery_refused(waiting).await;
orchestrator.discovery_refused(reading).await;
orchestrator.discovery_refused([0x53u8; 32]).await;

assert!(matches!(
orchestrator.get_session_status(&waiting).await,
Some(PairingState::Failed(_))
));
assert_eq!(
orchestrator.get_session_status(&reading).await,
Some(PairingState::ReadingIdentity)
);
assert_eq!(orchestrator.get_session_status(&[0x53u8; 32]).await, None);
tokio::time::timeout(std::time::Duration::from_secs(1), woken)
.await
.expect("the pairing loop was not woken");
}

/// A discovery the radio started leaves the attempt waiting for its
/// connection; one it refused frees the attempt once the retry delay has
/// passed. The delay the device waits is past the radio's 6 s gap between
/// scans.
#[tokio::test]
async fn a_refused_discovery_is_retried_after_the_scan_gap() {
assert_eq!(DISCOVERY_RETRY_DELAY, std::time::Duration::from_secs(7));
let orchestrator = PairingOrchestrator::new();
let started = [0x61u8; 32];
let refused = [0x62u8; 32];
{
let mut sessions = orchestrator.sessions.write().await;
for id in [started, refused] {
sessions.insert(
id,
PairingSession {
contact_device_id: id,
state: PairingState::WaitingForConnection,
ble_address: None,
peer_genesis_hash: None,
peer_chain_tip: None,
last_activity: Instant::now(),
},
);
}
}
let retry = std::time::Duration::from_millis(20);
orchestrator
.after_discovery_start(started, Ok(()), retry)
.await;
orchestrator
.after_discovery_start(refused, Err("scan refused".to_string()), retry)
.await;
assert_eq!(
orchestrator.get_session_status(&started).await,
Some(PairingState::WaitingForConnection)
);
assert!(matches!(
orchestrator.get_session_status(&refused).await,
Some(PairingState::Failed(_))
));
}

#[tokio::test]
async fn test_session_state_transitions() {
let orchestrator = PairingOrchestrator::new();
Expand Down
Loading
Loading