diff --git a/lib/Lacuna/DB/Result/Schedule.pm b/lib/Lacuna/DB/Result/Schedule.pm new file mode 100644 index 00000000..3cbcf719 --- /dev/null +++ b/lib/Lacuna/DB/Result/Schedule.pm @@ -0,0 +1,63 @@ +package Lacuna::DB::Result::Schedule; + +use Moose; +use utf8; +no warnings qw(uninitialized); +extends 'Lacuna::DB::Result'; +use Lacuna::Util qw(format_date); +use DateTime; +use Data::Dumper; + +__PACKAGE__->table('schedule'); +__PACKAGE__->add_columns( + queue => {data_type => 'varchar', size => 30, is_nullable => 0}, + delivery => {data_type => 'datetime', is_nullable => 0}, + priority => {data_type => 'int', size => 11, is_nullable => 0, default => 1000}, + parent_table => {data_type => 'varchar', size => 30, is_nullable => 0}, + parent_id => {data_type => 'int', size => 11, is_nullable => 0}, + task => {data_type => 'varchar', size => 30, is_nullable => 0}, + args => {data_type => 'medium_blob', is_nullable => 1, serializer_class => 'JSON'}, +); + +after 'insert' => sub { + my $self = shift; + +# my $earliest = DateTime->now->add( hours => 2); +# if ($self->delivery < $earliest) { + # If delivery is within the next couple of hours + # Then put it directly on the beanstalk queue + $self->queue_for_delivery; +# } + return $self; +}; + + +# Put this entry onto the beanstalk queue +# +sub queue_for_delivery { + my ($self) = @_; + + my $dur = $self->delivery->subtract_datetime_absolute(DateTime->now); + my $delay = int($dur->in_units('seconds')); + $delay = 0 if $delay < 0; + + my $queue = Lacuna->queue || 'default'; + my $priority = $self->priority || 1000; + + $queue->publish($self->queue, + { + id => $self->id, + parent_table => $self->parent_table, + parent_id => $self->parent_id, + task => $self->task, + args => $self->args, + },{ + delay => $delay, + priority => $priority, + } + ); +} + +no Moose; +__PACKAGE__->meta->make_immutable(inline_constructor => 0); + diff --git a/lib/Lacuna/Queue.pm b/lib/Lacuna/Queue.pm new file mode 100644 index 00000000..00b768f8 --- /dev/null +++ b/lib/Lacuna/Queue.pm @@ -0,0 +1,178 @@ +package Lacuna::Queue; + +use Moose; +use Beanstalk::Client; +use Data::Dumper; + +use Lacuna::Queue::Job; + +has '_beanstalk' => ( + is => 'ro', + isa => 'Beanstalk::Client', + lazy => 1, + builder => '__build_beanstalk', +); + +has 'max_timeouts' => ( + is => 'ro', + isa => 'Int', + lazy => 1, + default => 10, +); + +has 'max_reserves' => ( + is => 'ro', + isa => 'Int', + lazy => 1, + default => 10, +); + +has 'server' => ( + is => 'ro', + isa => 'Str', + lazy => 1, + default => 'localhost', +); + +has 'ttr' => ( + is => 'ro', + isa => 'Int', + lazy => 1, + default => 120, +); + +has 'debug' => ( + is => 'ro', + isa => 'Int', + lazy => 1, + default => 0, +); + + +sub __build_beanstalk { + my ($self) = @_; + + my $beanstalk = Beanstalk::Client->new({ + server => $self->server, + ttr => $self->ttr, + debug => $self->debug, + }); + return $beanstalk; +} + +sub publish { + my ($self, $queue, $payload, $options) = @_; + + my $beanstalk = $self->_beanstalk; + $options = defined $options ? $options : {}, + $beanstalk->use($queue); + + my $job = $beanstalk->put($options, $payload); + + return Lacuna::Queue::Job->new({job => $job}); +} + +sub peek { + my ($self, $job_id) = @_; + + my $beanstalk = $self->_beanstalk; + + my $job = $beanstalk->peek($job_id); + if ($job) { + return Lacuna::Queue::Job->new({job => $job}); + } + return; +} + +# DRY Principle +my $meta = __PACKAGE__->meta; + +foreach my $proc (qw(peek_buried peek_ready peek_delayed)) { + $meta->add_method($proc => sub { + my ($self) = @_; + + my $beanstalk = $self->_beanstalk; + my $job = $beanstalk->$proc; + if ($job) { + return Lacuna::Queue::Job->new({job => $job}); + } + return; + }); +} + +sub kick { + my ($self, $bound) = @_; + + $bound = $bound || 1; + + my $beanstalk = $self->_beanstalk; + my $kicked = $beanstalk->kick($bound); + + return $kicked; +} + +sub pause_tube { + my ($self, $tube, $seconds) = @_; + + $seconds = $seconds || 0; + + my $beanstalk = $self->_beanstalk; + my $ret = $beanstalk->pause_tube($tube, $seconds); +} + +sub stats { + my ($self) = @_; + + return $self->_beanstalk->stats; +} + +sub stats_tube { + my ($self, $tube) = @_; + + return $self->_beanstalk->stats_tube($tube); +} + +sub list_tubes { + my ($self) = @_; + + return $self->_beanstalk->list_tubes; +} + +sub consume { + my ($self,$tube) = @_; + + my $job; + my $beanstalk = $self->_beanstalk; + + RESERVE: + while (not $job) { + $beanstalk->watch_only($tube); + $job = $beanstalk->reserve; + + # Defend against undef jobs (most likely due to DEADLINE_SOON) + if (not $job) { + sleep 1; + redo RESERVE; + } + my $stats = $job->stats; + my $bury; + + if ($stats->timeouts > $self->max_timeouts) { + $bury = "timeouts"; + } + if ($stats->reserves > $self->max_reserves) { + $bury = "reserves"; + } + if ($bury) { + $job->bury; + undef $job; + } + } + return Lacuna::Queue::Job->new({job => $job}); +} + +__PACKAGE__->meta->make_immutable; + +1; + + diff --git a/lib/Lacuna/Queue/Job.pm b/lib/Lacuna/Queue/Job.pm new file mode 100644 index 00000000..a7d91121 --- /dev/null +++ b/lib/Lacuna/Queue/Job.pm @@ -0,0 +1,27 @@ +package Lacuna::Queue::Job; + +use Moose; +use YAML; + + +has 'job' => ( + is => 'ro', + isa => 'Beanstalk::Job', + required => 1, + handles => [qw(id buried reserved data error stats delete touch peek release bury args tube ttr priority)], +); + +sub payload { + my ($self) = @_; + + my $args = $self->job->args; + my $class = $args->{parent_table}; + my $id = $args->{parent_id}; + + my $thing = Lacuna->db->resultset($class)->find($id); + return $thing; +} + +__PACKAGE__->meta->make_immutable; +1; + diff --git a/lib/Lacuna/Queue/Job.pm.new b/lib/Lacuna/Queue/Job.pm.new new file mode 100644 index 00000000..a7d91121 --- /dev/null +++ b/lib/Lacuna/Queue/Job.pm.new @@ -0,0 +1,27 @@ +package Lacuna::Queue::Job; + +use Moose; +use YAML; + + +has 'job' => ( + is => 'ro', + isa => 'Beanstalk::Job', + required => 1, + handles => [qw(id buried reserved data error stats delete touch peek release bury args tube ttr priority)], +); + +sub payload { + my ($self) = @_; + + my $args = $self->job->args; + my $class = $args->{parent_table}; + my $id = $args->{parent_id}; + + my $thing = Lacuna->db->resultset($class)->find($id); + return $thing; +} + +__PACKAGE__->meta->make_immutable; +1; + diff --git a/lib/Lacuna/Role/Navigation.pm b/lib/Lacuna/Role/Navigation.pm new file mode 100644 index 00000000..cc944c0a --- /dev/null +++ b/lib/Lacuna/Role/Navigation.pm @@ -0,0 +1,37 @@ +package Lacuna::Role::Navigation; + +use Moose::Role; + +# Find a 'target' based on a number of methods +# +sub find_target { + my ($self, $target_params) = @_; + unless (ref $target_params eq 'HASH') { + confess [-32602, 'The target parameter should be a hash reference. For example { "star_id" : 9999 }.']; + } + my $target; + if (exists $target_params->{star_id}) { + $target = Lacuna->db->resultset('Map::Star')->find($target_params->{star_id}); + } + elsif (exists $target_params->{star_name}) { + $target = Lacuna->db->resultset('Map::Star')->search({ name => $target_params->{star_name} }, {rows=>1})->single; + } + if (exists $target_params->{body_id}) { + $target = Lacuna->db->resultset('Map::Body')->find($target_params->{body_id}); + } + elsif (exists $target_params->{body_name}) { + $target = Lacuna->db->resultset('Map::Body')->search({ name => $target_params->{body_name} }, {rows=>1})->single; + } + elsif (exists $target_params->{x}) { + $target = Lacuna->db->resultset('Map::Body')->search({ x => $target_params->{x}, y => $target_params->{y} }, {rows=>1})->single; + unless (defined $target) { + $target = Lacuna->db->resultset('Map::Star')->search({ x => $target_params->{x}, y => $target_params->{y} }, {rows=>1})->single; + } + } + unless (defined $target) { + confess [ 1002, 'Could not find the target.', $target]; + } + return $target; +} +1; + diff --git a/t/500_schedule.t b/t/500_schedule.t new file mode 100644 index 00000000..de2e9feb --- /dev/null +++ b/t/500_schedule.t @@ -0,0 +1,67 @@ +use lib '../lib'; + +use strict; +use warnings; + +use Test::More tests => 7; +use Test::Deep; +use Test::Memory::Cycle; +use Data::Dumper; +use 5.010; +use DateTime; +use Lacuna; +use TestHelper; + + +my $now = DateTime->now; +my $later = DateTime->now->add( seconds => 3); + +my $dur = $later->subtract_datetime_absolute($now); +my $seconds = $dur->in_units('seconds'); + +is($seconds, 3, "CPAN modules agree on seconds"); + +my $db = Lacuna->db; +my $thing = $db->resultset('ApiKey')->create({ + public_key => 'foo', + private_key => 'bar', + name => 'iain', + ip_address => '10.11.12.13', + email => 'iain@docherty.me', +}); + +my $schedule = $db->resultset('Schedule')->create({ + queue => 'foo', + delivery => $later, + parent_table => 'ApiKey', + parent_id => $thing->id, + task => 'bar', + args => {this => 'siht', that => 'taht'}, +}); + +isa_ok($schedule, 'Lacuna::DB::Result::Schedule', 'Correct class'); + +# Now test against beanstalk (it must be running) +# +my $queue = Lacuna::Queue->new; + +isa_ok($queue, 'Lacuna::Queue', 'Correct queue class'); + +my $job = $queue->consume('foo'); + +isa_ok($job, 'Lacuna::Queue::Job', 'Correct job class'); + +my $payload = $job->payload; +isa_ok($payload, 'Lacuna::DB::Result::ApiKey', 'Got back an ApiKey'); +is($payload->public_key,'foo', 'foo found'); +is($payload->name,'iain', 'iain found'); + +$now = DateTime->now; +diag("later = [$later] now = [$now]"); + +# Delete this job, we no longer need it +$job->delete; + + +1; +