frust-reactive 0.5.2

Reactive substrate for Frust: process-wide runtime, async executor, frame waker and shell event sources.
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
//! Process-wide deep-link source (`app_links` semantics): a shell delivers a
//! platform link (cold-start intent data, or a running app's
//! `onNewIntent`/`openURLContexts`) via [`push_deep_link`]; app
//! code reads the current state through [`deep_links`]/[`DeepLinks`] — the
//! `frust` facade re-exports both as `frust::deep_links()`/
//! `frust::DeepLinks`, so app code never names this crate directly.
//!
//! # `app_links` semantics
//!
//! Unlike the discontinued `uni_links` (an initial link fetched once, then a
//! separate stream for subsequent links), `app_links`-style delivery treats
//! the initial link as just the *first* element of one uniform stream: every
//! link (cold-start and warm) is written to the same [`RwSignal`]
//! ([`DeepLinks::latest`]), so a subscriber that only tracks `latest` sees
//! both uniformly. [`DeepLinks::initial`] additionally snapshots the very
//! first link ever pushed in this process — set at most once, readable any
//! number of times — for callers (the router glue) that need
//! cold-start precedence without setting up a subscription.
//!
//! **Documented limitation**: a push carries no "this is the cold-start link"
//! flag, so `initial` is simply "whichever link arrived first" in this
//! process. That is correct for the intended case (a real cold-start link
//! always arrives before the first app rebuild, per the shell's
//! queue-until-handle-exists contract), but a session with no real cold-start link whose
//! first-ever push happens to race ahead of the first [`deep_links`] call
//! would also see that push recorded as `initial`. Not a concern in practice
//! given the shell init order.
//!
//! # Thread contract
//!
//! [`push_deep_link`] must be called on the UI thread — the one
//! [`ReactiveRuntime::init`] ran on — mirroring `Executor::spawn_local`'s
//! contract (see `frust-reactive::executor`): the mobile shells' native
//! callbacks always run on the UI thread, so an off-thread call is a wiring
//! bug, not a runtime-data condition, and panics with the same message
//! convention `Executor::spawn_local` uses.
//!
//! A push that races ahead of [`ReactiveRuntime::init`] (an odd shell-
//! ordering edge case — the documented shell flow queues platform-side until
//! the native handle exists, so this should not happen through it) is
//! **dropped with a logged warning** rather than panicking or buffering
//! indefinitely: buffering would need to pick a bound and a flush point for a
//! path that isn't expected to be exercised, and a bare drop can never lose
//! the *cold-start* link specifically (that path is always delivered after
//! init, once the native handle exists).

use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::{Mutex, OnceLock};

use reactive_graph::signal::RwSignal;
use reactive_graph::traits::Set;

use crate::ReactiveRuntime;
use crate::executor::is_ui_thread;

/// One delivered deep link: the raw platform-provided URL/location string
/// (e.g. `myapp://profile/42`), unparsed — turning it into a route is the
/// router layer's job (the router/deep-link glue), not this crate's.
#[derive(Clone, Debug, PartialEq, Eq)]
pub struct DeepLink {
    pub url: String,
    /// Monotonic per-process delivery counter; compare sequences, not URLs, to
    /// apply each delivery exactly once — two identical URLs delivered twice
    /// are two deliveries.
    pub sequence: u64,
}

impl DeepLink {
    /// Wrap a raw URL/location string. Auto-assigns the next monotonic sequence.
    pub fn new(url: impl Into<String>) -> Self {
        let sequence = SEQUENCE.fetch_add(1, Ordering::SeqCst) + 1;
        Self {
            url: url.into(),
            sequence,
        }
    }
}

/// Process-wide monotonic sequence counter: incremented on each `DeepLink::new`.
static SEQUENCE: AtomicU64 = AtomicU64::new(0);

/// The app-facing deep-link read surface (see the module docs' `app_links`
/// semantics). Obtained via [`deep_links`] (`frust::deep_links()` at the
/// facade).
#[derive(Clone)]
pub struct DeepLinks {
    /// The first link ever pushed in this process, if any — a plain
    /// snapshot taken when [`deep_links`] is called, not itself reactive.
    /// The underlying value is set at most once (idempotent), so repeated
    /// calls to [`deep_links`] see the same `initial` once it exists.
    pub initial: Option<String>,
    /// The most recently pushed link — cold-start or warm, uniformly (see
    /// the module docs). Use [`DeepLink::sequence`] to deduplicate deliveries.
    /// Read/track it with the `Get`/`Track` traits (`frust::{Get, Track}`)
    /// the same way any other `RwSignal` is read.
    pub latest: RwSignal<Option<DeepLink>>,
}

/// The process-wide deep-link slot: the first-ever-pushed snapshot plus the
/// live `latest` signal. Lazily created (mirroring [`ReactiveRuntime`]'s own
/// process-wide, lazily-installed static) on first access once
/// [`ReactiveRuntime`] exists, without needing to modify
/// `ReactiveRuntime::init` itself.
struct Slot {
    initial: Mutex<Option<String>>,
    latest: RwSignal<Option<DeepLink>>,
}

static SLOT: OnceLock<Slot> = OnceLock::new();

/// Returns the process-wide slot, creating it (under the reactive root
/// [`Owner`](reactive_graph::owner::Owner)) on first access.
///
/// # Panics
///
/// Panics if [`ReactiveRuntime::init`] has not run yet — reading or pushing a
/// deep link before the reactive runtime exists across every other path
/// (`push_deep_link`'s pre-init case is handled separately, before this is
/// ever called) is a genuine wiring bug: the caller must initialize the
/// runtime first.
fn slot() -> &'static Slot {
    SLOT.get_or_init(|| {
        let rt = ReactiveRuntime::get().expect(
            "frust-reactive: deep_links() was called before ReactiveRuntime::init — an app \
             must run under the Frust facade's entry point (which initializes the reactive \
             runtime) before reading deep links",
        );
        Slot {
            initial: Mutex::new(None),
            latest: rt.with_owner(|| RwSignal::new(None)),
        }
    })
}

/// Test-only: clears the recorded `initial` URL so a test can observe a first
/// push regardless of which tests ran earlier in the same process. The slot is
/// shared process-wide, so callers must hold `WAKER_TEST_LOCK`.
#[cfg(test)]
fn reset_initial_for_test() {
    *slot()
        .initial
        .lock()
        .expect("frust-reactive: deep_link initial mutex poisoned") = None;
}

/// Deliver a platform deep link (cold-start or warm) into the process-wide
/// source. Called by a shell (the Android/iOS FFI glue) on the UI
/// thread; app code never calls this directly.
///
/// The first call in a process snapshots its URL into
/// [`DeepLinks::initial`] (idempotent — later calls do not overwrite it);
/// every call (including the first) also writes [`DeepLinks::latest`], so a
/// tracked reader observes both cold-start and warm links through the same
/// signal (see the module docs' `app_links` semantics). Each delivered link
/// is tagged with a monotonically increasing [`DeepLink::sequence`] — use it
/// to deduplicate deliveries rather than comparing URL text.
///
/// # Panics
///
/// Panics if called off the UI thread (see the module docs' thread
/// contract). A call before [`ReactiveRuntime::init`] does **not** panic — it
/// is dropped with a logged warning (see the module docs).
pub fn push_deep_link(url: impl Into<String>) {
    let url = url.into();

    if !is_ui_thread() {
        panic!(
            "frust-reactive: push_deep_link was called off the UI thread. Deep links can \
             only be pushed from the UI thread (the one `ReactiveRuntime::init` ran on) — this \
             is a wiring bug: route the platform delivery through the UI thread before pushing, \
             the same contract `Executor::spawn_local` enforces."
        );
    }

    if ReactiveRuntime::get().is_none() {
        eprintln!(
            "frust-reactive: push_deep_link(\"{url}\") dropped — ReactiveRuntime::init has \
             not run yet. A shell should queue a link platform-side until its native handle \
             exists; reaching this indicates an odd init-ordering race, not normal \
             operation."
        );
        return;
    }

    let slot = slot();
    {
        let mut initial = slot
            .initial
            .lock()
            .expect("frust-reactive: deep_link initial mutex poisoned");
        if initial.is_none() {
            *initial = Some(url.clone());
        }
    }
    slot.latest.set(Some(DeepLink::new(url)));
}

/// The current deep-link read surface: a snapshot of [`DeepLinks::initial`]
/// plus the live [`DeepLinks::latest`] signal. Call from a tracked context
/// (e.g. inside `Component::build`) to observe subsequent pushes as they
/// arrive.
pub fn deep_links() -> DeepLinks {
    let slot = slot();
    let initial = slot
        .initial
        .lock()
        .expect("frust-reactive: deep_link initial mutex poisoned")
        .clone();
    DeepLinks {
        initial,
        latest: slot.latest,
    }
}

#[cfg(test)]
mod tests {
    use super::*;
    use crate::{FrameWaker, TrackedScope};
    use reactive_graph::traits::{Get, GetUntracked};
    use std::sync::Arc;
    use std::sync::atomic::{AtomicUsize, Ordering};

    /// A recording waker: an `Arc<AtomicUsize>` bumped once per `wake()`.
    fn recording_waker() -> (FrameWaker, Arc<AtomicUsize>) {
        let counter = Arc::new(AtomicUsize::new(0));
        let seen = counter.clone();
        let waker: FrameWaker = Arc::new(move || {
            counter.fetch_add(1, Ordering::SeqCst);
        });
        (waker, seen)
    }

    /// Acceptance criterion (a): two pushes of the same URL yield sequences n
    /// and n+1 with equal URLs.
    #[test]
    fn sequence_increments_on_identical_urls() {
        let _guard = crate::WAKER_TEST_LOCK
            .lock()
            .unwrap_or_else(|e| e.into_inner());

        let (waker, _) = recording_waker();
        let _rt = ReactiveRuntime::init(waker);

        let url = "frust-test://identical".to_string();
        push_deep_link(url.clone());
        let link1 = deep_links().latest.get_untracked().unwrap();

        push_deep_link(url.clone());
        let link2 = deep_links().latest.get_untracked().unwrap();

        // Same URL, but strictly increasing sequences
        assert_eq!(link1.url, link2.url);
        assert_eq!(link1.sequence + 1, link2.sequence);
    }

    /// Acceptance criterion (b): `initial` (if constructed through
    /// DeepLink::new) has a sequence lower than any later push.
    #[test]
    fn initial_sequence_lower_than_later_push() {
        let _guard = crate::WAKER_TEST_LOCK
            .lock()
            .unwrap_or_else(|e| e.into_inner());

        let (waker, _) = recording_waker();
        let _rt = ReactiveRuntime::init(waker);

        let url1 = "frust-test://initial".to_string();
        push_deep_link(url1.clone());
        let links1 = deep_links();
        let initial_seq = links1
            .latest
            .get_untracked()
            .expect("initial push should have produced a link")
            .sequence;

        let url2 = "frust-test://later".to_string();
        push_deep_link(url2);
        let links2 = deep_links();
        let later_seq = links2
            .latest
            .get_untracked()
            .expect("second push should have produced a link")
            .sequence;

        assert!(
            initial_seq < later_seq,
            "initial sequence {} should be lower than later sequence {}",
            initial_seq,
            later_seq
        );
    }

    /// Acceptance criterion (c): sequences are strictly increasing across
    /// mixed URLs.
    #[test]
    fn sequences_strictly_increasing_across_mixed_urls() {
        let _guard = crate::WAKER_TEST_LOCK
            .lock()
            .unwrap_or_else(|e| e.into_inner());

        let (waker, _) = recording_waker();
        let _rt = ReactiveRuntime::init(waker);

        let urls = vec![
            "frust-test://a",
            "frust-test://b",
            "frust-test://a", // repeat
            "frust-test://c",
            "frust-test://b", // repeat
        ];

        let mut sequences = Vec::new();
        for url in urls {
            push_deep_link(url);
            let seq = deep_links()
                .latest
                .get_untracked()
                .expect("push should produce a link")
                .sequence;
            sequences.push(seq);
        }

        // Verify strictly increasing
        for i in 1..sequences.len() {
            assert!(
                sequences[i - 1] < sequences[i],
                "sequence {} should be strictly less than {}",
                sequences[i - 1],
                sequences[i]
            );
        }
    }

    /// Original waker/wake-bridge contract test, adapted for the sequence field.
    #[test]
    fn deep_link_push_and_wake_bridge() {
        let _guard = crate::WAKER_TEST_LOCK
            .lock()
            .unwrap_or_else(|e| e.into_inner());

        let (waker, wakes) = recording_waker();
        let _rt = ReactiveRuntime::init(waker);

        // Other tests push links into the same process-wide slot, so start
        // from a state where no link has been recorded as `initial`.
        reset_initial_for_test();

        // Criterion: push before any tracked read — a late subscriber sees it
        // immediately via BOTH the `initial` snapshot and the live `latest`
        // signal (app_links semantics: cold-start and warm links flow through
        // the same source).
        let url1 = "frust-test://a/1".to_string();
        push_deep_link(url1.clone());

        let links = deep_links();
        assert_eq!(links.initial.as_deref(), Some(url1.as_str()));
        let link1 = links
            .latest
            .get_untracked()
            .expect("should have pushed a link");
        assert_eq!(link1.url, url1);
        assert!(link1.sequence > 0, "sequence should be assigned");

        // A tracked scope reading `latest` after the push observes the
        // already-pushed link with no wake needed — an ordinary read.
        let scope = TrackedScope::new();
        let seen = scope.track(|| deep_links().latest.get());
        let seen_link = seen.expect("should have a link");
        assert_eq!(seen_link.url, url1);
        assert_eq!(
            seen_link.sequence, link1.sequence,
            "sequence should be stable"
        );
        assert!(!scope.is_dirty(), "a fresh track starts clean");

        // Criterion: push after a tracked read fires the signal-write -> wake
        // contract — the tracked scope re-dirties and the waker fires exactly
        // once (coalesced).
        let before = wakes.load(Ordering::SeqCst);
        let url2 = "frust-test://b/2".to_string();
        push_deep_link(url2.clone());
        assert!(
            scope.is_dirty(),
            "push_deep_link must dirty a scope tracking `latest`"
        );
        assert_eq!(
            wakes.load(Ordering::SeqCst) - before,
            1,
            "a push must fire the waker exactly once"
        );

        // `initial` is set once and never overwritten, even though `latest`
        // has since moved on to the second link.
        let links2 = deep_links();
        assert_eq!(
            links2.initial.as_deref(),
            Some(url1.as_str()),
            "initial is set once and stays"
        );
        let link2 = links2
            .latest
            .get_untracked()
            .expect("should have pushed a second link");
        assert_eq!(link2.url, url2);
        assert!(
            link2.sequence > link1.sequence,
            "second sequence should be higher"
        );
    }

    /// Criterion: the UI-thread contract is enforced (mirrors
    /// `Executor::spawn_local`'s off-thread panic). Also serializes on the
    /// waker lock since `ReactiveRuntime::init` swaps the process-wide waker.
    #[test]
    #[should_panic(expected = "wiring bug")]
    fn push_off_ui_thread_panics() {
        let _guard = crate::WAKER_TEST_LOCK
            .lock()
            .unwrap_or_else(|e| e.into_inner());
        let _rt = ReactiveRuntime::init(Arc::new(|| {}));

        std::thread::spawn(|| {
            push_deep_link("frust-test://off-thread");
        })
        .join()
        .unwrap_or_else(|e| std::panic::resume_unwind(e));
    }
}