From mboxrd@z Thu Jan 1 00:00:00 1970 Return-Path: Received: from gate001.proxmox.com (gate001.proxmox.com [IPv6:2a0f:8001:1:32::40]) by lore.proxmox.com (Postfix) with ESMTPS id 847FC1FF0A3 for ; Thu, 01 Oct 2026 14:11:39 +0200 (CEST) Received: from gate001.proxmox.com (localhost.localdomain [127.0.0.1]) by gate001.proxmox.com (Proxmox) with ESMTP id 724DB216DD; Thu, 01 Oct 2026 14:11:36 +0200 (CEST) Date: Thu, 1 Oct 2026 14:11:26 +0200 From: Wolfgang Bumiller To: Hannes Laimer Subject: Re: [PATCH pve-cluster 07/10] cfs: add perl client for the change notification socket Message-ID: References: <20260918144152.575163-1-h.laimer@proxmox.com> <20260918144152.575163-8-h.laimer@proxmox.com> MIME-Version: 1.0 Content-Type: text/plain; charset=utf-8 Content-Disposition: inline Content-Transfer-Encoding: 8bit In-Reply-To: <20260918144152.575163-8-h.laimer@proxmox.com> X-Bm-Milter-Handled: 55990f41-d878-4baa-be0a-ee34c49e34d2 X-Bm-Transport-Timestamp: 1790856687178 X-SPAM-LEVEL: Spam detection results: 0 AWL 0.510 Adjusted score from AWL reputation of From: address DMARC_MISSING 0.1 Missing DMARC policy KAM_DMARC_STATUS 0.01 Test Rule for DKIM or SPF Failure with Strict Alignment (newer systems) RCVD_IN_DNSWL_MED -2.3 Sender listed at https://www.dnswl.org/, medium trust SPF_HELO_NONE 0.001 SPF: HELO does not publish an SPF Record SPF_PASS -0.001 SPF: sender matches SPF record Message-ID-Hash: 4JC6L4J5ARYSV4RDAAPC6VQEDUEIS66J X-Message-ID-Hash: 4JC6L4J5ARYSV4RDAAPC6VQEDUEIS66J X-MailFrom: w.bumiller@proxmox.com X-Mailman-Rule-Misses: dmarc-mitigation; no-senders; approved; loop; banned-address; emergency; member-moderation; nonmember-moderation; administrivia; implicit-dest; max-recipients; max-size; news-moderation; no-subject; digests; suspicious-header CC: pve-devel@lists.proxmox.com X-Mailman-Version: 3.3.10 Precedence: list List-Id: Proxmox VE development discussion List-Help: List-Owner: List-Post: List-Subscribe: List-Unsubscribe: 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 > --- > 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 > > > > > --