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 07/10] cfs: add perl client for the change notification socket
Date: Thu, 1 Oct 2026 14:11:26 +0200 [thread overview]
Message-ID: <ubhltyu2gsrmvoxt4azdqtmyjxterzybz3zfhjzm2epjpj2b7r@hkpbnviz3ejx> (raw)
In-Reply-To: <20260918144152.575163-8-h.laimer@proxmox.com>
On Fri, Sep 18, 2026 at 04:41:49PM +0200, Hannes Laimer wrote:
> Daemons that want to react to cluster wide config changes all need to
> connect to the pmxcfs notification socket, subscribe with path
> templates, wait for events with a deadline, and recover when pmxcfs
> restarts. A shared client module handles that, so every consumer gets
> the reconnect and resync handling right instead of reimplementing it.
>
> The client resumes from the last sequence number it has, and only does
> a full resync when it has none or the daemon can no longer replay from
> it.
>
> A refused handshake, an error reply or a protocol mismatch is reported
> to the caller instead of retried, since a retry would not change any of
> them, while a lost connection is retried with backoff. A forked child
> opens its own connection instead of using the one it inherited, as the
> IPC client does, since two readers on one stream would corrupt the line
> framing.
>
> Signed-off-by: Hannes Laimer <h.laimer@proxmox.com>
> ---
> debian/pve-cluster.install | 1 +
> src/PVE/Cluster/Makefile | 2 +-
> src/PVE/Cluster/Watch.pm | 371 ++++++++++++++++++++++++++++++++++
> src/test/Makefile | 6 +-
> src/test/watch_client_test.pl | 294 +++++++++++++++++++++++++++
> 5 files changed, 672 insertions(+), 2 deletions(-)
> create mode 100644 src/PVE/Cluster/Watch.pm
> create mode 100644 src/test/watch_client_test.pl
>
> diff --git a/debian/pve-cluster.install b/debian/pve-cluster.install
> index f66cd06..77e3244 100644
> --- a/debian/pve-cluster.install
> +++ b/debian/pve-cluster.install
> @@ -5,4 +5,5 @@ usr/lib/
> usr/share/man/man8/pmxcfs.8
> usr/share/perl5/PVE/Cluster.pm
> usr/share/perl5/PVE/Cluster/IPCConst.pm
> +usr/share/perl5/PVE/Cluster/Watch.pm
> usr/share/perl5/PVE/IPCC.pm
> diff --git a/src/PVE/Cluster/Makefile b/src/PVE/Cluster/Makefile
> index 3f920cb..7beb976 100644
> --- a/src/PVE/Cluster/Makefile
> +++ b/src/PVE/Cluster/Makefile
> @@ -1,6 +1,6 @@
> PVEDIR=$(DESTDIR)/usr/share/perl5/PVE
>
> -SOURCES=IPCConst.pm Setup.pm
> +SOURCES=IPCConst.pm Setup.pm Watch.pm
>
> .PHONY: install
> install: $(SOURCES)
> diff --git a/src/PVE/Cluster/Watch.pm b/src/PVE/Cluster/Watch.pm
> new file mode 100644
> index 0000000..a41860b
> --- /dev/null
> +++ b/src/PVE/Cluster/Watch.pm
> @@ -0,0 +1,371 @@
> +package PVE::Cluster::Watch;
> +
> +use strict;
> +use warnings;
> +
> +use IO::Select;
> +use IO::Socket::UNIX;
> +use JSON;
> +use Socket qw(SOCK_STREAM MSG_NOSIGNAL);
> +use Time::HiRes qw(time sleep);
> +
> +use PVE::Tools;
> +
> +my $default_socket = '/run/pve-cluster/pmxcfs.sock';
> +my $protocol_version = 1;
> +my $reply_timeout = 10;
> +my $min_backoff = 1;
> +my $max_backoff = 30;
> +
> +sub new {
> + my ($class, %param) = @_;
> +
> + my $self = bless {
> + socket_path => $param{socket} // $default_socket,
> + patterns => {},
> + sock => undef,
> + buf => '',
> + queue => [],
> + backoff => $min_backoff,
> + next_connect => 0,
> + warned => 0,
> + last_seq => undef,
> + state => $param{state},
> + key => $param{key},
> + pid => $$,
> + }, $class;
> +
> + # resume point from the caller, or a previous process's checkpoint,
> + # numified since the daemon takes a JSON number and refuses a string
> + if (defined($param{since})) {
> + $self->{last_seq} = 0 + $param{since};
> + } elsif (defined($self->{state}) && -e $self->{state}) {
> + my $line = PVE::Tools::file_read_firstline($self->{state}) // '';
> + # a recorded point stands only for the key it was recorded under,
> + # anything else is as good as no state file
If we end up keeping this, maybe document this mechanism in this file &
commit message. I wouldn't really see a reason to have this without the
next commit (and maybe it should actually be handled there instead).
> + if (my ($seq, $key) = $line =~ m/^(\d+)(?:\s+(\S+))?$/) {
> + $self->{last_seq} = 0 + $seq if ($key // '') eq ($self->{key} // '');
> + }
> + }
> +
> + return $self;
> +}
> +
> +# A connection inherited across a fork belongs to the parent, so the child
> +# drops its copy without a shutdown and connects on its own, as IPCC does.
Mhh, worker tasks should probably get a way to have this closed eagerly,
since otherwise they also needlessly keep the connection open/connected.
But we'd probably need some more 'atfork' handling for this.
> +sub check_fork {
> + my ($self) = @_;
> +
> + return if $self->{pid} == $$;
> +
> + CORE::close($self->{sock}) if $self->{sock};
> + $self->{sock} = undef;
> + $self->{buf} = '';
> + $self->{queue} = [];
> + $self->{pid} = $$;
> + $self->{next_connect} = time();
> + $self->{backoff} = $min_backoff;
> + $self->{warned} = 0;
> + $self->{last_seq} = undef;
> + $self->{state} = undef;
^ could factor out this initialization for reuse in new()
> +
> + return;
> +}
> +
> +# Takes effect when connected and is replayed on every reconnect. A refusal
> +# is raised, a dead connection is dropped and comes back through the next wait.
We usually just say "dies if …" - it's perl ;-)
> +sub subscribe {
> + my ($self, %param) = @_;
> +
> + $self->check_fork();
> +
> + my $previous = $self->{patterns};
> + $self->{patterns} = { ($param{patterns} // {})->%* };
> +
> + return if !$self->{sock};
> +
> + my $refusal = eval { $self->send_subscription() };
^ That's an odd name to suddenly use here.
> + if (my $err = $@) {
> + chomp $err;
> + $self->disconnect($err) if $self->{sock};
> + return;
> + }
> + # the daemon keeps serving the previous subscription, so keep
> + # describing that one
> + if (defined($refusal)) {
> + $self->{patterns} = $previous;
> + die "subscription refused: $refusal\n";
> + }
> +
> + return;
> +}
> +
> +# The number the caller has processed up to, so a restart resumes there.
> +sub checkpoint {
> + my ($self, $seq) = @_;
> +
> + $self->check_fork();
> +
> + return if !defined($self->{state});
> +
> + if (!defined($seq)) {
> + unlink($self->{state});
> + return;
> + }
> +
> + my $key = defined($self->{key}) ? " $self->{key}" : '';
> + PVE::Tools::file_set_contents($self->{state}, "$seq$key\n");
> +
> + return;
> +}
> +
> +sub connected {
> + my ($self) = @_;
> +
> + $self->check_fork();
> +
> + return defined($self->{sock});
> +}
> +
> +# The descriptor changes on every reconnect and is undef while disconnected,
> +# so a select-loop caller refreshes it after each wait, and only wait reconnects.
> +sub fd {
> + my ($self) = @_;
> +
> + $self->check_fork();
> +
> + return $self->{sock} ? fileno($self->{sock}) : undef;
> +}
> +
> +sub close {
> + my ($self) = @_;
> +
> + $self->check_fork();
> +
> + $self->{sock}->close() if $self->{sock};
> + $self->{sock} = undef;
> + $self->{established} = 0;
> + $self->{buf} = '';
> +
> + return;
> +}
> +
> +# Returns the events that arrived within $timeout seconds, or none. Reconnect
> +# with handshake and reply timeout happens in here too, the one way a zero
> +# timeout can block. A first connect yields a resync, and a refused
> +# handshake dies, since it will not change on retry.
> +sub wait {
> + my ($self, $timeout) = @_;
> +
> + $self->check_fork();
> +
> + my $deadline = time() + ($timeout // 0);
> +
> + while (1) {
> + if (scalar($self->{queue}->@*)) {
> + my @events = $self->{queue}->@*;
> + $self->{queue} = [];
> + return @events;
> + }
> +
> + my $now = time();
> + if (!$self->{sock}) {
> + if ($now >= $self->{next_connect}) {
> + $self->connect();
> + next;
> + }
> + my $until = $self->{next_connect} < $deadline ? $self->{next_connect} : $deadline;
> + return if $until <= $now;
> + sleep($until - $now);
> + # a signal cut the sleep short, and a caller whose handler only
> + # sets a flag has to get back control to act on it
> + return if time() < $until;
> + next;
> + }
> +
> + my $remaining = $deadline - $now;
> + $remaining = 0 if $remaining < 0;
> +
> + return if !IO::Select->new($self->{sock})->can_read($remaining);
> + $self->read_messages();
> + }
> +}
> +
> +sub connect {
> + my ($self) = @_;
> +
> + # a daemon that cannot accept for a while leaves a blocking connect in
> + # its backlog, so the attempt carries a timeout and a failure goes to retry
> + my $sock = IO::Socket::UNIX->new(
> + Type => SOCK_STREAM,
> + Peer => $self->{socket_path},
> + Timeout => $reply_timeout,
> + );
> + if (!$sock) {
> + $self->connect_failed("connect to $self->{socket_path} failed: $!");
> + return;
> + }
> +
> + $self->{sock} = $sock;
> + $self->{buf} = '';
> + $self->{established} = 0;
> +
> + my $resume = defined($self->{last_seq});
> + my ($refusal, $hello_seq) = eval { $self->handshake() };
> + my $err = $@;
> + # a number learned in a handshake that did not go through is no resume
> + # point, the first connection that does still owes the caller a resync
> + $self->{last_seq} = undef if !$resume && ($err || defined($refusal));
> + if ($err) {
> + $self->close();
> + $self->connect_failed("handshake with $self->{socket_path} failed: $err");
> + return;
> + }
> + if (defined($refusal)) {
> + $self->close();
> + $self->delay_retry();
> + die "handshake with $self->{socket_path} refused: $refusal\n";
> + }
> +
> + $self->{established} = 1;
> + $self->{backoff} = $min_backoff;
> + $self->{warned} = 0;
> + # the resync carries the number the daemon handed out, so a
> + # consumer can record it once the resync is processed
> + unshift $self->{queue}->@*,
> + { type => 'resync', (defined($hello_seq) ? (seq => 0 + $hello_seq) : ()) }
> + if !$resume;
> +
> + return;
> +}
> +
> +sub handshake {
> + my ($self) = @_;
> +
> + my ($hello, $error) = $self->request('hello');
> + return $error if defined($error);
> +
> + my $version = ref($hello) eq 'HASH' ? $hello->{protocol} : undef;
> + return "unsupported protocol version " . ($version // 'unknown')
> + if ($version // -1) != $protocol_version;
> +
> + my $refusal = $self->send_subscription();
> +
> + # a client that saw no event yet still learns where the daemon stands, so
> + # a later reconnect resumes there instead of resyncing. The subscription
> + # goes out first, so a first connect still owes and gets its resync
> + $self->{last_seq} //= 0 + $hello->{seq} if defined($hello->{seq});
> +
> + return ($refusal, $hello->{seq});
> +}
> +
> +sub connect_failed {
> + my ($self, $msg) = @_;
> +
> + if (!$self->{warned}) {
> + chomp $msg;
> + warn "$msg, retrying\n";
> + $self->{warned} = 1;
> + }
> +
> + $self->delay_retry();
> +
> + return;
> +}
> +
> +sub delay_retry {
> + my ($self) = @_;
> +
> + $self->{next_connect} = time() + $self->{backoff};
> + $self->{backoff} *= 2;
> + $self->{backoff} = $max_backoff if $self->{backoff} > $max_backoff;
> +
> + return;
> +}
> +
> +# A connection lost before its handshake completed is the connect attempt's
> +# failure to report, so a daemon that accepts and closes does not warn twice.
> +sub disconnect {
> + my ($self, $reason) = @_;
> +
> + my $established = $self->{established};
> + $self->close();
> + return if !$established;
> +
> + warn "notification socket $reason, reconnecting\n";
> + $self->{next_connect} = time();
> + $self->{backoff} = $min_backoff;
> +
> + return;
> +}
> +
> +sub send_subscription {
> + my ($self) = @_;
> +
> + my $args = { patterns => $self->{patterns} };
> + $args->{since} = $self->{last_seq} if defined($self->{last_seq});
> + my (undef, $error) = $self->request('subscribe', $args);
> +
> + return $error;
> +}
> +
> +# Sends a request and returns its reply as (data, error), queueing any events
> +# that arrive meanwhile. Dies when no reply comes at all, meaning the
> +# connection is gone.
> +sub request {
> + my ($self, $command, $args) = @_;
> +
> + my $line = encode_json({ command => $command, ($args ? (args => $args) : ()) }) . "\n";
> + my $written = 0;
> + while ($written < length($line)) {
> + my $len = send($self->{sock}, substr($line, $written), MSG_NOSIGNAL);
> + next if !defined($len) && $!{EINTR};
> + die "write failed: $!\n" if !defined($len);
> + $written += $len;
> + }
> +
> + my $deadline = time() + $reply_timeout;
> + while (1) {
> + my $remaining = $deadline - time();
> + die "timeout waiting for reply to '$command'\n" if $remaining <= 0;
> + next if !IO::Select->new($self->{sock})->can_read($remaining);
^ Why this? read_messages() does a blocking read, after all. But it only
does one, so it may only have half a response.
> +
> + for my $msg ($self->read_messages()) {
> + return ($msg->{ok}, undef) if exists $msg->{ok};
> + return (undef, $msg->{error} // 'unknown error') if exists $msg->{error};
> + }
> + die "connection lost while waiting for reply to '$command'\n" if !$self->{sock};
> + }
> +}
> +
> +sub read_messages {
> + my ($self) = @_;
> +
> + my $len;
> + do {
> + $len = sysread($self->{sock}, $self->{buf}, 65536, length($self->{buf}));
When using sysread, this should append to an existing buffer, since it
can end with a partial line.
> + } while (!defined($len) && $!{EINTR});
> + if (!$len) {
> + $self->disconnect(defined($len) ? 'closed' : "read failed: $!");
> + return;
> + }
> +
> + my @replies;
> + while ($self->{buf} =~ s/^([^\n]*)\n//) {
> + my $msg = eval { decode_json($1) };
> + if (ref($msg) ne 'HASH') {
> + $self->disconnect('sent a malformed line');
> + return;
> + }
> + if (my $event = $msg->{event}) {
> + push $self->{queue}->@*, $event;
> + $self->{last_seq} = $event->{seq} if defined($event->{seq});
While I have already noted that I think we should drop this entire
sequence handling, still, as a review note here: We only use it for the
handshake, and it's the job of the user of this package to provide the
`checkpoint()` sub with a sequence number, so tracking `last_seq` beyond
the initial handshake is probably not worth it.
> + } else {
> + push @replies, $msg;
> + }
> + }
> +
> + return @replies;
> +}
> +
> +1;
> diff --git a/src/test/Makefile b/src/test/Makefile
> index cdd37d0..7b36fca 100644
> --- a/src/test/Makefile
> +++ b/src/test/Makefile
> @@ -4,7 +4,7 @@ cpgtest: cpgtest.c
> gcc -Wall cpgtest.c $(shell pkg-config --cflags --libs libcpg libqb) -o cpgtest
>
> .PHONY: check install clean distclean
> -check: corosync-parser-test test-mac-prefix
> +check: corosync-parser-test test-mac-prefix watch-client-test
>
> .PHONY: corosync-parser-test
> corosync-parser-test:
> @@ -14,5 +14,9 @@ corosync-parser-test:
> test-mac-prefix:
> perl test_mac_prefix.pl
>
> +.PHONY: watch-client-test
> +watch-client-test:
> + perl watch_client_test.pl
> +
> distclean: clean
> clean:
> diff --git a/src/test/watch_client_test.pl b/src/test/watch_client_test.pl
> new file mode 100644
> index 0000000..3ad3aa5
> --- /dev/null
> +++ b/src/test/watch_client_test.pl
> @@ -0,0 +1,294 @@
> +#!/usr/bin/perl
> +
> +use lib '..';
> +
> +use strict;
> +use warnings;
> +
> +use File::Temp qw(tempdir);
> +use IO::Socket::UNIX;
> +use JSON;
> +use Socket qw(SOCK_STREAM);
> +use Test::More;
> +use Time::HiRes qw(sleep time);
> +
> +use PVE::Cluster::Watch;
> +
> +my $dir = tempdir(CLEANUP => 1);
> +my $path = "$dir/sock";
> +
> +my $listener = IO::Socket::UNIX->new(Type => SOCK_STREAM, Local => $path, Listen => 1)
> + or die "listen failed: $!\n";
> +
> +# A stand-in for pmxcfs that serves fourteen connections in a row. The
> +# fourth refuses the subscription, the ninth hangs up after answering the
> +# hello, the eleventh refuses a second subscription on a live connection
> +# and the last three hang up before answering anything. The others answer
> +# the handshake, mirror the subscription back inside two events and hang
> +# up, which forces the client through its reconnect path. Like the daemon
> +# it takes a resume point only as a number.
> +my $server = fork() // die "fork failed: $!\n";
> +if (!$server) {
> + for my $round (1 .. 14) {
> + my $conn = $listener->accept() or die "accept failed: $!\n";
> + if ($round >= 12) {
> + close($conn);
> + next;
> + }
> + my $subscribed = 0;
> + while (my $line = <$conn>) {
> + my $req = decode_json($line);
> + if ($req->{command} eq 'hello') {
> + print $conn encode_json({ ok => { protocol => 1, seq => 3 } }), "\n";
> + } elsif ($round == 9) {
> + last;
> + } elsif ($req->{command} eq 'subscribe') {
> + if ($round == 4) {
> + print $conn encode_json({ error => "pattern 'guest': no good" }), "\n";
> + last;
> + }
> + if ($round == 11) {
> + if (!$subscribed++) {
> + print $conn encode_json({ ok => undef }), "\n";
> + next;
> + }
> + print $conn encode_json({ error => "pattern 'bad': no good" }), "\n";
> + last;
> + }
> + my $patterns = $req->{args}->{patterns};
> + my $since = $req->{args}->{since};
> + if (defined($since) && encode_json([$since]) !~ /^\[\d+\]$/) {
> + print $conn encode_json({ error => "'since' must be an unsigned integer" }),
> + "\n";
> + last;
> + }
> + print $conn encode_json({ ok => undef }), "\n";
> + print $conn encode_json({
> + event => {
> + seq => 1,
> + type => 'write',
> + path => 'x',
> + params =>
> + { map { $_ => [{ template => $patterns->{$_} }] } keys %$patterns },
> + },
> + }),
> + "\n";
> + print $conn encode_json({
> + event => {
> + seq => 2,
> + type => 'rename',
> + path => 'old',
> + to => join(',', sort keys %$patterns),
> + (defined($since) ? (since => $since) : ()),
> + },
> + }),
> + "\n";
> + last;
> + } else {
> + print $conn encode_json({ error => "unknown command '$req->{command}'" }), "\n";
> + }
> + }
> + close($conn);
> + }
> + exit(0);
> +}
> +close($listener);
> +
> +my $collect = sub {
> + my ($watch, $count, $timeout) = @_;
> + my @events;
> + my $deadline = time() + 10;
> + while (scalar(@events) < $count && time() < $deadline) {
> + push @events, $watch->wait($timeout);
> + sleep(0.05) if !$timeout;
> + }
> + return \@events;
> +};
> +
> +my $watch = PVE::Cluster::Watch->new(socket => $path);
> +ok(!$watch->connected(), 'not connected before the first wait');
> +
> +my $template = 'nodes/{node}/qemu-server/{vmid}.conf';
> +$watch->subscribe(patterns => { guest => $template });
> +
> +my $write = {
> + seq => 1,
> + type => 'write',
> + path => 'x',
> + params => { guest => [{ template => $template }] },
> +};
> +my $rename = { seq => 2, type => 'rename', path => 'old', to => 'guest' };
> +# a fresh client has nothing to resume from, so it subscribes without one
> +my $expected = [{ type => 'resync', seq => 3 }, $write, $rename];
> +my $resumed = [$write, { $rename->%*, since => 2 }];
> +
> +is_deeply($collect->($watch, 3, 1), $expected, 'connect gives a resync followed by the events');
> +ok(defined($watch->fd()), 'a descriptor is available while connected');
> +
> +$watch->{pid} = -1; # pretend the object was inherited across a fork
> +is_deeply(
> + $collect->($watch, 3, 0),
> + $expected,
> + 'a forked child gets a connection of its own, a zero timeout drains it',
> +);
> +
> +is_deeply(
> + $collect->($watch, 2, 1),
> + $resumed,
> + 'reconnect after hangup resumes after the last event and replays the subscription',
> +);
> +
> +# the server has hung up on that connection by now, writing to it must
> +# neither raise SIGPIPE nor die, the next wait reconnects
> +sleep(0.2);
> +$watch->subscribe(patterns => { guest => $template });
> +ok(!$watch->connected(), 'a subscription on a dead connection drops it');
> +
> +eval { $watch->wait(1) };
> +like($@, qr/refused: pattern 'guest': no good/, 'a refused subscription is raised, not retried');
> +ok(!$watch->connected(), 'and its connection is closed');
> +
> +# a client resuming from a state file, served by rounds five and six
> +my $state = "$dir/state";
> +open(my $fh, '>', $state) or die "open: $!";
> +print $fh "5\n";
> +close($fh);
> +my $resuming = PVE::Cluster::Watch->new(socket => $path, state => $state);
> +$resuming->subscribe(patterns => { guest => $template });
> +is_deeply(
> + $collect->($resuming, 2, 1),
> + [$write, { $rename->%*, since => 5 }],
> + 'a state file provides the resume point and replaces the first resync',
> +);
> +is_deeply(
> + $collect->($resuming, 2, 1),
> + [$write, { $rename->%*, since => 2 }],
> + 'from then on the last event received is the resume point',
> +);
> +$resuming->checkpoint(2);
> +open($fh, '<', $state) or die "open: $!";
> +my $stored = <$fh>;
> +close($fh);
> +is($stored, "2\n", 'checkpoint writes the state file');
> +is(PVE::Cluster::Watch->new(since => 7)->{last_seq}, 7, 'an explicit resume point is taken as is');
> +ok(eval { $watch->checkpoint(1); 1 }, 'checkpoint without a state file does nothing');
> +
> +# a resume point bound to a key, served by rounds seven and eight
> +my $key = 'a1b2c3';
> +open($fh, '>', $state) or die "open: $!";
> +print $fh "5 $key\n";
> +close($fh);
> +my $keyed = PVE::Cluster::Watch->new(socket => $path, state => $state, key => $key);
> +$keyed->subscribe(patterns => { guest => $template });
> +is_deeply(
> + $collect->($keyed, 2, 1),
> + [$write, { $rename->%*, since => 5 }],
> + 'a state file written under the same key provides the resume point',
> +);
> +$keyed->checkpoint(2);
> +open($fh, '<', $state) or die "open: $!";
> +$stored = <$fh>;
> +close($fh);
> +is($stored, "2 $key\n", 'which a checkpoint records next to the number');
> +
> +open($fh, '>', $state) or die "open: $!";
> +print $fh "5 d4e5f6\n";
> +close($fh);
> +my $rekeyed = PVE::Cluster::Watch->new(socket => $path, state => $state, key => $key);
> +$rekeyed->subscribe(patterns => { guest => $template });
> +is_deeply(
> + $collect->($rekeyed, 3, 1),
> + $expected,
> + 'a state file written under another key is no resume point',
> +);
> +
> +# a handshake cut after the hello, rounds nine and ten. The number the
> +# first attempt learned must not pass as a resume point on the second.
> +{
> + my @warnings;
> + local $SIG{__WARN__} = sub { push @warnings, $_[0] };
> + my $half = PVE::Cluster::Watch->new(socket => $path);
> + $half->subscribe(patterns => { guest => $template });
> + is_deeply(
> + $collect->($half, 3, 1),
> + $expected,
> + 'a first connect that is cut after the hello still gives a resync',
> + );
> + is(scalar(@warnings), 1, 'and the cut handshake is reported');
> +}
> +
> +# a refused subscription on a live connection, round eleven. The daemon keeps
> +# serving the previous one, so the client keeps describing that one.
> +{
> + my $live = PVE::Cluster::Watch->new(socket => $path);
> + $live->subscribe(patterns => { guest => $template });
> + is_deeply(
> + $collect->($live, 1, 1),
> + [{ type => 'resync', seq => 3 }],
> + 'connected with the first subscription',
> + );
> + eval { $live->subscribe(patterns => { bad => 'x/{y}' }) };
> + like($@, qr/subscription refused: pattern 'bad': no good/, 'a refusal is raised');
> + is_deeply($live->{patterns}, { guest => $template }, 'and the previous patterns stay');
> + ok($live->connected(), 'on a connection that stays up');
> +}
> +
> +# a daemon that accepts and hangs up immediately, rounds twelve to fourteen
> +{
> + my @warnings;
> + local $SIG{__WARN__} = sub { push @warnings, $_[0] };
> + my $cut = PVE::Cluster::Watch->new(socket => $path);
> + $cut->subscribe(patterns => { guest => $template });
> + my @none = $cut->wait(5);
> + is(scalar(@none), 0, 'nothing arrives from a daemon that hangs up immediately');
> + is(scalar(@warnings), 1, 'which is reported once and retried with backoff');
> + ok(!$cut->connected(), 'and leaves the client disconnected');
> +}
> +
> +# a signal that cuts a wait short, here while the client is waiting for its
> +# next connect attempt. A caller whose handler only raises a flag has to get
> +# back control to act on it, rather than at the end of the wait.
> +{
> + my @warnings;
> + local $SIG{__WARN__} = sub { push @warnings, $_[0] };
> + my $flagged = 0;
> + local $SIG{USR1} = sub { $flagged = 1 };
> +
> + my $parent = $$;
> + my $signaller = fork() // die "fork failed: $!\n";
> + if (!$signaller) {
> + sleep(0.3);
> + kill 'USR1', $parent;
> + exit(0);
> + }
> +
> + my $gone = PVE::Cluster::Watch->new(socket => "$dir/gone");
> + my $started = time();
> + my @none = $gone->wait(3);
> + my $elapsed = time() - $started;
> + waitpid($signaller, 0);
> +
> + is(scalar(@none), 0, 'a wait cut short by a signal returns no events');
> + ok($flagged, 'the handler of the caller has run');
> + # a signal that lands after the first backoff sleep ended only cuts the
> + # second one short, so the bound has to sit above one backoff and still
> + # below the deadline
> + ok(
> + $elapsed < 2.5,
> + sprintf('and the wait gave up after %.2fs, short of its deadline', $elapsed),
> + );
> + is(scalar(@warnings), 1, 'the socket that is not there is reported once');
> +}
> +
> +# a shortfall in connections must fail the test, not hang it
> +local $SIG{ALRM} = sub { kill 'KILL', $server; die "fake server did not finish\n" };
> +alarm(30);
> +waitpid($server, 0);
> +alarm(0);
> +is($?, 0, 'fake server ran all rounds cleanly');
> +
> +my @none = $watch->wait(1);
> +is(scalar(@none), 0, 'no events while the socket is gone');
> +ok(!$watch->connected(), 'disconnected after the server went away');
> +
> +done_testing();
> --
> 2.47.3
>
>
>
>
>
--
next prev parent reply other threads:[~2026-10-01 12:11 UTC|newest]
Thread overview: 25+ 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
2026-09-29 15:50 ` Wolfgang Bumiller
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-29 15:07 ` Wolfgang Bumiller
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-10-01 12:11 ` Wolfgang Bumiller [this message]
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
2026-10-01 7:31 ` Wolfgang Bumiller
2026-10-01 7:43 ` 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=ubhltyu2gsrmvoxt4azdqtmyjxterzybz3zfhjzm2epjpj2b7r@hkpbnviz3ejx \
--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