Beekeeper
view release on metacpan or search on metacpan
lib/Beekeeper/Worker.pm view on Meta::CPAN
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;
}
);
}
( run in 1.298 second using v1.01-cache-2.11-cpan-bbcb1afb8fc )