Skip to content
Merged
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
100 changes: 100 additions & 0 deletions src/client.rs
Original file line number Diff line number Diff line change
Expand Up @@ -2406,6 +2406,53 @@ impl Client {
});
}
}
"401" => {
// 401: unauthorized — session invalid, needs re-authentication.
// Matches WA Web's handling of unauthorized stream errors.
info!("Got 401 stream error (unauthorized). Logging out.");
self.expect_disconnect().await;
self.enable_auto_reconnect.store(false, Ordering::Relaxed);
self.core.event_bus.dispatch(&Event::LoggedOut(
crate::types::events::LoggedOut {
on_connect: false,
reason: ConnectFailureReason::LoggedOut,
},
));

let transport_opt = self.transport.lock().await.clone();
if let Some(transport) = transport_opt {
tokio::spawn(async move {
info!("Disconnecting transport after 401");
transport.disconnect().await;
});
}
}
"409" => {
// 409: conflict — another client instance connected.
// Same semantics as conflict child element but via code.
info!("Got 409 stream error (conflict). Another session replaced this one.");
self.expect_disconnect().await;
self.enable_auto_reconnect.store(false, Ordering::Relaxed);
self.core
.event_bus
.dispatch(&Event::StreamReplaced(crate::types::events::StreamReplaced));

let transport_opt = self.transport.lock().await.clone();
if let Some(transport) = transport_opt {
tokio::spawn(async move {
info!("Disconnecting transport after 409");
transport.disconnect().await;
});
}
}
"429" => {
// 429: rate limited — server is throttling connections.
// Auto-reconnect with extended backoff.
warn!(
"Got 429 stream error (rate limited). Will auto-reconnect with extended backoff."
);
self.auto_reconnect_errors.fetch_add(5, Ordering::Relaxed);
}
"503" => {
info!("Got 503 service unavailable, will auto-reconnect.");
}
Expand Down Expand Up @@ -4587,6 +4634,59 @@ mod tests {
);
}

// ── stream error tests ─────────────────────────────────────────────

#[tokio::test]
async fn test_stream_error_401_disables_reconnect() {
let client = create_offline_sync_test_client().await;
let node = NodeBuilder::new("stream:error").attr("code", "401").build();
client.handle_stream_error(&node).await;
assert!(
!client.enable_auto_reconnect.load(Ordering::Relaxed),
"401 should disable auto-reconnect"
);
}

#[tokio::test]
async fn test_stream_error_409_disables_reconnect() {
let client = create_offline_sync_test_client().await;
let node = NodeBuilder::new("stream:error").attr("code", "409").build();
client.handle_stream_error(&node).await;
assert!(
!client.enable_auto_reconnect.load(Ordering::Relaxed),
"409 should disable auto-reconnect"
);
}

#[tokio::test]
async fn test_stream_error_429_keeps_reconnect_with_backoff() {
let client = create_offline_sync_test_client().await;
let before = client.auto_reconnect_errors.load(Ordering::Relaxed);
let node = NodeBuilder::new("stream:error").attr("code", "429").build();
client.handle_stream_error(&node).await;
assert!(
client.enable_auto_reconnect.load(Ordering::Relaxed),
"429 should keep auto-reconnect enabled"
);
let after = client.auto_reconnect_errors.load(Ordering::Relaxed);
assert_eq!(
after,
before + 5,
"429 should increase backoff by exactly 5: before={before}, after={after}"
);
}

#[tokio::test]
async fn test_stream_error_503_keeps_reconnect() {
let client = create_offline_sync_test_client().await;
let node = NodeBuilder::new("stream:error").attr("code", "503").build();
client.handle_stream_error(&node).await;
assert!(
client.enable_auto_reconnect.load(Ordering::Relaxed),
"503 should keep auto-reconnect enabled"
);
}

#[tokio::test]
async fn test_custom_cache_config_is_respected() {
use crate::cache_config::{CacheConfig, CacheEntryConfig};
Expand Down
Loading