lib/PublicInbox/HTTP.pm | 18 ++++++++++++++++++ lib/PublicInbox/HTTPD/Async.pm | 2 ++ diff --git a/lib/PublicInbox/HTTP.pm b/lib/PublicInbox/HTTP.pm index f69056f87932980cb8dffe2d41ae9d1c3be4297b..d523bd42a3bd6571ec33d6074779a8dea7fe56ab 100644 --- a/lib/PublicInbox/HTTP.pm +++ b/lib/PublicInbox/HTTP.pm @@ -219,6 +219,24 @@ if (defined(my $body = $res->[2])) { if (ref $body eq 'ARRAY') { $write->($_) foreach @$body; $close->(); + } elsif ($body->can('async_pass')) { # HTTPD::Async + # prevent us from reading the body faster than we + # can write to the client + my $restart_read = sub { $body->watch_read(1) }; + $body->async_pass(sub { + local $/ = \8192; + my $buf = $body->getline; + if (defined $buf) { + $write->($buf); + if ($self->{write_buf}) { + $body->watch_read(0); + $self->write($restart_read); + } + return; # continue waiting + } + $body->close; + $close->(); + }); } else { my $pull; $pull = sub { diff --git a/lib/PublicInbox/HTTPD/Async.pm b/lib/PublicInbox/HTTPD/Async.pm index bedb397d0f9cdbf7b143615613fb117510549ea4..ceba738e3ce3f9793c314b16c7a3fe3733b55358 100644 --- a/lib/PublicInbox/HTTPD/Async.pm +++ b/lib/PublicInbox/HTTPD/Async.pm @@ -21,10 +21,12 @@ $self->watch_read(1); $self; } +sub async_pass { $_[0]->{cb} = $_[1] } sub event_read { $_[0]->{cb}->() } sub event_hup { $_[0]->{cb}->() } sub event_err { $_[0]->{cb}->() } sub sysread { shift->{sock}->sysread(@_) } +sub getline { $_[0]->{sock}->getline }; sub close { my $self = shift;