Created
August 30, 2010 23:02
-
-
Save redneb/558192 to your computer and use it in GitHub Desktop.
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| diff --git a/lib/Starman/InputStream.pm b/lib/Starman/InputStream.pm | |
| new file mode 100644 | |
| index 0000000..e0cfc38 | |
| --- /dev/null | |
| +++ b/lib/Starman/InputStream.pm | |
| @@ -0,0 +1,121 @@ | |
| +use strict; | |
| +use warnings; | |
| + | |
| +# ----------------------------------- | |
| + | |
| +package Starman::InputStream; | |
| + | |
| +sub finalize { | |
| + my $self = shift; | |
| + | |
| + # consume what's left | |
| + my $buf=''; | |
| + while (1) { | |
| + my $read = $self->read($buf, 65536); | |
| + return unless defined $read; | |
| + return 1 if $read == 0; | |
| + } | |
| +} | |
| + | |
| +# ----------------------------------- | |
| + | |
| +package Starman::InputStream::Empty; | |
| +use parent 'Starman::InputStream'; | |
| + | |
| +sub new { | |
| + bless \do {my $x}, $_[0]; | |
| +} | |
| + | |
| +sub read { | |
| + return 0; | |
| +} | |
| + | |
| +sub finalize { | |
| + return 1; | |
| +} | |
| + | |
| +# ----------------------------------- | |
| + | |
| +package Starman::InputStream::Identity; | |
| +use parent 'Starman::InputStream'; | |
| + | |
| +use constant { | |
| + READ_SOCK => 0, | |
| + LENGTH => 1, | |
| +}; | |
| + | |
| +sub new { | |
| + my($class, $readsocket, $length) = @_; | |
| + | |
| + bless [$readsocket, $length], $class; | |
| +} | |
| + | |
| +sub read { | |
| + my($self, $_buf, $len, $offset) = @_; | |
| + | |
| + # make sure that we won't read past the end of the message body | |
| + $len = $self->[LENGTH] if $len > $self->[LENGTH]; | |
| + | |
| + my $read = $self->[READ_SOCK]->read($_[1], $len, $offset); | |
| + $self->[LENGTH] -= $read if defined $read; | |
| + | |
| + return $read; | |
| +} | |
| + | |
| +# ----------------------------------- | |
| + | |
| +package Starman::InputStream::Chunked; | |
| +use parent 'Starman::InputStream'; | |
| + | |
| +use constant { | |
| + READ_SOCK => 0, | |
| + EOF => 1, | |
| + CHUNK_SIZE => 2, | |
| +}; | |
| + | |
| +sub new { | |
| + my($class, $readsocket) = @_; | |
| + | |
| + bless [$readsocket, 0, 0], $class; | |
| +} | |
| + | |
| +sub read { | |
| + my($self, $_buf, $len, $offset) = @_; | |
| + | |
| + if ($self->[EOF]) { | |
| + return 0 | |
| + } | |
| + else { | |
| + # $self->[CHUNK_SIZE] > 0 at this point means that we partially read the chunk in the previous run but we did not read all of it | |
| + # therefore there is no chunk header left to parse | |
| + | |
| + if ($self->[CHUNK_SIZE] == 0) { | |
| + my $chunk_header = $self->[READ_SOCK]->readrecord(qr/^[0-9a-fA-F]{1,8}.*?\r?\n/, 2048); | |
| + return unless defined $chunk_header; | |
| + $chunk_header =~ /^([0-9a-fA-F]{1,8}).*?\r?\n/; | |
| + $self->[CHUNK_SIZE] = hex $1; | |
| + } | |
| + | |
| + if ($self->[CHUNK_SIZE] == 0) { | |
| + $self->[EOF] = 1; | |
| + | |
| + # discard the trailer | |
| + $self->[READ_SOCK]->readrecord(qr/^(?:.+?\r?\n)*?\r?\n/, 2048); | |
| + | |
| + return 0; | |
| + } | |
| + else { | |
| + $len = $self->[CHUNK_SIZE] if $self->[CHUNK_SIZE] < $len; | |
| + my $read = $self->[READ_SOCK]->read($_[1], $len); | |
| + $self->[CHUNK_SIZE] -= $read if defined $read; | |
| + if ($self->[CHUNK_SIZE] == 0) { | |
| + # discard the CRLF at the end of the chunk | |
| + $self->[READ_SOCK]->readrecord(qr/^\r?\n/, 2) or return; | |
| + } | |
| + | |
| + return $read; | |
| + } | |
| + } | |
| +} | |
| + | |
| +1; | |
| diff --git a/lib/Starman/ReadSocket.pm b/lib/Starman/ReadSocket.pm | |
| new file mode 100644 | |
| index 0000000..4646706 | |
| --- /dev/null | |
| +++ b/lib/Starman/ReadSocket.pm | |
| @@ -0,0 +1,67 @@ | |
| +package Starman::ReadSocket; | |
| +use strict; | |
| +use warnings; | |
| + | |
| +use constant { | |
| + SOCKET => 0, | |
| + BUF => 1, | |
| + READSIZE => 64 * 1024, | |
| + READ_TIMEOUT => 5, | |
| +}; | |
| + | |
| +sub new { | |
| + bless [$_[1], ''], $_[0]; | |
| +} | |
| + | |
| +sub read { | |
| + my ($self, $_buf, $len, $offset) = @_; | |
| + $offset ||= 0; | |
| + | |
| + my $buflen = length $self->[BUF]; | |
| + if ($buflen == 0) { | |
| + my $read; | |
| + eval { | |
| + # if SIGALRM is received then we will die() and $read will be undef | |
| + # undef indicates read error, so this is fine | |
| + local $SIG{ALRM} = sub { die "\n" }; | |
| + alarm READ_TIMEOUT; | |
| + $read = sysread $self->[SOCKET], $_[1], $len, $offset; | |
| + alarm 0; | |
| + }; | |
| + return $read; | |
| + } | |
| + else { | |
| + $len = $buflen if $buflen < $len; | |
| + $_[1] = '' unless defined $_[1]; | |
| + substr $_[1], $offset, length($_[1]), (substr $self->[BUF], 0, $len, ''); | |
| + return $len; | |
| + } | |
| +} | |
| + | |
| +# this is like readline only that it accepts a regex as a record separator | |
| +# returns undef or a complete record (i.e. a string terminated with $sep) | |
| +sub readrecord { | |
| + my ($self, $sep, $maxlen) = @_; | |
| + | |
| + while (1) { | |
| + return substr $self->[BUF], 0, $+[0], '' if $self->[BUF] =~ $sep; | |
| + return if length($self->[BUF]) > $maxlen; | |
| + | |
| + # read some more | |
| + my $read; | |
| + eval { | |
| + local $SIG{ALRM} = sub { die "\n" }; | |
| + alarm READ_TIMEOUT; | |
| + $read = sysread $self->[SOCKET], $self->[BUF], READSIZE, length($self->[BUF]); | |
| + alarm 0; | |
| + }; | |
| + | |
| + unless ($read) { | |
| + # we reached the EOF or there was a read error or a time out | |
| + # a complete record has not been read so we must return undef | |
| + return; | |
| + } | |
| + } | |
| +} | |
| + | |
| +1; | |
| diff --git a/lib/Starman/Server.pm b/lib/Starman/Server.pm | |
| index 1ae048e..93b9fb4 100644 | |
| --- a/lib/Starman/Server.pm | |
| +++ b/lib/Starman/Server.pm | |
| @@ -11,13 +11,10 @@ use HTTP::Date qw(time2str); | |
| use Symbol; | |
| use Plack::Util; | |
| -use Plack::TempBuffer; | |
| +use Starman::ReadSocket; | |
| +use Starman::InputStream; | |
| use constant DEBUG => $ENV{STARMAN_DEBUG} || 0; | |
| -use constant CHUNKSIZE => 64 * 1024; | |
| -use constant READ_TIMEOUT => 5; | |
| - | |
| -my $null_io = do { open my $io, "<", \""; $io }; | |
| use Net::Server::SIG qw(register_sig); | |
| @@ -115,8 +112,6 @@ sub post_accept_hook { | |
| my $self = shift; | |
| $self->{client} = { | |
| - headerbuf => '', | |
| - inputbuf => '', | |
| keepalive => 1, | |
| }; | |
| } | |
| @@ -124,6 +119,8 @@ sub post_accept_hook { | |
| sub process_request { | |
| my $self = shift; | |
| my $conn = $self->{server}->{client}; | |
| + my $conn_read = Starman::ReadSocket->new($conn); | |
| + $self->{client}->{conn_read} = $conn_read; | |
| if ($conn->NS_proto eq 'TCP') { | |
| setsockopt($conn, IPPROTO_TCP, TCP_NODELAY, 1) | |
| @@ -134,7 +131,12 @@ sub process_request { | |
| last if !$conn->connected; | |
| # Read until we see all headers | |
| - last if !$self->_read_headers; | |
| + my $headerbuf = $conn_read->readrecord(qr/\r?\n\r?\n/, 2048); | |
| + | |
| + if ( !defined $headerbuf ) { | |
| + DEBUG && warn "[$$] Read error or client connection timed out\n"; | |
| + last; | |
| + } | |
| my $env = { | |
| REMOTE_ADDR => $self->{server}->{peeraddr}, | |
| @@ -151,11 +153,11 @@ sub process_request { | |
| 'psgi.multithread' => Plack::Util::FALSE, | |
| 'psgi.multiprocess' => Plack::Util::TRUE, | |
| 'psgix.io' => $conn, | |
| - 'psgix.input.buffered' => Plack::Util::TRUE, | |
| + 'psgix.input.buffered' => Plack::Util::FALSE, | |
| }; | |
| # Parse headers | |
| - my $reqlen = parse_http_request(delete $self->{client}->{headerbuf}, $env); | |
| + my $reqlen = parse_http_request($headerbuf, $env); | |
| if ( $reqlen == -1 ) { | |
| # Bad request | |
| DEBUG && warn "[$$] Bad request\n"; | |
| @@ -211,11 +213,14 @@ sub process_request { | |
| $self->{client}->{keepalive} = 0; | |
| } | |
| - $self->_prepare_env($env); | |
| + my $input_stream = $self->_prepare_env($env); | |
| # Run PSGI apps | |
| my $res = Plack::Util::run_app($self->{app}, $env); | |
| + # consume what's left in the input stream | |
| + last if !$input_stream->finalize; | |
| + | |
| if (ref $res eq 'CODE') { | |
| $res->(sub { $self->_finalize_response($env, $_[0]) }); | |
| } else { | |
| @@ -223,92 +228,11 @@ sub process_request { | |
| } | |
| DEBUG && warn "[$$] Request done\n"; | |
| - | |
| - if ( $self->{client}->{keepalive} ) { | |
| - # If we still have data in the input buffer it may be a pipelined request | |
| - if ( $self->{client}->{inputbuf} ) { | |
| - if ( $self->{client}->{inputbuf} =~ /^(?:GET|HEAD)/ ) { | |
| - if ( DEBUG ) { | |
| - warn "Pipelined GET/HEAD request in input buffer: " | |
| - . dump( $self->{client}->{inputbuf} ) . "\n"; | |
| - } | |
| - | |
| - # Continue processing the input buffer | |
| - next; | |
| - } | |
| - else { | |
| - # Input buffer just has junk, clear it | |
| - if ( DEBUG ) { | |
| - warn "Clearing junk from input buffer: " | |
| - . dump( $self->{client}->{inputbuf} ) . "\n"; | |
| - } | |
| - | |
| - $self->{client}->{inputbuf} = ''; | |
| - } | |
| - } | |
| - | |
| - DEBUG && warn "[$$] Waiting on previous connection for keep-alive request...\n"; | |
| - | |
| - my $sel = IO::Select->new($conn); | |
| - last unless $sel->can_read(1); | |
| - } | |
| } | |
| DEBUG && warn "[$$] Closing connection\n"; | |
| } | |
| -sub _read_headers { | |
| - my $self = shift; | |
| - | |
| - eval { | |
| - local $SIG{ALRM} = sub { die "Timed out\n"; }; | |
| - | |
| - alarm( READ_TIMEOUT ); | |
| - | |
| - while (1) { | |
| - # Do we have a full header in the buffer? | |
| - # This is before sysread so we don't read if we have a pipelined request | |
| - # waiting in the buffer | |
| - last if $self->{client}->{inputbuf} =~ /$CRLF$CRLF/s; | |
| - | |
| - # If not, read some data | |
| - my $read = sysread $self->{server}->{client}, my $buf, CHUNKSIZE; | |
| - | |
| - if ( !defined $read || $read == 0 ) { | |
| - die "Read error: $!\n"; | |
| - } | |
| - | |
| - if ( DEBUG ) { | |
| - warn "[$$] Read $read bytes: " . dump($buf) . "\n"; | |
| - } | |
| - | |
| - $self->{client}->{inputbuf} .= $buf; | |
| - } | |
| - }; | |
| - | |
| - alarm(0); | |
| - | |
| - if ( $@ ) { | |
| - if ( $@ =~ /Timed out/ ) { | |
| - DEBUG && warn "[$$] Client connection timed out\n"; | |
| - return; | |
| - } | |
| - | |
| - if ( $@ =~ /Read error/ ) { | |
| - DEBUG && warn "[$$] Read error: $!\n"; | |
| - return; | |
| - } | |
| - } | |
| - | |
| - # Pull out the complete header into a new buffer | |
| - $self->{client}->{headerbuf} = $self->{client}->{inputbuf}; | |
| - | |
| - # Save any left-over data, possibly body data or pipelined requests | |
| - $self->{client}->{inputbuf} =~ s/.*?$CRLF$CRLF//s; | |
| - | |
| - return 1; | |
| -} | |
| - | |
| sub _http_error { | |
| my ( $self, $code, $env ) = @_; | |
| @@ -328,64 +252,14 @@ sub _http_error { | |
| sub _prepare_env { | |
| my($self, $env) = @_; | |
| - my $get_chunk = sub { | |
| - if ($self->{client}->{inputbuf}) { | |
| - my $chunk = delete $self->{client}->{inputbuf}; | |
| - return ($chunk, length $chunk); | |
| - } | |
| - my $read = sysread $self->{server}->{client}, my($chunk), CHUNKSIZE; | |
| - return ($chunk, $read); | |
| - }; | |
| - | |
| my $chunked = do { no warnings; lc delete $env->{HTTP_TRANSFER_ENCODING} eq 'chunked' }; | |
| if (my $cl = $env->{CONTENT_LENGTH}) { | |
| - my $buf = Plack::TempBuffer->new($cl); | |
| - while ($cl > 0) { | |
| - my($chunk, $read) = $get_chunk->(); | |
| - | |
| - if ( !defined $read || $read == 0 ) { | |
| - die "Read error: $!\n"; | |
| - } | |
| - | |
| - $cl -= $read; | |
| - $buf->print($chunk); | |
| - } | |
| - $env->{'psgi.input'} = $buf->rewind; | |
| + return $env->{'psgi.input'} = Starman::InputStream::Identity->new($self->{client}->{conn_read}, $cl); | |
| } elsif ($chunked) { | |
| - my $buf = Plack::TempBuffer->new; | |
| - my $chunk_buffer = ''; | |
| - my $length; | |
| - | |
| - DECHUNK: | |
| - while (1) { | |
| - my($chunk, $read) = $get_chunk->(); | |
| - $chunk_buffer .= $chunk; | |
| - | |
| - while ( $chunk_buffer =~ s/^(([0-9a-fA-F]+).*\015\012)// ) { | |
| - my $trailer = $1; | |
| - my $chunk_len = hex $2; | |
| - | |
| - if ($chunk_len == 0) { | |
| - last DECHUNK; | |
| - } elsif (length $chunk_buffer < $chunk_len) { | |
| - $chunk_buffer = $trailer . $chunk_buffer; | |
| - last; | |
| - } | |
| - | |
| - $buf->print(substr $chunk_buffer, 0, $chunk_len, ''); | |
| - $chunk_buffer =~ s/^\015\012//; | |
| - | |
| - $length += $chunk_len; | |
| - } | |
| - | |
| - last unless $read && $read > 0; | |
| - } | |
| - | |
| - $env->{CONTENT_LENGTH} = $length; | |
| - $env->{'psgi.input'} = $buf->rewind; | |
| + return $env->{'psgi.input'} = Starman::InputStream::Chunked->new($self->{client}->{conn_read}); | |
| } else { | |
| - $env->{'psgi.input'} = $null_io; | |
| + return $env->{'psgi.input'} = Starman::InputStream::Empty->new; | |
| } | |
| } | |
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment