[svn:qpsmtpd] r643 - in trunk: . lib lib/Danga plugins
[email protected] Tue, 20 Jun 2006 06:51:33 -0700 (PDT)
| Newsgroups | perl.cvs.qpsmtpd |
|---|---|
| Message-ID | <[email protected]> |
Author: msergeant
Date: Tue Jun 20 06:51:32 2006
New Revision: 643
Modified:
trunk/lib/Danga/Client.pm
trunk/lib/Danga/DNS.pm
trunk/lib/Danga/Socket.pm
trunk/lib/Qpsmtpd.pm
trunk/plugins/check_earlytalker
trunk/qpsmtpd
Log:
Simplify qpsmtpd script (remove inetd and forking server)
Greatly simplify Danga::Client due to no more need for line mode client
Update to latest Danga::Socket
Fix check_earlytalker to use new API
Fix Danga::DNS to use new API
Modified: trunk/lib/Danga/Client.pm
==============================================================================
--- trunk/lib/Danga/Client.pm (original)
+++ trunk/lib/Danga/Client.pm Tue Jun 20 06:51:32 2006
@@ -2,7 +2,7 @@
package Danga::Client;
use base 'Danga::TimeoutSocket';
-use fields qw(line closing disable_read can_read_mode);
+use fields qw(line pause_count);
use Time::HiRes ();
# 30 seconds max timeout!
@@ -21,68 +21,14 @@
sub reset_for_next_message {
my Danga::Client $self = shift;
$self->{line} = '';
- $self->{disable_read} = 0;
- $self->{can_read_mode} = 0;
+ $self->{pause_count} = 0;
return $self;
}
-sub get_line {
- my Danga::Client $self = shift;
- if (!$self->have_line) {
- $self->SetPostLoopCallback(sub { $self->have_line ? 0 : 1 });
- #warn("get_line PRE\n");
- $self->EventLoop();
- #warn("get_line POST\n");
- $self->disable_read();
- }
- return if $self->{closing};
- # now have a line.
- $self->{alive_time} = time;
- $self->{line} =~ s/^(.*?\n)//;
- return $1;
-}
-
-sub can_read {
- my Danga::Client $self = shift;
- my ($timeout) = @_;
- my $end = Time::HiRes::time() + $timeout;
- # warn("Calling can-read\n");
- $self->{can_read_mode} = 1;
- if (!length($self->{line})) {
- $self->disable_read();
- # loop because any callback, not just ours, can make EventLoop return
- while( !(length($self->{line}) || (Time::HiRes::time > $end)) ) {
- $self->SetPostLoopCallback(sub { (length($self->{line}) ||
- (Time::HiRes::time > $end)) ? 0 : 1 });
- #warn("get_line PRE\n");
- $self->EventLoop();
- #warn("get_line POST\n");
- }
- $self->enable_read();
- }
- $self->{can_read_mode} = 0;
- $self->SetPostLoopCallback(undef);
- return if $self->{closing};
- $self->{alive_time} = time;
- # warn("can_read returning for '$self->{line}'\n");
- return 1 if length($self->{line});
- return;
-}
-
-sub have_line {
- my Danga::Client $self = shift;
- return 1 if $self->{closing};
- if ($self->{line} =~ /\n/) {
- return 1;
- }
- return 0;
-}
-
sub event_read {
my Danga::Client $self = shift;
my $bref = $self->read(8192);
return $self->close($!) unless defined $bref;
- # $self->watch_read(0);
$self->process_read_buf($bref);
}
@@ -90,8 +36,7 @@
my Danga::Client $self = shift;
my $bref = shift;
$self->{line} .= $$bref;
- return if ! $self->readable();
- return if $::LineMode;
+ return if $self->paused();
while ($self->{line} =~ s/^(.*?\n)//) {
my $line = $1;
@@ -99,34 +44,40 @@
my $resp = $self->process_line($line);
if ($::DEBUG > 1 and $resp) { print "$$:".($self+0)."S: $_\n" for split(/\n/, $resp) }
$self->write($resp) if $resp;
- $self->watch_read(0) if $self->{disable_read};
- last if ! $self->readable();
- }
- if($self->have_line) {
- $self->shift_back_read($self->{line});
- $self->{line} = '';
+ # $self->watch_read(0) if $self->{pause_count};
+ last if $self->paused();
}
}
-sub readable {
+sub has_data {
+ my Danga::Client $self = shift;
+ return length($self->{line}) ? 1 : 0;
+}
+
+sub clear_data {
+ my Danga::Client $self = shift;
+ $self->{line} = '';
+}
+
+sub paused {
my Danga::Client $self = shift;
- return 0 if $self->{disable_read} > 0;
- return 0 if $self->{closed} > 0;
- return 1;
+ return 1 if $self->{pause_count};
+ return 1 if $self->{closed};
+ return 0;
}
-sub disable_read {
+sub pause_read {
my Danga::Client $self = shift;
- $self->{disable_read}++;
- $self->watch_read(0);
+ $self->{pause_count}++;
+ # $self->watch_read(0);
}
-sub enable_read {
+sub continue_read {
my Danga::Client $self = shift;
- $self->{disable_read}--;
- if ($self->{disable_read} <= 0) {
- $self->{disable_read} = 0;
- $self->watch_read(1);
+ $self->{pause_count}--;
+ if ($self->{pause_count} <= 0) {
+ $self->{pause_count} = 0;
+ # $self->watch_read(1);
}
}
@@ -137,7 +88,6 @@
sub close {
my Danga::Client $self = shift;
- $self->{closing} = 1;
print "closing @_\n" if $::DEBUG;
$self->SUPER::close(@_);
}
Modified: trunk/lib/Danga/DNS.pm
==============================================================================
--- trunk/lib/Danga/DNS.pm (original)
+++ trunk/lib/Danga/DNS.pm Tue Jun 20 06:51:32 2006
@@ -25,7 +25,7 @@
$resolver ||= Danga::DNS::Resolver->new();
my $client = $options{client};
- $client->disable_read if $client;
+ $client->pause_read() if $client;
$self = fields::new($self) unless ref $self;
@@ -40,13 +40,13 @@
if ($options{type}) {
if ( ($options{type} eq 'A') || ($options{type} eq 'PTR') ) {
if (!$resolver->query($self, @{$self->{hosts}})) {
- $client->enable_read() if $client;
+ $client->continue_read() if $client;
return;
}
}
else {
if (!$resolver->query_type($self, $options{type}, @{$self->{hosts}})) {
- $client->enable_read() if $client;
+ $client->continue_read() if $client;
return;
}
# die "Unsupported DNS query type: $options{type}";
@@ -54,7 +54,7 @@
}
else {
if (!$resolver->query($self, @{$self->{hosts}})) {
- $client->enable_read() if $client;
+ $client->continue_read() if $client;
return;
}
}
@@ -84,7 +84,7 @@
$self->{callback}->("NXDOMAIN", $host);
}
}
- $self->{client}->enable_read if $self->{client};
+ $self->{client}->continue_read() if $self->{client};
if ($self->{finished}) {
$self->{finished}->();
}
Modified: trunk/lib/Danga/Socket.pm
==============================================================================
--- trunk/lib/Danga/Socket.pm (original)
+++ trunk/lib/Danga/Socket.pm Tue Jun 20 06:51:32 2006
@@ -2,16 +2,100 @@
=head1 NAME
-Danga::Socket - Event-driven async IO class
+Danga::Socket - Event loop and event-driven async socket base class
=head1 SYNOPSIS
+ package My::Socket
+ use Danga::Socket;
use base ('Danga::Socket');
+ use fields ('my_attribute');
+
+ sub new {
+ my My::Socket $self = shift;
+ $self = fields::new($self) unless ref $self;
+ $self->SUPER::new( @_ );
+
+ $self->{my_attribute} = 1234;
+ return $self;
+ }
+
+ sub event_err { ... }
+ sub event_hup { ... }
+ sub event_write { ... }
+ sub event_read { ... }
+ sub close { ... }
+
+ $my_sock->tcp_cork($bool);
+
+ # write returns 1 if all writes have gone through, or 0 if there
+ # are writes in queue
+ $my_sock->write($scalar);
+ $my_sock->write($scalarref);
+ $my_sock->write(sub { ... }); # run when previous data written
+ $my_sock->write(undef); # kick-starts
+
+ # read max $bytecount bytes, or undef on connection closed
+ $scalar_ref = $my_sock->read($bytecount);
+
+ # watch for writability. not needed with ->write(). write()
+ # will automatically turn on watch_write when you wrote too much
+ # and turn it off when done
+ $my_sock->watch_write($bool);
+
+ # watch for readability
+ $my_sock->watch_read($bool);
+
+ # if you read too much and want to push some back on
+ # readable queue. (not incredibly well-tested)
+ $my_sock->push_back_read($buf); # scalar or scalar ref
+
+ Danga::Socket->AddOtherFds(..);
+ Danga::Socket->SetLoopTimeout($millisecs);
+ Danga::Socket->DescriptorMap();
+ Danga::Socket->WatchedSockets(); # count of DescriptorMap keys
+ Danga::Socket->SetPostLoopCallback($code);
+ Danga::Socket->EventLoop();
=head1 DESCRIPTION
-This is an abstract base class which provides the basic framework for
-event-driven asynchronous IO.
+This is an abstract base class for objects backed by a socket which
+provides the basic framework for event-driven asynchronous IO,
+designed to be fast. Danga::Socket is both a base class for objects,
+and an event loop.
+
+Callers subclass Danga::Socket. Danga::Socket's constructor registers
+itself with the Danga::Socket event loop, and invokes callbacks on the
+object for readability, writability, errors, and other conditions.
+
+Because Danga::Socket uses the "fields" module, your subclasses must
+too.
+
+=head1 MORE INFO
+
+For now, see servers using Danga::Socket for guidance. For example:
+perlbal, mogilefsd, or ddlockd.
+
+=head1 AUTHORS
+
+Brad Fitzpatrick <[email protected]> - author
+
+Michael Granger <[email protected]> - docs, testing
+
+Mark Smith <[email protected]> - contributor, heavy user, testing
+
+Matt Sergeant <[email protected]> - kqueue support
+
+=head1 BUGS
+
+Not documented enough.
+
+tcp_cork only works on Linux for now. No BSD push/nopush support.
+
+=head1 LICENSE
+
+License is granted to use and distribute this module under the same
+terms as Perl itself.
=cut
@@ -19,53 +103,53 @@
package Danga::Socket;
use strict;
+use bytes;
+use POSIX ();
+use Time::HiRes ();
+
+my $opt_bsd_resource = eval "use BSD::Resource; 1;";
use vars qw{$VERSION};
-$VERSION = do { my @r = (q$Revision: 1.4 $ =~ /\d+/g); sprintf "%d."."%02d" x $#r, @r };
+$VERSION = "1.51";
-use fields qw(sock fd write_buf write_buf_offset write_buf_size
- read_push_back post_loop_callback
- peer_ip
- closed event_watch debug_level);
+use warnings;
+no warnings qw(deprecated);
-use Errno qw(EINPROGRESS EWOULDBLOCK EISCONN
- EPIPE EAGAIN EBADF ECONNRESET);
+use Sys::Syscall qw(:epoll);
-use Socket qw(IPPROTO_TCP);
-use Carp qw{croak confess};
-use POSIX ();
+use fields ('sock', # underlying socket
+ 'fd', # numeric file descriptor
+ 'write_buf', # arrayref of scalars, scalarrefs, or coderefs to write
+ 'write_buf_offset', # offset into first array of write_buf to start writing at
+ 'write_buf_size', # total length of data in all write_buf items
+ 'read_push_back', # arrayref of "pushed-back" read data the application didn't want
+ 'closed', # bool: socket is closed
+ 'corked', # bool: socket is corked
+ 'event_watch', # bitmask of events the client is interested in (POLLIN,OUT,etc.)
+ 'peer_ip', # cached stringified IP address of $sock
+ 'peer_port', # cached port number of $sock
+ 'local_ip', # cached stringified IP address of local end of $sock
+ 'local_port', # cached port number of local end of $sock
+ 'writer_func', # subref which does writing. must return bytes written (or undef) and set $! on errors
+ );
-use constant TCP_CORK => 3; # FIXME: not hard-coded (Linux-specific too)
+use Errno qw(EINPROGRESS EWOULDBLOCK EISCONN ENOTSOCK
+ EPIPE EAGAIN EBADF ECONNRESET ENOPROTOOPT);
+use Socket qw(IPPROTO_TCP);
+use Carp qw(croak confess);
+use constant TCP_CORK => ($^O eq "linux" ? 3 : 0); # FIXME: not hard-coded (Linux-specific too)
use constant DebugLevel => 0;
-# for epoll definitions:
-our $HAVE_SYSCALL_PH = eval { require 'syscall.ph'; 1 } || eval { require 'sys/syscall.ph'; 1 };
-our $HAVE_KQUEUE = eval { require IO::KQueue; 1 };
-
-# Explicitly define the poll constants, as either one set or the other won't be
-# loaded. They're also badly implemented in IO::Epoll:
-# The IO::Epoll module is buggy in that it doesn't export constants efficiently
-# (at least as of 0.01), so doing constants ourselves saves 13% of the user CPU
-# time
-use constant EPOLLIN => 1;
-use constant EPOLLOUT => 4;
-use constant EPOLLERR => 8;
-use constant EPOLLHUP => 16;
-use constant EPOLL_CTL_ADD => 1;
-use constant EPOLL_CTL_DEL => 2;
-use constant EPOLL_CTL_MOD => 3;
-
use constant POLLIN => 1;
use constant POLLOUT => 4;
use constant POLLERR => 8;
use constant POLLHUP => 16;
use constant POLLNVAL => 32;
-# keep track of active clients
+our $HAVE_KQUEUE = eval { require IO::KQueue; 1 };
+
our (
- $DoneInit, # if we've done the one-time module init yet
- $TryEpoll, # Whether epoll should be attempted to be used.
$HaveEpoll, # Flag -- is epoll available? initially undefined.
$HaveKQueue,
%DescriptorMap, # fd (num) -> Danga::Socket object
@@ -75,8 +159,14 @@
@ToClose, # sockets to close when event loop is done
%OtherFds, # A hash of "other" (non-Danga::Socket) file
# descriptors for the event loop to track.
- $PostLoopCallback, # subref to call at the end of each loop, if defined
- %PLCMap, # fd (num) -> PostLoopCallback
+
+ $PostLoopCallback, # subref to call at the end of each loop, if defined (global)
+ %PLCMap, # fd (num) -> PostLoopCallback (per-object)
+
+ $LoopTimeout, # timeout of event loop in milliseconds
+ $DoProfile, # if on, enable profiling
+ %Profiling, # what => [ utime, stime, calls ]
+ $DoneInit, # if we've done the one-time module init yet
@Timers, # timers
);
@@ -86,21 +176,27 @@
### C L A S S M E T H O D S
#####################################################################
-### (CLASS) METHOD: Reset()
-### Reset all state
+# (CLASS) method: reset all state
sub Reset {
%DescriptorMap = ();
%PushBackSet = ();
@ToClose = ();
%OtherFds = ();
+ $LoopTimeout = -1; # no timeout by default
+ $DoProfile = 0;
+ %Profiling = ();
+ @Timers = ();
+
$PostLoopCallback = undef;
%PLCMap = ();
- @Timers = ();
}
### (CLASS) METHOD: HaveEpoll()
### Returns a true value if this class will use IO::Epoll for async IO.
-sub HaveEpoll { $HaveEpoll };
+sub HaveEpoll {
+ _InitPoller();
+ return $HaveEpoll;
+}
### (CLASS) METHOD: WatchedSockets()
### Returns the number of file descriptors which are registered with the global
@@ -110,43 +206,95 @@
}
*watched_sockets = *WatchedSockets;
+### (CLASS) METHOD: EnableProfiling()
+### Turns profiling on, clearing current profiling data.
+sub EnableProfiling {
+ if ($opt_bsd_resource) {
+ %Profiling = ();
+ $DoProfile = 1;
+ return 1;
+ }
+ return 0;
+}
+
+### (CLASS) METHOD: DisableProfiling()
+### Turns off profiling, but retains data up to this point
+sub DisableProfiling {
+ $DoProfile = 0;
+}
+
+### (CLASS) METHOD: ProfilingData()
+### Returns reference to a hash of data in format above (see %Profiling)
+sub ProfilingData {
+ return \%Profiling;
+}
### (CLASS) METHOD: ToClose()
### Return the list of sockets that are awaiting close() at the end of the
### current event loop.
sub ToClose { return @ToClose; }
-
### (CLASS) METHOD: OtherFds( [%fdmap] )
### Get/set the hash of file descriptors that need processing in parallel with
### the registered Danga::Socket objects.
sub OtherFds {
my $class = shift;
- if ( @_ ) { %OtherFds = (%OtherFds, @_) }
+ if ( @_ ) { %OtherFds = @_ }
return wantarray ? %OtherFds : \%OtherFds;
}
+### (CLASS) METHOD: AddOtherFds( [%fdmap] )
+### Add fds to the OtherFds hash for processing.
+sub AddOtherFds {
+ my $class = shift;
+ %OtherFds = ( %OtherFds, @_ ); # FIXME investigate what happens on dupe fds
+ return wantarray ? %OtherFds : \%OtherFds;
+}
+
+### (CLASS) METHOD: SetLoopTimeout( $timeout )
+### Set the loop timeout for the event loop to some value in milliseconds.
+sub SetLoopTimeout {
+ return $LoopTimeout = $_[1] + 0;
+}
+
+### (CLASS) METHOD: DebugMsg( $format, @args )
+### Print the debugging message specified by the C<sprintf>-style I<format> and
+### I<args>
+sub DebugMsg {
+ my ( $class, $fmt, @args ) = @_;
+ chomp $fmt;
+ printf STDERR ">>> $fmt\n", @args;
+}
+
+### (CLASS) METHOD: AddTimer( $seconds, $coderef )
+### Add a timer to occur $seconds from now. $seconds may be fractional. Don't
+### expect this to be accurate though.
sub AddTimer {
my $class = shift;
my ($secs, $coderef) = @_;
- my $timeout = time + $secs;
-
- if (!@Timers || ($timeout >= $Timers[-1][0])) {
- push @Timers, [$timeout, $coderef];
+
+ my $fire_time = Time::HiRes::time() + $secs;
+
+ if (!@Timers || $fire_time >= $Timers[-1][0]) {
+ push @Timers, [$fire_time, $coderef];
return;
}
-
- # Now where do we insert...
+
+ # Now, where do we insert? (NOTE: this appears slow, algorithm-wise,
+ # but it was compared against calendar queues, heaps, naive push/sort,
+ # and a bunch of other versions, and found to be fastest with a large
+ # variety of datasets.)
for (my $i = 0; $i < @Timers; $i++) {
- if ($Timers[$i][0] > $timeout) {
- splice(@Timers, $i, 0, [$timeout, $coderef]);
+ if ($Timers[$i][0] > $fire_time) {
+ splice(@Timers, $i, 0, [$fire_time, $coderef]);
return;
}
}
-
- die "Shouldn't get here spank matt.";
+
+ die "Shouldn't get here.";
}
+
### (CLASS) METHOD: DescriptorMap()
### Get the hash of Danga::Socket objects keyed by the file descriptor they are
### wrapping.
@@ -156,11 +304,11 @@
*descriptor_map = *DescriptorMap;
*get_sock_ref = *DescriptorMap;
-sub init_poller
+sub _InitPoller
{
return if $DoneInit;
$DoneInit = 1;
-
+
if ($HAVE_KQUEUE) {
$KQueue = IO::KQueue->new();
$HaveKQueue = $KQueue >= 0;
@@ -168,14 +316,14 @@
*EventLoop = *KQueueEventLoop;
}
}
- elsif ($TryEpoll) {
+ elsif (Sys::Syscall::epoll_defined()) {
$Epoll = eval { epoll_create(1024); };
$HaveEpoll = defined $Epoll && $Epoll >= 0;
if ($HaveEpoll) {
*EventLoop = *EpollEventLoop;
}
}
-
+
if (!$HaveEpoll && !$HaveKQueue) {
require IO::Poll;
*EventLoop = *PollEventLoop;
@@ -187,7 +335,7 @@
sub EventLoop {
my $class = shift;
- init_poller();
+ _InitPoller();
if ($HaveEpoll) {
EpollEventLoop($class);
@@ -198,63 +346,55 @@
}
}
-### The kqueue-based event loop. Gets installed as EventLoop if IO::KQueue works
-### okay.
-sub KQueueEventLoop {
- my $class = shift;
-
- foreach my $fd (keys %OtherFds) {
- $KQueue->EV_SET($fd, IO::KQueue::EVFILT_READ(), IO::KQueue::EV_ADD());
+## profiling-related data/functions
+our ($Prof_utime0, $Prof_stime0);
+sub _pre_profile {
+ ($Prof_utime0, $Prof_stime0) = getrusage();
+}
+
+sub _post_profile {
+ # get post information
+ my ($autime, $astime) = getrusage();
+
+ # calculate differences
+ my $utime = $autime - $Prof_utime0;
+ my $stime = $astime - $Prof_stime0;
+
+ foreach my $k (@_) {
+ $Profiling{$k} ||= [ 0.0, 0.0, 0 ];
+ $Profiling{$k}->[0] += $utime;
+ $Profiling{$k}->[1] += $stime;
+ $Profiling{$k}->[2]++;
}
-
- while (1) {
- my $now = time;
- # Run expired timers
- while (@Timers && $Timers[0][0] <= $now) {
- my $to_run = shift(@Timers);
- $to_run->[1]->($now);
- }
-
- # Get next timeout
- my $timeout = @Timers ? ($Timers[0][0] - $now) : 1;
- # print STDERR "kevent($timeout)\n";
- my @ret = $KQueue->kevent($timeout * 1000);
-
- foreach my $kev (@ret) {
- my ($fd, $filter, $flags, $fflags) = @$kev;
-
- my Danga::Socket $pob = $DescriptorMap{$fd};
-
- # prioritise OtherFds first - likely to be accept() socks (?)
- if (!$pob) {
- if (my $code = $OtherFds{$fd}) {
- $code->($filter);
- }
- else {
- print STDERR "kevent() returned fd $fd for which we have no mapping. removing.\n";
- POSIX::close($fd); # close deletes the kevent entry
- }
- next;
- }
-
- DebugLevel >= 1 && $class->DebugMsg("Event: fd=%d (%s), flags=%d \@ %s\n",
- $fd, ref($pob), $flags, time);
-
- $pob->event_read if $filter == IO::KQueue::EVFILT_READ() && !$pob->{closed};
- $pob->event_write if $filter == IO::KQueue::EVFILT_WRITE() && !$pob->{closed};
- if ($flags == IO::KQueue::EV_EOF() && !$pob->{closed}) {
- if ($fflags) {
- $pob->event_err;
- } else {
- $pob->event_hup;
- }
- }
- }
-
- return unless PostEventLoop();
+}
+
+# runs timers and returns milliseconds for next one, or next event loop
+sub RunTimers {
+ return $LoopTimeout unless @Timers;
+
+ my $now = Time::HiRes::time();
+
+ # Run expired timers
+ while (@Timers && $Timers[0][0] <= $now) {
+ my $to_run = shift(@Timers);
+ $to_run->[1]->($now);
}
-
- exit(0);
+
+ return $LoopTimeout unless @Timers;
+
+ # convert time to an even number of milliseconds, adding 1
+ # extra, otherwise floating point fun can occur and we'll
+ # call RunTimers like 20-30 times, each returning a timeout
+ # of 0.0000212 seconds
+ my $timeout = int(($Timers[0][0] - $now) * 1000) + 1;
+
+ # -1 is an infinite timeout, so prefer a real timeout
+ return $timeout if $LoopTimeout == -1;
+
+ # otherwise pick the lower of our regular timeout and time until
+ # the next timer
+ return $LoopTimeout if $LoopTimeout < $timeout;
+ return $timeout;
}
### The epoll-based event loop. Gets installed as EventLoop if IO::Epoll loads
@@ -263,24 +403,18 @@
my $class = shift;
foreach my $fd ( keys %OtherFds ) {
- epoll_ctl($Epoll, EPOLL_CTL_ADD, $fd, EPOLLIN);
+ if (epoll_ctl($Epoll, EPOLL_CTL_ADD, $fd, EPOLLIN) == -1) {
+ warn "epoll_ctl(): failure adding fd=$fd; $! (", $!+0, ")\n";
+ }
}
while (1) {
- my $now = time;
- # Run expired timers
- while (@Timers && $Timers[0][0] <= $now) {
- my $to_run = shift(@Timers);
- $to_run->[1]->($now);
- }
-
- # Get next timeout
- my $timeout = @Timers ? ($Timers[0][0] - $now) : 1;
-
my @events;
my $i;
- my $evcount = epoll_wait($Epoll, 1000, $timeout * 1000, \@events);
-
+ my $timeout = RunTimers();
+
+ # get up to 1000 events
+ my $evcount = epoll_wait($Epoll, 1000, $timeout, \@events);
EVENT:
for ($i=0; $i<$evcount; $i++) {
my $ev = $events[$i];
@@ -298,10 +432,9 @@
if (! $pob) {
if (my $code = $OtherFds{$ev->[0]}) {
$code->($state);
- }
- else {
+ } else {
my $fd = $ev->[0];
- print STDERR "epoll() returned fd $fd w/ state $state for which we have no mapping. removing.\n";
+ warn "epoll() returned fd $fd w/ state $state for which we have no mapping. removing.\n";
POSIX::close($fd);
epoll_ctl($Epoll, EPOLL_CTL_DEL, $fd, 0);
}
@@ -311,12 +444,46 @@
DebugLevel >= 1 && $class->DebugMsg("Event: fd=%d (%s), state=%d \@ %s\n",
$ev->[0], ref($pob), $ev->[1], time);
+ if ($DoProfile) {
+ my $class = ref $pob;
+
+ # call profiling action on things that need to be done
+ if ($state & EPOLLIN && ! $pob->{closed}) {
+ _pre_profile();
+ $pob->event_read;
+ _post_profile("$class-read");
+ }
+
+ if ($state & EPOLLOUT && ! $pob->{closed}) {
+ _pre_profile();
+ $pob->event_write;
+ _post_profile("$class-write");
+ }
+
+ if ($state & (EPOLLERR|EPOLLHUP)) {
+ if ($state & EPOLLERR && ! $pob->{closed}) {
+ _pre_profile();
+ $pob->event_err;
+ _post_profile("$class-err");
+ }
+ if ($state & EPOLLHUP && ! $pob->{closed}) {
+ _pre_profile();
+ $pob->event_hup;
+ _post_profile("$class-hup");
+ }
+ }
+
+ next;
+ }
+
+ # standard non-profiling codepat
$pob->event_read if $state & EPOLLIN && ! $pob->{closed};
$pob->event_write if $state & EPOLLOUT && ! $pob->{closed};
- $pob->event_err if $state & EPOLLERR && ! $pob->{closed};
- $pob->event_hup if $state & EPOLLHUP && ! $pob->{closed};
+ if ($state & (EPOLLERR|EPOLLHUP)) {
+ $pob->event_err if $state & EPOLLERR && ! $pob->{closed};
+ $pob->event_hup if $state & EPOLLHUP && ! $pob->{closed};
+ }
}
-
return unless PostEventLoop();
}
exit 0;
@@ -330,16 +497,8 @@
my Danga::Socket $pob;
while (1) {
- my $now = time;
- # Run expired timers
- while (@Timers && $Timers[0][0] <= $now) {
- my $to_run = shift(@Timers);
- $to_run->[1]->($now);
- }
-
- # Get next timeout
- my $timeout = @Timers ? ($Timers[0][0] - $now) : 1;
-
+ my $timeout = RunTimers();
+
# the following sets up @poll as a series of ($poll,$event_mask)
# items, then uses IO::Poll::_poll, implemented in XS, which
# modifies the array in place with the even elements being
@@ -348,14 +507,23 @@
foreach my $fd ( keys %OtherFds ) {
push @poll, $fd, POLLIN;
}
- foreach my $fd ( keys %DescriptorMap ) {
- my Danga::Socket $sock = $DescriptorMap{$fd};
+ while ( my ($fd, $sock) = each %DescriptorMap ) {
push @poll, $fd, $sock->{event_watch};
}
- return 0 unless @poll;
-
- # print STDERR "Poll for $timeout secs\n";
- my $count = IO::Poll::_poll($timeout * 1000, @poll);
+
+ # if nothing to poll, either end immediately (if no timeout)
+ # or just keep calling the callback
+ unless (@poll) {
+ select undef, undef, undef, ($timeout / 1000);
+ return unless PostEventLoop();
+ next;
+ }
+
+ my $count = IO::Poll::_poll($timeout, @poll);
+ unless ($count) {
+ return unless PostEventLoop();
+ next;
+ }
# Fetch handles with read events
while (@poll) {
@@ -364,8 +532,10 @@
$pob = $DescriptorMap{$fd};
- if ( !$pob && (my $code = $OtherFds{$fd}) ) {
- $code->($state);
+ if (!$pob) {
+ if (my $code = $OtherFds{$fd}) {
+ $code->($state);
+ }
next;
}
@@ -381,8 +551,84 @@
exit 0;
}
-## PostEventLoop is called at the end of the event loop to process things
-# like close() calls.
+### The kqueue-based event loop. Gets installed as EventLoop if IO::KQueue works
+### okay.
+sub KQueueEventLoop {
+ my $class = shift;
+
+ foreach my $fd (keys %OtherFds) {
+ $KQueue->EV_SET($fd, IO::KQueue::EVFILT_READ(), IO::KQueue::EV_ADD());
+ }
+
+ while (1) {
+ my $timeout = RunTimers();
+ my @ret = $KQueue->kevent($timeout);
+ if (!@ret) {
+ foreach my $fd ( keys %DescriptorMap ) {
+ my Danga::Socket $sock = $DescriptorMap{$fd};
+ if ($sock->can('ticker')) {
+ $sock->ticker;
+ }
+ }
+ }
+
+ foreach my $kev (@ret) {
+ my ($fd, $filter, $flags, $fflags) = @$kev;
+ my Danga::Socket $pob = $DescriptorMap{$fd};
+ if (!$pob) {
+ if (my $code = $OtherFds{$fd}) {
+ $code->($filter);
+ } else {
+ warn "kevent() returned fd $fd for which we have no mapping. removing.\n";
+ POSIX::close($fd); # close deletes the kevent entry
+ }
+ next;
+ }
+
+ DebugLevel >= 1 && $class->DebugMsg("Event: fd=%d (%s), flags=%d \@ %s\n",
+ $fd, ref($pob), $flags, time);
+
+ $pob->event_read if $filter == IO::KQueue::EVFILT_READ() && !$pob->{closed};
+ $pob->event_write if $filter == IO::KQueue::EVFILT_WRITE() && !$pob->{closed};
+ if ($flags == IO::KQueue::EV_EOF() && !$pob->{closed}) {
+ if ($fflags) {
+ $pob->event_err;
+ } else {
+ $pob->event_hup;
+ }
+ }
+ }
+ return unless PostEventLoop();
+ }
+
+ exit(0);
+}
+
+### CLASS METHOD: SetPostLoopCallback
+### Sets post loop callback function. Pass a subref and it will be
+### called every time the event loop finishes. Return 1 from the sub
+### to make the loop continue, else it will exit. The function will
+### be passed two parameters: \%DescriptorMap, \%OtherFds.
+sub SetPostLoopCallback {
+ my ($class, $ref) = @_;
+
+ if (ref $class) {
+ # per-object callback
+ my Danga::Socket $self = $class;
+ if (defined $ref && ref $ref eq 'CODE') {
+ $PLCMap{$self->{fd}} = $ref;
+ } else {
+ delete $PLCMap{$self->{fd}};
+ }
+ } else {
+ # global callback
+ $PostLoopCallback = (defined $ref && ref $ref eq 'CODE') ? $ref : undef;
+ }
+}
+
+# Internal function: run the post-event callback, send read events
+# for pushed-back data, and close pending connections. returns 1
+# if event loop should continue, or 0 to shut it all down.
sub PostEventLoop {
# fire read events for objects with pushed-back read data
my $loop = 1;
@@ -390,6 +636,14 @@
$loop = 0;
foreach my $fd (keys %PushBackSet) {
my Danga::Socket $pob = $PushBackSet{$fd};
+
+ # a previous event_read invocation could've closed a
+ # connection that we already evaluated in "keys
+ # %PushBackSet", so skip ones that seem to have
+ # disappeared. this is expected.
+ next unless $pob;
+
+ die "ASSERT: the $pob socket has no read_push_back" unless @{$pob->{read_push_back}};
next unless (! $pob->{closed} &&
$pob->{event_watch} & POLLIN);
$loop = 1;
@@ -400,34 +654,38 @@
# now we can close sockets that wanted to close during our event processing.
# (we didn't want to close them during the loop, as we didn't want fd numbers
# being reused and confused during the event loop)
- foreach my $f (@ToClose) {
- close($f);
+ while (my $sock = shift @ToClose) {
+ my $fd = fileno($sock);
+
+ # close the socket. (not a Danga::Socket close)
+ $sock->close;
+
+ # and now we can finally remove the fd from the map. see
+ # comment above in _cleanup.
+ delete $DescriptorMap{$fd};
}
- @ToClose = ();
- # now we're at the very end, call per-connection callbacks if defined
- my $ret = 1; # use $ret so's to not starve some FDs; return 0 if any PLCs return 0
+
+ # by default we keep running, unless a postloop callback (either per-object
+ # or global) cancels it
+ my $keep_running = 1;
+
+ # per-object post-loop-callbacks
for my $plc (values %PLCMap) {
- $ret &&= $plc->(\%DescriptorMap, \%OtherFds);
+ $keep_running &&= $plc->(\%DescriptorMap, \%OtherFds);
}
- # now we're at the very end, call global callback if defined
+ # now we're at the very end, call callback if defined
if (defined $PostLoopCallback) {
- $ret &&= $PostLoopCallback->(\%DescriptorMap, \%OtherFds);
+ $keep_running &&= $PostLoopCallback->(\%DescriptorMap, \%OtherFds);
}
- return $ret;
-}
-
-### (CLASS) METHOD: DebugMsg( $format, @args )
-### Print the debugging message specified by the C<sprintf>-style I<format> and
-### I<args>
-sub DebugMsg {
- my ( $class, $fmt, @args ) = @_;
- chomp $fmt;
- printf STDERR ">>> $fmt\n", @args;
+ return $keep_running;
}
+#####################################################################
+### Danga::Socket-the-object code
+#####################################################################
### METHOD: new( $socket )
### Create a new Danga::Socket object for the given I<socket> which will react
@@ -440,17 +698,21 @@
$self->{sock} = $sock;
my $fd = fileno($sock);
+
+ Carp::cluck("undef sock and/or fd in Danga::Socket->new. sock=" . ($sock || "") . ", fd=" . ($fd || ""))
+ unless $sock && $fd;
+
$self->{fd} = $fd;
$self->{write_buf} = [];
$self->{write_buf_offset} = 0;
$self->{write_buf_size} = 0;
$self->{closed} = 0;
+ $self->{corked} = 0;
$self->{read_push_back} = [];
- $self->{post_loop_callback} = undef;
$self->{event_watch} = POLLERR|POLLHUP|POLLNVAL;
- init_poller();
+ _InitPoller();
if ($HaveEpoll) {
epoll_ctl($Epoll, EPOLL_CTL_ADD, $fd, $self->{event_watch})
@@ -464,12 +726,14 @@
IO::KQueue::EV_ADD() | IO::KQueue::EV_DISABLE());
}
+ Carp::cluck("Danga::Socket::new blowing away existing descriptor map for fd=$fd ($DescriptorMap{$fd})")
+ if $DescriptorMap{$fd};
+
$DescriptorMap{$fd} = $self;
return $self;
}
-
#####################################################################
### I N S T A N C E M E T H O D S
#####################################################################
@@ -477,54 +741,126 @@
### METHOD: tcp_cork( $boolean )
### Turn TCP_CORK on or off depending on the value of I<boolean>.
sub tcp_cork {
- my Danga::Socket $self = shift;
- my $val = shift;
+ my Danga::Socket $self = $_[0];
+ my $val = $_[1];
- # FIXME: Linux-specific.
- setsockopt($self->{sock}, IPPROTO_TCP, TCP_CORK,
- pack("l", $val ? 1 : 0)) || die "setsockopt: $!";
+ # make sure we have a socket
+ return unless $self->{sock};
+ return if $val == $self->{corked};
+
+ my $rv;
+ if (TCP_CORK) {
+ $rv = setsockopt($self->{sock}, IPPROTO_TCP, TCP_CORK,
+ pack("l", $val ? 1 : 0));
+ } else {
+ # FIXME: implement freebsd *PUSH sockopts
+ $rv = 1;
+ }
+
+ # if we failed, close (if we're not already) and warn about the error
+ if ($rv) {
+ $self->{corked} = $val;
+ } else {
+ if ($! == EBADF || $! == ENOTSOCK) {
+ # internal state is probably corrupted; warn and then close if
+ # we're not closed already
+ warn "setsockopt: $!";
+ $self->close('tcp_cork_failed');
+ } elsif ($! == ENOPROTOOPT) {
+ # TCP implementation doesn't support corking, so just ignore it
+ } else {
+ # some other error; we should never hit here, but if we do, die
+ die "setsockopt: $!";
+ }
+ }
}
-### METHOD: close( [$reason] )
-### Close the socket. The I<reason> argument will be used in debugging messages.
-sub close {
- my Danga::Socket $self = shift;
- my $reason = shift || "";
+### METHOD: steal_socket
+### Basically returns our socket and makes it so that we don't try to close it,
+### but we do remove it from epoll handlers. THIS CLOSES $self. It is the same
+### thing as calling close, except it gives you the socket to use.
+sub steal_socket {
+ my Danga::Socket $self = $_[0];
+ return if $self->{closed};
- my $fd = $self->{fd};
+ # cleanup does most of the work of closing this socket
+ $self->_cleanup();
+
+ # now undef our internal sock and fd structures so we don't use them
my $sock = $self->{sock};
- $self->{closed} = 1;
+ $self->{sock} = undef;
+ return $sock;
+}
- # we need to flush our write buffer, as there may
- # be self-referential closures (sub { $client->close })
- # preventing the object from being destroyed
- $self->{write_buf} = [];
+### METHOD: close( [$reason] )
+### Close the socket. The I<reason> argument will be used in debugging messages.
+sub close {
+ my Danga::Socket $self = $_[0];
+ return if $self->{closed};
+ # print out debugging info for this close
if (DebugLevel) {
my ($pkg, $filename, $line) = caller;
- print STDERR "Closing \#$fd due to $pkg/$filename/$line ($reason)\n";
- }
-
- if ($HaveEpoll) {
- if (epoll_ctl($Epoll, EPOLL_CTL_DEL, $fd, $self->{event_watch}) == 0) {
- DebugLevel >= 1 && $self->debugmsg("Client %d disconnected.\n", $fd);
- } else {
- DebugLevel >= 1 && $self->debugmsg("poll->remove failed on fd %d\n", $fd);
- }
+ my $reason = $_[1] || "";
+ warn "Closing \#$self->{fd} due to $pkg/$filename/$line ($reason)\n";
}
- delete $PLCMap{$fd};
- delete $DescriptorMap{$fd};
- delete $PushBackSet{$fd};
+ # this does most of the work of closing us
+ $self->_cleanup();
# defer closing the actual socket until the event loop is done
# processing this round of events. (otherwise we might reuse fds)
- push @ToClose, $sock;
+ if ($self->{sock}) {
+ push @ToClose, $self->{sock};
+ $self->{sock} = undef;
+ }
return 0;
}
+### METHOD: _cleanup()
+### Called by our closers so we can clean internal data structures.
+sub _cleanup {
+ my Danga::Socket $self = $_[0];
+ # we're effectively closed; we have no fd and sock when we leave here
+ $self->{closed} = 1;
+
+ # we need to flush our write buffer, as there may
+ # be self-referential closures (sub { $client->close })
+ # preventing the object from being destroyed
+ $self->{write_buf} = [];
+
+ # uncork so any final data gets sent. only matters if the person closing
+ # us forgot to do it, but we do it to be safe.
+ $self->tcp_cork(0);
+
+ # if we're using epoll, we have to remove this from our epoll fd so we stop getting
+ # notifications about it
+ if ($HaveEpoll && $self->{fd}) {
+ if (epoll_ctl($Epoll, EPOLL_CTL_DEL, $self->{fd}, $self->{event_watch}) != 0) {
+ # dump_error prints a backtrace so we can try to figure out why this happened
+ $self->dump_error("epoll_ctl(): failure deleting fd=$self->{fd} during _cleanup(); $! (" . ($!+0) . ")");
+ }
+ }
+
+ # now delete from mappings. this fd no longer belongs to us, so we don't want
+ # to get alerts for it if it becomes writable/readable/etc.
+ delete $PushBackSet{$self->{fd}};
+ delete $PLCMap{$self->{fd}};
+
+ # we explicitly don't delete from DescriptorMap here until we
+ # actually close the socket, as we might be in the middle of
+ # processing an epoll_wait/etc that returned hundreds of fds, one
+ # of which is not yet processed and is what we're closing. if we
+ # keep it in DescriptorMap, then the event harnesses can just
+ # looked at $pob->{closed} and ignore it. but if it's an
+ # un-accounted for fd, then it (understandably) freak out a bit
+ # and emit warnings, thinking their state got off.
+
+ # and finally get rid of our fd so we can't use it anywhere else
+ $self->{fd} = undef;
+}
### METHOD: sock()
### Returns the underlying IO::Handle for the object.
@@ -533,6 +869,12 @@
return $self->{sock};
}
+sub set_writer_func {
+ my Danga::Socket $self = shift;
+ my $wtr = shift;
+ Carp::croak("Not a subref") unless !defined $wtr || ref $wtr eq "CODE";
+ $self->{writer_func} = $wtr;
+}
### METHOD: write( $data )
### Write the specified data to the underlying handle. I<data> may be scalar,
@@ -587,6 +929,12 @@
shift @{$self->{write_buf}};
}
$bref->();
+
+ # code refs are just run and never get reenqueued
+ # (they're one-shot), so turn off the flag indicating the
+ # outstanding data needs queueing.
+ $need_queue = 0;
+
undef $bref;
next WRITE;
}
@@ -594,7 +942,12 @@
}
my $to_write = $len - $self->{write_buf_offset};
- my $written = syswrite($self->{sock}, $$bref, $to_write, $self->{write_buf_offset});
+ my $written;
+ if (my $wtr = $self->{writer_func}) {
+ $written = $wtr->($bref, $to_write, $self->{write_buf_offset});
+ } else {
+ $written = syswrite($self->{sock}, $$bref, $to_write, $self->{write_buf_offset});
+ }
if (! defined $written) {
if ($! == EPIPE) {
@@ -626,7 +979,7 @@
# interested in pending writes:
$self->{write_buf_offset} += $written;
$self->{write_buf_size} -= $written;
- $self->watch_write(1);
+ $self->on_incomplete_write;
return 0;
} elsif ($written == $to_write) {
DebugLevel >= 2 && $self->debugmsg("Wrote ALL %d bytes to %d (nq=%d)",
@@ -647,6 +1000,11 @@
}
}
+sub on_incomplete_write {
+ my Danga::Socket $self = shift;
+ $self->watch_write(1);
+}
+
### METHOD: push_back_read( $buf )
### Push back I<buf> (a scalar or scalarref) into the read stream
sub push_back_read {
@@ -656,17 +1014,6 @@
$PushBackSet{$self->{fd}} = $self;
}
-### METHOD: shift_back_read( $buf )
-### Shift back I<buf> (a scalar or scalarref) into the read stream
-### Use this instead of push_back_read() when you need to unread
-### something you just read.
-sub shift_back_read {
- my Danga::Socket $self = shift;
- my $buf = shift;
- unshift @{$self->{read_push_back}}, ref $buf ? $buf : \$buf;
- $PushBackSet{$self->{fd}} = $self;
-}
-
### METHOD: read( $bytecount )
### Read at most I<bytecount> bytes from the underlying handle; returns scalar
### ref on read, or undef on connection closed.
@@ -679,21 +1026,23 @@
if (@{$self->{read_push_back}}) {
$buf = shift @{$self->{read_push_back}};
my $len = length($$buf);
- if ($len <= $buf) {
- unless (@{$self->{read_push_back}}) {
- delete $PushBackSet{$self->{fd}};
- }
+
+ if ($len <= $bytes) {
+ delete $PushBackSet{$self->{fd}} unless @{$self->{read_push_back}};
return $buf;
} else {
# if the pushed back read is too big, we have to split it
my $overflow = substr($$buf, $bytes);
$buf = substr($$buf, 0, $bytes);
- unshift @{$self->{read_push_back}}, \$overflow,
+ unshift @{$self->{read_push_back}}, \$overflow;
return \$buf;
}
}
- my $res = sysread($sock, $buf, $bytes, 0);
+ # max 5MB, or perl quits(!!)
+ my $req_bytes = $bytes > 5242880 ? 5242880 : $bytes;
+
+ my $res = sysread($sock, $buf, $req_bytes, 0);
DebugLevel >= 2 && $self->debugmsg("sysread = %d; \$! = %d", $res, $!);
if (! $res && $! != EWOULDBLOCK) {
@@ -741,14 +1090,14 @@
### Turn 'readable' event notification on or off.
sub watch_read {
my Danga::Socket $self = shift;
- return if $self->{closed};
+ return if $self->{closed} || !$self->{sock};
my $val = shift;
my $event = $self->{event_watch};
-
+
$event &= ~POLLIN if ! $val;
$event |= POLLIN if $val;
-
+
# If it changed, set it
if ($event != $self->{event_watch}) {
if ($HaveKQueue) {
@@ -757,22 +1106,22 @@
}
elsif ($HaveEpoll) {
epoll_ctl($Epoll, EPOLL_CTL_MOD, $self->{fd}, $event)
- and print STDERR "couldn't modify epoll settings for $self->{fd} " .
- "($self) from $self->{event_watch} -> $event\n";
+ and $self->dump_error("couldn't modify epoll settings for $self->{fd} " .
+ "from $self->{event_watch} -> $event: $! (" . ($!+0) . ")");
}
$self->{event_watch} = $event;
}
}
-### METHOD: watch_read( $boolean )
+### METHOD: watch_write( $boolean )
### Turn 'writable' event notification on or off.
sub watch_write {
my Danga::Socket $self = shift;
- return if $self->{closed};
+ return if $self->{closed} || !$self->{sock};
my $val = shift;
my $event = $self->{event_watch};
-
+
$event &= ~POLLOUT if ! $val;
$event |= POLLOUT if $val;
@@ -784,13 +1133,28 @@
}
elsif ($HaveEpoll) {
epoll_ctl($Epoll, EPOLL_CTL_MOD, $self->{fd}, $event)
- and print STDERR "couldn't modify epoll settings for $self->{fd} " .
- "($self) from $self->{event_watch} -> $event\n";
+ and $self->dump_error("couldn't modify epoll settings for $self->{fd} " .
+ "from $self->{event_watch} -> $event: $! (" . ($!+0) . ")");
}
$self->{event_watch} = $event;
}
}
+# METHOD: dump_error( $message )
+# Prints to STDERR a backtrace with information about this socket and what lead
+# up to the dump_error call.
+sub dump_error {
+ my $i = 0;
+ my @list;
+ while (my ($file, $line, $sub) = (caller($i++))[1..3]) {
+ push @list, "\t$file:$line called $sub\n";
+ }
+
+ warn "ERROR: $_[1]\n" .
+ "\t$_[0] = " . $_[0]->as_string . "\n" .
+ join('', @list);
+}
+
### METHOD: debugmsg( $format, @args )
### Print the debugging message specified by the C<sprintf>-style I<format> and
@@ -809,12 +1173,16 @@
### Returns the string describing the peer's IP
sub peer_ip_string {
my Danga::Socket $self = shift;
- return $self->{peer_ip} if defined $self->{peer_ip};
- my $pn = getpeername($self->{sock}) or return undef;
+ return _undef("peer_ip_string undef: no sock") unless $self->{sock};
+ return $self->{peer_ip} if defined $self->{peer_ip};
+
+ my $pn = getpeername($self->{sock});
+ return _undef("peer_ip_string undef: getpeername") unless $pn;
+
my ($port, $iaddr) = Socket::sockaddr_in($pn);
- my $r = Socket::inet_ntoa($iaddr);
- $self->{peer_ip} = $r;
- return $r;
+ $self->{peer_port} = $port;
+
+ return $self->{peer_ip} = Socket::inet_ntoa($iaddr);
}
### METHOD: peer_addr_string()
@@ -822,16 +1190,43 @@
### object in form "ip:port"
sub peer_addr_string {
my Danga::Socket $self = shift;
- my $pn = getpeername($self->{sock}) or return undef;
+ my $ip = $self->peer_ip_string;
+ return $ip ? "$ip:$self->{peer_port}" : undef;
+}
+
+### METHOD: local_ip_string()
+### Returns the string describing the local IP
+sub local_ip_string {
+ my Danga::Socket $self = shift;
+ return _undef("local_ip_string undef: no sock") unless $self->{sock};
+ return $self->{local_ip} if defined $self->{local_ip};
+
+ my $pn = getsockname($self->{sock});
+ return _undef("local_ip_string undef: getsockname") unless $pn;
+
my ($port, $iaddr) = Socket::sockaddr_in($pn);
- return Socket::inet_ntoa($iaddr) . ":$port";
+ $self->{local_port} = $port;
+
+ return $self->{local_ip} = Socket::inet_ntoa($iaddr);
}
+### METHOD: local_addr_string()
+### Returns the string describing the local end of the socket which underlies this
+### object in form "ip:port"
+sub local_addr_string {
+ my Danga::Socket $self = shift;
+ my $ip = $self->local_ip_string;
+ return $ip ? "$ip:$self->{local_port}" : undef;
+}
+
+
### METHOD: as_string()
### Returns a string describing this socket.
sub as_string {
my Danga::Socket $self = shift;
- my $ret = ref($self) . ": " . ($self->{closed} ? "closed" : "open");
+ my $rw = "(" . ($self->{event_watch} & POLLIN ? 'R' : '') .
+ ($self->{event_watch} & POLLOUT ? 'W' : '') . ")";
+ my $ret = ref($self) . "$rw: " . ($self->{closed} ? "closed" : "open");
my $peer = $self->peer_addr_string;
if ($peer) {
$ret .= " to " . $self->peer_addr_string;
@@ -839,140 +1234,15 @@
return $ret;
}
-### CLASS METHOD: SetPostLoopCallback
-### Sets post loop callback function. Pass a subref and it will be
-### called every time the event loop finishes. Return 1 from the sub
-### to make the loop continue, else it will exit. The function will
-### be passed two parameters: \%DescriptorMap, \%OtherFds.
-sub SetPostLoopCallback {
- my ($class, $ref) = @_;
- if(ref $class) {
- my Danga::Socket $self = $class;
- if( defined $ref && ref $ref eq 'CODE' ) {
- $PLCMap{$self->{fd}} = $ref;
- }
- else {
- delete $PLCMap{$self->{fd}};
- }
- }
- else {
- $PostLoopCallback = (defined $ref && ref $ref eq 'CODE') ? $ref : undef;
- }
-}
-
-sub DESTROY {
- my Danga::Socket $self = shift;
- $self->close() if !$self->{closed};
-}
-
-#####################################################################
-### U T I L I T Y F U N C T I O N S
-#####################################################################
-
-our ($SYS_epoll_create, $SYS_epoll_ctl, $SYS_epoll_wait);
-
-if ($^O eq "linux") {
- my ($sysname, $nodename, $release, $version, $machine) = POSIX::uname();
-
- # whether the machine requires 64-bit numbers to be on 8-byte
- # boundaries.
- my $u64_mod_8 = 0;
-
- if ($machine =~ m/^i[3456]86$/) {
- $SYS_epoll_create = 254;
- $SYS_epoll_ctl = 255;
- $SYS_epoll_wait = 256;
- } elsif ($machine eq "x86_64") {
- $SYS_epoll_create = 213;
- $SYS_epoll_ctl = 233;
- $SYS_epoll_wait = 232;
- } elsif ($machine eq "ppc64") {
- $SYS_epoll_create = 236;
- $SYS_epoll_ctl = 237;
- $SYS_epoll_wait = 238;
- $u64_mod_8 = 1;
- } elsif ($machine eq "ppc") {
- $SYS_epoll_create = 236;
- $SYS_epoll_ctl = 237;
- $SYS_epoll_wait = 238;
- $u64_mod_8 = 1;
- } elsif ($machine eq "ia64") {
- $SYS_epoll_create = 1243;
- $SYS_epoll_ctl = 1244;
- $SYS_epoll_wait = 1245;
- $u64_mod_8 = 1;
- }
-
- if ($u64_mod_8) {
- *epoll_wait = \&epoll_wait_mod8;
- *epoll_ctl = \&epoll_ctl_mod8;
- } else {
- *epoll_wait = \&epoll_wait_mod4;
- *epoll_ctl = \&epoll_ctl_mod4;
- }
-
- # if syscall numbers have been defined (and this module has been
- # tested on) the arch above, then try to use it. try means see if
- # the syscall is implemented. it may well be that this is Linux
- # 2.4 and we don't even have it available.
- $TryEpoll = 1 if $SYS_epoll_create;
-}
-
-# epoll_create wrapper
-# ARGS: (size)
-sub epoll_create {
- my $epfd = eval { syscall($SYS_epoll_create, $_[0]) };
- return -1 if $@;
- return $epfd;
-}
-
-# epoll_ctl wrapper
-# ARGS: (epfd, op, fd, events_mask)
-sub epoll_ctl_mod4 {
- syscall($SYS_epoll_ctl, $_[0]+0, $_[1]+0, $_[2]+0, pack("LLL", $_[3], $_[2], 0));
-}
-sub epoll_ctl_mod8 {
- syscall($SYS_epoll_ctl, $_[0]+0, $_[1]+0, $_[2]+0, pack("LLLL", $_[3], 0, $_[2], 0));
-}
-
-# epoll_wait wrapper
-# ARGS: (epfd, maxevents, timeout (milliseconds), arrayref)
-# arrayref: values modified to be [$fd, $event]
-our $epoll_wait_events;
-our $epoll_wait_size = 0;
-sub epoll_wait_mod4 {
- # resize our static buffer if requested size is bigger than we've ever done
- if ($_[1] > $epoll_wait_size) {
- $epoll_wait_size = $_[1];
- $epoll_wait_events = "\0" x 12 x $epoll_wait_size;
- }
- my $ct = syscall($SYS_epoll_wait, $_[0]+0, $epoll_wait_events, $_[1]+0, $_[2]+0);
- for ($_ = 0; $_ < $ct; $_++) {
- @{$_[3]->[$_]}[1,0] = unpack("LL", substr($epoll_wait_events, 12*$_, 8));
- }
- return $ct;
-}
-
-sub epoll_wait_mod8 {
- # resize our static buffer if requested size is bigger than we've ever done
- if ($_[1] > $epoll_wait_size) {
- $epoll_wait_size = $_[1];
- $epoll_wait_events = "\0" x 16 x $epoll_wait_size;
- }
- my $ct = syscall($SYS_epoll_wait, $_[0]+0, $epoll_wait_events, $_[1]+0, $_[2]+0);
- for ($_ = 0; $_ < $ct; $_++) {
- # 16 byte epoll_event structs, with format:
- # 4 byte mask [idx 1]
- # 4 byte padding (we put it into idx 2, useless)
- # 8 byte data (first 4 bytes are fd, into idx 0)
- @{$_[3]->[$_]}[1,2,0] = unpack("LLL", substr($epoll_wait_events, 16*$_, 12));
- }
- return $ct;
+sub _undef {
+ return undef unless $ENV{DS_DEBUG};
+ my $msg = shift || "";
+ warn "Danga::Socket: $msg\n";
+ return undef;
}
1;
-
# Local Variables:
# mode: perl
# c-basic-indent: 4
Modified: trunk/lib/Qpsmtpd.pm
==============================================================================
--- trunk/lib/Qpsmtpd.pm (original)
+++ trunk/lib/Qpsmtpd.pm Tue Jun 20 06:51:32 2006
@@ -337,7 +337,7 @@
@r = $self->run_hook($hook, $code, @_);
next unless @r;
if ($r[0] == CONTINUATION) {
- $self->disable_read() if $self->isa('Danga::Client');
+ $self->pause_read() if $self->isa('Danga::Client');
$self->{_continuation} = [$hook, [@_], @local_hooks];
}
last unless $r[0] == DECLINED;
@@ -351,7 +351,7 @@
sub finish_continuation {
my ($self) = @_;
die "No continuation in progress" unless $self->{_continuation};
- $self->enable_read() if $self->isa('Danga::Client');
+ $self->continue_read() if $self->isa('Danga::Client');
my $todo = $self->{_continuation};
$self->{_continuation} = undef;
my $hook = shift @$todo || die "No hook in the continuation";
@@ -361,7 +361,7 @@
my $code = shift @$todo;
@r = $self->run_hook($hook, $code, @$args);
if ($r[0] == CONTINUATION) {
- $self->disable_read() if $self->isa('Danga::Client');
+ $self->pause_read() if $self->isa('Danga::Client');
$self->{_continuation} = [$hook, $args, @$todo];
return @r;
}
Modified: trunk/plugins/check_earlytalker
==============================================================================
--- trunk/plugins/check_earlytalker (original)
+++ trunk/plugins/check_earlytalker Tue Jun 20 06:51:32 2006
@@ -44,6 +44,15 @@
is to react at the SMTP greeting stage by issuing the apropriate response code
and terminating the SMTP connection.
+=item check-at [string: connect, data]
+
+Defines when to check for early talkers, either at connect time (pre-greet pause)
+or at DATA time (pause before sending "354 go ahead").
+
+The default is I<connect>.
+
+Note that defer-reject has no meaning if check-at is I<data>.
+
=back
=cut
@@ -61,23 +70,27 @@
'wait' => 1,
'action' => 'denysoft',
'defer-reject' => 0,
+ 'check-at' => 'connect',
@args,
};
+ print STDERR "Check at: ", $self->{_args}{'check-at'}, "\n";
if ($qp->isa('Qpsmtpd::Apache')) {
require APR::Const;
APR::Const->import(qw(POLLIN SUCCESS));
- $self->register_hook('connect', 'hook_connect_apr');
+ $self->register_hook($self->{_args}->{'check-at'}, 'check_talker_apr');
}
else {
- $self->register_hook('connect', 'hook_connect');
+ $self->register_hook($self->{_args}->{'check-at'}, 'check_talker_poll');
+ }
+ $self->register_hook($self->{_args}->{'check-at'}, 'check_talker_post');
+ if ($self->{_args}{'check-at'} eq 'connect') {
+ $self->register_hook('mail', 'hook_mail')
+ if $self->{_args}->{'defer-reject'};
}
- $self->register_hook('connect', 'hook_connect_post');
- $self->register_hook('mail', 'hook_mail')
- if $self->{_args}->{'defer-reject'};
1;
}
-sub hook_connect_apr {
+sub check_talker_apr {
my ($self, $transaction) = @_;
return DECLINED if ($self->qp->connection->notes('whitelistclient'));
@@ -104,29 +117,27 @@
return DECLINED;
}
-sub hook_connect {
+sub check_talker_poll {
my ($self, $transaction) = @_;
my $qp = $self->qp;
my $conn = $qp->connection;
- $qp->AddTimer($self->{_args}{'wait'}, sub { read_now($qp, $conn) });
+ $qp->AddTimer($self->{_args}{'wait'}, sub { read_now($qp, $conn, $self->{_args}{'check-at'}) });
return CONTINUATION;
}
sub read_now {
- my ($qp, $conn) = @_;
+ my ($qp, $conn, $phase) = @_;
- if (my $data = $qp->read(1024)) {
- if (length($$data)) {
+ if ($qp->has_data) {
$qp->log(LOGNOTICE, 'remote host started talking before we said hello');
- $qp->push_back_read($data);
+ $qp->clear_data if $phase eq 'data';
$conn->notes('earlytalker', 1);
- }
}
$qp->finish_continuation;
}
-sub hook_connect_post {
+sub check_talker_post {
my ($self, $transaction) = @_;
my $conn = $self->qp->connection;
Modified: trunk/qpsmtpd
==============================================================================
--- trunk/qpsmtpd (original)
+++ trunk/qpsmtpd Tue Jun 20 06:51:32 2006
@@ -35,7 +35,6 @@
my $PORT = 2525;
my $LOCALADDR = '0.0.0.0';
-my $LineMode = 0;
my $PROCS = 1;
my $MAXCONN = 15; # max simultaneous connections
my $USER = 'smtpd'; # user to suid to
@@ -54,7 +53,6 @@
-c, --limit-connections N : limit concurrent connections to N; default 15
-u, --user U : run as a particular user; defualt 'smtpd'
-m, --max-from-ip M : limit connections from a single IP; default 5
- -f, --forkmode : fork a child for each connection
-j, --procs J : spawn J processes; default 1
-a, --accept K : accept up to K conns per loop; default 20
-h, --help : this page
@@ -73,7 +71,6 @@
'l|listen-address=s' => \$LOCALADDR,
'j|procs=i' => \$PROCS,
'd|debug+' => \$DEBUG,
- 'f|forkmode' => \$LineMode,
'c|limit-connections=i' => \$MAXCONN,
'm|max-from-ip=i' => \$MAXCONNIP,
'u|user=s' => \$USER,
@@ -90,8 +87,6 @@
if ($PROCS =~ /^(\d+)$/) { $PROCS = $1 } else { &help }
if ($NUMACCEPT =~ /^(\d+)$/) { $NUMACCEPT = $1 } else { &help }
my $_NUMACCEPT = $NUMACCEPT;
-$::LineMode = $LineMode;
-$PROCS = 1 if $LineMode;
# This is a bit of a hack, but we get to approximate MAXCONN stuff when we
# have multiple children listening on the same socket.
$MAXCONN /= $PROCS;
@@ -102,7 +97,7 @@
$Danga::Socket::HaveKQueue = 0;
}
-Danga::Socket::init_poller();
+# Danga::Socket::init_poller();
my $POLL = "with " . ($Danga::Socket::HaveEpoll ? "epoll()" :
$Danga::Socket::HaveKQueue ? "kqueue()" : "poll()");
@@ -110,12 +105,6 @@
my $SERVER;
my $CONFIG_SERVER;
-# Code for inetd/tcpserver mode
-if ($ENV{REMOTE_HOST} or $ENV{TCPREMOTEHOST}) {
- run_as_inetd();
- exit(0);
-}
-
my %childstatus = ();
run_as_server();
@@ -165,8 +154,7 @@
print "child $child died\n";
delete $childstatus{$child};
}
- return if $LineMode;
- # restart a new child if in poll server mode
+ # restart a new child (assuming this one died)
spawn_child();
$SIG{CHLD} = \&sig_chld;
}
@@ -177,33 +165,6 @@
exit(0);
}
-sub run_as_inetd {
- $LineMode = $::LineMode = 1;
-
- my $insock = IO::Handle->new_from_fd(0, "r");
- IO::Handle::blocking($insock, 0);
-
- my $outsock = IO::Handle->new_from_fd(1, "w");
- IO::Handle::blocking($outsock, 0);
-
- my $client = Danga::Client->new($insock);
-
- my $out = Qpsmtpd::PollServer->new($outsock);
- $out->load_plugins;
- $out->input_sock($client);
- $client->push_back_read("Connect\n");
- # Cause poll/kevent/epoll to end quickly in first iteration
- Qpsmtpd::PollServer->AddTimer(1, sub { });
-
- while (1) {
- $client->enable_read;
- my $line = $client->get_line;
- last if !defined($line);
- my $output = $out->process_line($line);
- $out->write($output) if $output;
- }
-}
-
sub run_as_server {
local $::MAXconn = $MAXCONN;
# establish SERVER socket, bind and listen.
@@ -261,11 +222,7 @@
sleep while (1);
}
else {
- if ($LineMode) {
- $SIG{INT} = $SIG{TERM} = \&HUNTSMAN;
- }
- $plugin_loader->log(LOGDEBUG, "Listening on $PORT with single process $POLL" .
- ($LineMode ? " (forking server)" : ""));
+ $plugin_loader->log(LOGDEBUG, "Listening on $PORT with single process $POLL");
Qpsmtpd::PollServer->OtherFds(fileno($SERVER) => \&accept_handler,
fileno($CONFIG_SERVER) => \&config_handler,
);
@@ -298,13 +255,8 @@
# Accept all new connections
sub accept_handler {
my $running;
- if( $LineMode ) {
- $running = scalar keys %childstatus;
- }
- else {
- my $descriptors = Danga::Client->DescriptorMap;
- $running = scalar keys %$descriptors;
- }
+ my $descriptors = Danga::Client->DescriptorMap;
+ $running = scalar keys %$descriptors;
for (1 .. $NUMACCEPT) {
if ($running >= $MAXCONN) {
@@ -349,93 +301,43 @@
IO::Handle::blocking($csock, 0);
setsockopt($csock, IPPROTO_TCP, TCP_NODELAY, pack("l", 1)) or die;
- if (!$LineMode) {
- # multiplex mode
- my $client = Qpsmtpd::PollServer->new($csock);
- my $rem_ip = $client->peer_ip_string;
-
- if ($PAUSED) {
- $client->write("451 Sorry, this server is currently paused\r\n");
- $client->close;
- return 1;
- }
-
- if ($MAXCONNIP) {
- my $num_conn = 1; # seed with current value
+ # multiplex mode
+ my $client = Qpsmtpd::PollServer->new($csock);
+ my $rem_ip = $client->peer_ip_string;
- # If we for-loop directly over values %childstatus, a SIGCHLD
- # can call REAPER and slip $rip out from under us. Causes
- # "Use of freed value in iteration" under perl 5.8.4.
- my $descriptors = Danga::Client->DescriptorMap;
- my @obj = values %$descriptors;
- foreach my $obj (@obj) {
- local $^W;
- # This is a bit of a slow way to do this. Wish I could cache the method call.
- ++$num_conn if ($obj->peer_ip_string eq $rem_ip);
- }
-
- if ($num_conn > $MAXCONNIP) {
- $client->log(LOGINFO,"Too many connections from $rem_ip: "
- ."$num_conn > $MAXCONNIP. Denying connection.");
- $client->write("451 Sorry, too many connections from $rem_ip, try again later\r\n");
- $client->close;
- return 1;
- }
- $client->log(LOGINFO, "accepted connection $running/$MAXCONN ($num_conn/$MAXCONNIP) from $rem_ip");
- }
-
- $client->push_back_read("Connect\n");
- $client->watch_read(1);
+ if ($PAUSED) {
+ $client->write("451 Sorry, this server is currently paused\r\n");
+ $client->close;
return 1;
}
-
- # fork-per-connection mode
- my $rem_ip = $csock->sockhost();
if ($MAXCONNIP) {
my $num_conn = 1; # seed with current value
- my @rip = values %childstatus;
- foreach my $rip (@rip) {
- ++$num_conn if (defined $rip && $rip eq $rem_ip);
+ # If we for-loop directly over values %childstatus, a SIGCHLD
+ # can call REAPER and slip $rip out from under us. Causes
+ # "Use of freed value in iteration" under perl 5.8.4.
+ my $descriptors = Danga::Client->DescriptorMap;
+ my @obj = values %$descriptors;
+ foreach my $obj (@obj) {
+ local $^W;
+ # This is a bit of a slow way to do this. Wish I could cache the method call.
+ ++$num_conn if ($obj->peer_ip_string eq $rem_ip);
}
if ($num_conn > $MAXCONNIP) {
- ::log(LOGINFO,"Too many connections from $rem_ip: "
+ $client->log(LOGINFO,"Too many connections from $rem_ip: "
."$num_conn > $MAXCONNIP. Denying connection.");
- print $csock "451 Sorry, too many connections from $rem_ip, try again later\r\n";
- close $csock;
+ $client->write("451 Sorry, too many connections from $rem_ip, try again later\r\n");
+ $client->close;
return 1;
}
+ $client->log(LOGINFO, "accepted connection $running/$MAXCONN ($num_conn/$MAXCONNIP) from $rem_ip");
}
- if (my $pid = _fork) {
- $childstatus{$pid} = $rem_ip;
- return $csock->close();
- }
-
- $SERVER->close(); # make sure the child doesn't accept() new connections
-
- $SIG{$_} = 'DEFAULT' for keys %SIG;
-
- my $client = Qpsmtpd::PollServer->new($csock);
$client->push_back_read("Connect\n");
- # Cause poll/kevent/epoll to end quickly in first iteration
- Qpsmtpd::PollServer->AddTimer(0.1, sub { });
-
- while (1) {
- $client->enable_read;
- my $line = $client->get_line;
- last if !defined($line);
- my $resp = $client->process_line($line);
- $client->write($resp) if $resp;
- }
-
- $client->log(LOGDEBUG, "Finished with child %d.\n", fileno($csock))
- if $DEBUG;
- $client->close();
-
- exit;
+ $client->watch_read(1);
+ return 1;
}
########################################################################