Skip to content

Subscriptions

Every reactive primitive is watched the same way, so what is here is true of a field, a cell and a map alike.

let sub = state.port().subscribe(move |port| {
seen.lock().unwrap().push(*port);
});
state.port().set(9090)?;
assert_eq!(*heard.lock().unwrap(), [9090]);
drop(sub);
state.port().set(1234)?;
assert_eq!(*heard.lock().unwrap(), [9090]);

subscribe returns a guard and the callback lives exactly as long as it does. There is no separate unsubscribe: dropping the guard is what unregisters.

Which is where the trap is. A guard assigned to _ is dropped at the end of that same statement, and the callback never fires at all:

let ignored = Arc::clone(&heard);
let _ = state.port().subscribe(move |port| {
ignored.lock().unwrap().push(*port);
});
state.port().set(4321)?;
assert_eq!(*heard.lock().unwrap(), [9090]);

heard still holds [9090], from the first subscription. So every example here binds the guard to a name, and the place to keep it is beside whatever the callback writes into, so the two die together.

let mut scope = ReactiveScope::new();
state
.port()
.subscribe(|port| println!("port {port}"))
.watch(&mut scope);
state
.host()
.subscribe(|host| println!("host {host}"))
.watch(&mut scope);
scope.clear();

ReactiveScope is one owner for many guards. clear drops them all; dropping the scope does the same.

subscribe covers the common case. Anything past it goes through subscription_with(), whose links compose in any order:

linkwhat it does
.external()skip changes this handle made
.key(k)on a map, narrow to one entry
.register(f)finish, returning the guard
.register_with_source(f)the same, with who made the change
.stream()finish as a Stream rather than a callback

Taking the changes into a loop of your own

Section titled “Taking the changes into a loop of your own”

A callback must be Send + Sync, because a change made to the file outside the process is delivered from a watcher thread. That rules out Rc state and most GUI context handles.

.stream() finishes the subscription as a Stream instead. The value crosses the thread boundary and nothing else does, so what you do with it runs on the thread that drives the loop:

let mut ports = state.port().subscription_with().stream();
state.port().set(9090)?;
state.port().set(1234)?;
let mut heard = Vec::new();
futures::executor::block_on(async {
while let Some(port) = ports.next().await {
heard.push(port);
if port == 1234 {
break;
}
}
});
assert_eq!(heard, [9090, 1234]);

A stream yields every change rather than coalescing - it is a sequence, and coalescing downstream is your choice. Dropping it ends the subscription.

Every write carries the id of the handle that made it, and .external() is the filter that uses it:

let watcher = state.port().fork();
let _sub = state
.port()
.subscription_with()
.external()
.register(move |port| {
seen.lock().unwrap().push(*port);
});
state.port().set(8080)?;
watcher.set(9090)?;
assert_eq!(*heard.lock().unwrap(), [9090]);

A background thread writing while the UI reacts is the usual shape: the thread holds a fork, the UI subscribes .external(), and the UI does not redraw on its own writes.

A change made outside the process - the file edited by hand - has no id at all, so it is nobody’s own write and reaches external subscribers too.

Both give another handle onto the same value. They differ in one thing: whose writes they count as.

let port = state.port();
let same = port.clone();
let other = port.fork();
let _sub = port
.subscription_with()
.external()
.register(move |value| seen.lock().unwrap().push(*value));
same.set(1111)?;
other.set(2222)?;
assert_eq!(*heard.lock().unwrap(), [2222]);
assert_eq!(port.instance_id(), same.instance_id());
assert_ne!(port.instance_id(), other.instance_id());

clone keeps the id, so the original and the clone are one actor and neither hears the other’s writes through external. fork takes a new id, so the two are separate actors and each hears the other.

instance_id() is that id, and it is on a Field and on a ReactiveMap alike. Ask a handle for it when you need to hold the answer rather than filter on it - to name the actor in a log, or to compare it against the id a change arrives with, which is the next section.

let _sub = state
.port()
.subscription_with()
.register_with_source(move |port, who| {
seen.lock().unwrap().push((*port, who));
});
state.port().set(9090)?;
let (port, who) = heard.lock().unwrap()[0];
assert_eq!(port, 9090);
assert_eq!(who, Some(state.port().instance_id()));

register_with_source hands the callback the id beside the value, for deciding per change rather than filtering wholesale.

.external() filters Update and nothing else - Insert, Remove and Clear reach everyone including whoever caused them. The reasoning, and what that implies for insert, is on ReactiveMap.