Beekeeper
view release on metacpan or search on metacpan
lib/Beekeeper/Worker.pm view on Meta::CPAN
my $bus_id = $self->{_WORKER}->{bus_id};
my $config = $self->{_WORKER}->{bus_config}->{$bus_id};
my $client = Beekeeper::Client->new(
%$config,
timeout => 0, # retry forever
on_error => sub {
my $errmsg = $_[0] || ""; $errmsg =~ s/\s+/ /sg;
log_fatal "Connection to $bus_id failed: $errmsg";
$self->stop_working;
},
);
$self->{_CLIENT} = $client->{_CLIENT};
$self->{_BUS} = $client->{_BUS};
$Beekeeper::Client::singleton = $self;
}
sub __init_worker {
my $self = shift;
$self->on_startup;
$self->__report_status;
AnyEvent->now_update;
$self->{_WORKER}->{report_status_timer} = AnyEvent->timer(
after => rand( $REPORT_STATUS_PERIOD ),
interval => $REPORT_STATUS_PERIOD,
cb => sub { $self->__report_status },
);
}
sub on_startup {
# Placeholder, intended to be overrided
my $class = ref $_[0];
log_fatal "Worker class $class doesn't define on_startup() method";
}
sub on_shutdown {
# Placeholder, can be overrided
}
sub authorize_request {
# Placeholder, must to be overrided
my $class = ref $_[0];
log_fatal "Worker class $class doesn't define authorize_request() method";
return undef; # do not authorize
}
sub accept_notifications {
my ($self, %args) = @_;
my $worker = $self->{_WORKER};
my $callbacks = $worker->{callbacks};
my ($file, $line) = (caller)[1,2];
my $at = "at $file line $line\n";
foreach my $fq_meth (keys %args) {
$fq_meth =~ m/^ ( [\w-]+ (?: \.[\w-]+ )* )
\. ( [\w-]+ | \* ) $/x or die "Invalid notification method '$fq_meth' $at";
my ($service, $method) = ($1, $2);
my $callback = $self->__get_cb_coderef($fq_meth, $args{$fq_meth});
die "Already accepting notifications '$fq_meth' $at" if exists $callbacks->{"msg.$fq_meth"};
$callbacks->{"msg.$fq_meth"} = $callback;
my $local_bus = $self->{_BUS}->{bus_role};
my $topic = "msg/$local_bus/$service/$method";
$topic =~ tr|.*|/#|;
$self->{_BUS}->subscribe(
topic => $topic,
on_publish => sub {
# ($payload_ref, $properties) = @_;
# Enqueue notification
push @{$worker->{task_queue_high}}, [ @_ ];
unless ($worker->{queued_tasks}) {
$worker->{queued_tasks} = 1;
AnyEvent::postpone { $self->__drain_task_queue };
}
},
on_suback => sub {
my ($success, $prop) = @_;
die "Could not subscribe to topic '$topic' $at" unless $success;
}
);
}
}
sub __get_cb_coderef {
my ($self, $method, $callback) = @_;
if (ref $callback eq 'CODE') {
# Already a coderef
return $callback;
}
elsif (!ref($callback) && $self->can($callback)) {
# Return a reference to given method
no strict 'refs';
my $class = ref $self;
return \&{"${class}::${callback}"};
}
else {
my ($file, $line) = (caller(1))[1,2];
my $at = "at $file line $line\n";
die "Invalid handler '$callback' for '$method' $at";
}
}
sub accept_remote_calls {
my ($self, %args) = @_;
my $worker = $self->{_WORKER};
my $callbacks = $worker->{callbacks};
my %subscribed_to;
my ($file, $line) = (caller)[1,2];
my $at = "at $file line $line\n";
foreach my $fq_meth (keys %args) {
$fq_meth =~ m/^ ( [\w-]+ (?: \.[\w-]+ )* )
\. ( [\w-]+ | \* ) $/x or die "Invalid remote call method '$fq_meth' $at";
my ($service, $method) = ($1, $2);
my $callback = $self->__get_cb_coderef($fq_meth, $args{$fq_meth});
die "Already accepting remote calls '$fq_meth' $at" if exists $callbacks->{"req.$fq_meth"};
$callbacks->{"req.$fq_meth"} = $callback;
next if $subscribed_to{$service};
$subscribed_to{$service} = 1;
if (keys %subscribed_to == 2) {
log_warn "Running multiple services within a single worker hurts load balancing $at";
}
my $local_bus = $self->{_BUS}->{bus_role};
my $topic = "\$share/BKPR/req/$local_bus/$service";
$topic =~ tr|.*|/#|;
$self->{_BUS}->subscribe(
topic => $topic,
maximum_qos => 1,
on_publish => sub {
# ($payload_ref, $mqtt_properties) = @_;
# Enqueue request
push @{$worker->{task_queue_low}}, [ @_ ];
unless ($worker->{queued_tasks}) {
$worker->{queued_tasks} = 1;
AnyEvent::postpone { $self->__drain_task_queue };
}
},
on_suback => sub {
my ($success, $prop) = @_;
die "Could not subscribe to topic '$topic' $at" unless $success;
}
);
}
}
my $_TASK_QUEUE_DEPTH = 0;
sub __drain_task_queue {
my $self = shift;
# Ensure that draining does not recurse
Carp::confess "Unexpected task queue processing recursion" if $_TASK_QUEUE_DEPTH;
$_TASK_QUEUE_DEPTH++;
my $timing_tasks;
unless (defined $BUSY_SINCE) {
lib/Beekeeper/Worker.pm view on Meta::CPAN
# fire & forget calls doesn't expect responses
return unless defined $request->{id};
unless (defined $BUSY_SINCE) {
$BUSY_SINCE = Time::HiRes::time;
$timing_tasks = 1;
}
if (blessed($result) && $result->isa('Beekeeper::JSONRPC::Error')) {
# Explicit error response
$response = $result;
$self->{_WORKER}->{error_count}++;
}
else {
# Build a success response
$response = {
jsonrpc => '2.0',
result => $result,
};
}
$response->{id} = $request->{id};
local $@;
my $json = eval { $JSON->encode( $response ) };
if ($@) {
# Probably response contains blessed references
log_error "Couldn't serialize response as JSON: $@";
$response = Beekeeper::JSONRPC::Error->server_error;
$response->{id} = $request->{id};
$json = $JSON->encode( $response );
$self->{_WORKER}->{error_count}++;
}
if ($request->{_deflate_response} && length($json) > $request->{_deflate_response}) {
my $compressed_json;
$DEFLATE->deflate(\$json, $compressed_json);
$DEFLATE->flush(\$compressed_json);
$DEFLATE->deflateReset;
$json = $compressed_json;
}
$self->{_BUS}->publish(
topic => $request->{_mqtt_properties}->{'response_topic'},
addr => $request->{_mqtt_properties}->{'addr'},
payload => \$json,
);
if (defined $timing_tasks) {
$BUSY_TIME += Time::HiRes::time - $BUSY_SINCE;
undef $BUSY_SINCE;
}
}
sub stop_accepting_notifications {
my ($self, @methods) = @_;
my ($file, $line) = (caller)[1,2];
my $at = "at $file line $line\n";
die "No method specified $at" unless @methods;
foreach my $fq_meth (@methods) {
$fq_meth =~ m/^ ( [\w-]+ (?: \.[\w-]+ )* )
\. ( [\w-]+ | \* ) $/x or die "Invalid method '$fq_meth' $at";
my ($service, $method) = ($1, $2);
my $worker = $self->{_WORKER};
unless (defined $worker->{callbacks}->{"msg.$fq_meth"}) {
log_warn "Not previously accepting notifications '$fq_meth' $at";
next;
}
my $local_bus = $self->{_BUS}->{bus_role};
my $topic = "msg/$local_bus/$service/$method";
$topic =~ tr|.*|/#|;
# Cannot remove callbacks right now, as new notifications could be in flight or be
# already queued. We must wait for unsubscription completion, and then until the
# notification queue is empty to ensure that all received ones were processed. And
# even then wait a bit more, as some brokers may send messages *after* unsubscription.
my $postpone = sub {
my $unsub_tmr; $unsub_tmr = AnyEvent->timer(
after => $UNSUBSCRIBE_LINGER, cb => sub {
delete $worker->{callbacks}->{"msg.$fq_meth"};
undef $unsub_tmr;
}
);
};
$self->{_BUS}->unsubscribe(
topic => $topic,
on_unsuback => sub {
my ($success, $prop) = @_;
log_error "Could not unsubscribe from topic '$topic' $at" unless $success;
my $postponed = $worker->{postponed} ||= [];
push @$postponed, $postpone;
AnyEvent::postpone { $self->__drain_task_queue };
}
);
}
}
sub stop_accepting_calls {
my ($self, @methods) = @_;
my ($file, $line) = (caller)[1,2];
my $at = "at $file line $line\n";
die "No method specified $at" unless @methods;
foreach my $fq_meth (@methods) {
$fq_meth =~ m/^ ( [\w-]+ (?: \.[\w-]+ )* )
\. ( [\w-]+ | \* ) $/x or die "Invalid remote call method '$fq_meth' $at";
my ($service, $method) = ($1, $2);
unless ($method eq '*') {
# Known limitation. As all calls for an entire service class are received
# through a single MQTT subscription (in order to load balance them), it is
# not possible to reject a single method. A workaround is to use a different
# class for each method that need to be individually rejected.
die "Cannot stop accepting individual methods, only '$service.*' is allowed $at";
}
my $worker = $self->{_WORKER};
my $callbacks = $worker->{callbacks};
my @cb_keys = grep { $_ =~ m/^req.\Q$service\E\b/ } keys %$callbacks;
unless (@cb_keys) {
log_warn "Not previously accepting remote calls '$fq_meth' $at";
next;
}
my $local_bus = $self->{_BUS}->{bus_role};
my $topic = "\$share/BKPR/req/$local_bus/$service";
$topic =~ tr|.*|/#|;
# Cannot remove callbacks right now, as new requests could be in flight or be already
# queued. We must wait for unsubscription completion, and then until the task queue
# is empty to ensure that all received requests were processed. And even then wait a
# bit more, as some brokers may send requests *after* unsubscription.
my $postpone = sub {
$worker->{stop_cv}->begin;
my $unsub_tmr; $unsub_tmr = AnyEvent->timer(
after => $UNSUBSCRIBE_LINGER, cb => sub {
delete $worker->{callbacks}->{$_} foreach @cb_keys;
delete $worker->{subscriptions}->{$service};
undef $unsub_tmr;
return unless $worker->{shutting_down};
if ($worker->{in_progress} > 0) {
# The task queue is empty now, but an asynchronous method handler is
# still busy processing some requests received previously. Wait for
# these requests to be completed before telling _work_forever to stop
my $wait_time = 60;
$worker->{stop_cv}->begin;
( run in 3.043 seconds using v1.01-cache-2.11-cpan-84de2e75c66 )