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}