From 7edb6bbfa0c2f5b74fb88c2c7db0f72ddbc5c423 Mon Sep 17 00:00:00 2001 From: Bruno Tavares Date: Wed, 3 Sep 2014 10:19:11 +0100 Subject: [PATCH 1/7] Connection routing registry --- lib/AnyEvent/NSQ/Reader.pm | 39 ++++++++++++++++++++++++++++++++++++++ 1 file changed, 39 insertions(+) diff --git a/lib/AnyEvent/NSQ/Reader.pm b/lib/AnyEvent/NSQ/Reader.pm index 9076e91..b288d5c 100644 --- a/lib/AnyEvent/NSQ/Reader.pm +++ b/lib/AnyEvent/NSQ/Reader.pm @@ -10,6 +10,15 @@ use Carp 'croak'; use parent 'AnyEvent::NSQ::Client'; +sub new { + my $class = shift; + my $self = $class->SUPER::new(@_); + + # To keep registry of the connection bound to each message + $self->{routing} = {}; + + return $self; +} #### Parameter parsing @@ -46,15 +55,45 @@ sub _identified { ## FIXME: bless it with the future AnyEvent::NSQ::Message my $msg = $_[1]; + ## Keep the connection in the registry in case the user + ## wants to issue FIN/REQs inside the callback + $self->{routing}->{$msg->{message_id}} = $conn; + my $action = $self->{message_cb}->($self, $msg); ## Action below -1 does nothing, we assume the user took care of it himself if (not defined $action) { $conn->mark_as_done_msg($_[1]) } elsif ($action >= -1) { $conn->requeue_msg($_[1], $action) } + + delete($self->{routing}->{$msg->{message_id}}); } ); $conn->ready($self->{ready_count} || int(($info->{max_rdy_count} || 2000) / 10)); } +sub mark_as_done_msg { + my $self = shift; + my $message = shift; +} + +sub requeue_msg { + my $self = shift; + my $message = shift; +} + +sub touch_message { + my $self = shift; + my $message = shift; +} + +sub _find_message_connection { + my $self = shift; + my $message = shift; + + my $message_id = ref($message) ? $message->{message_id} : $message; + + return $self->{routing}->{$message_id}; +} + 1; From ee0b61582a05d3f800b2534f19c1e4b5f296e11d Mon Sep 17 00:00:00 2001 From: Bruno Tavares Date: Wed, 3 Sep 2014 11:02:54 +0100 Subject: [PATCH 2/7] allow to issue commands with message id ( and message structs ) --- lib/AnyEvent/NSQ/Connection.pm | 16 +++++++++++----- 1 file changed, 11 insertions(+), 5 deletions(-) diff --git a/lib/AnyEvent/NSQ/Connection.pm b/lib/AnyEvent/NSQ/Connection.pm index 1f00362..a80e9ad 100644 --- a/lib/AnyEvent/NSQ/Connection.pm +++ b/lib/AnyEvent/NSQ/Connection.pm @@ -167,7 +167,9 @@ sub mark_as_done_msg { my ($self, $msg) = @_; return unless my $hdl = $self->{handle}; - $hdl->push_write("FIN $msg->{message_id}\012"); + my $id = ref($msg) ? $msg->{message_id} : $msg; + + $hdl->push_write("FIN $id\012"); return; } @@ -176,10 +178,13 @@ sub requeue_msg { my ($self, $msg, $delay) = @_; return unless my $hdl = $self->{handle}; - $delay = 0 unless defined $delay; - $delay = $msg->{attempts} * $self->{requeue_delay} if $delay < 0; + my $id = ref($msg) ? $msg->{message_id} : $msg; + my $attempts = ref($msg) ? $msg->{attempts} : 1; - $hdl->push_write("REQ $msg->{message_id} $delay\012"); + $delay = 0 unless defined $delay; + $delay = $attempts * $self->{requeue_delay} if $delay < 0; + + $hdl->push_write("REQ $id $delay\012"); return; } @@ -188,7 +193,8 @@ sub touch_msg { my ($self, $msg) = @_; return unless my $hdl = $self->{handle}; - $hdl->push_write("TOUCH $msg->{message_id}\012"); + my $id = ref($msg) ? $msg->{message_id} : $msg; + $hdl->push_write("TOUCH $id\012"); return; } From 73801c8488b30d87e4f1336ba28f80062cf28156 Mon Sep 17 00:00:00 2001 From: Bruno Tavares Date: Wed, 3 Sep 2014 11:03:14 +0100 Subject: [PATCH 3/7] interface methods to message commands on the connection --- lib/AnyEvent/NSQ/Reader.pm | 40 ++++++++++++++++++++++++++++---------- 1 file changed, 30 insertions(+), 10 deletions(-) diff --git a/lib/AnyEvent/NSQ/Reader.pm b/lib/AnyEvent/NSQ/Reader.pm index b288d5c..a70043a 100644 --- a/lib/AnyEvent/NSQ/Reader.pm +++ b/lib/AnyEvent/NSQ/Reader.pm @@ -73,27 +73,47 @@ sub _identified { } sub mark_as_done_msg { - my $self = shift; - my $message = shift; + my ($self, $msg) = @_; + + my $conn = $self->_find_message_connection($msg); + return 0 unless $conn; + + $conn->mark_as_done_msg($msg); + return 1; } sub requeue_msg { - my $self = shift; - my $message = shift; + my ($self, $msg, $delay) = @_; + + my $conn = $self->_find_message_connection($msg); + return 0 unless $conn; + + $conn->requeue_msg($msg, $delay); + return 1; } sub touch_message { - my $self = shift; - my $message = shift; + my ($self, $msg) = @_; + + my $conn = $self->_find_message_connection($msg); + return 0 unless $conn; + + $conn->touch_msg($msg); + return 1; } sub _find_message_connection { - my $self = shift; - my $message = shift; + my ($self, $msg) = @_; - my $message_id = ref($message) ? $message->{message_id} : $message; + my $id = ref($msg) ? $msg->{message_id} : $msg; - return $self->{routing}->{$message_id}; + my $conn = $self->{routing}->{$id}; + + if ( !$conn ) { + warn "WARN: Could not find the connection to route msg $id"; + } + + return $conn; } 1; From 78c4ca959fa62adad9490c44252473b793f26ace Mon Sep 17 00:00:00 2001 From: Bruno Tavares Date: Wed, 3 Sep 2014 18:22:22 +0100 Subject: [PATCH 4/7] Make the connection references weak references. --- lib/AnyEvent/NSQ/Reader.pm | 2 ++ 1 file changed, 2 insertions(+) diff --git a/lib/AnyEvent/NSQ/Reader.pm b/lib/AnyEvent/NSQ/Reader.pm index a70043a..6cfed56 100644 --- a/lib/AnyEvent/NSQ/Reader.pm +++ b/lib/AnyEvent/NSQ/Reader.pm @@ -7,6 +7,7 @@ package AnyEvent::NSQ::Reader; use strict; use warnings; use Carp 'croak'; +use Scalar::Util qw( weaken ); use parent 'AnyEvent::NSQ::Client'; @@ -58,6 +59,7 @@ sub _identified { ## Keep the connection in the registry in case the user ## wants to issue FIN/REQs inside the callback $self->{routing}->{$msg->{message_id}} = $conn; + weaken( $self->{routing}->{$msg->{message_id}} ); my $action = $self->{message_cb}->($self, $msg); From 24b133500065da8d40feedadc28f58cdbb019802 Mon Sep 17 00:00:00 2001 From: Bruno Tavares Date: Tue, 9 Sep 2014 11:52:56 +0100 Subject: [PATCH 5/7] Better routing clean up. Now requires explicit message handling --- examples/consumer.pl | 16 ++++++++++++---- lib/AnyEvent/NSQ/Reader.pm | 20 ++++++-------------- t/10-reader.t | 1 + 3 files changed, 19 insertions(+), 18 deletions(-) diff --git a/examples/consumer.pl b/examples/consumer.pl index be4751e..ae1df27 100755 --- a/examples/consumer.pl +++ b/examples/consumer.pl @@ -18,9 +18,6 @@ usage("topic and channel are required parameters") unless $topic and $channel; my $cv = AE::cv; -## return undef => mark_as_done_msg() -my $message_cb = $print ? sub { print "$_[1]{message}\n"; return } : sub {return}; - my $t = my $p = 0; my $r = AnyEvent::NSQ::Reader->new( topic => $topic, @@ -28,7 +25,7 @@ my $r = AnyEvent::NSQ::Reader->new( nsqd_tcp_addresses => '127.0.0.1', client_id => "${channel}_consumer/pid_$$", - message_cb => sub { $t++; $p++; $message_cb->(@_) }, + message_cb => \&message_handler, error_cb => sub { warn "$_[1]\n" if $verbose }, disconnect_cb => sub { warn "Disconnected after $t total messages... exiting...\n" if $verbose; $cv->send }, @@ -85,3 +82,14 @@ Usage: consumer.pl [--help|-h] [--print|-p] topic channel exit(1); } + +sub message_handler { + use Data::Dumper; + my ($reader, $message) = @_; + + if ($print) { + print $message->{message}."\n"; + } + + $reader->mark_as_done_msg($message); +} diff --git a/lib/AnyEvent/NSQ/Reader.pm b/lib/AnyEvent/NSQ/Reader.pm index 6cfed56..76d7068 100644 --- a/lib/AnyEvent/NSQ/Reader.pm +++ b/lib/AnyEvent/NSQ/Reader.pm @@ -63,11 +63,6 @@ sub _identified { my $action = $self->{message_cb}->($self, $msg); - ## Action below -1 does nothing, we assume the user took care of it himself - if (not defined $action) { $conn->mark_as_done_msg($_[1]) } - elsif ($action >= -1) { $conn->requeue_msg($_[1], $action) } - - delete($self->{routing}->{$msg->{message_id}}); } ); @@ -77,8 +72,7 @@ sub _identified { sub mark_as_done_msg { my ($self, $msg) = @_; - my $conn = $self->_find_message_connection($msg); - return 0 unless $conn; + my $conn = $self->_find_and_delete_message_connection($msg); $conn->mark_as_done_msg($msg); return 1; @@ -87,8 +81,7 @@ sub mark_as_done_msg { sub requeue_msg { my ($self, $msg, $delay) = @_; - my $conn = $self->_find_message_connection($msg); - return 0 unless $conn; + my $conn = $self->_find_and_delete_message_connection($msg); $conn->requeue_msg($msg, $delay); return 1; @@ -97,22 +90,21 @@ sub requeue_msg { sub touch_message { my ($self, $msg) = @_; - my $conn = $self->_find_message_connection($msg); - return 0 unless $conn; + my $conn = $self->_find_and_delete_message_connection($msg); $conn->touch_msg($msg); return 1; } -sub _find_message_connection { +sub _find_and_delete_message_connection { my ($self, $msg) = @_; my $id = ref($msg) ? $msg->{message_id} : $msg; - my $conn = $self->{routing}->{$id}; + my $conn = delete($self->{routing}->{$id}); if ( !$conn ) { - warn "WARN: Could not find the connection to route msg $id"; + croak "WARN: Could not find the connection to route msg $id"; } return $conn; diff --git a/t/10-reader.t b/t/10-reader.t index 1545b29..b686640 100755 --- a/t/10-reader.t +++ b/t/10-reader.t @@ -17,6 +17,7 @@ subtest 'basic connection' => sub { message_cb => sub { print STDERR "!!!! GOT MESSAGE '$_[1]{message}\n"; + $_[0]->mark_as_done_msg($_[1]); return; }, ); From 05bea0dcb8cdeaeed1850d372d884e443c396160 Mon Sep 17 00:00:00 2001 From: Bruno Tavares Date: Tue, 9 Sep 2014 11:59:43 +0100 Subject: [PATCH 6/7] Lost debugging lines. --- examples/consumer.pl | 1 - 1 file changed, 1 deletion(-) diff --git a/examples/consumer.pl b/examples/consumer.pl index ae1df27..51539dd 100755 --- a/examples/consumer.pl +++ b/examples/consumer.pl @@ -84,7 +84,6 @@ Usage: consumer.pl [--help|-h] [--print|-p] topic channel } sub message_handler { - use Data::Dumper; my ($reader, $message) = @_; if ($print) { From f2cf93d8e573f96314a6e4d5a10b225c74815a6d Mon Sep 17 00:00:00 2001 From: Bruno Tavares Date: Wed, 10 Sep 2014 17:17:06 +0100 Subject: [PATCH 7/7] no need to allocate a variable that has no use --- lib/AnyEvent/NSQ/Reader.pm | 3 +-- 1 file changed, 1 insertion(+), 2 deletions(-) diff --git a/lib/AnyEvent/NSQ/Reader.pm b/lib/AnyEvent/NSQ/Reader.pm index 76d7068..4e0dec1 100644 --- a/lib/AnyEvent/NSQ/Reader.pm +++ b/lib/AnyEvent/NSQ/Reader.pm @@ -61,8 +61,7 @@ sub _identified { $self->{routing}->{$msg->{message_id}} = $conn; weaken( $self->{routing}->{$msg->{message_id}} ); - my $action = $self->{message_cb}->($self, $msg); - + $self->{message_cb}->($self, $msg); } );