Skip to content

Observer in Rust: Events and Shutdown

An Observer design needs more than a callback list: define who owns subscribers, when removal takes effect and what shutdown means. Rust closures and traits can represent notification behaviour, but they do not choose those policies for you. This lesson tests a synchronous callback design against a matching C++ implementation.

Requirements: synchronous Observer delivery with explicit lifetime

Publish signed integer readings, including zero and negatives, to callbacks in registration order. Each callback registered at the beginning of a successful publish runs once for that event. A callback may return Remove to unsubscribe itself after receiving that event; later subscribers still receive it. Explicit unsubscribe before publishing prevents delivery. Removing an already removed ID returns false.

The bus owns its callback objects. IDs are scoped to the bus that issued them, are never reused there, and must not be passed to another bus. Exhaustion rejects registration rather than wrapping an ID. These local IDs are not globally authenticated handles; cross-container identity is a separate design requirement in the next lesson.

Close is explicit and idempotent: first close releases all stored callbacks, later close reports false, and publishing or registering after close returns Closed. There is no background work, queue or final shutdown event. Callbacks must return normally and must not re-enter or mutate the bus. Exception/panic recovery, asynchronous delivery and concurrent publishers are outside this fixture's contract.

C++ Observer using owned callables

Save as observer.cpp. The shared log lets the program inspect delivery after a callback is removed; it is not the subscription itself.

#include <algorithm>
#include <cassert>
#include <cstddef>
#include <cstdint>
#include <functional>
#include <iostream>
#include <limits>
#include <memory>
#include <string>
#include <utility>
#include <variant>
#include <vector>

enum class Action { Keep, Remove };
enum class Error { Closed, Exhausted, Empty };
class Bus {
    struct Entry {
        std::uint64_t id;
        std::function<Action(int)> callback;
        bool remove = false;
    };
    std::vector<Entry> entries_;
    std::uint64_t next_ = 1;
    bool closed_ = false;
public:
    Bus() = default;
    Bus(const Bus&) = delete;
    Bus& operator=(const Bus&) = delete;
    Bus(Bus&&) = delete;
    Bus& operator=(Bus&&) = delete;
    std::variant<std::uint64_t, Error> subscribe(std::function<Action(int)> callback) {
        if (closed_) return Error::Closed;
        if (!callback) return Error::Empty;
        if (next_ == std::numeric_limits<std::uint64_t>::max()) return Error::Exhausted;
        auto id = next_;
        entries_.push_back(Entry{id, std::move(callback)});
        ++next_;
        return id;
    }
    bool unsubscribe(std::uint64_t id) {
        auto before = entries_.size();
        std::erase_if(entries_, [id](const Entry& entry) { return entry.id == id; });
        return before != entries_.size();
    }
    std::variant<std::size_t, Error> publish(int reading) {
        if (closed_) return Error::Closed;
        std::size_t delivered = 0;
        for (auto& entry : entries_) {
            ++delivered;
            entry.remove = entry.callback(reading) == Action::Remove;
        }
        std::erase_if(entries_, [](const Entry& entry) { return entry.remove; });
        return delivered;
    }
    bool close() {
        if (closed_) return false;
        closed_ = true;
        entries_.clear();
        return true;
    }
};

int main() {
    auto log = std::make_shared<std::vector<std::string>>();
    Bus bus;
    auto a = std::get<std::uint64_t>(bus.subscribe([log](int value) {
        log->push_back("A:" + std::to_string(value));
        return Action::Keep;
    }));
    auto b = std::get<std::uint64_t>(bus.subscribe([log](int value) {
        log->push_back("B:" + std::to_string(value));
        return Action::Remove;
    }));
    auto c = std::get<std::uint64_t>(bus.subscribe([log](int value) {
        log->push_back("C:" + std::to_string(value));
        return Action::Keep;
    }));
    auto first = std::get<std::size_t>(bus.publish(0));
    assert(first == 3 && !bus.unsubscribe(b));
    assert(bus.unsubscribe(a) && !bus.unsubscribe(a));
    auto second = std::get<std::size_t>(bus.publish(-500));
    assert(second == 1);
    assert((*log == std::vector<std::string>{"A:0", "B:0", "C:0", "C:-500"}));
    std::cout << "first=" << first << ", second=" << second << '\n';
    std::string joined;
    for (const auto& entry : *log) {
        if (!joined.empty()) joined += '|';
        joined += entry;
    }
    std::cout << "log=" << joined << '\n';
    bool closed = bus.close();
    bool repeated = bus.close();
    assert(closed && !repeated && !bus.unsubscribe(c));
    auto rejected = bus.publish(46700);
    assert(std::get<Error>(rejected) == Error::Closed);
    assert(std::get<Error>(bus.subscribe([](int) { return Action::Keep; })) == Error::Closed);
    assert(log->size() == 4);
    std::cout << "closed=" << closed << ", repeated=" << repeated << '\n';
    std::cout << "after close=Closed\n";
    Bus empty;
    assert(std::get<Error>(empty.subscribe({})) == Error::Empty);
    assert(std::get<std::size_t>(empty.publish(0)) == 0);
    auto held = std::make_shared<int>(7);
    std::weak_ptr<int> observed = held;
    Bus owner;
    auto id = std::get<std::uint64_t>(owner.subscribe([held](int) { return Action::Keep; }));
    held.reset();
    assert(!observed.expired());
    assert(owner.unsubscribe(id) && observed.expired());
    auto closing = std::make_shared<int>(0);
    std::weak_ptr<int> closing_view = closing;
    auto ignored = owner.subscribe([closing](int) { return Action::Keep; });
    assert(std::holds_alternative<std::uint64_t>(ignored));
    closing.reset();
    assert(!closing_view.expired() && owner.close() && closing_view.expired());
    auto dropping = std::make_shared<int>(0);
    std::weak_ptr<int> dropping_view = dropping;
    int calls = 0;
    {
        Bus scoped;
        assert(std::holds_alternative<std::uint64_t>(scoped.subscribe([dropping, &calls](int) {
            ++calls;
            return Action::Keep;
        }))); // The returned ID is not retained.
        dropping.reset();
        assert(!dropping_view.expired());
        assert(std::get<std::size_t>(scoped.publish(0)) == 1);
    }
    assert(dropping_view.expired() && calls == 1); // No final event on destruction.
}
1
2
3
4
g++ -std=c++20 -Wall -Wextra -Wpedantic -O0 observer.cpp -o observer
./observer
g++ -std=c++20 -Wall -Wextra -Wpedantic -O2 observer.cpp -o observer
./observer

Both programs produce:

1
2
3
4
first=3, second=1
log=A:0|B:0|C:0|C:-500
closed=1, repeated=0
after close=Closed

std::function owns a callable wrapper, not every object a callable refers to. Capturing a reference still requires the referenced object to outlive its use. This fixture captures shared ownership of the log rather than a dangling reference. A virtual Observer interface is another option for named subscriber behaviour; it would still need the same lifetime and removal policies.

An empty std::function is representable, so the C++ boundary rejects it with Empty. Rust's generic FnMut parameter requires an actual callable; it has no corresponding empty-wrapper value in this API. This is a difference in representable input, not a different delivery policy.

Removal is recorded during the C++ traversal and applied afterward, so erasing an entry does not invalidate the active traversal. Re-entrant registration or unsubscribe is prohibited, including through a captured pointer to the bus. The type system does not enforce that C++ restriction. The assertions do not invoke prohibited or undefined behaviour.

Rust Observer using boxed FnMut callbacks

cargo new --edition 2024 observer_demo
cd observer_demo

Replace src/main.rs:

use std::cell::RefCell;
use std::rc::Rc;

#[derive(Debug, PartialEq, Eq)]
enum Action {
    Keep,
    Remove,
}
#[derive(Debug, PartialEq, Eq)]
enum Error {
    Closed,
    Exhausted,
}

struct Entry {
    id: u64,
    callback: Box<dyn FnMut(i32) -> Action>,
}
struct Bus {
    entries: Vec<Entry>,
    next: u64,
    closed: bool,
}
impl Bus {
    fn new() -> Self {
        Self {
            entries: Vec::new(),
            next: 1,
            closed: false,
        }
    }
    fn subscribe<F>(&mut self, callback: F) -> Result<u64, Error>
    where
        F: FnMut(i32) -> Action + 'static,
    {
        if self.closed {
            return Err(Error::Closed);
        }
        let next = self.next.checked_add(1).ok_or(Error::Exhausted)?;
        let id = self.next;
        self.entries.push(Entry {
            id,
            callback: Box::new(callback),
        });
        self.next = next;
        Ok(id)
    }
    fn unsubscribe(&mut self, id: u64) -> bool {
        let before = self.entries.len();
        self.entries.retain(|entry| entry.id != id);
        before != self.entries.len()
    }
    fn publish(&mut self, reading: i32) -> Result<usize, Error> {
        if self.closed {
            return Err(Error::Closed);
        }
        let mut delivered = 0;
        self.entries.retain_mut(|entry| {
            delivered += 1;
            (entry.callback)(reading) == Action::Keep
        });
        Ok(delivered)
    }
    fn close(&mut self) -> bool {
        if self.closed {
            return false;
        }
        self.closed = true;
        self.entries.clear();
        true
    }
}

fn main() {
    let log = Rc::new(RefCell::new(Vec::<String>::new()));
    let mut bus = Bus::new();
    let a_log = Rc::clone(&log);
    let a = bus
        .subscribe(move |value| {
            a_log.borrow_mut().push(format!("A:{value}"));
            Action::Keep
        })
        .unwrap();
    let b_log = Rc::clone(&log);
    let b = bus
        .subscribe(move |value| {
            b_log.borrow_mut().push(format!("B:{value}"));
            Action::Remove
        })
        .unwrap();
    let c_log = Rc::clone(&log);
    let c = bus
        .subscribe(move |value| {
            c_log.borrow_mut().push(format!("C:{value}"));
            Action::Keep
        })
        .unwrap();
    let first = bus.publish(0).unwrap();
    assert_eq!(first, 3);
    assert!(!bus.unsubscribe(b));
    assert!(bus.unsubscribe(a));
    assert!(!bus.unsubscribe(a));
    let second = bus.publish(-500).unwrap();
    assert_eq!(second, 1);
    println!("first={first}, second={second}");
    println!("log={}", log.borrow().join("|"));
    let closed = bus.close();
    let repeated = bus.close();
    assert!(closed && !repeated && !bus.unsubscribe(c));
    assert_eq!(bus.publish(46700), Err(Error::Closed));
    assert_eq!(bus.subscribe(|_| Action::Keep), Err(Error::Closed));
    assert_eq!(log.borrow().len(), 4);
    println!(
        "closed={}, repeated={}",
        u8::from(closed),
        u8::from(repeated)
    );
    println!("after close=Closed");
}

#[cfg(test)]
mod tests {
    use super::*;
    use std::cell::Cell;

    #[test]
    fn self_removal_preserves_order_and_remaining_delivery() {
        let log = Rc::new(RefCell::new(Vec::new()));
        let mut bus = Bus::new();
        for (label, action) in [
            ("A", Action::Keep),
            ("B", Action::Remove),
            ("C", Action::Keep),
        ] {
            let sink = Rc::clone(&log);
            bus.subscribe(move |value| {
                sink.borrow_mut().push(format!("{label}:{value}"));
                if action == Action::Remove {
                    Action::Remove
                } else {
                    Action::Keep
                }
            })
            .unwrap();
        }
        assert_eq!(bus.publish(0), Ok(3));
        assert_eq!(bus.publish(-500), Ok(2));
        assert_eq!(*log.borrow(), ["A:0", "B:0", "C:0", "A:-500", "C:-500"]);
    }

    #[test]
    fn explicit_unsubscribe_prevents_calls_and_does_not_reuse_ids() {
        let mut bus = Bus::new();
        let removed = bus.subscribe(|_| panic!("removed callback ran")).unwrap();
        assert!(bus.unsubscribe(removed));
        assert!(!bus.unsubscribe(removed));
        let replacement = bus.subscribe(|_| Action::Keep).unwrap();
        assert_ne!(removed, replacement);
        assert!(!bus.unsubscribe(removed));
        assert_eq!(bus.publish(0), Ok(1));
    }

    #[test]
    fn shutdown_rejects_new_work_without_delivering_a_final_event() {
        let calls = Rc::new(Cell::new(0));
        let seen = Rc::clone(&calls);
        let mut bus = Bus::new();
        bus.subscribe(move |_| {
            seen.set(seen.get() + 1);
            Action::Keep
        })
        .unwrap();
        assert!(bus.close());
        assert!(!bus.close());
        assert_eq!(bus.publish(0), Err(Error::Closed));
        assert_eq!(bus.subscribe(|_| Action::Keep), Err(Error::Closed));
        assert_eq!(calls.get(), 0);
    }

    struct DropProbe(Rc<Cell<usize>>);
    impl Drop for DropProbe {
        fn drop(&mut self) {
            self.0.set(self.0.get() + 1);
        }
    }

    #[test]
    fn removal_and_close_release_owned_callback_state() {
        let drops = Rc::new(Cell::new(0));
        let mut bus = Bus::new();
        let first = DropProbe(Rc::clone(&drops));
        let id = bus
            .subscribe(move |_| {
                let _alive = &first;
                Action::Keep
            })
            .unwrap();
        assert_eq!(drops.get(), 0);
        assert!(bus.unsubscribe(id));
        assert_eq!(drops.get(), 1);
        let second = DropProbe(Rc::clone(&drops));
        bus.subscribe(move |_| {
            let _alive = &second;
            Action::Keep
        })
        .unwrap();
        assert!(bus.close());
        assert_eq!(drops.get(), 2);
    }

    #[test]
    fn empty_delivery_and_id_exhaustion_are_explicit() {
        let mut bus = Bus::new();
        assert_eq!(bus.publish(0), Ok(0));
        bus.next = u64::MAX - 1;
        assert_eq!(bus.subscribe(|_| Action::Keep), Ok(u64::MAX - 1));
        assert_eq!(bus.subscribe(|_| Action::Keep), Err(Error::Exhausted));
        assert_eq!(bus.publish(i32::MIN), Ok(1));
    }

    #[test]
    fn discarding_id_keeps_subscription_until_bus_drop() {
        let drops = Rc::new(Cell::new(0));
        let calls = Rc::new(Cell::new(0));
        {
            let mut bus = Bus::new();
            let owned = DropProbe(Rc::clone(&drops));
            let seen = Rc::clone(&calls);
            bus.subscribe(move |_| {
                let _alive = &owned;
                seen.set(seen.get() + 1);
                Action::Keep
            })
            .unwrap(); // Do not retain the returned ID.
            assert_eq!(bus.publish(0), Ok(1));
            assert_eq!(drops.get(), 0);
        }
        assert_eq!(drops.get(), 1);
        assert_eq!(calls.get(), 1); // Destruction did not send another event.
    }
}
1
2
3
4
5
6
cargo fmt --check
cargo check --offline
cargo test --offline
cargo test --offline --release
cargo run --offline --quiet
cargo run --offline --release --quiet

FnMut permits stateful, repeated calls. Box owns each closure; its captured data follows that closure's ownership. The 'static bound excludes captures borrowing short-lived local data, but does not mean the closure lives forever. Removed closures are released normally. Capturing an Rc clone keeps the log alive while preserving an external view.

Rc is single-threaded shared ownership, not thread-safe event delivery. The log's RefCell checks borrows at runtime; holding a conflicting log borrow while publishing can panic. The example takes only short borrows and releases them before the next operation.

retain_mut visits entries once in their original order and retains selected entries. Removing the current subscription therefore does not skip the next one. Unlike C++'s deferred erase, a removed Rust callback can release its captured state during this traversal; this example promises delivery order and removal before the next publish, not identical destructor timing between languages. Captured destructors must not perform bus re-entry.

Dropping a numeric subscription ID does not unsubscribe. The bus still owns the callback. Conversely, dropping the bus normally releases callbacks, but does not run close as an application protocol, send a final event, flush work or guarantee execution during process termination. See ownership and RAII for those boundaries.

Changed requirement: queues, borrowed observers and failure policy

A borrowed observer list can avoid owning subscriber objects, but then lifetimes restrict how long the bus may retain them. Rc/Weak or shared_ptr/weak_ptr can instead model externally owned observers; define what happens when the observer expires. Strong references in both directions can create a cycle, so shared ownership alone is not an unsubscribe strategy.

If notifications move to a queue or another thread, event ownership, delivery acknowledgement, bounded capacity, slow subscribers and termination become new requirements. Rust's standard mpsc channels are multiple-producer, single-consumer channels, not automatic broadcast to all observers. Fan-out requires explicit per-subscriber routing or another abstraction. Successfully enqueueing a value is not proof that a subscriber processed it.

Define whether shutdown drains or discards queued values, how producers stop, what happens while sender clones remain, and whether a worker must be joined. Adding a channel does not implement those choices. This synchronous fixture has nothing to drain and does not claim asynchronous or exactly-once durable delivery.

If callbacks may fail, decide whether later subscribers still run and how failures are returned. The current Keep/Remove result is a subscription action, not an error or acknowledgement. A panic or exception halfway through this example is outside its successful-publish guarantee; do not promise transactional rollback or resume after failure without designing and testing it.

Intentional failures: capture contracts differ

This independent Rust example borrows a local String in a closure required to be 'static:

1
2
3
4
5
6
7
fn register<F: FnMut(i32) + 'static>(mut callback: F) {
    callback(0);
}
fn main() {
    let label = String::from("Pi 4B");
    register(|_| assert_eq!(label, "Pi 4B"));
}

The ownership repair is a move closure: the String then belongs to the callable. A borrowed registration API would be a different lifetime contract, not permission to outlive the borrowed label.

C++20 std::function instead requires a copy-constructible target. This independent move-only capture cannot be stored in it:

1
2
3
4
5
6
7
#include <functional>
#include <memory>
#include <utility>
int main() {
    auto owner = std::make_unique<int>(0);
    std::function<void()> callback = [owner = std::move(owner)] { (void)*owner; };
}

An intentionally shared capture is one repair if shared ownership fits the requirements. A move-only callable wrapper is another design, but is not this C++20 std::function example. Rust's boxed FnMut does not require Clone, and the DropProbe test stores non-Clone owned state. Do not add sharing solely to silence a compiler error without deciding who should own the subscriber state.

Decision criteria and exercises

Requirement Candidate Policy still needed
Synchronous local delivery Owned callbacks or named observer trait/interface Order, re-entry, failure and removal timing
Externally owned subscribers Borrowed references or weak ownership Expiry and lifetime rules
Subscription scoped to an owner Explicit unsubscribe or designed subscription guard Whether dropping a token removes it, and how it reaches the bus safely
Deferred or cross-thread processing Owned event queues and appropriate thread-safe bounds Backpressure, fan-out, acknowledgements and shutdown
  1. Add a callback with an internal call count that removes itself on the second event. Verify order and absence on the third event.
  2. Change the Rust removal predicate to always keep entries. Retain the self-removal test: this compiles but must fail the delivery contract.
  3. Repair the Rust capture with move and the C++ target with an explicitly shared owner. Confirm the captured value remains valid when called.
  4. Compare explicit close with ordinary destruction. Which application actions occur in neither example, and which resources are simply released?
  5. Before replacing the list with a channel, write the drain/discard, slow-subscriber and acknowledgement policies. Explain why one mpsc receiver is not a broadcast list.

Verification and next step

The exact examples were verified on Raspberry Pi 4B with rustc/cargo 1.99.0, GCC 14.2.0 and kernel 6.18.50+rpt-rpi-v8. Rust passed six tests in debug and release, exact-source formatting and both output checks. The second-event removal variant passed seven tests in both profiles, and the owned-capture repair passed. The borrowed 'static capture failed with E0373; deliberately retaining removed subscribers compiled but failed the delivery test. C++20 passed output/assertions and second-event/shared-capture repairs at -O0 and -O2, with assertions enabled; std::function rejected the move-only target. Normal destruction, explicit close, ownership release and discarded IDs were checked without invoking prohibited re-entry. These tests are not timing measurements or a proof for all callback failures and thread schedules.

Next: ownership, IDs and ECS design, identity and processing boundaries based on requirements rather than benchmarks.

Previous: Visitor and extension boundaries · Course overview

Donate