Acme-Sort-Sleep
view release on metacpan or search on metacpan
local/lib/perl5/IO/Async/Stream.pm view on Meta::CPAN
sub _flush_one_write
{
my $self = shift;
my $writequeue = $self->{writequeue};
my $head;
while( $head = $writequeue->[0] and ref $head->data ) {
if( ref $head->data eq "CODE" ) {
my $data = $head->data->( $self );
if( !defined $data ) {
$head->on_flush->( $self ) if $head->on_flush;
shift @$writequeue;
return 1;
}
if( !ref $data and my $encoding = $self->{encoding} ) {
$data = $encoding->encode( $data );
}
unshift @$writequeue, my $new = Writer(
$data, $head->writelen, $head->on_write, undef, undef, 0
);
next;
}
elsif( blessed $head->data and $head->data->isa( "Future" ) ) {
my $f = $head->data;
if( !$f->is_ready ) {
return 0 if $head->watching;
$f->on_ready( sub { $self->_flush_one_write } );
$head->watching++;
return 0;
}
my $data = $f->get;
if( !ref $data and my $encoding = $self->{encoding} ) {
$data = $encoding->encode( $data );
}
$head->data = $data;
next;
}
else {
die "Unsure what to do with reference ".ref($head->data)." in write queue";
}
}
my $second;
while( $second = $writequeue->[1] and
!ref $second->data and
$head->writelen == $second->writelen and
!$head->on_write and !$second->on_write and
!$head->on_flush ) {
$head->data .= $second->data;
$head->on_write = $second->on_write;
$head->on_flush = $second->on_flush;
splice @$writequeue, 1, 1, ();
}
die "TODO: head data does not contain a plain string" if ref $head->data;
if( $IO::Async::Debug::DEBUG > 1 ) {
my $data = substr $head->data, 0, $head->writelen;
$self->debug_printf( "WRITE len=%d", length $data );
IO::Async::Debug::log_hexdump( $data ) if $IO::Async::Debug::DEBUG_FLAGS{Sw};
}
my $writer = $self->{writer};
my $len = $self->$writer( $self->write_handle, $head->data, $head->writelen );
if( !defined $len ) {
my $errno = $!;
if( $errno == EAGAIN or $errno == EWOULDBLOCK ) {
$self->maybe_invoke_event( on_writeable_stop => ) if $self->{writeable};
$self->{writeable} = 0;
}
return 0 if _nonfatal_error( $errno );
if( $errno == EPIPE ) {
$self->{write_eof} = 1;
$self->maybe_invoke_event( on_write_eof => );
}
$head->on_error->( $self, $errno ) if $head->on_error;
$self->maybe_invoke_event( on_write_error => $errno )
or $self->close_now;
return 0;
}
if( my $on_write = $head->on_write ) {
$on_write->( $self, $len );
}
if( !length $head->data ) {
$head->on_flush->( $self ) if $head->on_flush;
shift @{ $self->{writequeue} };
}
return 1;
}
sub write
{
my $self = shift;
my ( $data, %params ) = @_;
carp "Cannot write data to a Stream that is closing" and return if $self->{stream_closing};
# Allow writes without a filehandle if we're not yet in a Loop, just don't
# try to flush them
my $handle = $self->write_handle;
croak "Cannot write data to a Stream with no write_handle" if !$handle and $self->loop;
if( !ref $data and my $encoding = $self->{encoding} ) {
$data = $encoding->encode( $data );
}
my $on_write = delete $params{on_write};
my $on_flush = delete $params{on_flush};
my $on_error = delete $params{on_error};
local/lib/perl5/IO/Async/Stream.pm view on Meta::CPAN
$self->invoke_event( on_read_low_watermark => length $self->{readbuff} );
}
if( ref $ret eq "CODE" ) {
# Replace the top CODE, or add it if there was none
$readqueue->[0] = Reader( $ret, undef );
return 1;
}
elsif( @$readqueue and !defined $ret ) {
shift @$readqueue;
return 1;
}
else {
return $ret && ( length( $self->{readbuff} ) > 0 || $eof );
}
}
sub _sysread
{
my $self = shift;
my ( $handle, undef, $len ) = @_;
return $handle->sysread( $_[1], $len );
}
sub on_read_ready
{
my $self = shift;
$self->_do_read if $self->{want} & WANT_READ_FOR_READ;
$self->_do_write if $self->{want} & WANT_READ_FOR_WRITE;
}
sub _do_read
{
my $self = shift;
my $handle = $self->read_handle;
my $reader = $self->{reader};
while(1) {
my $data;
my $len = $self->$reader( $handle, $data, $self->{read_len} );
if( !defined $len ) {
my $errno = $!;
return if _nonfatal_error( $errno );
$self->maybe_invoke_event( on_read_error => $errno )
or $self->close_now;
foreach ( @{ $self->{readqueue} } ) {
$_->future->fail( "read failed: $errno", sysread => $errno ) if $_->future;
}
undef @{ $self->{readqueue} };
return;
}
if( $IO::Async::Debug::DEBUG > 1 ) {
$self->debug_printf( "READ len=%d", $len );
IO::Async::Debug::log_hexdump( $data ) if $IO::Async::Debug::DEBUG_FLAGS{Sr};
}
my $eof = $self->{read_eof} = ( $len == 0 );
if( my $encoding = $self->{encoding} ) {
my $bytes = defined $self->{bytes_remaining} ? $self->{bytes_remaining} . $data : $data;
$data = $encoding->decode( $bytes, STOP_AT_PARTIAL );
$self->{bytes_remaining} = $bytes;
}
$self->{readbuff} .= $data if !$eof;
1 while $self->_flush_one_read( $eof );
if( $eof ) {
$self->maybe_invoke_event( on_read_eof => );
$self->close_now if $self->{close_on_read_eof};
foreach ( @{ $self->{readqueue} } ) {
$_->future->done( undef ) if $_->future;
}
undef @{ $self->{readqueue} };
return;
}
last unless $self->{read_all};
}
if( defined $self->{read_high_watermark} and length $self->{readbuff} >= $self->{read_high_watermark} ) {
$self->{at_read_high_watermark} or
$self->invoke_event( on_read_high_watermark => length $self->{readbuff} );
$self->{at_read_high_watermark} = 1;
}
}
sub on_read_high_watermark
{
my $self = shift;
$self->want_readready_for_read( 0 );
}
sub on_read_low_watermark
{
my $self = shift;
$self->want_readready_for_read( 1 );
}
=head2 push_on_read
$stream->push_on_read( $on_read )
Pushes a new temporary C<on_read> handler to the end of the queue. This queue,
if non-empty, is used to provide C<on_read> event handling code in preference
to using the object's main event handler or method. New handlers can be
supplied at any time, and they will be used in first-in first-out (FIFO)
order.
As with the main C<on_read> event handler, each can return a (defined) boolean
to indicate if they wish to be invoked again or not, another C<CODE> reference
( run in 0.612 second using v1.01-cache-2.11-cpan-6de40a662fe )