Skip to main content

weida_runtime/
registry.rs

1//! A namespace of bound names with a byte budget.
2
3use std::collections::HashMap;
4use std::sync::Mutex;
5
6use tokio::sync::{Notify, mpsc};
7use weida_core::Error;
8
9/// A namespace of names that can be bound, dialled and unbound, generic over
10/// what a bound name hands its acceptor.
11///
12/// **What it is for.** Every protocol with an in-process transport needs the
13/// same object: a set of names, one owner per name, a queue from the dialler
14/// to the owner, and a ceiling on how long a name may be. weida's
15/// `weida+inproc://<bus>/<path>` buses and ZeroMQ's `inproc://` endpoints are
16/// that object twice, down to the 256-byte budget, which libzmq set first
17/// ([decisions/0010](../../../docs/decisions/0010-local-transport.md) §4.8,
18/// `docs/research/zeromq.md` §11).
19///
20/// `T` is whatever binding a name entitles its owner to receive — one side of
21/// a connection, a pipe pair, a socket. The registry never looks at it.
22///
23/// **Scope is the holder's choice.** A `static` registry is a process-wide
24/// namespace, which is what a library with a process-global transport wants;
25/// one owned by a context object is a per-context namespace, which is what
26/// libzmq's `inproc://` actually is. Neither needs a different type.
27///
28/// **The byte budget is remote input's bound.** A name may come from a peer,
29/// an address or a configuration file, so its length is capped at
30/// construction time and every entry point checks it
31/// (`docs/INVARIANTS.md`). Control bytes are refused for the same reason a
32/// name is a name: it ends up in a log line, an error message and a
33/// comparison.
34pub struct NameRegistry<T> {
35    max_name_bytes: usize,
36    entries: Mutex<HashMap<String, mpsc::UnboundedSender<T>>>,
37    /// Woken on every bind, for a dialler waiting for a name to appear.
38    bound: Notify,
39}
40
41impl<T> NameRegistry<T> {
42    /// A registry whose names may be at most `max_name_bytes` bytes long.
43    pub fn new(max_name_bytes: usize) -> NameRegistry<T> {
44        NameRegistry {
45            max_name_bytes,
46            entries: Mutex::new(HashMap::new()),
47            bound: Notify::new(),
48        }
49    }
50
51    /// The byte budget names in this registry live under.
52    pub const fn max_name_bytes(&self) -> usize {
53        self.max_name_bytes
54    }
55
56    /// Checks `name` against the budget and against the control bytes.
57    ///
58    /// A consumer whose own grammar forbids more than this — a separator
59    /// byte, say, because the name is one field of an address — checks that
60    /// itself and calls this for the rest.
61    pub fn validate(&self, name: &str) -> Result<(), Error> {
62        if name.is_empty() || name.len() > self.max_name_bytes {
63            return Err(Error::InvalidAddress(format!(
64                "name must be 1..={} bytes: {name:?}",
65                self.max_name_bytes
66            )));
67        }
68        if name.bytes().any(|b| b < 0x20) {
69            return Err(Error::InvalidAddress(format!(
70                "invalid byte in name: {name:?}"
71            )));
72        }
73        Ok(())
74    }
75
76    /// Binds `name` and returns the queue of whatever is dialled to it.
77    ///
78    /// Fails with [`Error::InvalidAddress`] for a name that
79    /// [`NameRegistry::validate`] refuses and with
80    /// [`Error::AlreadyRegistered`] for a name somebody already holds: one
81    /// owner per name, and the second caller is told rather than silently
82    /// displacing the first.
83    pub fn bind(&self, name: &str) -> Result<mpsc::UnboundedReceiver<T>, Error> {
84        self.validate(name)?;
85        let mut entries = self.entries.lock().expect("name registry poisoned");
86        if entries.contains_key(name) {
87            return Err(Error::AlreadyRegistered);
88        }
89        let (tx, rx) = mpsc::unbounded_channel();
90        entries.insert(name.to_owned(), tx);
91        drop(entries);
92        self.bound.notify_waiters();
93        Ok(rx)
94    }
95
96    /// Removes `name`. The binding's own drop is the caller; unbinding a name
97    /// nobody holds is not an error, because a drop cannot fail.
98    pub fn unbind(&self, name: &str) {
99        self.entries
100            .lock()
101            .expect("name registry poisoned")
102            .remove(name);
103    }
104
105    /// The sender for `name`, or `None` when nothing is bound there.
106    ///
107    /// `None` is the in-process equivalent of a dial to a closed port, and
108    /// the caller reports it in its own vocabulary: this crate has no opinion
109    /// on what "nobody is listening" means to a protocol.
110    pub fn lookup(&self, name: &str) -> Option<mpsc::UnboundedSender<T>> {
111        self.entries
112            .lock()
113            .expect("name registry poisoned")
114            .get(name)
115            .cloned()
116    }
117
118    /// Resolves once `name` is bound, at once if it already is.
119    ///
120    /// The in-process counterpart of redialling a socket: a bus is back
121    /// exactly when its name is registered again, so a dialler waits on the
122    /// registry rather than on a clock. Registered before the check, so a
123    /// bind between the check and the wait is not missed.
124    pub async fn wait_bound(&self, name: &str) {
125        loop {
126            let bound = self.bound.notified();
127            tokio::pin!(bound);
128            bound.as_mut().enable();
129            if self.lookup(name).is_some() {
130                return;
131            }
132            bound.await;
133        }
134    }
135}
136
137impl<T> std::fmt::Debug for NameRegistry<T> {
138    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
139        f.debug_struct("NameRegistry")
140            .field("max_name_bytes", &self.max_name_bytes)
141            .finish_non_exhaustive()
142    }
143}
144
145#[cfg(test)]
146mod tests {
147    use super::*;
148
149    /// Claim: a name is bounded at both ends and carries no control bytes,
150    /// and the budget is the one the registry was built with.
151    #[test]
152    fn a_name_is_bounded_by_the_registry_budget() {
153        let registry: NameRegistry<u8> = NameRegistry::new(8);
154        assert_eq!(registry.max_name_bytes(), 8);
155        assert!(registry.validate("orders").is_ok());
156        assert!(registry.validate(&"x".repeat(8)).is_ok());
157        assert!(registry.validate(&"x".repeat(9)).is_err());
158        assert!(registry.validate("").is_err());
159        assert!(registry.validate("has\ncontrol").is_err());
160    }
161
162    /// Claim: one owner per name — the second bind is refused rather than
163    /// displacing the first — and unbinding frees the name again.
164    #[test]
165    fn one_owner_per_name() {
166        let registry: NameRegistry<u8> = NameRegistry::new(16);
167        let first = registry.bind("orders").expect("first bind");
168        assert!(matches!(
169            registry.bind("orders"),
170            Err(Error::AlreadyRegistered)
171        ));
172        drop(first);
173        // Dropping the receiver does not free the name: the binding's own
174        // drop is what unbinds, so a dropped acceptor cannot be replaced
175        // behind the binding's back.
176        assert!(matches!(
177            registry.bind("orders"),
178            Err(Error::AlreadyRegistered)
179        ));
180        registry.unbind("orders");
181        assert!(registry.bind("orders").is_ok());
182    }
183
184    /// Claim: what a dialler hands over reaches the name's owner, and a name
185    /// nobody holds resolves to nothing at all.
186    #[test]
187    fn a_dial_reaches_the_owner_and_an_unbound_name_reaches_nobody() {
188        let registry: NameRegistry<u8> = NameRegistry::new(16);
189        assert!(registry.lookup("orders").is_none());
190        let mut incoming = registry.bind("orders").expect("bind");
191        registry
192            .lookup("orders")
193            .expect("bound name has a sender")
194            .send(9)
195            .expect("the owner is still listening");
196        assert_eq!(incoming.try_recv().expect("delivered"), 9);
197
198        registry.unbind("orders");
199        assert!(registry.lookup("orders").is_none());
200    }
201
202    /// Claim: a name over the budget is refused by `bind` and `validate`
203    /// alike, so the check cannot be skipped by going through the registry.
204    #[test]
205    fn an_oversized_name_cannot_be_bound() {
206        let registry: NameRegistry<u8> = NameRegistry::new(4);
207        let err = registry.bind("toolong").unwrap_err();
208        assert!(matches!(err, Error::InvalidAddress(_)), "{err:?}");
209    }
210}