From: Wolfgang Bumiller <w.bumiller@proxmox.com>
To: Hannes Laimer <h.laimer@proxmox.com>
Cc: pve-devel@lists.proxmox.com
Subject: Re: [PATCH pve-cluster 02/10] rust: notify: add change notification socket server
Date: Fri, 25 Sep 2026 17:18:48 +0200 [thread overview]
Message-ID: <6k5xskft7jh6a4vz7eb2673pmlq2t6dbgxog27kzx47y5jqdmc@6a4oe6hae4p7> (raw)
In-Reply-To: <20260918144152.575163-3-h.laimer@proxmox.com>
Review not done yet - but need to go for now.
On Fri, Sep 18, 2026 at 04:41:44PM +0200, Hannes Laimer wrote:
> Daemons currently learn about /etc/pve changes by polling the version
> vector over IPC on every loop iteration, and inotify only sees writes
> made on the local node. Add a small server that clients reach through a
> Unix stream socket. A client subscribes with named path templates such
> as nodes/{node}/qemu-server/{vmid}.conf. Each matching mutation then
> arrives as one pushed line, carrying the values the placeholders
> captured. Requests, replies and events are all JSON objects, one per
> line, so a client needs a single parser and the stream stays readable
> as text.
(We may want to consider `json-seq` format (record separator + data +
newline))
>
> Matching has to stay bounded, since it runs for every mutation and
> every connection. A template allows one placeholder per path component
> and a spanning one only as the last component, which makes a match a
> single pass over the path with no regex engine behind it, and the
> captures come out of that same pass for free.
>
> The daemon must never block on a client. The last few thousand
> mutations live in a fixed size ring and each connection keeps a cursor
> into it. The mutating thread only appends and wakes the delivery
> thread, which does every send, matches on the way out, and finishes a
> partially written line before the next one. A connection whose socket
> is full is caught up once it drains, and a cursor that fell off the
> ring turns into a single resync event. A reconnecting client resumes
> from the last number it saw. The hello reply hands out that number too,
> so a client that has seen no event yet can resume as well. The sequence
> number is the memdb version, so it means the same on every node. Should
> the delivery thread fail, it closes every connection on its way out, so
> a client always learns that it has to reconnect instead of staying
> attached to a socket that went silent.
>
> Connections authorized by group membership never see events about
> private paths, the same rule the IPC and FUSE side apply to content.
>
> Signed-off-by: Hannes Laimer <h.laimer@proxmox.com>
> ---
> src/rust/pmxcfs-notify/Cargo.toml | 5 +
> src/rust/pmxcfs-notify/src/conn.rs | 317 ++++++++++
> src/rust/pmxcfs-notify/src/lib.rs | 42 ++
> src/rust/pmxcfs-notify/src/protocol.rs | 237 ++++++++
> src/rust/pmxcfs-notify/src/registry.rs | 405 +++++++++++++
> src/rust/pmxcfs-notify/src/ring.rs | 171 ++++++
> src/rust/pmxcfs-notify/src/server.rs | 780 +++++++++++++++++++++++++
> src/rust/pmxcfs-notify/src/template.rs | 254 ++++++++
> 8 files changed, 2211 insertions(+)
> create mode 100644 src/rust/pmxcfs-notify/src/conn.rs
> create mode 100644 src/rust/pmxcfs-notify/src/protocol.rs
> create mode 100644 src/rust/pmxcfs-notify/src/registry.rs
> create mode 100644 src/rust/pmxcfs-notify/src/ring.rs
> create mode 100644 src/rust/pmxcfs-notify/src/server.rs
> create mode 100644 src/rust/pmxcfs-notify/src/template.rs
>
> diff --git a/src/rust/pmxcfs-notify/Cargo.toml b/src/rust/pmxcfs-notify/Cargo.toml
> index 3de30b5..5a67507 100644
> --- a/src/rust/pmxcfs-notify/Cargo.toml
> +++ b/src/rust/pmxcfs-notify/Cargo.toml
> @@ -9,3 +9,8 @@ homepage.workspace = true
> rust-version.workspace = true
>
> [dependencies]
> +anyhow.workspace = true
> +log.workspace = true
> +nix.workspace = true
> +serde.workspace = true
> +serde_json.workspace = true
> diff --git a/src/rust/pmxcfs-notify/src/conn.rs b/src/rust/pmxcfs-notify/src/conn.rs
> new file mode 100644
> index 0000000..e5df526
> --- /dev/null
> +++ b/src/rust/pmxcfs-notify/src/conn.rs
> @@ -0,0 +1,317 @@
> +use std::io::{self, Read};
> +use std::os::fd::{AsRawFd, RawFd};
> +use std::os::unix::net::UnixStream;
> +
> +use nix::errno::Errno;
> +use nix::sys::socket::{MsgFlags, Shutdown, send, shutdown};
> +
> +use crate::protocol::Params;
> +use crate::registry::is_private;
> +use crate::template::Template;
> +
> +const MAX_LINE: usize = 64 * 1024;
> +const MAX_OUTBUF: usize = 256 * 1024;
> +
> +pub(crate) struct Conn {
> + stream: UnixStream,
> + privileged: bool,
> + pub(crate) templates: Vec<(String, Template)>,
> + pub(crate) subscribed: bool,
> + pub(crate) dead: bool,
> + pub(crate) cursor: u64,
> + pub(crate) resync_pending: bool,
> + outbuf: Vec<u8>,
> + inbuf: Vec<u8>,
> +}
> +
> +#[derive(Debug)]
> +pub(crate) enum ReadError {
> + Eof,
> + LineTooLong,
> + Io(io::Error),
> +}
> +
> +impl Conn {
> + pub(crate) fn new(stream: UnixStream, privileged: bool) -> io::Result<Self> {
> + stream.set_nonblocking(true)?;
> + Ok(Self {
> + stream,
> + privileged,
> + templates: Vec::new(),
> + subscribed: false,
> + dead: false,
> + cursor: 0,
> + resync_pending: false,
> + outbuf: Vec::new(),
> + inbuf: Vec::new(),
> + })
> + }
> +
> + pub(crate) fn fd(&self) -> RawFd {
> + self.stream.as_raw_fd()
> + }
> +
> + pub(crate) fn lagging(&self) -> bool {
nit: `is_lagging`? Not sure about the term itself even. `.is_blocked()`?
> + !self.outbuf.is_empty()
> + }
> +
> + pub(crate) fn matches<'a>(
> + &'a self,
> + path: Option<&'a str>,
> + to: Option<&'a str>,
> + ) -> Option<Params<'a>> {
> + if !self.subscribed {
> + return None;
> + }
> + if !self.privileged && (path.is_some_and(is_private) || to.is_some_and(is_private)) {
> + return None;
> + }
> + let mut params = Params::new();
> + for (name, template) in &self.templates {
> + let mut sets = Vec::new();
> + for candidate in [path, to].into_iter().flatten() {
↑ This is a bit awkward to read.
Maybe just skip the combinator and use just `[path, to]`
and use
if let Some(captures) = candidate.and_then(|c| template.matches(c))
but meh, both works.
My main issue here is that it wasn't immediately obvious to obvious to
see, given the way it is written, is that if they're both `None` we do a
lot of nothing - since this happens inside another loop - which should
be caught at the top:
if path.is_none() && to.is_none() {
return;
}
> + if let Some(captures) = template.matches(candidate)
> + && !sets.contains(&captures)
↑ For the `sets.contains()` it might be worth a comment that sets has at
most 1 element at this point. Otherwise one might be tempted to turn it
into a `BTreeSet` (which would still compile when also changing the
`Params` value type to a set.
> + {
> + sets.push(captures);
> + }
> + }
> + if !sets.is_empty() {
> + params.insert(name.as_str(), sets);
> + }
> + }
> + (!params.is_empty()).then_some(params)
> + }
> +
> + // Bytes the kernel did not take stay here and go out first once the
> + // socket has room again, so a line is never torn. A client that keeps
> + // sending commands without reading the replies is cut at the cap.
> + pub(crate) fn queue(&mut self, line: &str) {
> + if self.dead {
> + return;
> + }
> + self.outbuf.extend_from_slice(line.as_bytes());
> + if self.outbuf.len() > MAX_OUTBUF {
^ So MAX_OUTBUF is not a "hard cap" - the buffer can still grow beyond,
but that's probably fine given what we are producing as output.
> + self.mark_dead();
> + }
> + }
> +
> + /// Returns whether everything queued went out.
> + pub(crate) fn flush(&mut self) -> bool {
> + while !self.outbuf.is_empty() && !self.dead {
> + let flags = MsgFlags::MSG_NOSIGNAL | MsgFlags::MSG_DONTWAIT;
> + match send(self.fd(), &self.outbuf, flags) {
> + Ok(0) | Err(Errno::EAGAIN) => return false,
> + Ok(n) => {
> + self.outbuf.drain(..n);
> + }
> + Err(Errno::EINTR) => continue,
> + Err(_) => self.mark_dead(),
> + }
> + }
> + !self.dead
> + }
> +
> + pub(crate) fn mark_dead(&mut self) {
> + self.dead = true;
> + self.outbuf.clear();
> + let _ = shutdown(self.fd(), Shutdown::Both);
> + }
> +
> + pub(crate) fn read_lines(&mut self) -> Result<Vec<String>, ReadError> {
^ This should mention a bit about how it works.
This performs a single read, looping on EINTR, ending on EAGAIN.
If the first read does not provide any lines,
> + let mut buf = [0u8; 4096];
> + loop {
> + match self.stream.read(&mut buf) {
> + Ok(0) => return Err(ReadError::Eof),
> + Ok(n) => {
> + self.inbuf.extend_from_slice(&buf[..n]);
> + break;
> + }
> + Err(err) if err.kind() == io::ErrorKind::WouldBlock => return Ok(Vec::new()),
But if we do a non-blocking read and already collect multiple lines, we
might as well just `break` here, and not `break` in the `Ok(n)` case.
Thena gain, if we expect "problematic" clients, doing a single pass
might be the better choice. 🤷
> + Err(err) if err.kind() == io::ErrorKind::Interrupted => continue,
> + Err(err) => return Err(ReadError::Io(err)),
> + }
> + }
> + let mut lines = Vec::new();
> + while let Some(pos) = self.inbuf.iter().position(|&b| b == b'\n') {
> + let line = String::from_utf8_lossy(&self.inbuf[..pos]).into_owned();
Let's not go with lossy here. If a client sends garbage, *disconnect
them*!
> + self.inbuf.drain(..=pos);
> + lines.push(line);
> + }
> + if self.inbuf.len() > MAX_LINE {
> + return Err(ReadError::LineTooLong);
> + }
> + Ok(lines)
> + }
> +}
> +
> +#[cfg(test)]
> +mod tests {
> + use std::collections::BTreeMap;
> + use std::io::Write;
> +
> + use super::*;
> +
> + fn pair(privileged: bool) -> (Conn, UnixStream) {
> + let (server, client) = UnixStream::pair().unwrap();
> + (Conn::new(server, privileged).unwrap(), client)
> + }
> +
> + fn templates(specs: &[(&str, &str)]) -> Vec<(String, Template)> {
> + specs
> + .iter()
> + .map(|(name, text)| ((*name).to_owned(), Template::parse(text).unwrap()))
> + .collect()
> + }
> +
> + #[test]
> + fn reassembles_partial_lines() {
> + let (mut conn, mut client) = pair(true);
> + client.write_all(b"{\"command\":").unwrap();
> + assert!(conn.read_lines().unwrap().is_empty());
> + client
> + .write_all(b"\"hello\"}\n{\"command\":\"unsubscribe\"}\n")
> + .unwrap();
> + let lines = conn.read_lines().unwrap();
> + assert_eq!(
> + lines,
> + vec!["{\"command\":\"hello\"}", "{\"command\":\"unsubscribe\"}"]
> + );
> + assert!(conn.read_lines().unwrap().is_empty());
> + }
> +
> + #[test]
> + fn reports_eof_and_overlong_lines() {
> + let (mut conn, mut client) = pair(true);
> + client.write_all(&vec![b'x'; MAX_LINE + 1]).unwrap();
> + let mut result = conn.read_lines();
> + while let Ok(lines) = &result {
> + assert!(lines.is_empty());
> + result = conn.read_lines();
> + }
> + assert!(matches!(result, Err(ReadError::LineTooLong)));
> +
> + let (mut conn, client) = pair(true);
> + drop(client);
> + assert!(matches!(conn.read_lines(), Err(ReadError::Eof)));
> + }
> +
> + #[test]
> + fn short_writes_are_finished_before_the_next_line() {
> + let (mut conn, mut client) = pair(true);
> + client.set_nonblocking(true).unwrap();
> + let line = format!("{}\n", "x".repeat(3000));
> +
> + let mut queued = 0;
> + while !conn.lagging() {
> + conn.queue(&line);
> + conn.flush();
> + queued += 1;
> + assert!(queued < 100_000, "socket never filled");
> + }
> +
> + let mut received = Vec::new();
> + let mut buf = [0u8; 65536];
> + let mut read_available = |client: &mut UnixStream, received: &mut Vec<u8>| {
> + loop {
> + match client.read(&mut buf) {
> + Ok(0) => break,
> + Ok(n) => received.extend_from_slice(&buf[..n]),
> + Err(err) if err.kind() == io::ErrorKind::WouldBlock => break,
> + Err(err) => panic!("read failed: {err}"),
> + }
> + }
> + };
> + loop {
> + read_available(&mut client, &mut received);
> + if !conn.lagging() {
> + break;
> + }
> + conn.flush();
> + }
> + read_available(&mut client, &mut received);
> +
> + assert_eq!(received.len(), queued * line.len());
> + assert!(
> + received
> + .chunks(line.len())
> + .all(|chunk| chunk == line.as_bytes())
> + );
> + assert!(!conn.dead);
> + }
> +
> + #[test]
> + fn matching_needs_a_subscription() {
> + let (mut conn, _client) = pair(true);
> + conn.templates = templates(&[
> + ("guest", "nodes/{node}/qemu-server/{vmid}.conf"),
> + ("dc", "datacenter.cfg"),
> + ]);
> + let guest = "nodes/n1/qemu-server/100.conf";
> + assert!(conn.matches(Some(guest), None).is_none());
> +
> + conn.subscribed = true;
> + let params = conn.matches(Some(guest), None).unwrap();
> + assert_eq!(params.len(), 1);
> + assert_eq!(
> + params["guest"],
> + vec![BTreeMap::from([("node", "n1"), ("vmid", "100")])]
> + );
> + assert_eq!(
> + conn.matches(Some("datacenter.cfg"), None).unwrap()["dc"],
> + vec![BTreeMap::new()]
> + );
> + assert!(conn.matches(Some("datacenter.cfg.tmp.1"), None).is_none());
> + assert!(conn.matches(Some("other"), None).is_none());
> + assert!(conn.matches(None, None).is_none());
> + }
> +
> + #[test]
> + fn rename_matches_on_either_side() {
> + let (mut conn, _client) = pair(true);
> + conn.templates = templates(&[("guest", "nodes/{node}/qemu-server/{vmid}.conf")]);
> + conn.subscribed = true;
> + let n1 = "nodes/n1/qemu-server/100.conf";
> + let n2 = "nodes/n2/qemu-server/100.conf";
> +
> + let params = conn.matches(Some(n1), Some(n2)).unwrap();
> + assert_eq!(params["guest"].len(), 2);
> + assert_eq!(params["guest"][0]["node"], "n1");
> + assert_eq!(params["guest"][1]["node"], "n2");
> +
> + let params = conn
> + .matches(Some("nodes/n2/qemu-server/100.conf.tmp.1"), Some(n2))
> + .unwrap();
> + assert_eq!(params["guest"].len(), 1);
> + assert_eq!(params["guest"][0]["node"], "n2");
> +
> + assert_eq!(conn.matches(Some(n1), Some(n1)).unwrap()["guest"].len(), 1);
> + }
> +
> + #[test]
> + fn private_paths_need_privileges() {
> + let (mut conn, _client) = pair(false);
> + conn.templates = templates(&[("all", "{path...}")]);
> + conn.subscribed = true;
> + assert!(conn.matches(Some("datacenter.cfg"), None).is_some());
> + assert!(conn.matches(Some("priv/lock/x"), None).is_none());
> + assert!(conn.matches(Some("nodes/n1/priv/ssl.key"), None).is_none());
> + assert!(
> + conn.matches(Some("datacenter.cfg"), Some("priv/moved"))
> + .is_none()
> + );
> + assert!(
> + conn.matches(Some("priv/old"), Some("datacenter.cfg"))
> + .is_none()
> + );
> +
> + let (mut conn, _client) = pair(true);
> + conn.templates = templates(&[("all", "{path...}")]);
> + conn.subscribed = true;
> + assert!(conn.matches(Some("priv/lock/x"), None).is_some());
> + assert!(
> + conn.matches(Some("priv/old"), Some("datacenter.cfg"))
> + .is_some()
> + );
> + }
> +}
> diff --git a/src/rust/pmxcfs-notify/src/lib.rs b/src/rust/pmxcfs-notify/src/lib.rs
> index 8b13789..d363efb 100644
> --- a/src/rust/pmxcfs-notify/src/lib.rs
> +++ b/src/rust/pmxcfs-notify/src/lib.rs
> @@ -1 +1,43 @@
> +//! Change notification socket for pmxcfs.
> +//!
> +//! Clients connect to a Unix stream socket and exchange newline
> +//! delimited JSON objects. A request is `{"command": "...", "args":
> +//! {...}}` and is answered with `{"ok": <json>}` or `{"error": "..."}`.
> +//! Once a client has sent `subscribe` with named path templates, every
> +//! matching memdb mutation is pushed as `{"event": {...}}` with the
> +//! memdb version as sequence number, the event type, the path, the
> +//! rename target, and per subscription name the parameter sets captured
> +//! from the path and the rename target. A client only ever needs to look
> +//! at which key an object has.
> +//!
> +//! A template is a memdb path whose components may each hold one
> +//! `{name}` placeholder standing for part of that component, and whose
> +//! last component may be `{name...}` for the rest of the path, so
> +//! `nodes/{node}/qemu-server/{vmid}.conf` captures node and vmid and
> +//! `{path...}` matches everything. Connections authorized by group
> +//! membership rather than as root never receive events that mention a
> +//! private path, matching what the IPC and FUSE side hide from them.
> +//!
> +//! The daemon keeps the last mutations in a fixed size ring and each
> +//! connection a cursor into it, so a client whose socket is full is
> +//! caught up once it reads again, and one that fell off the ring gets a
> +//! single resync event instead. A reconnecting client passes the last
> +//! sequence number it saw as `since` and receives what it missed, or a
> +//! resync when that is no longer available. The reply to `hello` carries
> +//! the protocol version and the newest number, so a client that has not
> +//! seen an event yet can still resume later.
>
> +mod conn;
> +mod protocol;
> +mod registry;
> +mod ring;
> +mod server;
> +mod template;
> +
> +pub use protocol::EventKind;
> +pub use server::{Config, Server};
> +
> +/// memdb paths carry no leading slash.
> +pub(crate) fn normalize_path(path: &str) -> &str {
> + path.trim_start_matches('/')
> +}
> diff --git a/src/rust/pmxcfs-notify/src/protocol.rs b/src/rust/pmxcfs-notify/src/protocol.rs
> new file mode 100644
> index 0000000..917cec6
> --- /dev/null
> +++ b/src/rust/pmxcfs-notify/src/protocol.rs
> @@ -0,0 +1,237 @@
> +use std::collections::BTreeMap;
> +
> +use anyhow::{Error, bail, format_err};
> +use serde::{Deserialize, Serialize};
> +use serde_json::Value;
> +
> +use crate::template::{Captures, Template};
> +
> +pub(crate) const PROTOCOL_VERSION: u32 = 1;
> +const MAX_PATTERNS: usize = 256;
> +
> +/// One entry per matching side of the event, keyed by subscription name.
> +pub(crate) type Params<'a> = BTreeMap<&'a str, Vec<Captures<'a>>>;
> +
> +#[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize)]
> +#[serde(rename_all = "lowercase")]
> +pub enum EventKind {
> + Create,
> + Write,
> + Mtime,
> + Rename,
> + Delete,
> + Mkdir,
> + Resync,
> +}
> +
> +#[derive(Serialize)]
> +pub(crate) struct Event<'a> {
> + pub seq: u64,
> + #[serde(rename = "type")]
> + pub kind: EventKind,
> + #[serde(skip_serializing_if = "Option::is_none")]
> + pub path: Option<&'a str>,
> + #[serde(skip_serializing_if = "Option::is_none")]
> + pub to: Option<&'a str>,
> + #[serde(skip_serializing_if = "Option::is_none")]
> + pub params: Option<Params<'a>>,
> +}
> +
> +/// Everything the server sends. The variant name is the single key of
> +/// the JSON object on the wire, so a client dispatches on that key.
> +#[derive(Serialize)]
> +#[serde(rename_all = "lowercase")]
> +pub(crate) enum Message<'a> {
> + Ok(Value),
> + Error(String),
> + Event(Event<'a>),
> +}
> +
> +impl Message<'_> {
> + pub(crate) fn line(&self) -> String {
> + let mut line = serde_json::to_string(self).expect("message serialization cannot fail");
> + line.push('\n');
> + line
> + }
> +}
> +
> +#[derive(Deserialize)]
> +struct Request {
> + command: String,
> + #[serde(default)]
> + args: Value,
> +}
> +
> +#[derive(Debug, PartialEq, Eq)]
How about we just make `Template` deserializable and then drop the
entire manual parsing below:
If you derive Deserialize and just add
#[serde(tag = "command", content = "args", rename_all = "lowercase")]
here (and turn the patterns into a HashMap), then you get the exact
protocl you're manually parsing with a single call to
serde_json::from_str().
That said - why didn't you design the protocol itself *by* writing the
structs out?
Never just handroll parsers for structural data in rust!
Anyway, you then still need to
let patterns = patterns.into_iter().collect();
when you call `registry.subscribe`. It wouldn't be more costly, since
the `Value` you use *now* already produces a map anyway.
(Or you turn the end users of this Vec into just sticking with the
HashMap as well…, that would also make it cheap enough to remove a
single entry by name if we ever want that later)
Then `Request`, `patterns()`, `since()` and `parse_command()` cann all
just be dropped.
> +pub(crate) enum Command {
> + Hello,
> + Subscribe {
> + patterns: Vec<(String, Template)>,
> + since: Option<u64>,
> + },
> + Unsubscribe,
> +}
> +
> +fn patterns(args: &Value) -> Result<Vec<(String, Template)>, Error> {
> + let entries = match args.get("patterns") {
> + None | Some(Value::Null) => return Ok(Vec::new()),
> + Some(Value::Object(entries)) => entries,
> + Some(_) => bail!("'patterns' must map names to path templates"),
> + };
> + if entries.len() > MAX_PATTERNS {
> + bail!("at most {MAX_PATTERNS} patterns per subscription");
> + }
> + entries
> + .iter()
> + .map(|(name, text)| {
> + let text = text
> + .as_str()
> + .ok_or_else(|| format_err!("pattern '{name}' must be a string"))?;
> + let template =
> + Template::parse(text).map_err(|err| format_err!("pattern '{name}': {err}"))?;
> + Ok((name.clone(), template))
> + })
> + .collect()
> +}
> +
> +fn since(args: &Value) -> Result<Option<u64>, Error> {
> + match args.get("since") {
> + None | Some(Value::Null) => Ok(None),
> + Some(value) => value
> + .as_u64()
> + .map(Some)
> + .ok_or_else(|| format_err!("'since' must be an unsigned integer")),
> + }
> +}
> +
> +pub(crate) fn parse_command(line: &str) -> Result<Command, Error> {
> + let request: Request = serde_json::from_str(line)?;
> + match request.command.as_str() {
> + "hello" => Ok(Command::Hello),
> + "subscribe" => Ok(Command::Subscribe {
> + patterns: patterns(&request.args)?,
> + since: since(&request.args)?,
> + }),
> + "unsubscribe" => Ok(Command::Unsubscribe),
> + other => bail!("unknown command '{other}'"),
> + }
> +}
> +
> +#[cfg(test)]
> +mod tests {
> + use super::*;
> +
> + #[test]
> + fn parses_commands() {
> + assert_eq!(
> + parse_command(r#"{"command":"hello"}"#).unwrap(),
> + Command::Hello
> + );
> + assert_eq!(
> + parse_command(
> + r#"{"command":"subscribe","args":{"patterns":{"guest":"nodes/{node}/qemu-server/{vmid}.conf","dc":"/datacenter.cfg"}}}"#
> + )
> + .unwrap(),
> + Command::Subscribe {
> + patterns: vec![
> + ("dc".to_owned(), Template::parse("datacenter.cfg").unwrap()),
> + (
> + "guest".to_owned(),
> + Template::parse("nodes/{node}/qemu-server/{vmid}.conf").unwrap()
> + ),
> + ],
(And you can then just extend the `vec![]` with
`.into_iter().collect()` here in the test.
> + since: None,
> + }
> + );
> + assert_eq!(
> + parse_command(r#"{"command":"subscribe"}"#).unwrap(),
> + Command::Subscribe {
> + patterns: Vec::new(),
> + since: None,
> + }
> + );
> + assert_eq!(
> + parse_command(r#"{"command":"subscribe","args":{"since":42}}"#).unwrap(),
> + Command::Subscribe {
> + patterns: Vec::new(),
> + since: Some(42),
> + }
> + );
> + assert_eq!(
> + parse_command(r#"{"command":"unsubscribe"}"#).unwrap(),
> + Command::Unsubscribe
> + );
> + }
> +
> + #[test]
> + fn rejects_bad_requests() {
> + assert!(parse_command("not json").is_err());
> + assert!(parse_command(r#"{"command":"subscribe","args":{"patterns":["a.cfg"]}}"#).is_err());
> + assert!(parse_command(r#"{"command":"subscribe","args":{"patterns":{"x":1}}}"#).is_err());
> + assert!(parse_command(r#"{"command":"subscribe","args":{"since":"x"}}"#).is_err());
> + assert!(parse_command(r#"{"command":"subscribe","args":{"since":-1}}"#).is_err());
> + let err = parse_command(r#"{"command":"subscribe","args":{"patterns":{"x":"{a}{b}"}}}"#)
> + .unwrap_err();
> + assert_eq!(
> + err.to_string(),
> + "pattern 'x': only one placeholder per component in '{a}{b}'"
> + );
> + let err = parse_command(r#"{"command":"bogus"}"#).unwrap_err();
> + assert_eq!(err.to_string(), "unknown command 'bogus'");
> + }
> +
> + #[test]
> + fn formats_events() {
> + let params = Params::from([(
> + "guest",
> + vec![Captures::from([("node", "n1"), ("vmid", "100")])],
> + )]);
> + let write = Message::Event(Event {
> + seq: 7,
> + kind: EventKind::Write,
> + path: Some("nodes/n1/qemu-server/100.conf"),
> + to: None,
> + params: Some(params),
> + });
> + assert_eq!(
> + write.line(),
> + "{\"event\":{\"seq\":7,\"type\":\"write\",\"path\":\"nodes/n1/qemu-server/100.conf\",\"params\":{\"guest\":[{\"node\":\"n1\",\"vmid\":\"100\"}]}}}\n"
> + );
> + let params = Params::from([("dc", vec![Captures::new()])]);
> + let rename = Message::Event(Event {
> + seq: 8,
> + kind: EventKind::Rename,
> + path: Some("a"),
> + to: Some("datacenter.cfg"),
> + params: Some(params),
> + });
> + assert_eq!(
> + rename.line(),
> + "{\"event\":{\"seq\":8,\"type\":\"rename\",\"path\":\"a\",\"to\":\"datacenter.cfg\",\"params\":{\"dc\":[{}]}}}\n"
> + );
> + let resync = Message::Event(Event {
> + seq: 9,
> + kind: EventKind::Resync,
> + path: None,
> + to: None,
> + params: None,
> + });
> + assert_eq!(
> + resync.line(),
> + "{\"event\":{\"seq\":9,\"type\":\"resync\"}}\n"
> + );
> + }
> +
> + #[test]
> + fn formats_replies() {
> + assert_eq!(Message::Ok(Value::Null).line(), "{\"ok\":null}\n");
> + assert_eq!(
> + Message::Ok(serde_json::json!({"protocol": 1})).line(),
> + "{\"ok\":{\"protocol\":1}}\n"
> + );
> + assert_eq!(
> + Message::Error("multi\nline".to_owned()).line(),
> + "{\"error\":\"multi\\nline\"}\n"
> + );
> + }
> +}
> diff --git a/src/rust/pmxcfs-notify/src/registry.rs b/src/rust/pmxcfs-notify/src/registry.rs
> new file mode 100644
> index 0000000..29fe4dd
> --- /dev/null
> +++ b/src/rust/pmxcfs-notify/src/registry.rs
> @@ -0,0 +1,405 @@
> +use std::collections::{HashMap, HashSet};
> +use std::os::fd::RawFd;
> +use std::sync::{Mutex, MutexGuard};
> +
> +use crate::conn::Conn;
> +use crate::protocol::{Event, EventKind, Message};
> +use crate::ring::{Entry, Ring};
> +use crate::template::Template;
> +
> +// Entries copied out of the ring per lock hold. Bounds how long the
> +// mutating thread can wait behind a delivery walk.
> +const BATCH: usize = 256;
> +
> +// Batches one connection may take per delivery round before the others get a
> +// turn, so a client catching up a large backlog cannot starve them.
> +const WALK_ROUNDS: usize = 8;
> +
> +pub(crate) fn lock_ring(ring: &Mutex<Ring>) -> MutexGuard<'_, Ring> {
> + ring.lock().unwrap_or_else(|poisoned| poisoned.into_inner())
> +}
> +
> +#[derive(Default)]
> +pub(crate) struct Registry {
> + conns: HashMap<u64, Conn>,
> + next_id: u64,
IMO Ids are worth a proper type. Makes it easier to adapt in the future.
#[derive(Clone, Copy, Debug, Eq, PartialEq, Ord, PartialOrd, Hash)]
pub struct Id(u64);
impl Id {
pub fn as_raw(self) -> u64;
}
> + #[cfg(test)]
> + pub(crate) panic_on_delivery: bool,
> +}
> +
> +impl Registry {
> + pub(crate) fn add(&mut self, mut conn: Conn, head: u64) -> u64 {
> + conn.cursor = head;
> + let id = self.next_id;
> + self.next_id += 1;
> + self.conns.insert(id, conn);
> + id
> + }
> +
> + pub(crate) fn remove(&mut self, id: u64) -> Option<Conn> {
> + self.conns.remove(&id)
> + }
> +
> + pub(crate) fn get_mut(&mut self, id: u64) -> Option<&mut Conn> {
> + self.conns.get_mut(&id)
> + }
> +
> + pub(crate) fn drain_dead(&mut self) -> Vec<Conn> {
> + self.conns
> + .extract_if(|_, conn| conn.dead)
> + .map(|(_, conn)| conn)
> + .collect()
> + }
> +
> + pub(crate) fn fds(&self) -> Vec<(u64, RawFd, bool)> {
↑ Doc comment for the type please.
> + self.conns
> + .iter()
> + .map(|(id, conn)| (*id, conn.fd(), conn.lagging()))
> + .collect()
> + }
> +
> + pub(crate) fn len(&self) -> usize {
> + self.conns.len()
> + }
> +
> + pub(crate) fn clear(&mut self) {
> + self.conns.clear();
> + }
> +
> + // A fresh subscription starts at the tail, or right after `since` when
> + // the ring still holds it, and with a resync when it does not. On a live
> + // subscription the cursor stays put, what the client asks for may still
> + // be in flight to it.
> + pub(crate) fn subscribe(
> + &mut self,
> + id: u64,
> + patterns: Vec<(String, Template)>,
> + since: Option<u64>,
> + ring: &Mutex<Ring>,
> + ) {
> + let Some(conn) = self.conns.get_mut(&id) else {
> + return;
> + };
> + let fresh = !conn.subscribed;
> + conn.templates = patterns;
> + conn.subscribed = true;
> + if !fresh {
> + return;
> + }
> + let ring = lock_ring(ring);
> + conn.cursor = ring.head();
> + conn.resync_pending = false;
> + let Some(seq) = since else {
> + return;
> + };
> + match ring.resume(seq) {
> + Some(pos) => conn.cursor = pos,
> + None => conn.resync_pending = true,
> + }
> + }
> +
> + pub(crate) fn unsubscribe(&mut self, id: u64, head: u64) {
> + let Some(conn) = self.conns.get_mut(&id) else {
> + return;
> + };
> + conn.templates.clear();
> + conn.subscribed = false;
> + conn.resync_pending = false;
> + conn.cursor = head;
> + }
> +
> + /// Sends to every connection that can take data, a lagging one only
> + /// when its socket was reported writable. Returns whether a connection
> + /// stopped at its batch cap with more still to send.
> + pub(crate) fn deliver(&mut self, ring: &Mutex<Ring>, writable: &HashSet<u64>) -> bool {
> + let mut more = false;
> + for (id, conn) in &mut self.conns {
> + if conn.dead || (conn.lagging() && !writable.contains(id)) {
> + continue;
> + }
> + more |= walk(ring, conn);
> + }
> + more
> + }
> +}
> +
> +fn resync_line(seq: u64) -> String {
> + Message::Event(Event {
> + seq,
> + kind: EventKind::Resync,
> + path: None,
> + to: None,
> + params: None,
> + })
> + .line()
> +}
> +
> +// Moves a connection's cursor toward the head, queueing the lines it
> +// subscribed to, until the socket is full, nothing is left, or it has taken
> +// its batch cap for this round. A cursor that fell off the ring turns into
> +// one resync at the head. The ring is only locked to copy a batch, never
> +// while a line is built or sent. Returns whether the cap stopped it with more
> +// still to send.
> +fn walk(ring: &Mutex<Ring>, conn: &mut Conn) -> bool {
> + for _ in 0..WALK_ROUNDS {
> + if !conn.flush() {
> + return false;
> + }
> + let (batch, last_seq) = {
> + let ring = lock_ring(ring);
> + if !conn.subscribed {
> + conn.cursor = ring.head();
> + return false;
> + }
> + if conn.cursor < ring.oldest() {
> + conn.resync_pending = true;
> + }
> + // the resync stands for everything up to the head as of now, so
> + // nothing before it may follow it
> + if conn.resync_pending {
> + conn.cursor = ring.head();
> + }
> + let batch: Vec<Entry> = (conn.cursor..ring.head())
> + .take(BATCH)
> + .filter_map(|pos| ring.get(pos).cloned())
> + .collect();
> + (batch, ring.last_seq())
> + };
> + if conn.resync_pending {
> + conn.resync_pending = false;
> + conn.queue(&resync_line(last_seq.unwrap_or(0)));
> + continue;
> + }
> + if batch.is_empty() {
> + return false;
> + }
> + for entry in &batch {
> + conn.cursor += 1;
> + let line = if entry.kind == EventKind::Resync {
> + resync_line(entry.seq)
> + } else {
> + match conn.matches(entry.path.as_deref(), entry.to.as_deref()) {
> + Some(params) => Message::Event(Event {
> + seq: entry.seq,
> + kind: entry.kind,
> + path: entry.path.as_deref(),
> + to: entry.to.as_deref(),
> + params: Some(params),
> + })
> + .line(),
> + None => continue,
> + }
> + };
> + conn.queue(&line);
> + if !conn.flush() {
> + return false;
> + }
> + }
> + }
> + true
> +}
> +
> +// Same rule as path_is_private in the daemon, which also hides these paths
> +// from group authorized IPC and FUSE clients.
> +pub(crate) fn is_private(path: &str) -> bool {
> + fn priv_component(path: &str) -> bool {
> + path == "priv" || path.starts_with("priv/")
> + }
> + let path = crate::normalize_path(path);
> + priv_component(path)
> + || path
> + .strip_prefix("nodes/")
> + .and_then(|rest| rest.split_once('/'))
> + .is_some_and(|(_, below)| priv_component(below))
> +}
> +
> +#[cfg(test)]
> +mod tests {
> + use std::io::{BufRead, BufReader};
> + use std::os::unix::net::UnixStream;
> + use std::time::Duration;
> +
> + use super::*;
> +
> + fn entry(seq: u64) -> Entry {
> + Entry {
> + seq,
> + kind: EventKind::Write,
> + path: Some(format!("f{seq}").into()),
> + to: None,
> + }
> + }
> +
> + fn all() -> Vec<(String, Template)> {
> + vec![("all".to_owned(), Template::parse("{path...}").unwrap())]
> + }
> +
> + fn client(registry: &mut Registry, ring: &Mutex<Ring>) -> (u64, BufReader<UnixStream>) {
> + let (server, client) = UnixStream::pair().unwrap();
> + client
> + .set_read_timeout(Some(Duration::from_millis(300)))
> + .unwrap();
> + let head = lock_ring(ring).head();
> + (
> + registry.add(Conn::new(server, true).unwrap(), head),
> + BufReader::new(client),
> + )
> + }
> +
> + fn lines(reader: &mut BufReader<UnixStream>) -> Vec<String> {
> + let mut lines = Vec::new();
> + let mut line = String::new();
> + while reader.read_line(&mut line).is_ok_and(|n| n > 0) {
> + lines.push(line.trim_end().to_owned());
> + line.clear();
> + }
> + lines
> + }
> +
> + #[test]
> + fn resume_positions_the_cursor() {
> + let ring = Mutex::new(Ring::new(4));
> + let mut registry = Registry::default();
> + for seq in 1..=3 {
> + lock_ring(&ring).push(entry(seq));
> + }
> + let (id, _reader) = client(&mut registry, &ring);
> +
> + registry.subscribe(id, all(), None, &ring);
> + assert_eq!(registry.get_mut(id).unwrap().cursor, 3);
> +
> + registry.unsubscribe(id, 3);
> + registry.subscribe(id, all(), Some(1), &ring);
> + assert_eq!(registry.get_mut(id).unwrap().cursor, 1);
> + assert!(!registry.get_mut(id).unwrap().resync_pending);
> +
> + registry.unsubscribe(id, 3);
> + registry.subscribe(id, all(), Some(3), &ring);
> + assert_eq!(registry.get_mut(id).unwrap().cursor, 3);
> +
> + registry.unsubscribe(id, 3);
> + registry.subscribe(id, all(), Some(9), &ring);
> + assert!(registry.get_mut(id).unwrap().resync_pending);
> +
> + // a live subscription keeps its cursor whatever it asks for
> + registry.get_mut(id).unwrap().resync_pending = false;
> + registry.subscribe(id, all(), Some(1), &ring);
> + assert_eq!(registry.get_mut(id).unwrap().cursor, 3);
> + assert!(!registry.get_mut(id).unwrap().resync_pending);
> + }
> +
> + #[test]
> + fn walk_delivers_or_resyncs() {
> + let ring = Mutex::new(Ring::new(3));
> + let mut registry = Registry::default();
> + let (id, mut reader) = client(&mut registry, &ring);
> + registry.subscribe(id, all(), None, &ring);
> +
> + lock_ring(&ring).push(entry(1));
> + registry.deliver(&ring, &HashSet::new());
> + assert_eq!(
> + lines(&mut reader),
> + vec![
> + "{\"event\":{\"seq\":1,\"type\":\"write\",\"path\":\"f1\",\"params\":{\"all\":[{\"path\":\"f1\"}]}}}"
> + ]
> + );
> +
> + for seq in 2..=6 {
> + lock_ring(&ring).push(entry(seq));
> + }
> + registry.deliver(&ring, &HashSet::new());
> + assert_eq!(
> + lines(&mut reader),
> + vec!["{\"event\":{\"seq\":6,\"type\":\"resync\"}}"],
> + "a cursor that fell off gets one resync"
> + );
> +
> + lock_ring(&ring).push(entry(7));
> + registry.deliver(&ring, &HashSet::new());
> + assert_eq!(lines(&mut reader).len(), 1);
> + }
> +
> + #[test]
> + fn a_resync_at_the_resume_number_is_delivered() {
> + let ring = Mutex::new(Ring::new(4));
> + let mut registry = Registry::default();
> + lock_ring(&ring).push(Entry::marker(5));
> + lock_ring(&ring).push(Entry {
> + seq: 5,
> + kind: EventKind::Resync,
> + path: None,
> + to: None,
> + });
> + let (id, mut reader) = client(&mut registry, &ring);
> + registry.subscribe(id, all(), Some(5), &ring);
> + registry.deliver(&ring, &HashSet::new());
> + assert_eq!(
> + lines(&mut reader),
> + vec!["{\"event\":{\"seq\":5,\"type\":\"resync\"}}"],
> + "the state was replaced at the number the client resumes from"
> + );
> + }
> +
> + #[test]
> + fn a_pending_resync_swallows_what_arrived_before_the_walk() {
> + let ring = Mutex::new(Ring::new(3));
> + let mut registry = Registry::default();
> + for seq in 1..=3 {
> + lock_ring(&ring).push(entry(seq));
> + }
> + let (id, mut reader) = client(&mut registry, &ring);
> + registry.subscribe(id, all(), Some(0), &ring);
> + assert!(registry.get_mut(id).unwrap().resync_pending);
> +
> + lock_ring(&ring).push(entry(4));
> + registry.deliver(&ring, &HashSet::new());
> + assert_eq!(
> + lines(&mut reader),
> + vec!["{\"event\":{\"seq\":4,\"type\":\"resync\"}}"],
> + "the entry pushed between subscribe and walk is covered by the resync"
> + );
> +
> + lock_ring(&ring).push(entry(5));
> + registry.deliver(&ring, &HashSet::new());
> + assert_eq!(lines(&mut reader).len(), 1);
> + }
> +
> + #[test]
> + fn walks_hand_the_ring_back_between_batches() {
> + let ring = Mutex::new(Ring::new(2000));
> + let mut registry = Registry::default();
> + let (id, mut reader) = client(&mut registry, &ring);
> + registry.subscribe(id, all(), None, &ring);
> + for seq in 1..=1000 {
> + lock_ring(&ring).push(entry(seq));
> + }
> + // the pair's buffer holds a few hundred lines, so the walk stops
> + // lagging and is resumed as the poll loop would on writability
> + let mut received = 0;
> + for _ in 0..20 {
> + registry.deliver(&ring, &HashSet::from([id]));
> + received += lines(&mut reader).len();
> + if received == 1000 {
> + break;
> + }
> + }
> + assert_eq!(received, 1000, "every entry arrives across batches");
> + assert!(
> + !lock_ring(&ring).entries_in_use(),
> + "no batch outlives its walk"
> + );
> + }
> +
> + #[test]
> + fn private_rules() {
> + assert!(is_private("priv"));
> + assert!(is_private("priv/lock/x"));
> + assert!(is_private("/priv/tfa.cfg"));
> + assert!(is_private("nodes/n1/priv"));
> + assert!(is_private("nodes/n1/priv/ssl.key"));
> + assert!(!is_private("private"));
> + assert!(!is_private("nodes/n1/qemu-server/100.conf"));
> + assert!(!is_private("nodes/priv"));
> + assert!(!is_private("datacenter.cfg"));
> + }
> +}
> diff --git a/src/rust/pmxcfs-notify/src/ring.rs b/src/rust/pmxcfs-notify/src/ring.rs
> new file mode 100644
> index 0000000..f6f4594
> --- /dev/null
> +++ b/src/rust/pmxcfs-notify/src/ring.rs
> @@ -0,0 +1,171 @@
> +use std::collections::VecDeque;
> +use std::sync::Arc;
> +
> +use crate::protocol::EventKind;
> +
> +// Reference-counted strings, so copying an entry out of the ring for
> +// delivery clones only a pointer.
> +#[derive(Clone)]
> +pub(crate) struct Entry {
> + pub seq: u64,
> + pub kind: EventKind,
> + pub path: Option<Arc<str>>,
> + pub to: Option<Arc<str>>,
> +}
> +
> +impl Entry {
> + /// Stands for the last mutation before the server started. It carries
> + /// that mutation's number and no path, so a resuming client at that
> + /// number is current, while nothing ever matches it.
> + pub(crate) fn marker(seq: u64) -> Self {
> + Self {
> + seq,
> + kind: EventKind::Write,
> + path: None,
> + to: None,
> + }
> + }
> +}
> +
> +// Positions are absolute append counts, so a cursor stays valid while the
> +// ring wraps, and a position below the oldest one means its reader fell
> +// off and has to resync.
> +pub(crate) struct Ring {
> + entries: VecDeque<Entry>,
> + capacity: usize,
↑ Could consider dropping this.
> + head: u64,
> +}
> +
> +impl Ring {
> + pub(crate) fn new(capacity: usize) -> Self {
> + Self {
> + entries: VecDeque::new(),
and use `.with_capacity(capacity)` here
(Although I think this technically uses capacity as a lower bound, but
that's probably not an issue)
> + capacity: capacity.max(1),
> + head: 0,
> + }
> + }
> +
> + pub(crate) fn push(&mut self, entry: Entry) {
> + if self.entries.len() == self.capacity {
and use `self.entries.len() == self.entries.capacity()` here
> + self.entries.pop_front();
> + }
> + self.entries.push_back(entry);
> + self.head += 1;
This should be
self.head = self.head.wrapping_add(1);
> + }
> +
> + pub(crate) fn head(&self) -> u64 {
> + self.head
> + }
> +
> + pub(crate) fn oldest(&self) -> u64 {
> + self.head - self.entries.len() as u64
And this should use `.wrapping_sub()`.
While it is unlikely, it's cheap, and means that wrapping around
u64::MAX would "just work".
> + }
> +
> + pub(crate) fn get(&self, pos: u64) -> Option<&Entry> {
> + if pos < self.oldest() || pos >= self.head {
> + return None;
> + }
If we first turn pos into an index, the above condition becomes `if
index >= capacity`:
let index = pos.wrapping_sub(self.oldest());
if index > self.capacity as u64 {
return None;
}
self.entries.get(index as usize)
(^ I prepared & tested the wrapping changes in a branch if you want to
take a look)
> + self.entries.get((pos - self.oldest()) as usize)
> + }
> +
> + pub(crate) fn last_seq(&self) -> Option<u64> {
> + self.entries.back().map(|entry| entry.seq)
> + }
> +
> + /// Whether any entry is shared with a copy outside the ring.
> + #[cfg(test)]
> + pub(crate) fn entries_in_use(&self) -> bool {
> + self.entries.iter().any(|entry| {
> + entry
> + .path
> + .as_ref()
> + .is_some_and(|p| Arc::strong_count(p) > 1)
> + || entry.to.as_ref().is_some_and(|t| Arc::strong_count(t) > 1)
> + })
> + }
> +
> + /// Where a client that saw `seq` last continues. Right after the newest
> + /// entry carrying it, unless that entry is a resync, which the client
> + /// then has to see since the state was replaced at that number.
> + pub(crate) fn resume(&self, seq: u64) -> Option<u64> {
> + self.entries
> + .iter()
> + .rposition(|entry| entry.seq == seq)
> + .map(|index| {
> + let pos = self.oldest() + index as u64;
> + if self.entries[index].kind == EventKind::Resync {
> + pos
> + } else {
> + pos + 1
> + }
> + })
> + }
> +}
> +
> +#[cfg(test)]
> +mod tests {
> + use super::*;
> +
> + fn entry(seq: u64) -> Entry {
> + Entry {
> + seq,
> + kind: EventKind::Write,
> + path: Some(format!("f{seq}").into()),
> + to: None,
> + }
> + }
> +
> + #[test]
> + fn positions_survive_wrapping() {
> + let mut ring = Ring::new(3);
> + assert_eq!(ring.head(), 0);
> + assert_eq!(ring.oldest(), 0);
> + assert_eq!(ring.last_seq(), None);
> + assert!(ring.get(0).is_none());
> +
> + for seq in 10..15 {
> + ring.push(entry(seq));
> + }
> + assert_eq!(ring.head(), 5);
> + assert_eq!(ring.oldest(), 2);
> + assert_eq!(ring.last_seq(), Some(14));
> + assert!(ring.get(1).is_none());
> + assert_eq!(ring.get(2).unwrap().seq, 12);
> + assert_eq!(ring.get(4).unwrap().seq, 14);
> + assert!(ring.get(5).is_none());
> + }
> +
> + #[test]
> + fn resume_positions() {
> + let mut ring = Ring::new(4);
> + for seq in [1, 2, 2, 3] {
> + ring.push(entry(seq));
> + }
> + assert_eq!(ring.resume(1), Some(1));
> + assert_eq!(ring.resume(2), Some(3), "the newest duplicate wins");
> + assert_eq!(ring.resume(3), Some(4));
> + assert_eq!(ring.resume(0), None);
> + ring.push(entry(4));
> + assert_eq!(ring.resume(1), None, "fell off");
> + assert_eq!(ring.resume(4), Some(5));
> + }
> +
> + #[test]
> + fn a_resync_at_the_resume_number_is_not_skipped() {
> + let mut ring = Ring::new(4);
> + ring.push(Entry::marker(5));
> + ring.push(Entry {
> + seq: 5,
> + kind: EventKind::Resync,
> + path: None,
> + to: None,
> + });
> + assert_eq!(
> + ring.resume(5),
> + Some(1),
> + "the state was replaced at that number"
> + );
> + ring.push(entry(6));
> + assert_eq!(ring.resume(6), Some(3));
> + }
> +}
> diff --git a/src/rust/pmxcfs-notify/src/server.rs b/src/rust/pmxcfs-notify/src/server.rs
> new file mode 100644
> index 0000000..65d5744
> --- /dev/null
> +++ b/src/rust/pmxcfs-notify/src/server.rs
> @@ -0,0 +1,780 @@
> +use std::any::Any;
> +use std::collections::HashSet;
> +use std::fs;
> +use std::io;
> +use std::os::fd::{AsFd, BorrowedFd, RawFd};
> +use std::os::unix::fs::PermissionsExt;
> +use std::os::unix::net::{UnixListener, UnixStream};
> +use std::panic::{AssertUnwindSafe, catch_unwind};
> +use std::path::PathBuf;
> +use std::sync::atomic::{AtomicBool, Ordering};
> +use std::sync::{Arc, Mutex, MutexGuard};
> +use std::thread::JoinHandle;
> +use std::time::Duration;
> +
> +use anyhow::{Context, Error};
> +use nix::errno::Errno;
> +use nix::poll::{PollFd, PollFlags, PollTimeout, poll};
> +use nix::sys::eventfd::{EfdFlags, EventFd};
> +use nix::sys::socket::{getsockopt, setsockopt, sockopt};
> +use serde_json::{Value, json};
> +
> +use crate::conn::{Conn, ReadError};
> +use crate::normalize_path;
> +use crate::protocol::{Command, EventKind, Message, PROTOCOL_VERSION, parse_command};
> +use crate::registry::{Registry, lock_ring};
> +use crate::ring::{Entry, Ring};
> +
> +const MAX_CONNECTIONS: usize = 128;
> +const SNDBUF: usize = 2 * 1024 * 1024;
> +
> +pub struct Config {
> + pub path: PathBuf,
> + pub gid: u32,
> + /// Mutations kept for catching up lagging and reconnecting clients.
> + pub ring: usize,
> + /// The number of the last mutation before the server starts, so a
> + /// client resuming at it after a restart is current.
> + pub seq: u64,
> +}
> +
> +pub struct Server {
> + #[cfg_attr(not(test), allow(dead_code))]
> + registry: Arc<Mutex<Registry>>,
> + ring: Arc<Mutex<Ring>>,
> + stop: Arc<AtomicBool>,
> + wake: Arc<EventFd>,
> + thread: Option<JoinHandle<()>>,
> + path: PathBuf,
> +}
> +
> +fn lock(registry: &Mutex<Registry>) -> MutexGuard<'_, Registry> {
> + registry
> + .lock()
> + .unwrap_or_else(|poisoned| poisoned.into_inner())
> +}
> +
> +impl Server {
> + pub fn start(config: Config) -> Result<Server, Error> {
> + match fs::remove_file(&config.path) {
> + Ok(()) => {}
> + Err(err) if err.kind() == io::ErrorKind::NotFound => {}
> + Err(err) => {
> + return Err(err)
> + .with_context(|| format!("failed to remove stale {}", config.path.display()));
> + }
> + }
> + let listener = UnixListener::bind(&config.path)
> + .with_context(|| format!("failed to bind {}", config.path.display()))?;
> + listener.set_nonblocking(true)?;
> + std::os::unix::fs::chown(&config.path, None, Some(config.gid))?;
> + fs::set_permissions(&config.path, fs::Permissions::from_mode(0o660))?;
> +
> + let wake = Arc::new(EventFd::from_value_and_flags(
> + 0,
> + EfdFlags::EFD_NONBLOCK | EfdFlags::EFD_CLOEXEC,
> + )?);
> + let registry = Arc::new(Mutex::new(Registry::default()));
> + let mut seeded = Ring::new(config.ring);
> + seeded.push(Entry::marker(config.seq));
> + let ring = Arc::new(Mutex::new(seeded));
> + let stop = Arc::new(AtomicBool::new(false));
> + let worker = Worker {
> + listener,
> + path: config.path.clone(),
> + wake: Arc::clone(&wake),
> + stop: Arc::clone(&stop),
> + registry: Arc::clone(®istry),
> + ring: Arc::clone(&ring),
> + gid: config.gid,
> + accept_failing: false,
> + };
> + let thread = std::thread::Builder::new()
> + .name("pmxcfs-notify".to_owned())
> + .spawn(move || worker.run())?;
> +
> + Ok(Server {
> + registry,
> + ring,
> + stop,
> + wake,
> + thread: Some(thread),
> + path: config.path,
> + })
> + }
> +
> + /// Records a mutation and wakes the delivery thread. It never touches
> + /// a client socket and only takes the ring lock, which the delivery
> + /// thread holds for one batch copy at a time, so it is safe on the
> + /// mutating thread.
> + pub fn emit(&self, seq: u64, kind: EventKind, path: &str, to: Option<&str>) {
> + if self.stop.load(Ordering::SeqCst) {
> + return;
> + }
> + lock_ring(&self.ring).push(Entry {
> + seq,
> + kind,
> + path: Some(normalize_path(path).into()),
> + to: to.map(|to| normalize_path(to).into()),
> + });
> + let _ = self.wake.write(1);
> + }
> +
> + pub fn resync(&self, seq: u64) {
> + if self.stop.load(Ordering::SeqCst) {
> + return;
> + }
> + lock_ring(&self.ring).push(Entry {
> + seq,
> + kind: EventKind::Resync,
> + path: None,
> + to: None,
> + });
> + let _ = self.wake.write(1);
> + }
> +
> + pub fn shutdown(mut self) {
> + self.stop_worker();
> + }
> +
> + fn stop_worker(&mut self) {
> + let Some(thread) = self.thread.take() else {
> + return;
> + };
> + self.stop.store(true, Ordering::SeqCst);
> + let _ = self.wake.write(1);
> + let _ = thread.join();
> + let _ = fs::remove_file(&self.path);
> + }
> +
> + #[cfg(test)]
> + pub(crate) fn connection_count(&self) -> usize {
> + lock(&self.registry).len()
> + }
> +}
> +
> +impl Drop for Server {
> + fn drop(&mut self) {
> + self.stop_worker();
> + }
> +}
> +
> +struct Worker {
> + listener: UnixListener,
> + path: PathBuf,
> + wake: Arc<EventFd>,
> + stop: Arc<AtomicBool>,
> + registry: Arc<Mutex<Registry>>,
> + ring: Arc<Mutex<Ring>>,
> + gid: u32,
> + accept_failing: bool,
> +}
> +
> +fn panic_message(payload: &(dyn Any + Send)) -> String {
> + payload
> + .downcast_ref::<&str>()
> + .map(|s| s.to_string())
> + .or_else(|| payload.downcast_ref::<String>().cloned())
> + .unwrap_or_default()
> +}
> +
> +impl Worker {
> + // Whatever ends the loop, the connections are closed and the socket
> + // path removed on the way out, so clients notice, reconnect and log
> + // rather than staying attached to a thread that no longer serves them.
> + fn run(mut self) {
> + while !self.stop.load(Ordering::SeqCst) {
> + match catch_unwind(AssertUnwindSafe(|| self.iteration())) {
> + Ok(Ok(())) => {}
> + Ok(Err(err)) => {
> + log::error!("poll loop failed, stopping: {err}");
> + break;
> + }
> + Err(payload) => {
> + log::error!("poll loop panicked, stopping: {}", panic_message(&*payload));
> + break;
> + }
> + }
> + }
> + self.stop.store(true, Ordering::SeqCst);
> + lock(&self.registry).clear();
> + let _ = fs::remove_file(&self.path);
> + }
> +
> + fn iteration(&mut self) -> Result<(), Error> {
> + let (dead, clients): (Vec<Conn>, Vec<(u64, RawFd, bool)>) = {
> + let mut registry = lock(&self.registry);
> + (registry.drain_dead(), registry.fds())
> + };
> + drop(dead);
> +
> + let mut fds: Vec<PollFd> = Vec::with_capacity(clients.len() + 2);
> + fds.push(PollFd::new(self.listener.as_fd(), PollFlags::POLLIN));
> + fds.push(PollFd::new(self.wake.as_fd(), PollFlags::POLLIN));
> + for (_, fd, lagging) in &clients {
> + // Client fds are closed by this thread only, so the raw fds
> + // collected under the lock stay valid for the whole poll.
> + let fd = unsafe { BorrowedFd::borrow_raw(*fd) };
> + let mut flags = PollFlags::POLLIN;
> + if *lagging {
> + flags |= PollFlags::POLLOUT;
> + }
> + fds.push(PollFd::new(fd, flags));
> + }
> +
> + match poll(&mut fds, PollTimeout::NONE) {
> + Ok(_) => {}
> + Err(Errno::EINTR) => return Ok(()),
> + // out of memory is transient, so back off and re-poll rather than
> + // tear the notifier down for the rest of the daemon's life
> + Err(Errno::ENOMEM) => {
> + log::warn!("poll transiently out of memory, retrying");
> + std::thread::sleep(Duration::from_millis(100));
> + return Ok(());
> + }
> + Err(err) => return Err(err.into()),
> + }
> + let revents: Vec<PollFlags> = fds
> + .iter()
> + .map(|fd| fd.revents().unwrap_or(PollFlags::empty()))
> + .collect();
> + drop(fds);
> +
> + if revents[1].intersects(PollFlags::POLLIN) {
> + let _ = self.wake.read();
> + }
> + if self.stop.load(Ordering::SeqCst) {
> + return Ok(());
> + }
> + if revents[0].intersects(PollFlags::POLLIN) {
> + self.accept_all();
> + }
> + let readable =
> + PollFlags::POLLIN | PollFlags::POLLHUP | PollFlags::POLLERR | PollFlags::POLLNVAL;
> + let mut writable = HashSet::new();
> + for (index, (id, _, _)) in clients.iter().enumerate() {
> + let flags = revents[index + 2];
> + if flags.intersects(readable) {
> + self.service(*id);
> + }
> + if flags.intersects(PollFlags::POLLOUT) {
> + writable.insert(*id);
> + }
> + }
> + let mut registry = lock(&self.registry);
> + #[cfg(test)]
> + if registry.panic_on_delivery {
> + panic!("injected");
> + }
> + let more = registry.deliver(&self.ring, &writable);
> + drop(registry);
> + // a connection that stopped at its batch cap still has data queued, so
> + // wake the loop to resume it once the other fds have had a turn
> + if more {
> + let _ = self.wake.write(1);
> + }
> + Ok(())
> + }
> +
> + fn accept_all(&mut self) {
> + loop {
> + match self.listener.accept() {
> + Ok((stream, _)) => {
> + self.accept_failing = false;
> + self.admit(stream);
> + }
> + Err(err) if err.kind() == io::ErrorKind::WouldBlock => break,
> + Err(err) if err.kind() == io::ErrorKind::Interrupted => continue,
> + Err(err) => {
> + // the pending connection stays in the backlog and keeps
> + // the listener readable, so do not spin on it. a cause
> + // that lasts, like exhausted descriptors, would flood
> + // the log with one warning per round
> + if !self.accept_failing {
> + log::warn!("accept failed: {err}");
> + self.accept_failing = true;
> + }
> + std::thread::sleep(Duration::from_millis(100));
> + break;
> + }
> + }
> + }
> + }
> +
> + fn admit(&self, stream: UnixStream) {
> + let cred = match getsockopt(&stream, sockopt::PeerCredentials) {
> + Ok(cred) => cred,
> + Err(err) => {
> + log::warn!("failed to get peer credentials: {err}");
> + return;
> + }
> + };
> + let privileged = cred.uid() == 0 && cred.gid() == 0;
> + if !privileged && cred.gid() != self.gid {
> + log::warn!(
> + "connection from bad user {}/{} rejected",
> + cred.uid(),
> + cred.gid()
> + );
> + return;
> + }
> + if lock(&self.registry).len() >= MAX_CONNECTIONS {
> + log::warn!("connection limit of {MAX_CONNECTIONS} reached, refusing");
> + return;
> + }
> + // root may raise the send buffer past wmem_max, which gives a
> + // lagging client room for thousands of events instead of hundreds
> + if setsockopt(&stream, sockopt::SndBufForce, &SNDBUF).is_err()
> + && setsockopt(&stream, sockopt::SndBuf, &SNDBUF).is_err()
> + {
> + log::debug!("could not raise the send buffer of a connection");
> + }
> + match Conn::new(stream, privileged) {
> + Ok(conn) => {
> + let head = lock_ring(&self.ring).head();
> + lock(&self.registry).add(conn, head);
> + }
> + Err(err) => log::warn!("failed to set up connection: {err}"),
> + }
> + }
> +
> + fn service(&self, id: u64) {
> + let mut registry = lock(&self.registry);
> + let lines = match registry.get_mut(id).map(Conn::read_lines) {
> + None => return,
> + Some(Ok(lines)) => lines,
> + Some(Err(err)) => {
> + match err {
> + ReadError::Eof => {}
> + ReadError::LineTooLong => log::warn!("closing connection, line too long"),
> + ReadError::Io(err) => log::debug!("closing connection: {err}"),
> + }
> + registry.remove(id);
> + return;
> + }
> + };
> + for line in lines {
> + let reply = match parse_command(&line) {
> + Ok(Command::Hello) => {
> + // the newest number lets a client that saw no event yet
> + // resume after a restart instead of resyncing
> + let seq = lock_ring(&self.ring).last_seq().unwrap_or(0);
> + Message::Ok(json!({ "protocol": PROTOCOL_VERSION, "seq": seq })).line()
> + }
> + Ok(Command::Subscribe { patterns, since }) => {
> + registry.subscribe(id, patterns, since, &self.ring);
> + Message::Ok(Value::Null).line()
> + }
> + Ok(Command::Unsubscribe) => {
> + let head = lock_ring(&self.ring).head();
> + registry.unsubscribe(id, head);
> + Message::Ok(Value::Null).line()
> + }
> + Err(err) => Message::Error(err.to_string()).line(),
> + };
> + let Some(conn) = registry.get_mut(id) else {
> + return;
> + };
> + conn.queue(&reply);
> + conn.flush();
> + if conn.dead {
> + return;
> + }
> + }
> + }
> +}
> +
> +#[cfg(test)]
> +mod tests {
> + use std::io::{BufRead, BufReader, Write};
> + use std::os::unix::fs::MetadataExt;
> + use std::sync::atomic::AtomicUsize;
> + use std::time::Instant;
> +
> + use super::*;
> +
> + static COUNTER: AtomicUsize = AtomicUsize::new(0);
> +
> + struct TempDir(PathBuf);
> +
> + impl Drop for TempDir {
> + fn drop(&mut self) {
> + let _ = fs::remove_dir_all(&self.0);
> + }
> + }
> +
> + fn start(ring: usize) -> (Server, PathBuf, TempDir) {
> + start_at(ring, 0)
> + }
> +
> + fn start_at(ring: usize, seq: u64) -> (Server, PathBuf, TempDir) {
> + let n = COUNTER.fetch_add(1, Ordering::SeqCst);
> + let dir = std::env::temp_dir().join(format!("pmxcfs-notify-{}-{n}", std::process::id()));
> + fs::create_dir_all(&dir).unwrap();
> + let gid = fs::metadata(&dir).unwrap().gid();
> + let path = dir.join("sock");
> + let server = Server::start(Config {
> + path: path.clone(),
> + gid,
> + ring,
> + seq,
> + })
> + .unwrap();
> + (server, path, TempDir(dir))
> + }
> +
> + fn connect(path: &PathBuf) -> (UnixStream, BufReader<UnixStream>) {
> + let stream = UnixStream::connect(path).unwrap();
> + stream
> + .set_read_timeout(Some(Duration::from_secs(5)))
> + .unwrap();
> + let reader = BufReader::new(stream.try_clone().unwrap());
> + (stream, reader)
> + }
> +
> + fn request(stream: &mut UnixStream, reader: &mut BufReader<UnixStream>, line: &str) -> String {
> + stream.write_all(line.as_bytes()).unwrap();
> + stream.write_all(b"\n").unwrap();
> + let mut reply = String::new();
> + reader.read_line(&mut reply).unwrap();
> + reply
> + }
> +
> + fn subscribe(
> + stream: &mut UnixStream,
> + reader: &mut BufReader<UnixStream>,
> + patterns: &str,
> + since: Option<u64>,
> + ) {
> + let since = since.map(|s| format!(",\"since\":{s}")).unwrap_or_default();
> + let line =
> + format!("{{\"command\":\"subscribe\",\"args\":{{\"patterns\":{patterns}{since}}}}}");
> + assert_eq!(request(stream, reader, &line), "{\"ok\":null}\n");
> + }
> +
> + fn event(reader: &mut BufReader<UnixStream>) -> Value {
> + let mut line = String::new();
> + reader.read_line(&mut line).unwrap();
> + serde_json::from_str::<Value>(&line).unwrap()["event"].clone()
> + }
> +
> + /// The next event if one arrives shortly, None when the line stays quiet.
> + fn pending(reader: &mut BufReader<UnixStream>) -> Option<Value> {
> + let stream = reader.get_ref();
> + stream
> + .set_read_timeout(Some(Duration::from_millis(300)))
> + .unwrap();
> + let mut line = String::new();
> + let result = reader.read_line(&mut line);
> + reader
> + .get_ref()
> + .set_read_timeout(Some(Duration::from_secs(5)))
> + .unwrap();
> + match result {
> + Ok(n) if n > 0 => Some(serde_json::from_str::<Value>(&line).unwrap()["event"].clone()),
> + _ => None,
> + }
> + }
> +
> + fn wait_for(mut condition: impl FnMut() -> bool) {
> + let deadline = Instant::now() + Duration::from_secs(5);
> + while !condition() {
> + assert!(Instant::now() < deadline, "condition not met in time");
> + std::thread::sleep(Duration::from_millis(10));
> + }
> + }
> +
> + #[test]
> + fn subscribe_and_receive() {
> + let (server, path, _dir) = start(64);
> + let meta = fs::metadata(&path).unwrap();
> + assert_eq!(meta.mode() & 0o777, 0o660);
> +
> + let (mut stream, mut reader) = connect(&path);
> + assert_eq!(
> + request(&mut stream, &mut reader, r#"{"command":"hello"}"#),
> + "{\"ok\":{\"protocol\":1,\"seq\":0}}\n"
> + );
> + assert_eq!(server.connection_count(), 1);
> + assert_eq!(
> + request(&mut stream, &mut reader, r#"{"command":"bogus"}"#),
> + "{\"error\":\"unknown command 'bogus'\"}\n"
> + );
> + server.emit(1, EventKind::Write, "datacenter.cfg", None);
> + subscribe(
> + &mut stream,
> + &mut reader,
> + r#"{"guest":"/nodes/{node}/qemu-server/{vmid}.conf","dc":"datacenter.cfg"}"#,
> + None,
> + );
> + // a rejected subscription leaves the previous one in place
> + assert_eq!(
> + request(
> + &mut stream,
> + &mut reader,
> + r#"{"command":"subscribe","args":{"patterns":{"bad":"{a}{b}"}}}"#
> + ),
> + "{\"error\":\"pattern 'bad': only one placeholder per component in '{a}{b}'\"}\n"
> + );
> +
> + server.emit(2, EventKind::Write, "storage.cfg", None);
> + server.emit(3, EventKind::Create, "/nodes/n1/qemu-server/100.conf", None);
> + server.emit(
> + 4,
> + EventKind::Write,
> + "nodes/n1/qemu-server/100.conf.tmp.1",
> + None,
> + );
> + server.emit(
> + 5,
> + EventKind::Rename,
> + "nodes/n1/qemu-server/100.conf.tmp.1",
> + Some("datacenter.cfg"),
> + );
> + server.resync(6);
> +
> + let create = event(&mut reader);
> + assert_eq!(create["seq"], 3);
> + assert_eq!(create["type"], "create");
> + assert_eq!(create["path"], "nodes/n1/qemu-server/100.conf");
> + assert!(create.get("to").is_none());
> + assert_eq!(
> + create["params"],
> + json!({ "guest": [{ "node": "n1", "vmid": "100" }] })
> + );
> +
> + let rename = event(&mut reader);
> + assert_eq!(rename["seq"], 5);
> + assert_eq!(rename["type"], "rename");
> + assert_eq!(rename["to"], "datacenter.cfg");
> + assert_eq!(rename["params"], json!({ "dc": [{}] }));
> +
> + assert_eq!(event(&mut reader), json!({ "seq": 6, "type": "resync" }));
> +
> + assert_eq!(
> + request(&mut stream, &mut reader, r#"{"command":"unsubscribe"}"#),
> + "{\"ok\":null}\n"
> + );
> + server.emit(7, EventKind::Write, "datacenter.cfg", None);
> + server.resync(8);
> + // a fresh subscription without a resume point starts at the tail
> + subscribe(&mut stream, &mut reader, r#"{"dc":"datacenter.cfg"}"#, None);
> + server.emit(9, EventKind::Write, "datacenter.cfg", None);
> + assert_eq!(event(&mut reader)["seq"], 9);
> + assert!(pending(&mut reader).is_none());
> +
> + server.shutdown();
> + let mut line = String::new();
> + assert_eq!(reader.read_line(&mut line).unwrap(), 0);
> + assert!(!path.exists());
> + }
> +
> + #[test]
> + fn private_paths_follow_credentials() {
> + let (server, path, _dir) = start(64);
> + let meta = fs::metadata(&path).unwrap();
> + let privileged = meta.uid() == 0 && meta.gid() == 0;
> +
> + let (mut stream, mut reader) = connect(&path);
> + subscribe(&mut stream, &mut reader, r#"{"all":"{path...}"}"#, None);
> + server.emit(1, EventKind::Write, "nodes/n1/priv/ssl.key", None);
> + server.emit(2, EventKind::Write, "nodes/n1/pve-ssl.pem", None);
> +
> + if privileged {
> + assert_eq!(event(&mut reader)["path"], "nodes/n1/priv/ssl.key");
> + }
> + let public = event(&mut reader);
> + assert_eq!(public["path"], "nodes/n1/pve-ssl.pem");
> + assert_eq!(
> + public["params"],
> + json!({ "all": [{ "path": "nodes/n1/pve-ssl.pem" }] })
> + );
> + }
> +
> + #[test]
> + fn panic_in_the_loop_closes_the_connections() {
> + let (server, path, _dir) = start(64);
> + let (mut stream, mut reader) = connect(&path);
> + subscribe(&mut stream, &mut reader, r#"{"all":"{path...}"}"#, None);
> + lock(&server.registry).panic_on_delivery = true;
> + server.emit(1, EventKind::Write, "a", None);
> +
> + let mut line = String::new();
> + assert_eq!(reader.read_line(&mut line).unwrap(), 0, "client sees EOF");
> + wait_for(|| server.thread.as_ref().is_some_and(|t| t.is_finished()));
> + assert!(!path.exists());
> + assert!(UnixStream::connect(&path).is_err());
> + server.shutdown();
> + }
> +
> + #[test]
> + fn client_close_is_noticed() {
> + let (server, path, _dir) = start(64);
> + let (stream, mut reader) = connect(&path);
> + let mut stream = stream;
> + request(&mut stream, &mut reader, r#"{"command":"hello"}"#);
> + assert_eq!(server.connection_count(), 1);
> + drop(reader);
> + drop(stream);
> + wait_for(|| server.connection_count() == 0);
> + }
> +
> + #[test]
> + fn lagging_client_is_caught_up_or_resynced() {
> + let (server, path, _dir) = start(100);
> + let (mut stream, mut reader) = connect(&path);
> + subscribe(&mut stream, &mut reader, r#"{"all":"{path...}"}"#, None);
> + server.emit(1, EventKind::Write, "warm", None);
> + assert_eq!(event(&mut reader)["seq"], 1);
> +
> + // delivery interleaves with the burst, so where the cursor falls
> + // off the ring and how much still goes out depends on scheduling.
> + // Fixed are at least one resync, no event twice, and after the last
> + // resync a gapless run of events up to the newest one.
> + let big = "x".repeat(3000);
> + for seq in 2..=1201 {
> + server.emit(seq, EventKind::Write, &big, None);
> + }
> +
> + let mut events = vec![event(&mut reader)];
> + while let Some(ev) = pending(&mut reader) {
> + events.push(ev);
> + }
> + let writes: Vec<u64> = events
> + .iter()
> + .filter(|ev| ev["type"] == "write")
> + .map(|ev| ev["seq"].as_u64().unwrap())
> + .collect();
> + let resyncs: Vec<u64> = events
> + .iter()
> + .filter(|ev| ev["type"] == "resync")
> + .map(|ev| ev["seq"].as_u64().unwrap())
> + .collect();
> + assert!(!resyncs.is_empty(), "the cursor never fell off");
> + assert!(writes.len() < 1200, "delivered {}", writes.len());
> + assert!(
> + writes.windows(2).all(|w| w[0] < w[1]),
> + "an event arrived twice"
> + );
> + let last_resync = events
> + .iter()
> + .rposition(|ev| ev["type"] == "resync")
> + .unwrap();
> + let tail: Vec<u64> = events[last_resync + 1..]
> + .iter()
> + .map(|ev| ev["seq"].as_u64().unwrap())
> + .collect();
> + assert!(
> + tail.windows(2).all(|w| w[1] == w[0] + 1),
> + "gap after the last resync"
> + );
> + let newest = tail.last().or(resyncs.last()).copied();
> + assert_eq!(newest, Some(1201), "the newest event never arrived");
> +
> + server.emit(2000, EventKind::Write, "after", None);
> + let after = event(&mut reader);
> + assert_eq!(after["seq"], 2000);
> + assert_eq!(after["path"], "after");
> + assert_eq!(server.connection_count(), 1);
> + }
> +
> + #[test]
> + fn a_client_at_the_persisted_version_is_current() {
> + let (server, path, _dir) = start_at(100, 41);
> + let all = r#"{"all":"{path...}"}"#;
> +
> + let (mut stream, mut reader) = connect(&path);
> + assert_eq!(
> + request(&mut stream, &mut reader, r#"{"command":"hello"}"#),
> + "{\"ok\":{\"protocol\":1,\"seq\":41}}\n",
> + "hello tells a client where the daemon stands"
> + );
> + subscribe(&mut stream, &mut reader, all, Some(41));
> + assert!(pending(&mut reader).is_none(), "nothing to catch up on");
> + server.emit(42, EventKind::Write, "a", None);
> + assert_eq!(event(&mut reader)["seq"], 42);
> + assert!(
> + pending(&mut reader).is_none(),
> + "the marker is never delivered"
> + );
> + drop(stream);
> + drop(reader);
> +
> + for since in [40, 43] {
> + let (mut stream, mut reader) = connect(&path);
> + subscribe(&mut stream, &mut reader, all, Some(since));
> + assert_eq!(
> + event(&mut reader),
> + json!({ "seq": 42, "type": "resync" }),
> + "a number the ring cannot vouch for gets a resync"
> + );
> + drop(stream);
> + drop(reader);
> + }
> +
> + let (mut stream, mut reader) = connect(&path);
> + subscribe(&mut stream, &mut reader, all, None);
> + assert!(pending(&mut reader).is_none());
> + server.emit(43, EventKind::Write, "b", None);
> + assert_eq!(event(&mut reader)["seq"], 43);
> + }
> +
> + #[test]
> + fn since_resumes_where_the_client_left_off() {
> + let (server, path, _dir) = start(100);
> + let all = r#"{"all":"{path...}"}"#;
> +
> + // nothing recorded yet, so a resume point cannot be honored
> + let (mut stream, mut reader) = connect(&path);
> + subscribe(&mut stream, &mut reader, all, Some(1));
> + assert_eq!(event(&mut reader), json!({ "seq": 0, "type": "resync" }));
> + drop(stream);
> + drop(reader);
> +
> + let (mut stream, mut reader) = connect(&path);
> + subscribe(&mut stream, &mut reader, all, None);
> + for (seq, name) in [(1, "a"), (2, "b"), (3, "c")] {
> + server.emit(seq, EventKind::Write, name, None);
> + assert_eq!(event(&mut reader)["seq"], seq);
> + }
> + drop(stream);
> + drop(reader);
> + wait_for(|| server.connection_count() == 0);
> +
> + server.emit(4, EventKind::Write, "d", None);
> + server.emit(5, EventKind::Write, "e", None);
> +
> + let (mut stream, mut reader) = connect(&path);
> + subscribe(&mut stream, &mut reader, all, Some(3));
> + assert_eq!(event(&mut reader)["path"], "d");
> + assert_eq!(event(&mut reader)["path"], "e");
> + assert!(pending(&mut reader).is_none());
> + server.emit(6, EventKind::Write, "f", None);
> + assert_eq!(event(&mut reader)["seq"], 6);
> + drop(stream);
> + drop(reader);
> +
> + let (mut stream, mut reader) = connect(&path);
> + subscribe(&mut stream, &mut reader, all, Some(6));
> + assert!(pending(&mut reader).is_none());
> + server.emit(7, EventKind::Write, "g", None);
> + assert_eq!(event(&mut reader)["seq"], 7);
> + drop(stream);
> + drop(reader);
> +
> + let (mut stream, mut reader) = connect(&path);
> + subscribe(&mut stream, &mut reader, all, Some(99));
> + assert_eq!(event(&mut reader), json!({ "seq": 7, "type": "resync" }));
> + assert!(pending(&mut reader).is_none());
> + drop(stream);
> + drop(reader);
> +
> + for seq in 8..=200 {
> + server.emit(seq, EventKind::Write, "h", None);
> + }
> + let (mut stream, mut reader) = connect(&path);
> + subscribe(&mut stream, &mut reader, all, Some(3));
> + assert_eq!(event(&mut reader), json!({ "seq": 200, "type": "resync" }));
> + assert!(pending(&mut reader).is_none());
> + }
> +}
> diff --git a/src/rust/pmxcfs-notify/src/template.rs b/src/rust/pmxcfs-notify/src/template.rs
> new file mode 100644
> index 0000000..a7367eb
> --- /dev/null
> +++ b/src/rust/pmxcfs-notify/src/template.rs
> @@ -0,0 +1,254 @@
> +use std::collections::BTreeMap;
> +
> +use anyhow::{Error, bail};
> +
> +use crate::normalize_path;
> +
> +pub(crate) type Captures<'a> = BTreeMap<&'a str, &'a str>;
> +
> +#[derive(Debug, PartialEq, Eq)]
> +enum Component {
> + Literal(String),
> + Param {
> + prefix: String,
> + name: String,
> + suffix: String,
> + },
> + Rest(String),
> +}
> +
> +// One placeholder per component and the spanning form only at the end let
> +// a match run in a single pass over the path without backtracking. It runs for
> +// every entry, template and connection on the delivery thread.
> +#[derive(Debug, PartialEq, Eq)]
> +pub(crate) struct Template {
> + components: Vec<Component>,
> +}
> +
> +fn valid_name(name: &str) -> bool {
> + let mut chars = name.chars();
> + chars
> + .next()
> + .is_some_and(|c| c.is_ascii_alphabetic() || c == '_')
> + && chars.all(|c| c.is_ascii_alphanumeric() || c == '_')
> +}
> +
> +impl Template {
> + pub(crate) fn parse(text: &str) -> Result<Self, Error> {
Is there a reason to make this crate-private? We could just implement
`FromStr` here - as well as `Deserialize` (as mentioned above).
> + let text = normalize_path(text);
> + if text.is_empty() {
> + bail!("empty template");
> + }
> + let parts: Vec<&str> = text.split('/').collect();
> + let last = parts.len() - 1;
> + let mut components = Vec::with_capacity(parts.len());
> + let mut names: Vec<&str> = Vec::new();
> + for (index, part) in parts.into_iter().enumerate() {
> + if part.is_empty() {
> + bail!("empty path component");
> + }
> + let Some((prefix, rest)) = part.split_once('{') else {
> + if part.contains('}') {
> + bail!("malformed placeholder in '{part}'");
> + }
> + components.push(Component::Literal(part.to_owned()));
> + continue;
> + };
> + let Some((body, suffix)) = rest.split_once('}') else {
> + bail!("unterminated placeholder in '{part}'");
> + };
> + if suffix.contains('{') {
> + bail!("only one placeholder per component in '{part}'");
> + }
> + if prefix.contains('}') || body.contains('{') || suffix.contains('}') {
> + bail!("malformed placeholder in '{part}'");
> + }
> + let (name, spans) = match body.strip_suffix("...") {
> + Some(name) => (name, true),
> + None => (body, false),
> + };
> + if !valid_name(name) {
> + bail!("invalid placeholder name '{name}'");
> + }
> + if names.contains(&name) {
> + bail!("placeholder '{name}' used twice");
> + }
> + names.push(name);
> + if spans {
> + if index != last || !prefix.is_empty() || !suffix.is_empty() {
> + bail!("'{{{name}...}}' must be the whole last component");
> + }
> + components.push(Component::Rest(name.to_owned()));
> + } else {
> + components.push(Component::Param {
> + prefix: prefix.to_owned(),
> + name: name.to_owned(),
> + suffix: suffix.to_owned(),
> + });
> + }
> + }
> + Ok(Self { components })
> + }
> +
> + pub(crate) fn matches<'a>(&'a self, path: &'a str) -> Option<Captures<'a>> {
> + let mut captures = Captures::new();
> + let mut rest = path;
> + let last = self.components.len() - 1;
> + for (index, component) in self.components.iter().enumerate() {
> + if let Component::Rest(name) = component {
> + if rest.is_empty() {
> + return None;
> + }
> + captures.insert(name.as_str(), rest);
> + return Some(captures);
> + }
> + let (head, tail) = match rest.split_once('/') {
> + Some((head, tail)) => (head, Some(tail)),
> + None => (rest, None),
> + };
> + match component {
> + Component::Literal(literal) if head == literal.as_str() => {}
> + Component::Param {
> + prefix,
> + name,
> + suffix,
> + } => {
> + let value = head
> + .strip_prefix(prefix.as_str())?
> + .strip_suffix(suffix.as_str())?;
> + if value.is_empty() {
> + return None;
> + }
> + captures.insert(name.as_str(), value);
> + }
> + _ => return None,
> + }
> + match (tail, index == last) {
> + (None, true) => return Some(captures),
> + (Some(tail), false) => rest = tail,
> + _ => return None,
> + }
> + }
> + None
> + }
> +}
> +
> +#[cfg(test)]
> +mod tests {
> + use super::*;
> +
> + fn captures<'a>(pairs: &[(&'a str, &'a str)]) -> Captures<'a> {
> + pairs.iter().copied().collect()
> + }
> +
> + #[test]
> + fn parses_templates() {
> + let literal = |text: &str| Component::Literal(text.to_owned());
> + let param = |prefix: &str, name: &str, suffix: &str| Component::Param {
> + prefix: prefix.to_owned(),
> + name: name.to_owned(),
> + suffix: suffix.to_owned(),
> + };
> + assert_eq!(
> + Template::parse("/datacenter.cfg").unwrap(),
> + Template {
> + components: vec![literal("datacenter.cfg")]
> + }
> + );
> + assert_eq!(
> + Template::parse("nodes/{node}/qemu-server/{vmid}.conf").unwrap(),
> + Template {
> + components: vec![
> + literal("nodes"),
> + param("", "node", ""),
> + literal("qemu-server"),
> + param("", "vmid", ".conf"),
> + ]
> + }
> + );
> + assert_eq!(
> + Template::parse("{path...}").unwrap(),
> + Template {
> + components: vec![Component::Rest("path".to_owned())]
> + }
> + );
> + assert_eq!(
> + Template::parse("priv/lock/{path...}").unwrap(),
> + Template {
> + components: vec![
> + literal("priv"),
> + literal("lock"),
> + Component::Rest("path".to_owned())
> + ]
> + }
> + );
> + }
> +
> + #[test]
> + fn rejects_bad_templates() {
> + for (text, message) in [
> + ("", "empty template"),
> + ("/", "empty template"),
> + ("a//b", "empty path component"),
> + ("a/", "empty path component"),
> + ("{a}{b}", "only one placeholder per component"),
> + ("{a", "unterminated placeholder"),
> + ("a}", "malformed placeholder"),
> + ("{a}}", "malformed placeholder"),
> + ("{a{b}", "malformed placeholder"),
> + ("{}", "invalid placeholder name ''"),
> + ("{1a}", "invalid placeholder name '1a'"),
> + ("{a-b}", "invalid placeholder name 'a-b'"),
> + ("a/{x}/{x}", "placeholder 'x' used twice"),
> + ("{a...}/b", "must be the whole last component"),
> + ("x{a...}", "must be the whole last component"),
> + ("{a...}.conf", "must be the whole last component"),
> + ] {
> + let err = Template::parse(text).unwrap_err().to_string();
> + assert!(err.contains(message), "{text:?}: {err}");
> + }
> + }
> +
> + #[test]
> + fn matches_paths() {
> + let guest = Template::parse("nodes/{node}/qemu-server/{vmid}.conf").unwrap();
> + assert_eq!(
> + guest.matches("nodes/n1/qemu-server/100.conf"),
> + Some(captures(&[("node", "n1"), ("vmid", "100")]))
> + );
> + assert_eq!(guest.matches("nodes/n1/qemu-server/100.conf.tmp.1"), None);
> + assert_eq!(guest.matches("nodes/n1/qemu-server/.conf"), None);
> + assert_eq!(guest.matches("nodes/n1/lxc/100.conf"), None);
> + assert_eq!(guest.matches("nodes/n1/qemu-server"), None);
> + assert_eq!(guest.matches("nodes/n1/qemu-server/100.conf/x"), None);
> + assert_eq!(guest.matches("nodes//qemu-server/100.conf"), None);
> + assert_eq!(guest.matches(""), None);
> +
> + let host = Template::parse("nodes/{node}/host-{id}.cfg").unwrap();
> + assert_eq!(
> + host.matches("nodes/n1/host-7.cfg"),
> + Some(captures(&[("node", "n1"), ("id", "7")]))
> + );
> + assert_eq!(host.matches("nodes/n1/host-.cfg"), None);
> + assert_eq!(host.matches("nodes/n1/host-7.cfgx"), None);
> +
> + let exact = Template::parse("datacenter.cfg").unwrap();
> + assert_eq!(exact.matches("datacenter.cfg"), Some(captures(&[])));
> + assert_eq!(exact.matches("datacenter.cfg/x"), None);
> + assert_eq!(exact.matches("xdatacenter.cfg"), None);
> +
> + let all = Template::parse("{path...}").unwrap();
> + assert_eq!(all.matches("a"), Some(captures(&[("path", "a")])));
> + assert_eq!(all.matches("a/b/c"), Some(captures(&[("path", "a/b/c")])));
> + assert_eq!(all.matches(""), None);
> +
> + let locks = Template::parse("priv/lock/{path...}").unwrap();
> + assert_eq!(
> + locks.matches("priv/lock/file-storage_cfg/a"),
> + Some(captures(&[("path", "file-storage_cfg/a")]))
> + );
> + assert_eq!(locks.matches("priv/lock"), None);
> + assert_eq!(locks.matches("priv/lock/"), None);
> + assert_eq!(locks.matches("priv/locks/a"), None);
> + }
> +}
> --
> 2.47.3
>
>
>
>
>
--
next prev parent reply other threads:[~2026-09-25 15:27 UTC|newest]
Thread overview: 20+ messages / expand[flat|nested] mbox.gz Atom feed top
2026-09-18 14:41 [RFC cluster/manager 00/10] pmxcfs: add a change notification socket Hannes Laimer
2026-09-18 14:41 ` [PATCH pve-cluster 01/10] buildsys: add rust workspace under src/rust Hannes Laimer
2026-09-25 12:48 ` Wolfgang Bumiller
2026-09-18 14:41 ` [PATCH pve-cluster 02/10] rust: notify: add change notification socket server Hannes Laimer
2026-09-25 15:18 ` Wolfgang Bumiller [this message]
2026-09-18 14:41 ` [PATCH pve-cluster 03/10] rust: ffi: add C ABI staticlib for pmxcfs Hannes Laimer
2026-09-21 9:57 ` Robert Obkircher
2026-09-18 14:41 ` [PATCH pve-cluster 04/10] pmxcfs: memdb: add change notification hook Hannes Laimer
2026-09-18 14:41 ` [PATCH pve-cluster 05/10] buildsys: link pmxcfs against the rust notify staticlib Hannes Laimer
2026-09-18 14:41 ` [PATCH pve-cluster 06/10] pmxcfs: notify: emit change events over the notification socket Hannes Laimer
2026-09-18 14:41 ` [PATCH pve-cluster 07/10] cfs: add perl client for the change " Hannes Laimer
2026-09-18 14:41 ` [PATCH pve-cluster 08/10] cfs: add hook registry for change notification consumers Hannes Laimer
2026-09-18 14:41 ` [PATCH pve-manager 09/10] hooks: add runner executing cluster change hooks in children Hannes Laimer
2026-09-18 14:41 ` [PATCH pve-manager 10/10] pvescheduler: run cluster change hooks from a listener child Hannes Laimer
2026-09-25 12:01 ` [RFC cluster/manager 00/10] pmxcfs: add a change notification socket Wolfgang Bumiller
2026-09-25 12:18 ` Hannes Laimer
2026-09-25 12:38 ` Wolfgang Bumiller
2026-09-25 12:23 ` Hannes Laimer
2026-09-25 12:43 ` Wolfgang Bumiller
2026-09-25 12:48 ` Hannes Laimer
Reply instructions:
You may reply publicly to this message via plain-text email
using any one of the following methods:
* Save the following mbox file, import it into your mail client,
and reply-to-all from there: mbox
Avoid top-posting and favor interleaved quoting:
https://en.wikipedia.org/wiki/Posting_style#Interleaved_style
* Reply using the --to, --cc, and --in-reply-to
switches of git-send-email(1):
git send-email \
--in-reply-to=6k5xskft7jh6a4vz7eb2673pmlq2t6dbgxog27kzx47y5jqdmc@6a4oe6hae4p7 \
--to=w.bumiller@proxmox.com \
--cc=h.laimer@proxmox.com \
--cc=pve-devel@lists.proxmox.com \
/path/to/YOUR_REPLY
https://kernel.org/pub/software/scm/git/docs/git-send-email.html
* If your mail client supports setting the In-Reply-To header
via mailto: links, try the mailto: link
Be sure your reply has a Subject: header at the top and a blank line
before the message body.
This is a public inbox, see mirroring instructions
for how to clone and mirror all data and code used for this inbox