From 7c341d42e5dd4a69b7853ca117a821f62274689e Mon Sep 17 00:00:00 2001 From: Edison Jwa Date: Mon, 8 Jun 2026 20:15:50 +0900 Subject: [PATCH] fix(core,protocol): bound disconnect shutdown --- core/chanora_core/Cargo.toml | 2 +- core/chanora_core/src/events.rs | 285 +++++++++++ core/chanora_core/src/lib.rs | 486 ++++++------------- core/chanora_core/src/network_diagnostics.rs | 72 +++ crates/chanora_protocol/src/adapter.rs | 217 +++++++-- 5 files changed, 674 insertions(+), 388 deletions(-) create mode 100644 core/chanora_core/src/events.rs create mode 100644 core/chanora_core/src/network_diagnostics.rs diff --git a/core/chanora_core/Cargo.toml b/core/chanora_core/Cargo.toml index 5722f35..3003815 100644 --- a/core/chanora_core/Cargo.toml +++ b/core/chanora_core/Cargo.toml @@ -18,7 +18,7 @@ chanora_diagnostics = { path = "../../crates/chanora_diagnostics" } chanora_prefetch = { path = "../../crates/chanora_prefetch" } thiserror.workspace = true tracing.workspace = true -tokio = { version = "1", features = ["sync", "rt", "macros"] } +tokio = { version = "1", features = ["sync", "rt", "macros", "time"] } [dev-dependencies] # Used by integration tests to inspect the bookmark DB row layout diff --git a/core/chanora_core/src/events.rs b/core/chanora_core/src/events.rs new file mode 100644 index 0000000..20201c2 --- /dev/null +++ b/core/chanora_core/src/events.rs @@ -0,0 +1,285 @@ +use chanora_audio::{AudioRoute, PttBackendDescriptor}; +use chanora_protocol::MessageTarget; + +/// Privacy-safe snapshot of the active PTT capability. +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct PttDescriptorSnapshot { + /// Stable capability level name. + pub level: String, + /// Stable backend identifier. + pub backend_id: String, + /// Coarse bound input class; empty when no binding is active. + pub bound_input_class: String, +} + +impl From for PttDescriptorSnapshot { + fn from(desc: PttBackendDescriptor) -> Self { + Self { + level: desc.level.as_str().to_string(), + backend_id: desc.backend_id.to_string(), + bound_input_class: desc.bound_input_class.unwrap_or("").to_string(), + } + } +} + +/// Persisted PTT binding state exposed to callers. +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct PersistedPttBinding { + /// Stable input category string (`""`, `"keyboard"`, or + /// `"mouse-side-button"`). + pub input_class: String, + /// Display-only key label; empty when no binding is active. + pub key_label: String, +} + +impl PersistedPttBinding { + pub(crate) fn empty() -> Self { + Self { + input_class: String::new(), + key_label: String::new(), + } + } +} + +/// High-level lifecycle event surfaced to subscribers. +/// +/// This is the minimal set needed for A.6 (reconnect banner). The +/// full event catalogue lands in A.4. +#[derive(Debug, Clone)] +pub enum SessionEvent { + /// Initial connect succeeded, or reconnect attempt succeeded. + Connected { + /// Server name reported in the snapshot. + server_name: String, + }, + /// Connection lost; the supervisor will retry. + Lost { + /// Reason classification from the protocol layer. + reason: String, + }, + /// Supervisor is sleeping before its next reconnect attempt. + Reconnecting { + /// 1-based attempt counter for the current outage. + attempt: u32, + /// Seconds the supervisor will sleep before this attempt. + delay_secs: u32, + }, + /// Supervisor gave up after `attempt` failed retries (or the + /// user explicitly disconnected mid-outage). + Disconnected { + /// Reason classification from the protocol layer. + reason: String, + }, + /// Audio engine started (e.g. after a successful reconnect with + /// reattachment). + AudioStarted, + /// Audio engine stopped (e.g. before a reconnect cycle, or by + /// explicit user action). + AudioStopped, + /// Detected desktop Push-to-Talk capability (gen2 v0.9.3 / + /// DEC-023..028). Published when the audio engine starts or + /// when the active backend transitions (for example macOS + /// permission state change). Carries only the privacy-safe + /// descriptor — capability level, backend identifier, bound + /// input class — per SRS-202 / DEC-027. + PttCapability { + /// Stable level name from `PttCapabilityLevel::as_str()`. + level: String, + /// Stable backend identifier (e.g. `"focused"`). + backend_id: String, + /// Coarse bound input class (e.g. `"keyboard"`); empty when + /// no binding is active. + bound_input_class: String, + }, + /// Voice subsystem state snapshot (SDD-094). Emitted on + /// `voice_join` / `voice_leave`, transmit-mode changes, + /// hard-mute toggles, and release-tail edits. + VoiceState { + /// True when the user has joined a voice channel via + /// `voice_join` and the audio engine is running. + in_channel: bool, + /// Active transmit mode encoded as + /// [`chanora_audio::TransmitMode::as_u8`]. + transmit_mode: u8, + /// True when the hard-mute clamp is engaged. + mute: bool, + /// Current release-tail in milliseconds (0..=500). + release_tail_ms: u32, + /// Last confirmed authoritative channel id from the + /// `channel_join` reducer projection. + current_channel_id: Option, + /// Non-authoritative pending target channel id from the + /// reducer projection. + pending_target_channel_id: Option, + /// Whether the reducer currently allows a new join intent. + can_join: bool, + /// Whether the reducer currently allows leave intent. + can_leave: bool, + /// Join projection synchronization state. + join_sync_state: VoiceJoinSyncState, + /// Last stable sanitized join error code, if any. + join_error_code: Option, + }, + /// iOS audio-session interruption state (SDD-101). Emitted when + /// interruption begins and when it ends (with the platform hint + /// indicating whether audio should resume). + InterruptionState { + /// True when interruption began, false when interruption ended. + began: bool, + /// Platform-provided resume hint. For begin events this is false. + should_resume: bool, + }, + /// A text message was received from the server. + ChatMessage { + /// Client id of the sender. + sender_id: u64, + /// Nickname of the sender. + sender_name: String, + /// Message content. + message: String, + /// Target scope (server/channel/private/poke). + target: MessageTarget, + }, + /// Human-readable TeamSpeak-style server activity. + ServerActivity { + /// Activity line text. + message: String, + }, + /// Audio route changed (speaker/earpiece/BT/wired headset). + AudioRouteChanged { + /// New audio output route. + route: AudioRoute, + }, + /// A client moved to a different channel. + ClientMoved { + /// Unique client identifier. + client_id: u64, + /// Destination channel. + new_channel_id: u64, + }, + /// A new client connected. + ClientJoined { + /// Unique client identifier. + client_id: u64, + /// Channel the client joined. + channel_id: u64, + /// Display nickname. + name: String, + /// Microphone muted state. + input_muted: bool, + /// Speaker muted state. + output_muted: bool, + /// Whether this is a server query (bot) client. + is_server_query: bool, + /// Client's talk power value. + talk_power: i32, + /// Whether the server granted temporary talk power. + talk_power_granted: bool, + }, + /// A client disconnected. + ClientLeft { + /// Unique client identifier. + client_id: u64, + /// Display nickname at time of disconnect. + name: String, + }, + /// Client properties changed. + ClientUpdated { + /// Unique client identifier. + client_id: u64, + /// Microphone muted state. + input_muted: bool, + /// Speaker muted state. + output_muted: bool, + /// Whether this is a server query (bot) client. + is_server_query: bool, + /// Client's talk power value. + talk_power: i32, + /// Whether the server granted temporary talk power. + talk_power_granted: bool, + }, + /// A new channel appeared. + ChannelAdded { + /// Unique channel identifier. + id: u64, + /// Parent channel ID. + parent: u64, + /// Channel name. + name: String, + /// Predecessor channel ID within the same parent (TeamSpeak + /// linked-list ordering hint). Zero means first child. + order: i64, + /// Whether the channel requires a password. + has_password: bool, + /// Talk power required to speak, or `None` when unrestricted. + needed_talk_power: Option, + }, + /// A channel was deleted. + ChannelRemoved { + /// Channel identifier. + id: u64, + }, + /// Channel properties changed. + ChannelUpdated { + /// Unique channel identifier. + id: u64, + /// Channel name. + name: String, + /// Whether the channel requires a password. + has_password: bool, + /// Talk power required to speak, or `None` when unrestricted. + needed_talk_power: Option, + }, +} + +/// Bridge-safe mirror of channel-join projection sync state. +#[derive(Debug, Clone, Copy)] +pub enum VoiceJoinSyncState { + /// Reducer is ready to accept channel actions. + Ready, + /// Reducer is synchronizing against an initial snapshot. + SynchronizingInitialSnapshot, + /// Reducer is synchronizing after reconnect. + SynchronizingReconnect, +} + +/// Bridge-safe mirror of stable channel-join error codes. +#[derive(Debug, Clone, Copy)] +pub enum VoiceJoinErrorCode { + /// Duplicate same-target join intent was coalesced. + DuplicateSameTargetCoalesced, + /// A different target was requested while one is already pending. + JoinAlreadyPendingDifferentTarget, + /// Join denied by server policy/permission. + JoinDenied, + /// Join failed due to protocol-level error. + JoinProtocolFailure, + /// Join failed due to transport/network error. + JoinNetworkFailure, + /// Join timed out awaiting confirmation. + JoinTimeout, + /// Pending join was superseded by user leave. + JoinSupersededByLeave, + /// Stale join outcome was ignored. + JoinStaleOutcomeIgnored, + /// Authoritative membership reconciled to different channel. + JoinReconciledDifferentChannel, + /// Join command was rejected before send acceptance. + JoinCommandRejectedBeforeSend, + /// Join intent rejected while reducer synchronizing. + JoinCannotStartWhileSynchronizing, +} + +/// Coarse OS-reported network state. Populated by the Flutter side +/// via `connectivity_plus`; on platforms where no signal is wired +/// we stay at `Unknown` forever and the supervisor falls back to +/// pure watchdog/backoff behaviour. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum NetworkState { + /// No signal seen yet — treat as ambiguous; don't change behaviour. + Unknown, + /// OS reports at least one network with internet capability. + Online, + /// OS reports no networks available. + Offline, +} diff --git a/core/chanora_core/src/lib.rs b/core/chanora_core/src/lib.rs index 900e4ea..b18017a 100644 --- a/core/chanora_core/src/lib.rs +++ b/core/chanora_core/src/lib.rs @@ -52,6 +52,8 @@ use chanora_state::channel_join::{ ConnectionEpoch, JoinFailureKind, }; +mod events; +mod network_diagnostics; pub mod ptt; pub use chanora_audio::{ @@ -69,46 +71,11 @@ pub use chanora_protocol::{ MessageTarget, ProtocolError, ServerActivity, ServerSnapshot, }; pub use chanora_storage::{Bookmark, BookmarkRepository, IdentityFileStore}; - -/// Privacy-safe snapshot of the active PTT capability. -#[derive(Debug, Clone, PartialEq, Eq)] -pub struct PttDescriptorSnapshot { - /// Stable capability level name. - pub level: String, - /// Stable backend identifier. - pub backend_id: String, - /// Coarse bound input class; empty when no binding is active. - pub bound_input_class: String, -} - -impl From for PttDescriptorSnapshot { - fn from(desc: PttBackendDescriptor) -> Self { - Self { - level: desc.level.as_str().to_string(), - backend_id: desc.backend_id.to_string(), - bound_input_class: desc.bound_input_class.unwrap_or("").to_string(), - } - } -} - -/// Persisted PTT binding state exposed to callers. -#[derive(Debug, Clone, PartialEq, Eq)] -pub struct PersistedPttBinding { - /// Stable input category string (`""`, `"keyboard"`, or - /// `"mouse-side-button"`). - pub input_class: String, - /// Display-only key label; empty when no binding is active. - pub key_label: String, -} - -impl PersistedPttBinding { - fn empty() -> Self { - Self { - input_class: String::new(), - key_label: String::new(), - } - } -} +pub use events::{ + NetworkState, PersistedPttBinding, PttDescriptorSnapshot, SessionEvent, VoiceJoinErrorCode, + VoiceJoinSyncState, +}; +use network_diagnostics::NetworkDiagnostics; /// Errors that can arise during top-level orchestration. #[derive(Debug, Error)] @@ -147,291 +114,11 @@ pub enum CoreError { Ptt(#[from] ptt::PttControllerError), } -/// High-level lifecycle event surfaced to subscribers. -/// -/// This is the minimal set needed for A.6 (reconnect banner). The -/// full event catalogue lands in A.4. -#[derive(Debug, Clone)] -pub enum SessionEvent { - /// Initial connect succeeded, or reconnect attempt succeeded. - Connected { - /// Server name reported in the snapshot. - server_name: String, - }, - /// Connection lost; the supervisor will retry. - Lost { - /// Reason classification from the protocol layer. - reason: String, - }, - /// Supervisor is sleeping before its next reconnect attempt. - Reconnecting { - /// 1-based attempt counter for the current outage. - attempt: u32, - /// Seconds the supervisor will sleep before this attempt. - delay_secs: u32, - }, - /// Supervisor gave up after `attempt` failed retries (or the - /// user explicitly disconnected mid-outage). - Disconnected { - /// Reason classification from the protocol layer. - reason: String, - }, - /// Audio engine started (e.g. after a successful reconnect with - /// reattachment). - AudioStarted, - /// Audio engine stopped (e.g. before a reconnect cycle, or by - /// explicit user action). - AudioStopped, - /// Detected desktop Push-to-Talk capability (gen2 v0.9.3 / - /// DEC-023..028). Published when the audio engine starts or - /// when the active backend transitions (for example macOS - /// permission state change). Carries only the privacy-safe - /// descriptor — capability level, backend identifier, bound - /// input class — per SRS-202 / DEC-027. - PttCapability { - /// Stable level name from `PttCapabilityLevel::as_str()`. - level: String, - /// Stable backend identifier (e.g. `"focused"`). - backend_id: String, - /// Coarse bound input class (e.g. `"keyboard"`); empty when - /// no binding is active. - bound_input_class: String, - }, - /// Voice subsystem state snapshot (SDD-094). Emitted on - /// `voice_join` / `voice_leave`, transmit-mode changes, - /// hard-mute toggles, and release-tail edits. - VoiceState { - /// True when the user has joined a voice channel via - /// `voice_join` and the audio engine is running. - in_channel: bool, - /// Active transmit mode encoded as - /// [`chanora_audio::TransmitMode::as_u8`]. - transmit_mode: u8, - /// True when the hard-mute clamp is engaged. - mute: bool, - /// Current release-tail in milliseconds (0..=500). - release_tail_ms: u32, - /// Last confirmed authoritative channel id from the - /// `channel_join` reducer projection. - current_channel_id: Option, - /// Non-authoritative pending target channel id from the - /// reducer projection. - pending_target_channel_id: Option, - /// Whether the reducer currently allows a new join intent. - can_join: bool, - /// Whether the reducer currently allows leave intent. - can_leave: bool, - /// Join projection synchronization state. - join_sync_state: VoiceJoinSyncState, - /// Last stable sanitized join error code, if any. - join_error_code: Option, - }, - /// iOS audio-session interruption state (SDD-101). Emitted when - /// interruption begins and when it ends (with the platform hint - /// indicating whether audio should resume). - InterruptionState { - /// True when interruption began, false when interruption ended. - began: bool, - /// Platform-provided resume hint. For begin events this is false. - should_resume: bool, - }, - /// A text message was received from the server. - ChatMessage { - /// Client id of the sender. - sender_id: u64, - /// Nickname of the sender. - sender_name: String, - /// Message content. - message: String, - /// Target scope (server/channel/private/poke). - target: MessageTarget, - }, - /// Human-readable TeamSpeak-style server activity. - ServerActivity { - /// Activity line text. - message: String, - }, - /// Audio route changed (speaker/earpiece/BT/wired headset). - AudioRouteChanged { - /// New audio output route. - route: AudioRoute, - }, - /// A client moved to a different channel. - ClientMoved { - /// Unique client identifier. - client_id: u64, - /// Destination channel. - new_channel_id: u64, - }, - /// A new client connected. - ClientJoined { - /// Unique client identifier. - client_id: u64, - /// Channel the client joined. - channel_id: u64, - /// Display nickname. - name: String, - /// Microphone muted state. - input_muted: bool, - /// Speaker muted state. - output_muted: bool, - /// Whether this is a server query (bot) client. - is_server_query: bool, - /// Client's talk power value. - talk_power: i32, - /// Whether the server granted temporary talk power. - talk_power_granted: bool, - }, - /// A client disconnected. - ClientLeft { - /// Unique client identifier. - client_id: u64, - /// Display nickname at time of disconnect. - name: String, - }, - /// Client properties changed. - ClientUpdated { - /// Unique client identifier. - client_id: u64, - /// Microphone muted state. - input_muted: bool, - /// Speaker muted state. - output_muted: bool, - /// Whether this is a server query (bot) client. - is_server_query: bool, - /// Client's talk power value. - talk_power: i32, - /// Whether the server granted temporary talk power. - talk_power_granted: bool, - }, - /// A new channel appeared. - ChannelAdded { - /// Unique channel identifier. - id: u64, - /// Parent channel ID. - parent: u64, - /// Channel name. - name: String, - /// Predecessor channel ID within the same parent (TeamSpeak - /// linked-list ordering hint). Zero means first child. - order: i64, - /// Whether the channel requires a password. - has_password: bool, - /// Talk power required to speak, or `None` when unrestricted. - needed_talk_power: Option, - }, - /// A channel was deleted. - ChannelRemoved { - /// Channel identifier. - id: u64, - }, - /// Channel properties changed. - ChannelUpdated { - /// Unique channel identifier. - id: u64, - /// Channel name. - name: String, - /// Whether the channel requires a password. - has_password: bool, - /// Talk power required to speak, or `None` when unrestricted. - needed_talk_power: Option, - }, -} - -/// Bridge-safe mirror of channel-join projection sync state. -#[derive(Debug, Clone, Copy)] -pub enum VoiceJoinSyncState { - /// Reducer is ready to accept channel actions. - Ready, - /// Reducer is synchronizing against an initial snapshot. - SynchronizingInitialSnapshot, - /// Reducer is synchronizing after reconnect. - SynchronizingReconnect, -} - -/// Bridge-safe mirror of stable channel-join error codes. -#[derive(Debug, Clone, Copy)] -pub enum VoiceJoinErrorCode { - /// Duplicate same-target join intent was coalesced. - DuplicateSameTargetCoalesced, - /// A different target was requested while one is already pending. - JoinAlreadyPendingDifferentTarget, - /// Join denied by server policy/permission. - JoinDenied, - /// Join failed due to protocol-level error. - JoinProtocolFailure, - /// Join failed due to transport/network error. - JoinNetworkFailure, - /// Join timed out awaiting confirmation. - JoinTimeout, - /// Pending join was superseded by user leave. - JoinSupersededByLeave, - /// Stale join outcome was ignored. - JoinStaleOutcomeIgnored, - /// Authoritative membership reconciled to different channel. - JoinReconciledDifferentChannel, - /// Join command was rejected before send acceptance. - JoinCommandRejectedBeforeSend, - /// Join intent rejected while reducer synchronizing. - JoinCannotStartWhileSynchronizing, -} - -/// Coarse OS-reported network state. Populated by the Flutter side -/// via `connectivity_plus`; on platforms where no signal is wired -/// we stay at `Unknown` forever and the supervisor falls back to -/// pure watchdog/backoff behaviour. -#[derive(Debug, Clone, Copy, PartialEq, Eq)] -pub enum NetworkState { - /// No signal seen yet — treat as ambiguous; don't change behaviour. - Unknown, - /// OS reports at least one network with internet capability. - Online, - /// OS reports no networks available. - Offline, -} - /// Channel capacity for the broadcast events. Generous because /// reconnect cycles emit several events per attempt; if subscribers /// fall behind we'd rather skip than block the supervisor. const EVENT_CHANNEL_CAPACITY: usize = 64; -/// Network diagnostics snapshot collected across connection lifetimes. -#[derive(Debug, Clone, Default)] -struct NetworkDiagnostics { - /// Total count of connects (including the initial one). - connect_count: u64, - /// Count of disconnects (graceful + loss). - disconnect_count: u64, - /// Recent loss reasons (last 8, ring buffer). - loss_reasons: Vec, -} - -impl NetworkDiagnostics { - fn record_connect(&mut self) { - self.connect_count = self.connect_count.saturating_add(1); - } - fn record_loss(&mut self, reason: &str) { - self.disconnect_count = self.disconnect_count.saturating_add(1); - if self.loss_reasons.len() >= 8 { - self.loss_reasons.remove(0); - } - self.loss_reasons.push(reason.to_string()); - } - fn summary(&self) -> String { - let mut s = format!( - "connects: {}\ndisconnects: {}\n", - self.connect_count, self.disconnect_count - ); - if !self.loss_reasons.is_empty() { - s.push_str(&format!( - "loss_reasons: [{}]\n", - self.loss_reasons.join(", ") - )); - } - s - } -} - struct SupervisorInner { /// Optional cached AudioEngineConfig — set when start_audio is /// first called, used to re-create the engine after a reconnect. @@ -477,6 +164,10 @@ struct ConnectedState { local_output_muted: bool, } +async fn take_disconnect_state(inner: &Arc>>) -> Option { + inner.lock().await.take() +} + fn normalize_channel_password(password: Option) -> Option { password .map(|p| p.trim().to_string()) @@ -1878,8 +1569,7 @@ impl ChanoraSession { /// Disconnect from the server. No-op if not connected. pub async fn disconnect(&self) -> Result<(), CoreError> { - let mut guard = self.inner.lock().await; - if let Some(mut state) = guard.take() { + if let Some(mut state) = take_disconnect_state(&self.inner).await { // Signal the supervisor to exit (cancels any backoff sleep). if let Some(tx) = state.cancel_tx.take() { let _ = tx.send(()); @@ -1895,7 +1585,7 @@ impl ChanoraSession { // Wait for the supervisor to wind down so we don't race // a redial against the explicit disconnect. if let Some(handle) = state.supervisor.take() { - let _ = handle.await; + await_supervisor_shutdown(handle, SUPERVISOR_SHUTDOWN_TIMEOUT).await; } let _ = self.events_tx.send(SessionEvent::Disconnected { reason: "user requested".to_string(), @@ -1934,6 +1624,22 @@ const WATCHDOG_PROBE_TIMEOUT: Duration = Duration::from_secs(4); /// Number of consecutive watchdog failures before the supervisor /// declares the connection lost. const WATCHDOG_MAX_MISSES: u32 = 3; +const SUPERVISOR_SHUTDOWN_TIMEOUT: Duration = Duration::from_secs(1); + +async fn await_supervisor_shutdown(mut handle: JoinHandle<()>, timeout_duration: Duration) { + if tokio::time::timeout(timeout_duration, &mut handle) + .await + .is_err() + { + warn!( + target: "chanora_core", + timeout_ms = timeout_duration.as_millis() as u64, + "supervisor did not stop before shutdown timeout" + ); + handle.abort(); + let _ = handle.await; + } +} struct SupervisorContext { state_arc: Arc>>, @@ -1989,27 +1695,77 @@ fn spawn_event_forwarders( let mut rx = delta_rx; while let Some(delta) = rx.recv().await { let event = match delta { - ProtocolDelta::ClientMoved { client_id, new_channel_id } => { - SessionEvent::ClientMoved { client_id, new_channel_id } - } - ProtocolDelta::ClientJoined { client_id, channel_id, name, input_muted, output_muted, is_server_query, talk_power, talk_power_granted } => { - SessionEvent::ClientJoined { client_id, channel_id, name, input_muted, output_muted, is_server_query, talk_power, talk_power_granted } - } + ProtocolDelta::ClientMoved { + client_id, + new_channel_id, + } => SessionEvent::ClientMoved { + client_id, + new_channel_id, + }, + ProtocolDelta::ClientJoined { + client_id, + channel_id, + name, + input_muted, + output_muted, + is_server_query, + talk_power, + talk_power_granted, + } => SessionEvent::ClientJoined { + client_id, + channel_id, + name, + input_muted, + output_muted, + is_server_query, + talk_power, + talk_power_granted, + }, ProtocolDelta::ClientLeft { client_id, name } => { SessionEvent::ClientLeft { client_id, name } } - ProtocolDelta::ClientUpdated { client_id, input_muted, output_muted, is_server_query, talk_power, talk_power_granted } => { - SessionEvent::ClientUpdated { client_id, input_muted, output_muted, is_server_query, talk_power, talk_power_granted } - } - ProtocolDelta::ChannelAdded { id, parent, name, order, has_password, needed_talk_power } => { - SessionEvent::ChannelAdded { id, parent, name, order, has_password, needed_talk_power } - } - ProtocolDelta::ChannelRemoved { id } => { - SessionEvent::ChannelRemoved { id } - } - ProtocolDelta::ChannelUpdated { id, name, has_password, needed_talk_power } => { - SessionEvent::ChannelUpdated { id, name, has_password, needed_talk_power } - } + ProtocolDelta::ClientUpdated { + client_id, + input_muted, + output_muted, + is_server_query, + talk_power, + talk_power_granted, + } => SessionEvent::ClientUpdated { + client_id, + input_muted, + output_muted, + is_server_query, + talk_power, + talk_power_granted, + }, + ProtocolDelta::ChannelAdded { + id, + parent, + name, + order, + has_password, + needed_talk_power, + } => SessionEvent::ChannelAdded { + id, + parent, + name, + order, + has_password, + needed_talk_power, + }, + ProtocolDelta::ChannelRemoved { id } => SessionEvent::ChannelRemoved { id }, + ProtocolDelta::ChannelUpdated { + id, + name, + has_password, + needed_talk_power, + } => SessionEvent::ChannelUpdated { + id, + name, + has_password, + needed_talk_power, + }, }; let _ = ev_tx.send(event); } @@ -2633,6 +2389,54 @@ mod tests { s.disconnect().await.unwrap(); } + #[tokio::test] + async fn disconnect_state_take_releases_inner_lock_before_teardown() { + let inner = Arc::new(Mutex::new(Some(()))); + + let state = super::take_disconnect_state(&inner).await; + + assert_eq!(state, Some(())); + assert!(inner.try_lock().is_ok()); + } + + #[tokio::test] + async fn supervisor_join_returns_after_shutdown_timeout() { + let handle = tokio::spawn(async { + std::future::pending::<()>().await; + }); + let start = std::time::Instant::now(); + + super::await_supervisor_shutdown(handle, Duration::from_millis(10)).await; + + assert!(start.elapsed() < Duration::from_millis(100)); + } + + #[tokio::test] + async fn supervisor_shutdown_timeout_aborts_pending_task() { + struct DropNotice(Option>); + + impl Drop for DropNotice { + fn drop(&mut self) { + if let Some(tx) = self.0.take() { + let _ = tx.send(()); + } + } + } + + let (dropped_tx, dropped_rx) = tokio::sync::oneshot::channel(); + let handle = tokio::spawn(async move { + let _notice = DropNotice(Some(dropped_tx)); + std::future::pending::<()>().await; + }); + + super::await_supervisor_shutdown(handle, Duration::from_millis(10)).await; + + tokio::time::timeout(Duration::from_millis(100), dropped_rx) + .await + .expect("pending supervisor task should be aborted") + .expect("drop notice should be delivered"); + } + #[test] fn signature_detects_in_channel_move() { use chanora_protocol::{ChannelInfo, ClientInfo}; diff --git a/core/chanora_core/src/network_diagnostics.rs b/core/chanora_core/src/network_diagnostics.rs new file mode 100644 index 0000000..faf6fb2 --- /dev/null +++ b/core/chanora_core/src/network_diagnostics.rs @@ -0,0 +1,72 @@ +use std::collections::VecDeque; + +/// Network diagnostics snapshot collected across connection lifetimes. +#[derive(Debug, Clone, Default)] +pub(crate) struct NetworkDiagnostics { + /// Total count of connects (including the initial one). + connect_count: u64, + /// Count of disconnects (graceful + loss). + disconnect_count: u64, + /// Recent loss reasons (last 8, ring buffer). + loss_reasons: VecDeque, +} + +impl NetworkDiagnostics { + pub(crate) fn record_connect(&mut self) { + self.connect_count = self.connect_count.saturating_add(1); + } + + pub(crate) fn record_loss(&mut self, reason: &str) { + self.disconnect_count = self.disconnect_count.saturating_add(1); + if self.loss_reasons.len() >= 8 { + self.loss_reasons.pop_front(); + } + self.loss_reasons.push_back(reason.to_string()); + } + + pub(crate) fn summary(&self) -> String { + let mut s = format!( + "connects: {}\ndisconnects: {}\n", + self.connect_count, self.disconnect_count + ); + if !self.loss_reasons.is_empty() { + s.push_str(&format!( + "loss_reasons: [{}]\n", + self.loss_reasons + .iter() + .map(String::as_str) + .collect::>() + .join(", ") + )); + } + s + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn network_diagnostics_keeps_last_eight_loss_reasons() { + let mut diagnostics = NetworkDiagnostics::default(); + + for i in 0..10 { + diagnostics.record_loss(&format!("loss-{i}")); + } + + assert_eq!(diagnostics.disconnect_count, 10); + assert_eq!(diagnostics.loss_reasons.len(), 8); + assert_eq!( + diagnostics.loss_reasons.front().map(String::as_str), + Some("loss-2") + ); + assert_eq!( + diagnostics.loss_reasons.back().map(String::as_str), + Some("loss-9") + ); + assert!(diagnostics.summary().contains( + "loss_reasons: [loss-2, loss-3, loss-4, loss-5, loss-6, loss-7, loss-8, loss-9]" + )); + } +} diff --git a/crates/chanora_protocol/src/adapter.rs b/crates/chanora_protocol/src/adapter.rs index 0d3b2bd..10628fe 100644 --- a/crates/chanora_protocol/src/adapter.rs +++ b/crates/chanora_protocol/src/adapter.rs @@ -46,6 +46,9 @@ use crate::ProtocolError; const SPEAKING_ACTIVITY_WINDOW: Duration = Duration::from_millis(750); const INBOUND_VOICE_SEND_TIMEOUT: Duration = Duration::from_millis(40); const PROFILE_REFRESH_RESULT_TIMEOUT: Duration = Duration::from_secs(3); +const OUTBOUND_VOICE_PACKETS_PER_TICK: usize = 8; +const DISCONNECT_REPLY_TIMEOUT: Duration = Duration::from_secs(1); +const DISCONNECT_EVENT_DRAIN_TIMEOUT: Duration = Duration::from_millis(500); type PendingMoves = HashMap< MessageHandle, @@ -84,6 +87,30 @@ async fn send_with_timeout( } } +fn drain_voice_packets_for_tick( + voice_out_rx: &mut mpsc::Receiver, + max_packets: usize, + mut send: impl FnMut(T) -> Result<(), E>, +) -> usize { + let mut drained = 0; + for _ in 0..max_packets { + let packet = match voice_out_rx.try_recv() { + Ok(packet) => packet, + Err(_) => break, + }; + let _ = send(packet); + drained += 1; + } + drained +} + +async fn bounded_drain_stream(stream: S, timeout_duration: Duration) +where + S: futures::Stream, +{ + let _ = tokio::time::timeout(timeout_duration, stream.for_each(|_| future::ready(()))).await; +} + /// Pick the TeamSpeak `client_version`/platform/signature triple /// (sourced from `ReSpeak/tsdeclarations/Versions.csv`, baked into /// `tsproto-types` at vendor-time) that best matches the *runtime* @@ -341,8 +368,21 @@ impl ProtocolClient { /// Disconnect cleanly. Blocks until the task exits. pub async fn disconnect(self) { let (tx, rx) = oneshot::channel(); - if self.tx.send(Request::Disconnect(tx)).await.is_ok() { - let _ = rx.await; + let request_path = async { + if self.tx.send(Request::Disconnect(tx)).await.is_ok() { + let _ = rx.await; + } + }; + + if tokio::time::timeout(DISCONNECT_REPLY_TIMEOUT, request_path) + .await + .is_err() + { + warn!( + target: "chanora_protocol", + timeout_ms = DISCONNECT_REPLY_TIMEOUT.as_millis() as u64, + "disconnect request did not complete before timeout" + ); } } @@ -676,12 +716,14 @@ async fn connection_task( // Main loop: pump events, service requests, forward voice. loop { - // 1. Drain any outbound voice packets first — they're time-sensitive. - while let Ok(pkt) = voice_out_rx.try_recv() { + // 1. Send a bounded batch of outbound voice packets first — they're + // time-sensitive, but control requests must still make progress. + drain_voice_packets_for_tick(&mut voice_out_rx, OUTBOUND_VOICE_PACKETS_PER_TICK, |pkt| { if let Err(e) = con.send_audio(pkt) { warn!(target: "chanora_protocol", error = %e, "send_audio failed"); } - } + Ok::<(), ()>(()) + }); // 2. Advance event stream by at most one event with a small timeout. let pump = async { @@ -689,21 +731,19 @@ async fn connection_task( tokio::time::timeout(Duration::from_millis(20), ev_stream.next()).await }; match pump.await { - Ok(Some(Ok(item))) => { - match item { - StreamItem::Audio(buf) => { - handle_audio_stream_item(&channels.voice_in, &mut voice_activity, buf).await; - } - other => handle_non_audio_stream_item( - &con, - other, - &channels.chat, - &channels.activity, - &channels.delta, - &mut pending_moves, - ), + Ok(Some(Ok(item))) => match item { + StreamItem::Audio(buf) => { + handle_audio_stream_item(&channels.voice_in, &mut voice_activity, buf).await; } - } + other => handle_non_audio_stream_item( + &con, + other, + &channels.chat, + &channels.activity, + &channels.delta, + &mut pending_moves, + ), + }, Ok(Some(Err(e))) => { warn!(target: "chanora_protocol", error = %e, "event error"); // Some errors are transient; treat persistent ones @@ -807,7 +847,7 @@ async fn connection_task( } Ok(Request::Disconnect(reply)) => { let _ = con.disconnect(DisconnectOptions::new()); - con.events().for_each(|_| future::ready(())).await; + bounded_drain_stream(con.events(), DISCONNECT_EVENT_DRAIN_TIMEOUT).await; let _ = reply.send(()); info!(target: "chanora_protocol", "clean disconnect"); exit!(DisconnectReason::UserRequested); @@ -815,7 +855,7 @@ async fn connection_task( Err(mpsc::error::TryRecvError::Empty) => {} Err(mpsc::error::TryRecvError::Disconnected) => { let _ = con.disconnect(DisconnectOptions::new()); - con.events().for_each(|_| future::ready(())).await; + bounded_drain_stream(con.events(), DISCONNECT_EVENT_DRAIN_TIMEOUT).await; info!(target: "chanora_protocol", "handle dropped; implicit disconnect"); exit!(DisconnectReason::UserRequested); } @@ -927,7 +967,9 @@ fn handle_non_audio_stream_item( let mapped = match target { tsclientlib::MessageTarget::Server => MessageTarget::Server, tsclientlib::MessageTarget::Channel => MessageTarget::Channel, - tsclientlib::MessageTarget::Client(id) => MessageTarget::Client(id.0 as u64), + tsclientlib::MessageTarget::Client(id) => { + MessageTarget::Client(id.0 as u64) + } tsclientlib::MessageTarget::Poke(id) => MessageTarget::Poke(id.0 as u64), }; let _ = chat_tx.try_send(ChatMessage { @@ -1125,7 +1167,15 @@ async fn fetch_client_profile( ) -> Result { let target_id = TsClientId(client_id as u16); - let (database_id, uid_b64, has_optional, has_connection, is_own, needs_server_groups, needs_channel_groups) = { + let ( + database_id, + uid_b64, + has_optional, + has_connection, + is_own, + needs_server_groups, + needs_channel_groups, + ) = { let state = con .get_state() .map_err(|e| ProtocolError::Backend(format!("get_state: {e}")))?; @@ -1214,15 +1264,9 @@ async fn fetch_client_profile( } let db_info = if refresh_plan.needs_client_db_info { - request_client_db_info( - con, - database_id, - channels, - pending_moves, - voice_activity, - ) - .await - .ok() + request_client_db_info(con, database_id, channels, pending_moves, voice_activity) + .await + .ok() } else { None }; @@ -1289,10 +1333,18 @@ async fn fetch_client_profile( .or_else(|| db_info.as_ref().map(|info| info.created.unix_timestamp())), last_connected_unix_seconds: optional .map(|info| info.last_connected.unix_timestamp()) - .or_else(|| db_info.as_ref().map(|info| info.last_connected.unix_timestamp())), + .or_else(|| { + db_info + .as_ref() + .map(|info| info.last_connected.unix_timestamp()) + }), connections_total: optional .map(|info| u64::from(info.connections_total)) - .or_else(|| db_info.as_ref().map(|info| u64::from(info.connections_total))), + .or_else(|| { + db_info + .as_ref() + .map(|info| u64::from(info.connections_total)) + }), online_seconds: connection .and_then(|info| info.connected_time.map(|duration| duration.whole_seconds())), idle_milliseconds: connection.map(|info| duration_millis(info.idle_time)), @@ -1325,14 +1377,10 @@ async fn fetch_client_profile( .or_else(|| db_info.as_ref().map(|info| info.bytes_uploaded_total)), packet_loss_client_to_server_total: net_stats .map(|s| s.get_packetloss()) - .or_else(|| { - connection.map(|info| info.client_to_server_packetloss_total) - }), + .or_else(|| connection.map(|info| info.client_to_server_packetloss_total)), packet_loss_server_to_client_total: net_stats .map(|s| s.get_packetloss_s2c_total()) - .or_else(|| { - connection.and_then(|info| info.server_to_client_packetloss_total) - }), + .or_else(|| connection.and_then(|info| info.server_to_client_packetloss_total)), }) } @@ -1876,10 +1924,12 @@ const _: () = { #[cfg(test)] mod tests { use super::{ - client_profile_refresh_plan, is_server_query_client_type, send_with_timeout, - server_socket_from_config, sort_channels_tree_by, std_duration_millis, - ConnectConfig, SendTimeoutError, + bounded_drain_stream, client_profile_refresh_plan, drain_voice_packets_for_tick, + is_server_query_client_type, send_with_timeout, server_socket_from_config, + sort_channels_tree_by, std_duration_millis, ConnectConfig, ProtocolClient, Request, + SendTimeoutError, DISCONNECT_REPLY_TIMEOUT, }; + use futures::stream; use std::time::Duration; use tokio::sync::mpsc; use tsproto_types::ClientType; @@ -2122,6 +2172,83 @@ mod tests { assert_eq!(result, Err(SendTimeoutError::Timeout(2))); } + + #[tokio::test] + async fn disconnect_request_send_is_bounded_when_request_channel_is_full() { + let (tx, _rx) = mpsc::channel(1); + let (reply_tx, _reply_rx) = tokio::sync::oneshot::channel(); + tx.send(Request::Snapshot(reply_tx)) + .await + .expect("seed first request"); + + let (disconnect_tx, _disconnect_rx) = tokio::sync::oneshot::channel(); + let result = send_with_timeout( + &tx, + Request::Disconnect(disconnect_tx), + Duration::from_millis(10), + ) + .await; + + assert!(matches!(result, Err(SendTimeoutError::Timeout(_)))); + } + + #[tokio::test] + async fn protocol_client_disconnect_returns_when_request_channel_is_full() { + let (tx, _rx) = mpsc::channel(1); + let (snapshot_tx, _snapshot_rx) = tokio::sync::oneshot::channel(); + tx.send(Request::Snapshot(snapshot_tx)) + .await + .expect("seed first request"); + let (voice_out_tx, _voice_out_rx) = mpsc::channel(1); + let (_voice_in_tx, voice_in_rx) = mpsc::channel(1); + let (_lost_tx, lost_rx) = tokio::sync::oneshot::channel(); + let (_chat_tx, chat_rx) = mpsc::channel(1); + let (_activity_tx, activity_rx) = mpsc::channel(1); + let (_delta_tx, delta_rx) = mpsc::channel(1); + let client = ProtocolClient { + tx, + voice_out_tx, + voice_in_rx: std::sync::Mutex::new(Some(voice_in_rx)), + lost_rx: std::sync::Mutex::new(Some(lost_rx)), + chat_rx: std::sync::Mutex::new(Some(chat_rx)), + activity_rx: std::sync::Mutex::new(Some(activity_rx)), + delta_rx: std::sync::Mutex::new(Some(delta_rx)), + }; + + tokio::time::timeout( + DISCONNECT_REPLY_TIMEOUT + Duration::from_millis(100), + client.disconnect(), + ) + .await + .expect("disconnect should not wait indefinitely for request channel capacity"); + } + + #[tokio::test] + async fn voice_drain_stops_at_per_tick_budget() { + let (tx, mut rx) = mpsc::channel(8); + for value in 0_u8..5 { + tx.send(value).await.expect("seed voice packet"); + } + + let mut sent = Vec::new(); + let drained = drain_voice_packets_for_tick(&mut rx, 2, |value| { + sent.push(value); + Ok::<(), ()>(()) + }); + + assert_eq!(drained, 2); + assert_eq!(sent, vec![0, 1]); + assert_eq!(rx.len(), 3); + } + + #[tokio::test] + async fn disconnect_stream_drain_returns_after_timeout() { + let start = tokio::time::Instant::now(); + + bounded_drain_stream(stream::pending::<()>(), Duration::from_millis(10)).await; + + assert!(start.elapsed() < Duration::from_millis(100)); + } } fn forward_delta( @@ -2204,9 +2331,7 @@ fn forward_delta( old: PropertyValue::Channel(channel), .. } => { - let _ = delta_tx.try_send(ProtocolDelta::ChannelRemoved { - id: channel.id.0, - }); + let _ = delta_tx.try_send(ProtocolDelta::ChannelRemoved { id: channel.id.0 }); } Event::PropertyChanged { id: PropertyId::Channel(channel_id),