ViewVC Help
View File | Revision Log | Show Annotations | Download File
/cvs/cvsroot/AnyEvent-Fork-RPC/RPC/Async.pm
Revision: 1.4
Committed: Sat Apr 27 23:49:01 2013 UTC (13 years, 5 months ago) by root
Branch: MAIN
CVS Tags: rel-1_1
Changes since 1.3: +12 -10 lines
Log Message:
*** empty log message ***

File Contents

# Content
1 package AnyEvent::Fork::RPC::Async;
2
3 use common::sense; # actually required to avoid spurious warnings...
4
5 use Errno ();
6
7 use AnyEvent;
8
9 # declare only
10 sub AnyEvent::Fork::RPC::event;
11
12 sub run {
13 my ($function, $init, $serialiser) = splice @_, -3, 3,;
14
15 my $rfh = shift;
16 my $wfh = fileno $rfh ? $rfh : *STDOUT;
17
18 {
19 package main;
20 &$init if length $init;
21 $function = \&$function; # resolve function early for extra speed
22 }
23
24 my $busy = 1; # exit when == 0
25
26 my ($f, $t) = eval $serialiser; die $@ if $@;
27 my ($wbuf, $ww);
28
29 my $wcb = sub {
30 my $len = syswrite $wfh, $wbuf;
31
32 unless (defined $len) {
33 if ($! != Errno::EAGAIN && $! != Errno::EWOULDBLOCK) {
34 undef $ww;
35 die "AnyEvent::Fork::RPC: write error ($!), parent gone?\n";
36 }
37 }
38
39 substr $wbuf, 0, $len, "";
40
41 unless (length $wbuf) {
42 undef $ww;
43 unless ($busy) {
44 shutdown $wfh, 1;
45 exit;
46 }
47 }
48 };
49
50 my $write = sub {
51 $wbuf .= $_[0];
52 $ww ||= AE::io $wfh, 1, $wcb;
53 };
54
55 *AnyEvent::Fork::RPC::event = sub {
56 $write->(pack "NN/a*", 0, &$f);
57 };
58
59 my ($rlen, $rbuf, $rw) = 512 - 16;
60
61 my $len;
62
63 $rw = AE::io $rfh, 0, sub {
64 $rlen = $rlen * 2 + 16 if $rlen - 128 < length $rbuf;
65 $len = sysread $rfh, $rbuf, $rlen - length $rbuf, length $rbuf;
66
67 if ($len) {
68 while (8 <= length $rbuf) {
69 (my $id, $len) = unpack "NN", $rbuf;
70 8 + $len <= length $rbuf
71 or last;
72
73 my @r = $t->(substr $rbuf, 8, $len);
74 substr $rbuf, 0, 8 + $len, "";
75
76 ++$busy;
77 $function->(sub {
78 --$busy;
79 $write->(pack "NN/a*", $id, &$f);
80 }, @r);
81 }
82 } elsif (defined $len) {
83 undef $rw;
84 --$busy;
85 $ww ||= AE::io $wfh, 1, $wcb;
86 } elsif ($! != Errno::EAGAIN && $! != Errno::EWOULDBLOCK) {
87 undef $rw;
88 die "AnyEvent::Fork::RPC: read error in child: $!\n";
89 }
90 };
91
92 $AnyEvent::MODEL eq "EV"
93 ? EV::loop ()
94 : AE::cv->recv;
95 }
96
97 1
98