ViewVC Help
View File | Revision Log | Show Annotations | Download File
/cvs/AnyEvent/lib/AnyEvent/Handle.pm
(Generate patch)

Comparing AnyEvent/lib/AnyEvent/Handle.pm (file contents):
Revision 1.166 by root, Tue Jul 28 02:07:18 2009 UTC vs.
Revision 1.179 by root, Wed Aug 12 15:50:44 2009 UTC

1package AnyEvent::Handle;
2
3use Scalar::Util ();
4use Carp ();
5use Errno qw(EAGAIN EINTR);
6
7use AnyEvent (); BEGIN { AnyEvent::common_sense }
8use AnyEvent::Util qw(WSAEWOULDBLOCK);
9
10=head1 NAME 1=head1 NAME
11 2
12AnyEvent::Handle - non-blocking I/O on file handles via AnyEvent 3AnyEvent::Handle - non-blocking I/O on file handles via AnyEvent
13
14=cut
15
16our $VERSION = 4.88;
17 4
18=head1 SYNOPSIS 5=head1 SYNOPSIS
19 6
20 use AnyEvent; 7 use AnyEvent;
21 use AnyEvent::Handle; 8 use AnyEvent::Handle;
59C<on_error> callback. 46C<on_error> callback.
60 47
61All callbacks will be invoked with the handle object as their first 48All callbacks will be invoked with the handle object as their first
62argument. 49argument.
63 50
51=cut
52
53package AnyEvent::Handle;
54
55use Scalar::Util ();
56use List::Util ();
57use Carp ();
58use Errno qw(EAGAIN EINTR);
59
60use AnyEvent (); BEGIN { AnyEvent::common_sense }
61use AnyEvent::Util qw(WSAEWOULDBLOCK);
62
63our $VERSION = $AnyEvent::VERSION;
64
64=head1 METHODS 65=head1 METHODS
65 66
66=over 4 67=over 4
67 68
68=item $handle = B<new> AnyEvent::TLS fh => $filehandle, key => value... 69=item $handle = B<new> AnyEvent::TLS fh => $filehandle, key => value...
216memory and push it into the queue, but instead only read more data from 217memory and push it into the queue, but instead only read more data from
217the file when the write queue becomes empty. 218the file when the write queue becomes empty.
218 219
219=item timeout => $fractional_seconds 220=item timeout => $fractional_seconds
220 221
222=item rtimeout => $fractional_seconds
223
224=item wtimeout => $fractional_seconds
225
221If non-zero, then this enables an "inactivity" timeout: whenever this many 226If non-zero, then these enables an "inactivity" timeout: whenever this
222seconds pass without a successful read or write on the underlying file 227many seconds pass without a successful read or write on the underlying
223handle, the C<on_timeout> callback will be invoked (and if that one is 228file handle (or a call to C<timeout_reset>), the C<on_timeout> callback
224missing, a non-fatal C<ETIMEDOUT> error will be raised). 229will be invoked (and if that one is missing, a non-fatal C<ETIMEDOUT>
230error will be raised).
231
232There are three variants of the timeouts that work fully independent
233of each other, for both read and write, just read, and just write:
234C<timeout>, C<rtimeout> and C<wtimeout>, with corresponding callbacks
235C<on_timeout>, C<on_rtimeout> and C<on_wtimeout>, and reset functions
236C<timeout_reset>, C<rtimeout_reset>, and C<wtimeout_reset>.
225 237
226Note that timeout processing is also active when you currently do not have 238Note that timeout processing is also active when you currently do not have
227any outstanding read or write requests: If you plan to keep the connection 239any outstanding read or write requests: If you plan to keep the connection
228idle then you should disable the timout temporarily or ignore the timeout 240idle then you should disable the timout temporarily or ignore the timeout
229in the C<on_timeout> callback, in which case AnyEvent::Handle will simply 241in the C<on_timeout> callback, in which case AnyEvent::Handle will simply
438 delete $self->{_skip_drain_rbuf}; 450 delete $self->{_skip_drain_rbuf};
439 $self->_start; 451 $self->_start;
440 452
441 $self->{on_connect} 453 $self->{on_connect}
442 and $self->{on_connect}($self, $host, $port, sub { 454 and $self->{on_connect}($self, $host, $port, sub {
443 delete @$self{qw(fh _tw _ww _rw _eof _queue rbuf _wbuf tls _tls_rbuf _tls_wbuf)}; 455 delete @$self{qw(fh _tw _rtw _wtw _ww _rw _eof _queue rbuf _wbuf tls _tls_rbuf _tls_wbuf)};
444 $self->{_skip_drain_rbuf} = 1; 456 $self->{_skip_drain_rbuf} = 1;
445 &$retry; 457 &$retry;
446 }); 458 });
447 459
448 } else { 460 } else {
474sub _start { 486sub _start {
475 my ($self) = @_; 487 my ($self) = @_;
476 488
477 AnyEvent::Util::fh_nonblocking $self->{fh}, 1; 489 AnyEvent::Util::fh_nonblocking $self->{fh}, 1;
478 490
491 $self->{_activity} =
492 $self->{_ractivity} =
479 $self->{_activity} = AnyEvent->now; 493 $self->{_wactivity} = AE::now;
480 $self->_timeout; 494
495 $self->timeout (delete $self->{timeout} ) if $self->{timeout};
496 $self->rtimeout (delete $self->{rtimeout}) if $self->{rtimeout};
497 $self->wtimeout (delete $self->{wtimeout}) if $self->{wtimeout};
481 498
482 $self->no_delay (delete $self->{no_delay}) if exists $self->{no_delay}; 499 $self->no_delay (delete $self->{no_delay}) if exists $self->{no_delay};
483 500
484 $self->starttls (delete $self->{tls}, delete $self->{tls_ctx}) 501 $self->starttls (delete $self->{tls}, delete $self->{tls_ctx})
485 if $self->{tls}; 502 if $self->{tls};
489 $self->start_read 506 $self->start_read
490 if $self->{on_read} || @{ $self->{_queue} }; 507 if $self->{on_read} || @{ $self->{_queue} };
491 508
492 $self->_drain_wbuf; 509 $self->_drain_wbuf;
493} 510}
494
495#sub _shutdown {
496# my ($self) = @_;
497#
498# delete @$self{qw(_tw _rw _ww fh wbuf on_read _queue)};
499# $self->{_eof} = 1; # tell starttls et. al to stop trying
500#
501# &_freetls;
502#}
503 511
504sub _error { 512sub _error {
505 my ($self, $errno, $fatal, $message) = @_; 513 my ($self, $errno, $fatal, $message) = @_;
506 514
507 $! = $errno; 515 $! = $errno;
544 $_[0]{on_eof} = $_[1]; 552 $_[0]{on_eof} = $_[1];
545} 553}
546 554
547=item $handle->on_timeout ($cb) 555=item $handle->on_timeout ($cb)
548 556
549Replace the current C<on_timeout> callback, or disables the callback (but 557=item $handle->on_rtimeout ($cb)
550not the timeout) if C<$cb> = C<undef>. See the C<timeout> constructor
551argument and method.
552 558
553=cut 559=item $handle->on_wtimeout ($cb)
554 560
555sub on_timeout { 561Replace the current C<on_timeout>, C<on_rtimeout> or C<on_wtimeout>
556 $_[0]{on_timeout} = $_[1]; 562callback, or disables the callback (but not the timeout) if C<$cb> =
557} 563C<undef>. See the C<timeout> constructor argument and method.
564
565=cut
566
567# see below
558 568
559=item $handle->autocork ($boolean) 569=item $handle->autocork ($boolean)
560 570
561Enables or disables the current autocork behaviour (see C<autocork> 571Enables or disables the current autocork behaviour (see C<autocork>
562constructor argument). Changes will only take effect on the next write. 572constructor argument). Changes will only take effect on the next write.
602 612
603sub on_starttls { 613sub on_starttls {
604 $_[0]{on_stoptls} = $_[1]; 614 $_[0]{on_stoptls} = $_[1];
605} 615}
606 616
617=item $handle->rbuf_max ($max_octets)
618
619Configures the C<rbuf_max> setting (C<undef> disables it).
620
621=cut
622
623sub rbuf_max {
624 $_[0]{rbuf_max} = $_[1];
625}
626
607############################################################################# 627#############################################################################
608 628
609=item $handle->timeout ($seconds) 629=item $handle->timeout ($seconds)
610 630
631=item $handle->rtimeout ($seconds)
632
633=item $handle->wtimeout ($seconds)
634
611Configures (or disables) the inactivity timeout. 635Configures (or disables) the inactivity timeout.
612 636
613=cut 637=item $handle->timeout_reset
614 638
615sub timeout { 639=item $handle->rtimeout_reset
640
641=item $handle->wtimeout_reset
642
643Reset the activity timeout, as if data was received or sent.
644
645These methods are cheap to call.
646
647=cut
648
649for my $dir ("", "r", "w") {
650 my $timeout = "${dir}timeout";
651 my $tw = "_${dir}tw";
652 my $on_timeout = "on_${dir}timeout";
653 my $activity = "_${dir}activity";
654 my $cb;
655
656 *$on_timeout = sub {
657 $_[0]{$on_timeout} = $_[1];
658 };
659
660 *$timeout = sub {
616 my ($self, $timeout) = @_; 661 my ($self, $new_value) = @_;
617 662
618 $self->{timeout} = $timeout; 663 $self->{$timeout} = $new_value;
619 $self->_timeout; 664 delete $self->{$tw}; &$cb;
620} 665 };
621 666
667 *{"${dir}timeout_reset"} = sub {
668 $_[0]{$activity} = AE::now;
669 };
670
671 # main workhorse:
622# reset the timeout watcher, as neccessary 672 # reset the timeout watcher, as neccessary
623# also check for time-outs 673 # also check for time-outs
624sub _timeout { 674 $cb = sub {
625 my ($self) = @_; 675 my ($self) = @_;
626 676
627 if ($self->{timeout} && $self->{fh}) { 677 if ($self->{$timeout} && $self->{fh}) {
628 my $NOW = AnyEvent->now; 678 my $NOW = AE::now;
629 679
630 # when would the timeout trigger? 680 # when would the timeout trigger?
631 my $after = $self->{_activity} + $self->{timeout} - $NOW; 681 my $after = $self->{$activity} + $self->{$timeout} - $NOW;
632 682
633 # now or in the past already? 683 # now or in the past already?
634 if ($after <= 0) { 684 if ($after <= 0) {
635 $self->{_activity} = $NOW; 685 $self->{$activity} = $NOW;
636 686
637 if ($self->{on_timeout}) { 687 if ($self->{$on_timeout}) {
638 $self->{on_timeout}($self); 688 $self->{$on_timeout}($self);
639 } else { 689 } else {
640 $self->_error (Errno::ETIMEDOUT); 690 $self->_error (Errno::ETIMEDOUT);
691 }
692
693 # callback could have changed timeout value, optimise
694 return unless $self->{$timeout};
695
696 # calculate new after
697 $after = $self->{$timeout};
641 } 698 }
642 699
643 # callback could have changed timeout value, optimise 700 Scalar::Util::weaken $self;
644 return unless $self->{timeout}; 701 return unless $self; # ->error could have destroyed $self
645 702
646 # calculate new after 703 $self->{$tw} ||= AE::timer $after, 0, sub {
647 $after = $self->{timeout}; 704 delete $self->{$tw};
705 $cb->($self);
706 };
707 } else {
708 delete $self->{$tw};
648 } 709 }
649
650 Scalar::Util::weaken $self;
651 return unless $self; # ->error could have destroyed $self
652
653 $self->{_tw} ||= AnyEvent->timer (after => $after, cb => sub {
654 delete $self->{_tw};
655 $self->_timeout;
656 });
657 } else {
658 delete $self->{_tw};
659 } 710 }
660} 711}
661 712
662############################################################################# 713#############################################################################
663 714
711 my $len = syswrite $self->{fh}, $self->{wbuf}; 762 my $len = syswrite $self->{fh}, $self->{wbuf};
712 763
713 if (defined $len) { 764 if (defined $len) {
714 substr $self->{wbuf}, 0, $len, ""; 765 substr $self->{wbuf}, 0, $len, "";
715 766
716 $self->{_activity} = AnyEvent->now; 767 $self->{_activity} = $self->{_wactivity} = AE::now;
717 768
718 $self->{on_drain}($self) 769 $self->{on_drain}($self)
719 if $self->{low_water_mark} >= (length $self->{wbuf}) + (length $self->{_tls_wbuf}) 770 if $self->{low_water_mark} >= (length $self->{wbuf}) + (length $self->{_tls_wbuf})
720 && $self->{on_drain}; 771 && $self->{on_drain};
721 772
727 778
728 # try to write data immediately 779 # try to write data immediately
729 $cb->() unless $self->{autocork}; 780 $cb->() unless $self->{autocork};
730 781
731 # if still data left in wbuf, we need to poll 782 # if still data left in wbuf, we need to poll
732 $self->{_ww} = AnyEvent->io (fh => $self->{fh}, poll => "w", cb => $cb) 783 $self->{_ww} = AE::io $self->{fh}, 1, $cb
733 if length $self->{wbuf}; 784 if length $self->{wbuf};
734 }; 785 };
735} 786}
736 787
737our %WH; 788our %WH;
827Other languages could read single lines terminated by a newline and pass 878Other languages could read single lines terminated by a newline and pass
828this line into their JSON decoder of choice. 879this line into their JSON decoder of choice.
829 880
830=cut 881=cut
831 882
883sub json_coder() {
884 eval { require JSON::XS; JSON::XS->new->utf8 }
885 || do { require JSON; JSON->new->utf8 }
886}
887
832register_write_type json => sub { 888register_write_type json => sub {
833 my ($self, $ref) = @_; 889 my ($self, $ref) = @_;
834 890
835 require JSON; 891 my $json = $self->{json} ||= json_coder;
836 892
837 $self->{json} ? $self->{json}->encode ($ref) 893 $json->encode ($ref)
838 : JSON::encode_json ($ref)
839}; 894};
840 895
841=item storable => $reference 896=item storable => $reference
842 897
843Freezes the given reference using L<Storable> and writes it to the 898Freezes the given reference using L<Storable> and writes it to the
981 1036
982sub _drain_rbuf { 1037sub _drain_rbuf {
983 my ($self) = @_; 1038 my ($self) = @_;
984 1039
985 # avoid recursion 1040 # avoid recursion
986 return if exists $self->{_skip_drain_rbuf}; 1041 return if $self->{_skip_drain_rbuf};
987 local $self->{_skip_drain_rbuf} = 1; 1042 local $self->{_skip_drain_rbuf} = 1;
988
989 if (
990 defined $self->{rbuf_max}
991 && $self->{rbuf_max} < length $self->{rbuf}
992 ) {
993 $self->_error (Errno::ENOSPC, 1), return;
994 }
995 1043
996 while () { 1044 while () {
997 # we need to use a separate tls read buffer, as we must not receive data while 1045 # we need to use a separate tls read buffer, as we must not receive data while
998 # we are draining the buffer, and this can only happen with TLS. 1046 # we are draining the buffer, and this can only happen with TLS.
999 $self->{rbuf} .= delete $self->{_tls_rbuf} 1047 $self->{rbuf} .= delete $self->{_tls_rbuf}
1039 $self->{on_eof} 1087 $self->{on_eof}
1040 ? $self->{on_eof}($self) 1088 ? $self->{on_eof}($self)
1041 : $self->_error (0, 1, "Unexpected end-of-file"); 1089 : $self->_error (0, 1, "Unexpected end-of-file");
1042 1090
1043 return; 1091 return;
1092 }
1093
1094 if (
1095 defined $self->{rbuf_max}
1096 && $self->{rbuf_max} < length $self->{rbuf}
1097 ) {
1098 $self->_error (Errno::ENOSPC, 1), return;
1044 } 1099 }
1045 1100
1046 # may need to restart read watcher 1101 # may need to restart read watcher
1047 unless ($self->{_rw}) { 1102 unless ($self->{_rw}) {
1048 $self->start_read 1103 $self->start_read
1397=cut 1452=cut
1398 1453
1399register_read_type json => sub { 1454register_read_type json => sub {
1400 my ($self, $cb) = @_; 1455 my ($self, $cb) = @_;
1401 1456
1402 my $json = $self->{json} ||= 1457 my $json = $self->{json} ||= json_coder;
1403 eval { require JSON::XS; JSON::XS->new->utf8 }
1404 || do { require JSON; JSON->new->utf8 };
1405 1458
1406 my $data; 1459 my $data;
1407 my $rbuf = \$self->{rbuf}; 1460 my $rbuf = \$self->{rbuf};
1408 1461
1409 sub { 1462 sub {
1529 my ($self) = @_; 1582 my ($self) = @_;
1530 1583
1531 unless ($self->{_rw} || $self->{_eof}) { 1584 unless ($self->{_rw} || $self->{_eof}) {
1532 Scalar::Util::weaken $self; 1585 Scalar::Util::weaken $self;
1533 1586
1534 $self->{_rw} = AnyEvent->io (fh => $self->{fh}, poll => "r", cb => sub { 1587 $self->{_rw} = AE::io $self->{fh}, 0, sub {
1535 my $rbuf = \($self->{tls} ? my $buf : $self->{rbuf}); 1588 my $rbuf = \($self->{tls} ? my $buf : $self->{rbuf});
1536 my $len = sysread $self->{fh}, $$rbuf, $self->{read_size} || 8192, length $$rbuf; 1589 my $len = sysread $self->{fh}, $$rbuf, $self->{read_size} || 8192, length $$rbuf;
1537 1590
1538 if ($len > 0) { 1591 if ($len > 0) {
1539 $self->{_activity} = AnyEvent->now; 1592 $self->{_activity} = $self->{_ractivity} = AE::now;
1540 1593
1541 if ($self->{tls}) { 1594 if ($self->{tls}) {
1542 Net::SSLeay::BIO_write ($self->{_rbio}, $$rbuf); 1595 Net::SSLeay::BIO_write ($self->{_rbio}, $$rbuf);
1543 1596
1544 &_dotls ($self); 1597 &_dotls ($self);
1552 $self->_drain_rbuf; 1605 $self->_drain_rbuf;
1553 1606
1554 } elsif ($! != EAGAIN && $! != EINTR && $! != WSAEWOULDBLOCK) { 1607 } elsif ($! != EAGAIN && $! != EINTR && $! != WSAEWOULDBLOCK) {
1555 return $self->_error ($!, 1); 1608 return $self->_error ($!, 1);
1556 } 1609 }
1557 }); 1610 };
1558 } 1611 }
1559} 1612}
1560 1613
1561our $ERROR_SYSCALL; 1614our $ERROR_SYSCALL;
1562our $ERROR_WANT_READ; 1615our $ERROR_WANT_READ;
1722 Net::SSLeay::CTX_set_mode ($tls, 1|2); 1775 Net::SSLeay::CTX_set_mode ($tls, 1|2);
1723 1776
1724 $self->{_rbio} = Net::SSLeay::BIO_new (Net::SSLeay::BIO_s_mem ()); 1777 $self->{_rbio} = Net::SSLeay::BIO_new (Net::SSLeay::BIO_s_mem ());
1725 $self->{_wbio} = Net::SSLeay::BIO_new (Net::SSLeay::BIO_s_mem ()); 1778 $self->{_wbio} = Net::SSLeay::BIO_new (Net::SSLeay::BIO_s_mem ());
1726 1779
1780 Net::SSLeay::BIO_write ($self->{_rbio}, delete $self->{rbuf});
1781
1727 Net::SSLeay::set_bio ($tls, $self->{_rbio}, $self->{_wbio}); 1782 Net::SSLeay::set_bio ($tls, $self->{_rbio}, $self->{_wbio});
1728 1783
1729 $self->{_on_starttls} = sub { $_[0]{on_starttls}(@_) } 1784 $self->{_on_starttls} = sub { $_[0]{on_starttls}(@_) }
1730 if $self->{on_starttls}; 1785 if $self->{on_starttls};
1731 1786
1760 my ($self) = @_; 1815 my ($self) = @_;
1761 1816
1762 return unless $self->{tls}; 1817 return unless $self->{tls};
1763 1818
1764 $self->{tls_ctx}->_put_session (delete $self->{tls}) 1819 $self->{tls_ctx}->_put_session (delete $self->{tls})
1765 if ref $self->{tls}; 1820 if $self->{tls} > 0;
1766 1821
1767 delete @$self{qw(_rbio _wbio _tls_wbuf _on_starttls)}; 1822 delete @$self{qw(_rbio _wbio _tls_wbuf _on_starttls)};
1768} 1823}
1769 1824
1770sub DESTROY { 1825sub DESTROY {
1778 my $fh = delete $self->{fh}; 1833 my $fh = delete $self->{fh};
1779 my $wbuf = delete $self->{wbuf}; 1834 my $wbuf = delete $self->{wbuf};
1780 1835
1781 my @linger; 1836 my @linger;
1782 1837
1783 push @linger, AnyEvent->io (fh => $fh, poll => "w", cb => sub { 1838 push @linger, AE::io $fh, 1, sub {
1784 my $len = syswrite $fh, $wbuf, length $wbuf; 1839 my $len = syswrite $fh, $wbuf, length $wbuf;
1785 1840
1786 if ($len > 0) { 1841 if ($len > 0) {
1787 substr $wbuf, 0, $len, ""; 1842 substr $wbuf, 0, $len, "";
1788 } else { 1843 } else {
1789 @linger = (); # end 1844 @linger = (); # end
1790 } 1845 }
1791 }); 1846 };
1792 push @linger, AnyEvent->timer (after => $linger, cb => sub { 1847 push @linger, AE::timer $linger, 0, sub {
1793 @linger = (); 1848 @linger = ();
1794 }); 1849 };
1795 } 1850 }
1796} 1851}
1797 1852
1798=item $handle->destroy 1853=item $handle->destroy
1799 1854

Diff Legend

Removed lines
+ Added lines
< Changed lines
> Changed lines