ViewVC Help
View File | Revision Log | Show Annotations | Download File
/cvs/cvsroot/Net-IRC3/samples/JSONConnection.pm
Revision: 1.5
Committed: Sat Feb 17 13:01:38 2007 UTC (19 years, 7 months ago) by elmex
Branch: MAIN
CVS Tags: HEAD
Changes since 1.4: +0 -0 lines
State: FILE REMOVED
Log Message:
removed json examples,
fixed a few minor bugs and added connect/connect_error events
with improved network code.

File Contents

# Content
1 package SOMEConnection;
2 use strict;
3 use Socket;
4 use AnyEvent;
5 use IO::Socket::INET;
6
7 sub new {
8 my $this = shift;
9 my $class = ref($this) || $this;
10 my $self = {@_};
11 bless $self, $class;
12 return $self;
13 }
14
15 sub client_connect {
16 my ($self, $lid) = @_;
17 }
18
19 sub client_disconnect {
20 my ($self, $lid) = @_;
21 }
22
23 sub start_listener {
24 my ($self) = @_;
25
26 my $sock = IO::Socket::INET->new(
27 Listen => 5,
28 ReuseAddr => 1,
29 Reuse => 1,
30 LocalPort => 1236,
31 Proto => 'tcp'
32 );
33
34 $sock or die "Couldn't create listener: $!";
35
36 $self->{listener} =
37 AnyEvent->io (poll => 'r', fh => $sock, cb => sub {
38 my $cl = $sock->accept ()
39 or die "couldn't accept client: $!";
40 binmode $cl;
41 $cl->autoflush (1);
42 $self->handle_client ($cl);
43 });
44 }
45
46 sub write_data {
47 my ($self, $lid, $data) = @_;
48 return unless $self->{$lid . '_r'};
49
50 my $cl = $self->{$lid}->{socket};
51 $self->{$lid}->{write_buffer} .= $data;
52
53 unless ($self->{$lid . '_w'}) {
54 $self->{$lid . '_w'} =
55 AnyEvent->io (poll => 'w', fh => $cl, cb => sub {
56 if (my $data = $self->{$lid}->{write_buffer}) {
57 my $len = syswrite $cl, $data;
58 unless ($len) {
59 if (not defined $len) {
60 warn "error when writing data on $lid: $!";
61 return;
62 } else {
63 delete $self->{$lid . '_w'};
64 }
65 }
66
67 if ($len == length $self->{$lid}->{write_buffer}) {
68 delete $self->{$lid . '_w'};
69 }
70
71 $self->{$lid}->{write_buffer} = substr $self->{$lid}->{write_buffer}, $len;
72 }
73 });
74 }
75 }
76
77 sub handle_client {
78 my ($self, $cl) = @_;
79 my ($chost, $cport) = ($cl->peerhost (), $cl->peerport ());
80 my $lid = "$chost:$cport";
81
82 $self->{$lid}->{socket} = $cl;
83
84 $self->{$lid . '_r'} =
85 AnyEvent->io (poll => 'r', fh => $cl, cb => sub {
86 my $res = sysread $cl, my $data, 1024;
87 if ($res) {
88 $self->{$lid}->{read_buffer} .= $data;
89 $self->handle_data ($lid, \$self->{$lid}->{read_buffer});
90 } else {
91 if (not defined $res) {
92 warn "error when receiving data on $lid: $!";
93 } else {
94 warn "got eof on $lid: $!";
95 }
96 $self->close_client ($lid);
97 }
98 });
99
100 $self->client_connect ($lid);
101 }
102
103 sub handle_data {
104 my ($self, $lid, $buf) = @_;
105 die "implement";
106 }
107
108 sub close_client {
109 my ($self, $lid) = @_;
110 eval { delete $self->{$lid}->{socket} };
111 delete $self->{$lid . '_r'};
112 delete $self->{$lid . '_w'};
113 delete $self->{$lid};
114 $self->client_disconnect ($lid);
115 }
116
117 package JSONConnection;
118 use strict;
119 our @ISA = qw/SOMEConnection/;
120 use JSON::Syck;
121
122 our %CLIENTS;
123
124 sub client_connect {
125 my ($self, $lid) = @_;
126 $CLIENTS{$lid} = 1;
127 $self->{connect_cb}->($self, $lid);
128 }
129
130 sub client_disconnect {
131 my ($self, $lid) = @_;
132 delete $CLIENTS{$lid};
133 }
134
135 sub broadcast {
136 my ($self, $data) = @_;
137 for (keys %CLIENTS) {
138 $self->send_data ($_, $data);
139 }
140 }
141
142 sub send_data {
143 my ($self, $lid, $data) = @_;
144
145 # $JSON::Syck::ImplicitUnicode = 0;
146 my $dump = JSON::Syck::Dump ($data);
147 $self->write_data ($lid, (length $dump) . " " . $dump . "\015\012");
148 }
149
150 sub handle_data {
151 my ($self, $lid, $buf) = @_;
152
153 while ($$buf =~ m/^(\s*(\d+) )(.*)$/s) {
154 my ($prefix, $len, $rembuf) = ($1, $2, $3);
155 if ((length $rembuf) >= $len) {
156 my $data = substr $rembuf, 0, $len;
157 substr $$buf, 0, (length $prefix) + (length $data), '';
158 # $JSON::Syck::ImplicitUnicode = 1;
159 $self->{packet_cb}->($self, $lid, JSON::Syck::Load ($data));
160 } else {
161 return
162 }
163 }
164 }
165
166 package JSONClientConnection;
167 use strict;
168 use JSON::Syck;
169
170 sub new {
171 my $this = shift;
172 my $class = ref($this) || $this;
173 my $self = { disconnect_cb => sub {}, @_ };
174 bless $self, $class;
175 return $self;
176 }
177
178 sub connect {
179 my ($self, $host, $port) = @_;
180
181 $self->{socket}
182 and return;
183
184 my $sock = IO::Socket::INET->new (
185 PeerAddr => $host,
186 PeerPort => $port,
187 Proto => 'tcp',
188 Blocking => 0
189 );
190 die "couldn't connect to server '$host:$port': $!\n" unless $sock->connected;
191
192 $self->{socket} = $sock;
193 $self->{host} = $host;
194 $self->{port} = $port;
195
196 $self->{r} =
197 AnyEvent->io (poll => 'r', fh => $sock, cb => sub {
198 my $l = sysread $sock, my $data, 1024;
199
200 $self->{read_buffer} .= $data;
201 $self->handle_data (\$self->{read_buffer});
202
203 unless ($l) {
204 if (defined $l) {
205 $self->{disconnect_cb}->("EOF from json server '$host:$port'");
206 delete $self->{r};
207 delete $self->{socket};
208 return;
209
210 } else {
211 $self->{disconnect_cb}->("Error while reading from json server '$host:$port': $!");
212 delete $self->{socket};
213 delete $self->{r};
214 return;
215 }
216 }
217 });
218 }
219
220 sub handle_data {
221 my ($self, $buf) = @_;
222
223 while ($$buf =~ m/^(\s*(\d+) )(.*)$/s) {
224 my ($prefix, $len, $rembuf) = ($1, $2, $3);
225 if ((length $rembuf) >= $len) {
226 my $data = substr $rembuf, 0, $len;
227 substr $$buf, 0, (length $prefix) + (length $data), '';
228 # $JSON::Syck::ImplicitUnicode = 1;
229 $self->{packet_cb}->($self, JSON::Syck::Load ($data));
230 } else {
231 return
232 }
233 }
234 }
235
236 sub send_data {
237 my ($self, $data) = @_;
238 # $JSON::Syck::ImplicitUnicode = 0;
239 my $dump = JSON::Syck::Dump ($data);
240 $self->write_data ((length $dump) . " " . $dump . "\015\012");
241 }
242
243 sub write_data {
244 my ($self, $data) = @_;
245 return unless $self->{r};
246
247 my $cl = $self->{socket};
248 $self->{write_buffer} .= $data;
249
250 unless ($self->{w}) {
251 $self->{w} =
252 AnyEvent->io (poll => 'w', fh => $cl, cb => sub {
253 if (my $data = $self->{write_buffer}) {
254 my $len = syswrite $cl, $data;
255 unless ($len) {
256 if (not defined $len) {
257 warn "error when writing data on $self->{host}:$self->{port}: $!";
258 return;
259 } else {
260 delete $self->{w};
261 }
262 }
263
264 if ($len == length $self->{write_buffer}) {
265 delete $self->{w};
266 }
267
268 $self->{write_buffer} = substr $self->{write_buffer}, $len;
269 }
270 });
271 }
272 }
273
274 1;