ViewVC Help
View File | Revision Log | Show Annotations | Download File
/cvs/cvsroot/AnyEvent-Fork-RPC/RPC/Async.pm
Revision: 1.16
Committed: Thu Sep 12 15:15:49 2019 UTC (7 years ago) by root
Branch: MAIN
Changes since 1.15: +11 -2 lines
Log Message:
*** empty log message ***

File Contents

# User Rev Content
1 root 1.1 package AnyEvent::Fork::RPC::Async;
2    
3     use common::sense; # actually required to avoid spurious warnings...
4    
5 root 1.16 our $VERSION = 2; # protocol version
6    
7 root 1.1 use Errno ();
8    
9     use AnyEvent;
10    
11     # declare only
12     sub AnyEvent::Fork::RPC::event;
13    
14 root 1.9 sub do_exit { exit } # workaround for perl 5.14 and below
15    
16 root 1.1 sub run {
17 root 1.16 my %kv = splice @_, pop;
18    
19 root 1.4 my $rfh = shift;
20     my $wfh = fileno $rfh ? $rfh : *STDOUT;
21 root 1.1
22 root 1.16 my $function = delete $kv{function};
23     my $serialiser = delete $kv{serialiser};
24     my $rlen = delete $kv{rlen};
25     my $done = delete $kv{done};
26    
27 root 1.15 $0 =~ s/^(\d+).*$/$1 $function/s;
28 root 1.5
29 root 1.1 {
30     package main;
31 root 1.16 my $init = delete $kv{init};
32 root 1.1 &$init if length $init;
33     $function = \&$function; # resolve function early for extra speed
34     }
35    
36     my $busy = 1; # exit when == 0
37    
38 root 1.8 my ($f, $t) = eval $serialiser; AE::log fatal => $@ if $@;
39 root 1.1 my ($wbuf, $ww);
40    
41     my $wcb = sub {
42 root 1.4 my $len = syswrite $wfh, $wbuf;
43 root 1.1
44     unless (defined $len) {
45     if ($! != Errno::EAGAIN && $! != Errno::EWOULDBLOCK) {
46     undef $ww;
47 root 1.8 AE::log fatal => "AnyEvent::Fork::RPC: write error ($!), parent gone?";
48 root 1.1 }
49     }
50    
51     substr $wbuf, 0, $len, "";
52    
53     unless (length $wbuf) {
54     undef $ww;
55 root 1.3 unless ($busy) {
56 root 1.4 shutdown $wfh, 1;
57 root 1.6 @_ = (); goto &$done;
58 root 1.3 }
59 root 1.1 }
60     };
61    
62     my $write = sub {
63     $wbuf .= $_[0];
64 root 1.4 $ww ||= AE::io $wfh, 1, $wcb;
65 root 1.1 };
66    
67     *AnyEvent::Fork::RPC::event = sub {
68 root 1.4 $write->(pack "NN/a*", 0, &$f);
69 root 1.1 };
70    
71 root 1.16 my ($rbuf, $rw);
72 root 1.1
73 root 1.2 my $len;
74 root 1.1
75 root 1.4 $rw = AE::io $rfh, 0, sub {
76 root 1.1 $rlen = $rlen * 2 + 16 if $rlen - 128 < length $rbuf;
77 root 1.4 $len = sysread $rfh, $rbuf, $rlen - length $rbuf, length $rbuf;
78 root 1.1
79     if ($len) {
80     while (8 <= length $rbuf) {
81 root 1.4 (my $id, $len) = unpack "NN", $rbuf;
82 root 1.1 8 + $len <= length $rbuf
83     or last;
84    
85     my @r = $t->(substr $rbuf, 8, $len);
86     substr $rbuf, 0, 8 + $len, "";
87    
88     ++$busy;
89     $function->(sub {
90     --$busy;
91 root 1.4 $write->(pack "NN/a*", $id, &$f);
92 root 1.1 }, @r);
93     }
94 root 1.7 } elsif (defined $len or $! == Errno::EINVAL) { # EINVAL is for microshit windoze
95 root 1.1 undef $rw;
96     --$busy;
97 root 1.4 $ww ||= AE::io $wfh, 1, $wcb;
98 root 1.1 } elsif ($! != Errno::EAGAIN && $! != Errno::EWOULDBLOCK) {
99     undef $rw;
100 root 1.8 AE::log fatal => "AnyEvent::Fork::RPC: read error in child: $!";
101 root 1.1 }
102     };
103    
104 root 1.14 $AnyEvent::MODEL eq "AnyEvent::Impl::EV"
105 root 1.10 ? EV::run ()
106 root 1.1 : AE::cv->recv;
107     }
108    
109     1
110