Skip to content

Instantly share code, notes, and snippets.

@redneb
Created August 30, 2010 23:02
Show Gist options
  • Select an option

  • Save redneb/558192 to your computer and use it in GitHub Desktop.

Select an option

Save redneb/558192 to your computer and use it in GitHub Desktop.
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