]> Sergey Matveev's repositories - public-inbox.git/commitdiff
nntpd: support systemd FD inheritance + signals
authorEric Wong <e@80x24.org>
Sat, 19 Sep 2015 08:42:58 +0000 (08:42 +0000)
committerEric Wong <e@80x24.org>
Sun, 20 Sep 2015 02:59:01 +0000 (02:59 +0000)
Avoid depending on IO::Socket::INET if we can help it,
we do not need to bloat ourselves with lot of that
functionality.

lib/PublicInbox/NewsGroup.pm
public-inbox-nntpd
t/nntpd.t [new file with mode: 0644]

index 6cc3f248cd4b69efa64ede7bb87a341e7e1303d4..b8aed529a93fd773ed15ec8a11809438df6eb1b9 100644 (file)
@@ -36,7 +36,10 @@ sub gcf {
 }
 
 sub mm {
-       my ($self) = @_;
+       my ($self, $check_only) = @_;
+       if ($check_only) {
+               return eval { PublicInbox::Msgmap->new($self->{git_dir}) };
+       }
        $self->{mm} ||= eval {
                my $mm = PublicInbox::Msgmap->new($self->{git_dir});
 
index 0c221fa388ff263284db7f8f731f727c905adf2c..588efdd3cc32ea87e5bb6015a688f55459a0cf1b 100644 (file)
 # License: AGPLv3 or later (https://www.gnu.org/licenses/agpl-3.0.txt)
 use strict;
 use warnings;
+my @CMD = ($0, @ARGV);
 require Danga::Socket;
-use IO::Socket;
-use Socket qw(SO_KEEPALIVE IPPROTO_TCP TCP_NODELAY);
-require PublicInbox::NNTP;
+require IO::Handle;
+use Getopt::Long qw/:config gnu_getopt no_ignore_case auto_abbrev/;
 require PublicInbox::NewsGroup;
 my $nntpd = PublicInbox::NNTPD->new;
 my $refresh = sub { $nntpd->refresh_groups };
+$SIG{HUP} = $SIG{USR1} = $SIG{USR2} = $SIG{PIPE} =
+       $SIG{TTIN} = $SIG{TTOU} = $SIG{WINCH} = 'IGNORE';
 
+$refresh->();
+my (@cfg_listen, $stdout, $stderr);
+my $worker_processes = 0;
 my %opts = (
-       LocalAddr => '127.0.0.1:1133',
-       Type => SOCK_STREAM,
-       Proto => 'tcp',
-       Blocking => 0,
-       Reuse => 1,
-       Listen => 1024,
+       'l|listen=s' => \@cfg_listen,
+       '1|stdout=s' => \$stdout,
+       '2|stderr=s' => \$stderr,
+       'W|worker-processes=i' => \$worker_processes,
 );
-my $s = IO::Socket::INET->new(%opts) or die "Error creating socket: $@\n";
-$s->sockopt(SO_KEEPALIVE, 1);
-$s->setsockopt(IPPROTO_TCP, TCP_NODELAY, 1);
+GetOptions(%opts) or die "bad command-line args\n";
+my %pids;
+my %listener_names;
+my $reexec_pid;
+my @listeners = inherit();
 
-$SIG{PIPE} = 'IGNORE';
-$SIG{HUP} = $refresh;
-$refresh->();
+# default NNTP listener if no listeners
+push @cfg_listen, '0.0.0.0:119' unless (@listeners || @cfg_listen);
+
+foreach my $l (@cfg_listen) {
+       next if $listener_names{$l}; # already inherited
+       require IO::Socket::INET6; # works for IPv4, too
+       my %o = (
+               LocalAddr => $l,
+               ReuseAddr => 1,
+               Proto => 'tcp',
+       );
+       if (my $s = IO::Socket::INET6->new(%o)) {
+               $listener_names{sockname($s)} = $s;
+               push @listeners, $s;
+       } else {
+               warn "error binding $l: $!\n";
+       }
+}
+die 'No listeners bound' unless @listeners;
+open(STDIN, '+<', '/dev/null');
+
+if ($worker_processes > 0) {
+       # my ($p0, $p1, $r, $w);
+       pipe(my ($p0, $p1)) or die "failed to create parent-pipe: $!\n";
+       my %pwatch = ( fileno($p0) => sub { kill('TERM', $$) } );
+       pipe(my ($r, $w)) or die "failed to create self-pipe: $!\n";
+       IO::Handle::blocking($w, 0);
+       my $set_workers = $worker_processes;
+       my @caught;
+       my $master_pid = $$;
+       foreach my $s (qw(HUP CHLD QUIT INT TERM USR1 USR2 TTIN TTOU WINCH)) {
+               $SIG{$s} = sub {
+                       return if $$ != $master_pid;
+                       push @caught, $s;
+                       syswrite($w, '.');
+               };
+       }
+       reopen_logs();
+       # main loop
+       while (1) {
+               while (my $s = shift @caught) {
+                       if ($s eq 'USR1') {
+                               reopen_logs();
+                               kill_workers($s);
+                       } elsif ($s eq 'USR2') {
+                               upgrade();
+                       } elsif ($s =~ /\A(?:QUIT|TERM|INT)\z/) {
+                               # drops pipes and causes children to die
+                               exit
+                       } elsif ($s eq 'WINCH') {
+                               $worker_processes = 0;
+                       } elsif ($s eq 'HUP') {
+                               $worker_processes = $set_workers;
+                               $refresh->();
+                               kill_workers($s);
+                       } elsif ($s eq 'TTIN') {
+                               if ($set_workers > $worker_processes) {
+                                       ++$worker_processes;
+                               } else {
+                                       $worker_processes = ++$set_workers;
+                               }
+                       } elsif ($s eq 'TTOU') {
+                               if ($set_workers > 0) {
+                                       $worker_processes = --$set_workers;
+                               }
+                       } elsif ($s eq 'CHLD') {
+                               reap_children();
+                       }
+               }
+
+               my $n = scalar keys %pids;
+               if ($n > $worker_processes) {
+                       while (my ($k, $v) = each %pids) {
+                               kill('TERM', $k) if $v >= $worker_processes;
+                       }
+                       $n = $worker_processes;
+               }
+               foreach my $i ($n..($worker_processes - 1)) {
+                       my ($pid, $err) = do_fork();
+                       if (!defined $pid) {
+                               warn "failed to fork worker[$i]: $err\n";
+                       } elsif ($pid == 0) {
+                               close($_) for ($w, $r, $p1);
+                               Danga::Socket->AddOtherFds(%pwatch);
+                               goto worker;
+                       } else {
+                               warn "PID=$pid is worker[$i]\n";
+                               $pids{$pid} = $i;
+                       }
+               }
+               sysread($r, my $buf, 8);
+       }
+} else {
+worker:
+       # this calls epoll_create:
+       @listeners = map { PublicInbox::Listener->new($_) } @listeners;
+       reopen_logs();
+       $SIG{QUIT} = $SIG{INT} = $SIG{TERM} = *worker_quit;
+       $SIG{USR1} = *reopen_logs;
+       $SIG{HUP} = $refresh;
+       $_->watch_read(1) for @listeners;
+       Danga::Socket->EventLoop;
+}
+
+# end of main
+
+sub worker_quit {
+       # killing again terminates immediately:
+       exit unless @listeners;
 
-Danga::Socket->AddOtherFds(fileno($s) => sub {
-       while (my $c = $s->accept) {
-               $c->blocking(0); # no accept4 :<
+       $_->close for @listeners;
+       @listeners = ();
+
+       # drop idle connections and try to quit gracefully
+       Danga::Socket->SetPostLoopCallback(sub {
+               my ($dmap, undef) = @_;
+               my $n = 0;
+               foreach my $s (values %$dmap) {
+                       if ($s->{write_buf_size} || @{$s->{read_push_back}}) {
+                               ++$n;
+                       } else {
+                               $s->close;
+                       }
+               }
+               $n; # true: loop continues, false: loop breaks
+       });
+}
+
+sub reopen_logs {
+       if ($stdout) {
+               open STDOUT, '>>', $stdout or
+                       warn "failed to redirect stdout to $stdout: $!\n";
+       }
+       if ($stderr) {
+               open STDERR, '>>', $stderr or
+                       warn "failed to redirect stderr to $stderr: $!\n";
+       }
+}
+
+sub sockname {
+       my ($s) = @_;
+       my $n = getsockname($s) or return;
+       my ($port, $addr);
+       if (bytes::length($n) >= 28) {
+               require Socket6;
+               ($port, $addr) = Socket6::unpack_sockaddr_in6($n);
+       } else {
+               ($port, $addr) = Socket::sockaddr_in($n);
+       }
+       if (bytes::length($addr) == 4) {
+               $n = Socket::inet_ntoa($addr)
+       } else {
+               $n = '['.Socket6::inet_ntop(Socket6::AF_INET6(), $addr).']';
+       }
+       $n .= ":$port";
+}
+
+sub inherit {
+       return () if ($ENV{LISTEN_PID} || 0) != $$;
+       my $fds = $ENV{LISTEN_FDS} or return ();
+       my $end = $fds + 2; # LISTEN_FDS_START - 1
+       my @rv = ();
+       foreach my $fd (3..$end) {
+               my $s = IO::Handle->new;
+               $s->fdopen($fd, 'r');
+               if (my $k = sockname($s)) {
+                       $listener_names{$k} = $s;
+                       push @rv, $s;
+               } else {
+                       warn "failed to inherit fd=$fd (LISTEN_FDS=$fds)";
+               }
+       }
+       @rv
+}
+
+sub upgrade {
+       if ($reexec_pid) {
+               warn "upgrade in-progress: $reexec_pid\n";
+               return;
+       }
+       my ($pid, $err) = do_fork();
+       unless (defined $pid) {
+               warn "fork failed: $err\n";
+               return;
+       }
+       if ($pid == 0) {
+               use Fcntl qw(FD_CLOEXEC F_SETFD F_GETFD);
+               $ENV{LISTEN_FDS} = scalar @listeners;
+               $ENV{LISTEN_PID} = $$;
+               foreach my $s (@listeners) {
+                       my $fl = fcntl($s, F_GETFD, 0);
+                       fcntl($s, F_SETFD, $fl &= ~FD_CLOEXEC);
+               }
+               exec @CMD;
+               die "Failed to exec: $!\n";
+       }
+       $reexec_pid = $pid;
+}
+
+sub kill_workers {
+       my ($s) = @_;
+
+       while (my ($pid, $id) = each %pids) {
+               kill $s, $pid;
+       }
+}
+
+sub do_fork {
+       require POSIX;
+       my $new = POSIX::SigSet->new;
+       $new->fillset;
+       my $old = POSIX::SigSet->new;
+       POSIX::sigprocmask(&POSIX::SIG_BLOCK, $new, $old) or
+                               die "SIG_BLOCK: $!\n";
+       my $pid = fork;
+       my $err = $!;
+       POSIX::sigprocmask(&POSIX::SIG_SETMASK, $old) or
+                               die "SIG_SETMASK: $!\n";
+       ($pid, $err);
+}
+
+sub reap_children {
+       while (1) {
+               my $p = waitpid(-1, &POSIX::WNOHANG) or return;
+               if (defined $reexec_pid && $p == $reexec_pid) {
+                       $reexec_pid = undef;
+                       warn "reexec PID($p) died with: $?\n";
+               } elsif (defined(my $id = delete $pids{$p})) {
+                       warn "worker[$id] PID($p) died with: $?\n";
+               } elsif ($p > 0) {
+                       warn "unknown PID($p) reaped: $?\n";
+               } else {
+                       return;
+               }
+       }
+}
+
+1;
+package PublicInbox::Listener;
+use strict;
+use warnings;
+use base 'Danga::Socket';
+use Socket qw(SOL_SOCKET SO_KEEPALIVE IPPROTO_TCP TCP_NODELAY);
+use PublicInbox::NNTP;
+
+sub new ($$) {
+       my ($class, $s) = @_;
+       setsockopt($s, SOL_SOCKET, SO_KEEPALIVE, 1);
+       setsockopt($s, IPPROTO_TCP, TCP_NODELAY, 1);
+       listen($s, 1024);
+       IO::Handle::blocking($s, 0);
+       my $self = fields::new($class);
+       $self->SUPER::new($s);
+}
+
+sub event_read {
+       my ($self) = @_;
+       # no loop here, we want to fairly distribute clients
+       # between multiple processes sharing the same socket
+       if (accept(my $c, $self->{sock})) {
+               IO::Handle::blocking($c, 0); # no accept4 :<
                PublicInbox::NNTP->new($c, $nntpd);
        }
-});
-Danga::Socket->EventLoop();
+}
 
+1;
 package PublicInbox::NNTPD;
 use strict;
 use warnings;
@@ -70,7 +329,7 @@ sub refresh_groups {
                }
 
                # Only valid if Msgmap works
-               $new->{$g} = $ng if $ng->mm;
+               $new->{$g} = $ng if $ng->mm(1);
        }
        # this will destroy old groups that got deleted
        %{$self->{groups}} = %$new;
diff --git a/t/nntpd.t b/t/nntpd.t
new file mode 100644 (file)
index 0000000..527cfc2
--- /dev/null
+++ b/t/nntpd.t
@@ -0,0 +1,106 @@
+# Copyright (C) 2015 all contributors <meta@public-inbox.org>
+# License: AGPLv3 or later (https://www.gnu.org/licenses/agpl-3.0.txt)
+use strict;
+use warnings;
+use Test::More;
+eval { require PublicInbox::SearchIdx };
+plan skip_all => "Xapian missing for nntpd" if $@;
+eval { require PublicInbox::Msgmap };
+plan skip_all => "DBD::SQLite missing for nntpd" if $@;
+use Cwd;
+use Email::Simple;
+use IO::Socket;
+use Fcntl qw(FD_CLOEXEC F_SETFD F_GETFD);
+use Socket qw(SO_KEEPALIVE IPPROTO_TCP TCP_NODELAY);
+use File::Temp qw/tempdir/;
+use Net::NNTP;
+use IPC::Run qw(run);
+use Data::Dumper;
+
+my $tmpdir = tempdir(CLEANUP => 1);
+my $home = "$tmpdir/pi-home";
+my $err = "$tmpdir/stderr.log";
+my $out = "$tmpdir/stdout.log";
+my $pi_home = "$home/.public-inbox";
+my $pi_config = "$pi_home/config";
+my $maindir = "$tmpdir/main.git";
+my $main_bin = getcwd()."/t/main-bin";
+my $main_path = "$main_bin:$ENV{PATH}"; # for spamc ham mock
+my $group = 'test-nntpd';
+my $addr = $group . '@example.com';
+my $cfgpfx = "publicinbox.$group";
+my $failbox = "$home/fail.mbox";
+local $ENV{PI_EMERGENCY} = $failbox;
+my $mda = 'blib/script/public-inbox-mda';
+my $nntpd = 'blib/script/public-inbox-nntpd';
+my $init = 'blib/script/public-inbox-init';
+my $index = 'blib/script/public-inbox-index';
+
+my %opts = (
+       LocalAddr => '127.0.0.1',
+       ReuseAddr => 1,
+       Proto => 'tcp',
+       Type => SOCK_STREAM,
+       Listen => 1024,
+);
+my $sock = IO::Socket::INET->new(%opts);
+plan skip_all => 'sock fd!=3, cannot test nntpd integration' if fileno($sock) != 3;
+my $pid;
+END { kill 'TERM', $pid if defined $pid };
+{
+       local $ENV{HOME} = $home;
+       system($init, $group, $maindir, 'http://example.com/', $addr);
+
+       # ensure successful message delivery
+       {
+               local $ENV{ORIGINAL_RECIPIENT} = $addr;
+               my $simple = Email::Simple->new(<<EOF);
+From: Me <me\@example.com>
+To: You <you\@example.com>
+Cc: $addr
+Message-Id: <nntp\@example.com>
+Subject: hihi
+Date: Thu, 01 Jan 1970 00:00:00 +0000
+
+nntp
+EOF
+               my $in = $simple->as_string;
+               local $ENV{PATH} = $main_path;
+               IPC::Run::run([$mda], \$in);
+               is(0, $?, 'ran MDA correctly');
+               is(0, system($index, $maindir), 'indexed git dir');
+       }
+
+       ok($sock, 'sock created');
+       $! = 0;
+       my $fl = fcntl($sock, F_GETFD, 0);
+       ok(! $!, 'no error from fcntl(F_GETFD)');
+       is($fl, FD_CLOEXEC, 'cloexec set by default (Perl behavior)');
+       $pid = fork;
+       if ($pid == 0) {
+               # pretend to be systemd
+               fcntl($sock, F_SETFD, $fl &= ~FD_CLOEXEC);
+               $ENV{LISTEN_PID} = $$;
+               $ENV{LISTEN_FDS} = 1;
+               exec $nntpd, "--stdout=$out", "--stderr=$err";
+               die "FAIL: $!\n";
+       }
+       ok(defined $pid, 'forked nntpd process successfully');
+       $! = 0;
+       ok(! $!, 'no error from fcntl(F_SETFD)');
+       fcntl($sock, F_SETFD, $fl |= FD_CLOEXEC);
+       my $n = Net::NNTP->new($sock->sockhost . ':' . $sock->sockport);
+       my $list = $n->list;
+       is_deeply($list, { $group => [ qw(1 1 n) ] }, 'LIST works');
+       is_deeply([$n->group($group)], [ qw(0 1 1), $group ], 'GROUP works');
+
+       # TODO: upgrades and such
+
+       ok(kill('TERM', $pid), 'killed nntpd');
+       $pid = undef;
+       waitpid(-1, 0);
+}
+
+done_testing();
+
+1;