lib/PublicInbox/Daemon.pm | 41 ++++++++++++++++++++++++++++++-----------
lib/PublicInbox/ParentPipe.pm | 21 +++++++++++++++++++++
t/nntpd.t | 12 ++++++++----
diff --git a/lib/PublicInbox/Daemon.pm b/lib/PublicInbox/Daemon.pm
index c9594a371f17afde7c4400b44127634474b005c9..8de7ff24612881cd43babdc23f5c9120b2f2a702 100644
--- a/lib/PublicInbox/Daemon.pm
+++ b/lib/PublicInbox/Daemon.pm
@@ -14,6 +14,7 @@ STDERR->autoflush(1);
require Danga::Socket;
require POSIX;
require PublicInbox::Listener;
+require PublicInbox::ParentPipe;
my @CMD;
my $set_user;
my (@cfg_listen, $stdout, $stderr, $group, $user, $pid_file, $daemonize);
@@ -161,16 +162,18 @@ };
}
}
-sub worker_quit () {
+
+sub worker_quit {
+ my ($reason) = @_;
# killing again terminates immediately:
exit unless @listeners;
$_->close foreach @listeners; # call Danga::Socket::close
@listeners = ();
-
- # give slow clients 30s to finish reading/writing whatever
- Danga::Socket->AddTimer(30, sub { exit });
+ $reason->close if ref($reason) eq 'PublicInbox::ParentPipe';
+ my $proc_name;
+ my $warn = 0;
# drop idle connections and try to quit gracefully
Danga::Socket->SetPostLoopCallback(sub {
my ($dmap, undef) = @_;
@@ -178,12 +181,23 @@ my $n = 0;
foreach my $s (values %$dmap) {
if ($s->can('busy') && $s->busy) {
- $n = 1;
+ ++$n;
} else {
# close as much as possible, early as possible
$s->close;
}
}
+ if ($n) {
+ if (($warn + 5) < time) {
+ warn "$$ quitting, $n client(s) left\n";
+ $warn = time;
+ }
+ unless (defined $proc_name) {
+ $proc_name = (split(/\s+/, $0))[0];
+ $proc_name =~ s!\A.*?([^/]+)\z!$1!;
+ }
+ $0 = "$proc_name quitting, $n client(s) left";
+ }
$n; # true: loop continues, false: loop breaks
});
}
@@ -359,6 +373,7 @@ };
}
reopen_logs();
# main loop
+ my $quit = 0;
while (1) {
while (my $s = shift @caught) {
if ($s eq 'USR1') {
@@ -367,8 +382,8 @@ kill_workers($s);
} elsif ($s eq 'USR2') {
upgrade();
} elsif ($s =~ /\A(?:QUIT|TERM|INT)\z/) {
- # drops pipes and causes children to die
- exit
+ exit if $quit++;
+ kill_workers($s);
} elsif ($s eq 'WINCH') {
$worker_processes = 0;
} elsif ($s eq 'HUP') {
@@ -390,6 +405,11 @@ }
}
my $n = scalar keys %pids;
+ if ($quit) {
+ exit if $n == 0;
+ $set_workers = $worker_processes = $n = 0;
+ }
+
if ($n > $worker_processes) {
while (my ($k, $v) = each %pids) {
kill('TERM', $k) if $v >= $worker_processes;
@@ -419,13 +439,12 @@ my ($refresh, $post_accept) = @_;
my $parent_pipe;
if ($worker_processes > 0) {
$refresh->(); # preload by default
- $parent_pipe = master_loop(); # returns if in child process
- my $fd = fileno($parent_pipe);
- Danga::Socket->AddOtherFds($fd => *worker_quit);
+ my $fh = master_loop(); # returns if in child process
+ $parent_pipe = PublicInbox::ParentPipe->new($fh, *worker_quit);
} else {
reopen_logs();
$set_user->() if $set_user;
- $SIG{USR2} = sub { worker_quit() if upgrade() };
+ $SIG{USR2} = sub { worker_quit('USR2') if upgrade() };
$refresh->();
}
$uid = $gid = undef;
diff --git a/lib/PublicInbox/ParentPipe.pm b/lib/PublicInbox/ParentPipe.pm
new file mode 100644
index 0000000000000000000000000000000000000000..d2d054ce31f7ca12a16ca848a4f40490c6ce199b
--- /dev/null
+++ b/lib/PublicInbox/ParentPipe.pm
@@ -0,0 +1,21 @@
+# Copyright (C) 2016 all contributors
+# License: AGPL-3.0+
+# only for PublicInbox::Daemon
+package PublicInbox::ParentPipe;
+use strict;
+use warnings;
+use base qw(Danga::Socket);
+use fields qw(cb);
+
+sub new ($$$) {
+ my ($class, $pipe, $cb) = @_;
+ my $self = fields::new($class);
+ $self->SUPER::new($pipe);
+ $self->{cb} = $cb;
+ $self->watch_read(1);
+ $self;
+}
+
+sub event_read { $_[0]->{cb}->($_[0]) }
+
+1;
diff --git a/t/nntpd.t b/t/nntpd.t
index d597855b6fd49bef7687b439e86028af63864acb..d0332216957295160cbf9752cb7469a307f17b4b 100644
--- a/t/nntpd.t
+++ b/t/nntpd.t
@@ -185,7 +185,12 @@ is_deeply($n->xhdr(qw(list-id 1-)), {},
'XHDR on invalid header returns empty');
{
- syswrite($s, "HDR List-id 1-\r\n");
+ setsockopt($s, IPPROTO_TCP, TCP_NODELAY, 1);
+ syswrite($s, 'HDR List-id 1-');
+ select(undef, undef, undef, 0.15);
+ ok(kill('TERM', $pid), 'killed nntpd');
+ select(undef, undef, undef, 0.15);
+ syswrite($s, "\r\n");
$buf = '';
do {
sysread($s, $buf, 4096, length($buf));
@@ -196,9 +201,8 @@ 'got 5xx response for unoptimized HDR');
is(scalar @r, 1, 'only one response line');
}
- ok(kill('TERM', $pid), 'killed nntpd');
- $pid = undef;
- waitpid(-1, 0);
+ is($pid, waitpid($pid, 0), 'nntpd exited successfully');
+ is($?, 0, 'no error in exited process');
}
done_testing();