openrtc 1.0.2

OpenRTC: a Rust-first P2P runtime for device discovery, signaling, and iroh/QUIC networking.
Documentation
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
//! Phase 3 unit tests for `DriveGrantConnectionActor` + `DriveGrantActorRegistry`.
//!
//! These tests use a fake `DriveGrantTransport` so the actor can be
//! exercised without booting an iroh endpoint. They exercise the
//! lifecycle invariants the remediation plan calls out:
//!
//!   1. Concurrent `ensure_ready()` from N requesters yields exactly one
//!      underlying `transport.dial()` invocation.
//!   2. `shutdown()` is idempotent and prevents subsequent `ensure_ready()`
//!      from issuing more dials.
//!   3. The registry hands out one shared actor per key (idempotent
//!      `get_or_spawn`); different keys get different actors.
//!   4. WebRTC observation is non-blocking and reflects state transitions.
//!   5. `open_logical_channel` enforces ready ordering and surfaces the
//!      Phase 6 placeholder error today.

#![cfg(all(test, not(target_arch = "wasm32")))]

use std::sync::atomic::{AtomicU32, Ordering};
use std::sync::Arc;
use std::time::Duration;

use async_trait::async_trait;
use tokio::sync::Notify;
use tokio::time::timeout;

use super::auth_readiness::{AuthLeg, AuthLegState, AuthReadinessStore};
use super::drive_grant_actor::{
    ActorError, DriveGrantActorConfig, DriveGrantActorRegistry, DriveGrantConnectionActor,
    DriveGrantConnectionKey, DriveGrantTransport, WebRtcState,
};

/// Fake transport that counts dials. Test-controlled outcome via a closure.
struct FakeTransport {
    dial_count: Arc<AtomicU32>,
    close_count: Arc<AtomicU32>,
    /// Per-call outcome callback. Default success.
    outcome: Arc<dyn Fn(u32) -> Result<(), String> + Send + Sync>,
    /// Optional: notify each time a dial begins, so tests can synchronise
    /// concurrency carefully.
    on_dial_begin: Arc<Notify>,
}

impl FakeTransport {
    fn new(outcome: impl Fn(u32) -> Result<(), String> + Send + Sync + 'static) -> Arc<Self> {
        Arc::new(Self {
            dial_count: Arc::new(AtomicU32::new(0)),
            close_count: Arc::new(AtomicU32::new(0)),
            outcome: Arc::new(outcome),
            on_dial_begin: Arc::new(Notify::new()),
        })
    }

    fn dials(&self) -> u32 {
        self.dial_count.load(Ordering::SeqCst)
    }

    fn closes(&self) -> u32 {
        self.close_count.load(Ordering::SeqCst)
    }
}

#[async_trait]
impl DriveGrantTransport for FakeTransport {
    async fn dial(&self, _key: &DriveGrantConnectionKey) -> Result<(), String> {
        let n = self.dial_count.fetch_add(1, Ordering::SeqCst) + 1;
        self.on_dial_begin.notify_waiters();
        // Make the dial yield to the executor so concurrent waiters can
        // observe in-progress dials in tests that care.
        tokio::task::yield_now().await;
        (self.outcome)(n)
    }

    async fn close(&self, _key: &DriveGrantConnectionKey) {
        self.close_count.fetch_add(1, Ordering::SeqCst);
    }
}

fn key() -> DriveGrantConnectionKey {
    DriveGrantConnectionKey::new("drive-grant:test", "node-A")
}

fn fast_config() -> DriveGrantActorConfig {
    DriveGrantActorConfig {
        initial_retry: Duration::from_millis(1),
        max_retry: Duration::from_millis(4),
        max_consecutive_failures: 3,
        auth_readiness_wait_timeout: Some(Duration::from_millis(250)),
    }
}

#[tokio::test]
async fn concurrent_ensure_ready_coalesces_to_one_dial() {
    let transport = FakeTransport::new(|_| Ok(()));
    let registry = DriveGrantActorRegistry::with_config(transport.clone(), fast_config());

    let actor = registry.get_or_spawn(key()).await;

    // 16 concurrent requesters, all coming in before the first dial
    // resolves. Each should observe Ready, but the underlying transport
    // must see exactly one dial.
    let mut handles = Vec::new();
    for _ in 0..16 {
        let actor = actor.clone();
        handles.push(tokio::spawn(async move { actor.ensure_ready().await }));
    }
    for h in handles {
        h.await.expect("join").expect("ensure_ready");
    }
    assert_eq!(
        transport.dials(),
        1,
        "concurrent ensure_ready must coalesce to one dial",
    );
}

#[tokio::test]
async fn ensure_ready_is_memoized_after_first_success() {
    let transport = FakeTransport::new(|_| Ok(()));
    let registry = DriveGrantActorRegistry::with_config(transport.clone(), fast_config());
    let actor = registry.get_or_spawn(key()).await;

    actor.ensure_ready().await.expect("first");
    actor.ensure_ready().await.expect("second");
    actor.ensure_ready().await.expect("third");
    assert_eq!(transport.dials(), 1);
}

#[tokio::test]
async fn ensure_ready_retries_on_transient_failure_then_succeeds() {
    // Fail twice, then succeed. Within max_consecutive_failures (3).
    let transport = FakeTransport::new(|attempt| {
        if attempt < 3 {
            Err(format!("transient-{attempt}"))
        } else {
            Ok(())
        }
    });
    let registry = DriveGrantActorRegistry::with_config(transport.clone(), fast_config());
    let actor = registry.get_or_spawn(key()).await;

    actor.ensure_ready().await.expect("ready after retries");
    assert_eq!(transport.dials(), 3);

    // Second call must not redial.
    actor.ensure_ready().await.expect("memoized");
    assert_eq!(transport.dials(), 3);
}

#[tokio::test]
async fn ensure_ready_surfaces_dial_failed_after_budget() {
    let transport = FakeTransport::new(|_| Err("permanent".to_string()));
    let registry = DriveGrantActorRegistry::with_config(transport.clone(), fast_config());
    let actor = registry.get_or_spawn(key()).await;

    let result = actor.ensure_ready().await;
    match result {
        Err(ActorError::DialFailed(reason)) => {
            assert!(reason.contains("permanent"));
        }
        other => panic!("expected DialFailed, got {other:?}"),
    }
    // Budget = 3 → exactly 3 attempts before surfacing the error.
    assert_eq!(transport.dials(), 3);
}

#[tokio::test]
async fn shutdown_is_idempotent_and_blocks_subsequent_ensure_ready() {
    let transport = FakeTransport::new(|_| Ok(()));
    let registry = DriveGrantActorRegistry::with_config(transport.clone(), fast_config());
    let actor = registry.get_or_spawn(key()).await;

    actor.ensure_ready().await.expect("first ready");
    assert_eq!(transport.dials(), 1);

    actor.shutdown().await;
    actor.shutdown().await; // idempotent
    assert_eq!(transport.closes(), 1, "close fires once");

    let result = actor.ensure_ready().await;
    assert!(matches!(result, Err(ActorError::Shutdown)));
    // No further dials after shutdown.
    assert_eq!(transport.dials(), 1);
}

#[tokio::test]
async fn registry_returns_same_actor_for_same_key() {
    let transport = FakeTransport::new(|_| Ok(()));
    let registry = DriveGrantActorRegistry::with_config(transport.clone(), fast_config());

    let a = registry.get_or_spawn(key()).await;
    let b = registry.get_or_spawn(key()).await;
    assert!(Arc::ptr_eq(&a, &b), "same key → same actor");
}

#[tokio::test]
async fn registry_returns_distinct_actors_for_different_keys() {
    let transport = FakeTransport::new(|_| Ok(()));
    let registry = DriveGrantActorRegistry::with_config(transport.clone(), fast_config());

    let a = registry
        .get_or_spawn(DriveGrantConnectionKey::new("scope-a", "node-1"))
        .await;
    let b = registry
        .get_or_spawn(DriveGrantConnectionKey::new("scope-b", "node-1"))
        .await;
    let c = registry
        .get_or_spawn(DriveGrantConnectionKey::new("scope-a", "node-2"))
        .await;

    assert!(
        !Arc::ptr_eq(&a, &b),
        "different scope must yield different actor"
    );
    assert!(
        !Arc::ptr_eq(&a, &c),
        "different node must yield different actor"
    );
    assert_eq!(registry.len().await, 3);
}

#[tokio::test]
async fn registry_get_finds_existing_actor_without_spawning() {
    let transport = FakeTransport::new(|_| Ok(()));
    let registry = DriveGrantActorRegistry::with_config(transport.clone(), fast_config());

    assert!(registry.get(&key()).await.is_none());

    let _ = registry.get_or_spawn(key()).await;
    assert!(registry.get(&key()).await.is_some());
    assert_eq!(registry.len().await, 1);
}

#[tokio::test]
async fn registry_shutdown_all_tears_down_every_actor() {
    let transport = FakeTransport::new(|_| Ok(()));
    let registry = DriveGrantActorRegistry::with_config(transport.clone(), fast_config());

    let a = registry
        .get_or_spawn(DriveGrantConnectionKey::new("scope", "node-1"))
        .await;
    let b = registry
        .get_or_spawn(DriveGrantConnectionKey::new("scope", "node-2"))
        .await;
    a.ensure_ready().await.unwrap();
    b.ensure_ready().await.unwrap();

    registry.shutdown_all().await;
    assert_eq!(registry.len().await, 0);
    assert_eq!(transport.closes(), 2);
    assert!(matches!(a.ensure_ready().await, Err(ActorError::Shutdown)));
    assert!(matches!(b.ensure_ready().await, Err(ActorError::Shutdown)));
}

#[tokio::test]
async fn observe_webrtc_reflects_state_transitions_without_blocking_dial() {
    let transport = FakeTransport::new(|_| Ok(()));
    let registry = DriveGrantActorRegistry::with_config(transport.clone(), fast_config());
    let actor = registry.get_or_spawn(key()).await;

    let mut rx = actor.observe_webrtc();
    assert_eq!(*rx.borrow_and_update(), WebRtcState::Unknown);

    // ensure_ready must succeed even with WebRTC stuck on Unknown — Phase 2
    // invariant: drive-view correctness must not depend on WebRTC.
    timeout(Duration::from_millis(500), actor.ensure_ready())
        .await
        .expect("did not block on webrtc")
        .expect("ready");

    actor.set_webrtc_state(WebRtcState::Connecting);
    rx.changed().await.unwrap();
    assert_eq!(*rx.borrow_and_update(), WebRtcState::Connecting);

    actor.set_webrtc_state(WebRtcState::Connected);
    rx.changed().await.unwrap();
    assert_eq!(*rx.borrow_and_update(), WebRtcState::Connected);
}

#[tokio::test]
async fn observe_webrtc_signals_failed_on_shutdown() {
    let transport = FakeTransport::new(|_| Ok(()));
    let registry = DriveGrantActorRegistry::with_config(transport.clone(), fast_config());
    let actor = registry.get_or_spawn(key()).await;

    let mut rx = actor.observe_webrtc();
    actor.shutdown().await;
    rx.changed().await.unwrap();
    assert!(matches!(*rx.borrow(), WebRtcState::Failed(_)));
}

#[tokio::test]
async fn open_logical_channel_requires_ensure_ready_first() {
    let transport = FakeTransport::new(|_| Ok(()));
    let registry = DriveGrantActorRegistry::with_config(transport.clone(), fast_config());
    let actor = registry.get_or_spawn(key()).await;

    let err = actor
        .open_logical_channel("drive-view.list")
        .await
        .expect_err("must reject before ensure_ready");
    assert!(matches!(err, ActorError::NotReady));
}

#[tokio::test]
async fn open_logical_channel_returns_handle_for_known_labels() {
    let transport = FakeTransport::new(|_| Ok(()));
    let registry = DriveGrantActorRegistry::with_config(transport.clone(), fast_config());
    let actor = registry.get_or_spawn(key()).await;
    actor.ensure_ready().await.unwrap();

    // Phase 6: known labels return a real handle.
    let handle = actor
        .open_logical_channel("drive-view.list")
        .await
        .expect("known label should succeed");
    assert_eq!(handle.label(), "drive-view.list");

    let handle = actor
        .open_logical_channel("drive-view.buckets")
        .await
        .expect("known label should succeed");
    assert_eq!(handle.label(), "drive-view.buckets");

    // Unknown labels still error.
    let err = actor
        .open_logical_channel("unknown-channel")
        .await
        .expect_err("unknown label should fail");
    match err {
        ActorError::LogicalChannelNotImplemented(label) => {
            assert_eq!(label, "unknown-channel");
        }
        other => panic!("expected LogicalChannelNotImplemented, got {other:?}"),
    }

    // The test stub roundtrips the label so Phase 4 wiring can assert
    // observers really pass the right channel name through.
    let handle = actor
        .open_logical_channel_test_stub("drive-view.list")
        .await
        .unwrap();
    assert_eq!(handle.label(), "drive-view.list");
}

#[tokio::test]
async fn debug_state_tracks_failures_and_readiness() {
    let attempts = Arc::new(AtomicU32::new(0));
    let attempts_clone = attempts.clone();
    let transport = FakeTransport::new(move |attempt| {
        attempts_clone.store(attempt, Ordering::SeqCst);
        if attempt < 2 {
            Err("retry-me".to_string())
        } else {
            Ok(())
        }
    });
    let registry = DriveGrantActorRegistry::with_config(transport.clone(), fast_config());
    let actor = registry.get_or_spawn(key()).await;

    let initial = actor.debug_state().await;
    assert!(!initial.ready);
    assert!(!initial.ever_succeeded);
    assert!(!initial.shutdown);

    actor.ensure_ready().await.unwrap();

    let after = actor.debug_state().await;
    assert!(after.ready);
    assert!(after.ever_succeeded);
    assert_eq!(after.consecutive_failures, 0, "reset after success");

    actor.shutdown().await;
    let final_state = actor.debug_state().await;
    assert!(final_state.shutdown);
}

#[tokio::test]
async fn auth_readiness_gate_blocks_first_dial_until_ready() {
    let transport = FakeTransport::new(|_| Ok(()));
    let store = Arc::new(AuthReadinessStore::new());
    let registry = DriveGrantActorRegistry::with_auth_readiness(
        transport.clone(),
        fast_config(),
        store.clone(),
        vec![AuthLeg::RuntimeAuth, AuthLeg::Firestore],
    );
    let actor = registry.get_or_spawn(key()).await;

    let pending = {
        let actor = actor.clone();
        tokio::spawn(async move { actor.ensure_ready().await })
    };

    tokio::time::sleep(Duration::from_millis(25)).await;
    assert_eq!(transport.dials(), 0, "dial must wait for auth readiness");

    store.mark_pluto_rtc_auth(AuthLegState::Ready, None);
    store.mark_runtime_auth(AuthLegState::Ready, None);

    pending.await.expect("join").expect("ready");
    assert_eq!(transport.dials(), 1);
}

#[tokio::test]
async fn auth_readiness_timeout_surfaces_without_dialing() {
    let transport = FakeTransport::new(|_| Ok(()));
    let store = Arc::new(AuthReadinessStore::new());
    let registry = DriveGrantActorRegistry::with_auth_readiness(
        transport.clone(),
        fast_config(),
        store,
        vec![AuthLeg::RuntimeAuth, AuthLeg::Firestore],
    );
    let actor = registry.get_or_spawn(key()).await;

    let result = actor.ensure_ready().await;
    match result {
        Err(ActorError::AuthNotReady(reason)) => {
            assert!(reason.contains("RuntimeAuth"));
            assert!(reason.contains("Firestore"));
        }
        other => panic!("expected AuthNotReady, got {other:?}"),
    }
    assert_eq!(transport.dials(), 0);
}

/// Phase 3 lifecycle invariant: the actor's internal `select_loop` mutex
/// serializes the dial path so concurrent ensure_ready callers from many
/// threads never run dial bodies at the same time. (This is the
/// representative coverage for the "inbound replacement during in-flight
/// admission is queued, never racing the response writer" property — full
/// inbound-dial coverage lands in Phase 4 once the actor is wired into
/// `Client::connect_device`.)
#[tokio::test]
async fn select_loop_serializes_dials() {
    let in_flight = Arc::new(AtomicU32::new(0));
    let max_observed = Arc::new(AtomicU32::new(0));
    let in_flight_clone = in_flight.clone();
    let max_observed_clone = max_observed.clone();
    let transport = FakeTransport::new(move |_| {
        let now = in_flight_clone.fetch_add(1, Ordering::SeqCst) + 1;
        max_observed_clone.fetch_max(now, Ordering::SeqCst);
        // Yield several times so any leak in the serialization shows up
        // as max_observed > 1.
        std::thread::sleep(Duration::from_millis(2));
        in_flight_clone.fetch_sub(1, Ordering::SeqCst);
        Err("force-retry".to_string())
    });
    let registry = DriveGrantActorRegistry::with_config(
        transport.clone(),
        DriveGrantActorConfig {
            initial_retry: Duration::from_millis(1),
            max_retry: Duration::from_millis(2),
            max_consecutive_failures: 4,
            auth_readiness_wait_timeout: Some(Duration::from_millis(250)),
        },
    );
    let actor = registry.get_or_spawn(key()).await;

    let mut handles = Vec::new();
    for _ in 0..6 {
        let actor = actor.clone();
        handles.push(tokio::spawn(async move {
            // Don't unwrap; budget is small so we expect failure surface.
            let _ = actor.ensure_ready().await;
        }));
    }
    for h in handles {
        h.await.unwrap();
    }
    assert_eq!(
        max_observed.load(Ordering::SeqCst),
        1,
        "at most one dial in-flight at a time",
    );

    // ensure helper crate stays warm
    let _: Arc<DriveGrantConnectionActor> = actor;
}