From 699fc820213daf160e272943700b6333acf2d29a Mon Sep 17 00:00:00 2001 From: Ian Docherty Date: Thu, 13 Dec 2012 08:40:50 -0500 Subject: [PATCH] Added scheduler that uses beanstalk --- bin/schedule_daemon.pl | 105 ++++++++++++++++++++++ lib/Lacuna/DB/Result/Building.pm | 40 +++++++-- lib/Lacuna/DB/Result/Map/Body/Planet.pm | 79 +++++++++-------- lib/Lacuna/DB/Result/Schedule.pm | 9 +- lib/Lacuna/Queue.pm | 1 + t/500_schedule.t | 111 +++++++++++++++--------- t/510_bs_building.t | 69 +++++++++++++++ 7 files changed, 327 insertions(+), 87 deletions(-) create mode 100644 bin/schedule_daemon.pl create mode 100644 t/510_bs_building.t diff --git a/bin/schedule_daemon.pl b/bin/schedule_daemon.pl new file mode 100644 index 00000000..1e6be91e --- /dev/null +++ b/bin/schedule_daemon.pl @@ -0,0 +1,105 @@ +use 5.010; +use strict; +use feature "switch"; +use lib '/data/Lacuna-Server/lib'; +use Lacuna::DB; +use Lacuna; +use Lacuna::Util qw(randint format_date); +use Getopt::Long; +use App::Daemon qw(daemonize ); +use Data::Dumper; +use Try::Tiny; + +$|=1; + +# -------------------------------------------------------------------- +# command line arguments: +# +my $daemonise = 1; +my $loop = 1; +our $quiet = 1; + +GetOptions( + 'daemonise!' => \$daemonise, + 'loop!' => \$loop, + 'quiet!' => \$quiet, +); + +# Catch SIG INT +# +#my $sig_int = 0; +#local $SIG{'INT'} = sub { $sig_int = 1; }; +my $start = time; + +# -------------------------------------------------------------------- +# Daemonise + +if ($daemonise) { + daemonize(); + out('Running as a daemon'); +} +else { + out('Running in the foreground'); +} + +my $config = Lacuna->config; + +my $queue = Lacuna::Queue->new({ + max_timeouts => $config->get('beanstalk/max_timeouts'), + max_reserves => $config->get('beanstalk/max_reserves'), + server => $config->get('beanstalk/server'), + ttr => $config->get('beanstalk/ttr'), + debug => $config->get('beanstalk/debug'), +}); + +out("queue = $queue"); + +# -------------------------------------------------------------------- +# Main processing loop + +out('Started'); +do { + out('In Main Processing Loop'); + my $job = $queue->consume('default'); + my $args = $job->args; + my $task = $args->{task}; + my $task_args = $args->{args}; + + out('job received ['.$job->id.']'); + + my $payload = $job->payload; + + try { + # process the job + out("Process class=$payload task=$task"); + $payload->$task($task_args); + out("Processing done. Delete job ".$job->id); + $job->delete; + } + catch { + # bury the job, it failed + out("Job ".$job->id." failed: $_"); + $job->bury; + }; +# if ($sig_int) { +# out('Received INT signal, jumping out of polling loop'); +# undef $loop; +# } +} while ($loop); + +my $finish = time; +out('Finished'); +out(int(($finish - $start)/60)." minutes have elapsed"); +exit 0; + +############### +## SUBROUTINES +############### + +sub out { + my ($message) = @_; + if (not $quiet) { + say format_date(DateTime->now), " ", $message; + } +} + diff --git a/lib/Lacuna/DB/Result/Building.pm b/lib/Lacuna/DB/Result/Building.pm index f6624eea..f892b0b5 100755 --- a/lib/Lacuna/DB/Result/Building.pm +++ b/lib/Lacuna/DB/Result/Building.pm @@ -816,22 +816,30 @@ sub start_upgrade { $cost ||= $self->cost_to_upgrade; # set time to build, plus what's in the queue - my $time_to_build = $in_parallel ? DateTime->now : $body->get_existing_build_queue_time; + my $now = DateTime->now; + my $upgrade_ends = $in_parallel ? $now : $body->get_existing_build_queue_time; + if ($upgrade_ends < $now) { + $upgrade_ends = $now; + } + my $time_to_add = $body->isa('Lacuna::DB::Result::Map::Body::Planet::Station') ? 60 * 60 * 72 : $cost->{time}; - $time_to_build->add(seconds=>$time_to_add); +# print STDERR "start_upgrade, building=[$self] time_to_add=$time_to_add upgrade_ends=$upgrade_ends\n"; + $upgrade_ends->add(seconds=>$time_to_add); +# print STDERR "start_upgrade, building=[$self] new upgrade_ends=$upgrade_ends now=".DateTime->now."\n"; # add to queue $self->update({ is_upgrading => 1, upgrade_started => DateTime->now, - upgrade_ends => $time_to_build, + upgrade_ends => $upgrade_ends, }); my $schedule = Lacuna->db->resultset('Schedule')->create({ - delivery => $time_to_build, + delivery => $upgrade_ends, parent_table => 'Building', parent_id => $self->id, task => 'finish_upgrade', }); + return $self; } sub finish_upgrade { @@ -853,7 +861,13 @@ sub finish_upgrade { my %levels = (5=>'a quiet',10=>'an extravagant',15=>'a lavish',20=>'a magnificent',25=>'a historic',30=>'a magical'); $self->body->add_news($self->level*4,"In %s ceremony, %s unveiled its newly augmented %s.", $levels{$self->level}, $empire->name, $self->name); } - + my ($schedule) = Lacuna->db->resultset('Schedule')->search({ + parent_table => 'Building', + parent_id => $self->id, + task => 'finish_upgrade', + }); + $schedule->delete if defined $schedule; + return $self; } @@ -883,6 +897,15 @@ sub start_work { $self->work_started($now); $self->work_ends($now->clone->add(seconds=>$duration)); $self->work($work); + + # add to queue + my $schedule = Lacuna->db->resultset('Schedule')->create({ + delivery => $self->work_ends, + parent_table => 'Building', + parent_id => $self->id, + task => 'finish_work', + }); + return $self; } @@ -890,6 +913,13 @@ sub finish_work { my ($self) = @_; $self->is_working(0); $self->work({}); + + my ($schedule) = Lacuna->db->resultset('Schedule')->search({ + parent_table => 'Building', + parent_id => $self->id, + task => 'finish_work', + }); + $schedule->delete if defined $schedule; return $self; } diff --git a/lib/Lacuna/DB/Result/Map/Body/Planet.pm b/lib/Lacuna/DB/Result/Map/Body/Planet.pm index 2821fded..623f0674 100644 --- a/lib/Lacuna/DB/Result/Map/Body/Planet.pm +++ b/lib/Lacuna/DB/Result/Map/Body/Planet.pm @@ -790,43 +790,43 @@ sub has_room_in_build_queue { use constant operating_resource_names => qw(food_hour energy_hour ore_hour water_hour); has future_operating_resources => ( - is => 'rw', - clearer => 'clear_future_operating_resources', - lazy => 1, - default => sub { + is => 'rw', + clearer => 'clear_future_operating_resources', + lazy => 1, + default => sub { my $self = shift; -# get current + # get current my %future; foreach my $method ($self->operating_resource_names) { - $future{$method} = $self->$method; + $future{$method} = $self->$method; } -# adjust for what's already in build queue + # adjust for what's already in build queue my @queued_builds = @{$self->builds}; foreach my $build (@queued_builds) { - my $other = $build->stats_after_upgrade; - foreach my $method ($self->operating_resource_names) { - $future{$method} += $other->{$method} - $build->$method; - } + my $other = $build->stats_after_upgrade; + foreach my $method ($self->operating_resource_names) { + $future{$method} += $other->{$method} - $build->$method; + } } return \%future; - }, - ); + }, +); sub has_resources_to_operate { my ($self, $building) = @_; -# get future + # get future my $future = $self->future_operating_resources; -# get change for this building + # get change for this building my $after = $building->stats_after_upgrade; -# check our ability to sustain ourselves + # check our ability to sustain ourselves foreach my $method ($self->operating_resource_names) { my $delta = $after->{$method} - $building->$method; -# don't allow it if it sucks resources && its sucking more than we're producing + # don't allow it if it sucks resources && its sucking more than we're producing if ($delta < 0 && $future->{$method} + $delta < 0) { my $resource = $method; $resource =~ s/(\w+)_hour/$1/; @@ -839,12 +839,12 @@ sub has_resources_to_operate { sub has_resources_to_operate_after_building_demolished { my ($self, $building) = @_; -# get future + # get future my $planet = $self->future_operating_resources; -# check our ability to sustain ourselves + # check our ability to sustain ourselves foreach my $method ($self->operating_resource_names) { -# don't allow it if it sucks resources && its sucking more than we're producing + # don't allow it if it sucks resources && its sucking more than we're producing if ($planet->{$method} - $building->$method < 0) { my $resource = $method; $resource =~ s/(\w+)_hour/$1/; @@ -891,6 +891,9 @@ sub builds { sub get_existing_build_queue_time { my $self = shift; my ($building) = @{$self->builds(1)}; + +#print STDERR "GET_EXISTING_BUILD_QUEUE_TIME: building=[$building]\n"; + return (defined $building) ? $building->upgrade_ends : DateTime->now; } @@ -1319,26 +1322,26 @@ sub tick { my $i; # in case 2 things finish at exactly the same time # get building tasks - my @buildings = grep { - ($_->is_upgrading and $_->upgrade_ends->epoch <= $now_epoch) - or ($_->is_working and $_->work_ends->epoch <= $now_epoch) - } @{$self->building_cache}; - - foreach my $building (@buildings) { +# my @buildings = grep { +# ($_->is_upgrading and $_->upgrade_ends->epoch <= $now_epoch) +# or ($_->is_working and $_->work_ends->epoch <= $now_epoch) +# } @{$self->building_cache}; +# +# foreach my $building (@buildings) { # if ($building->is_upgrading && $building->upgrade_ends->epoch <= $now_epoch) { # $todo{format_date($building->upgrade_ends).$i} = { # object => $building, # type => 'building upgraded', # }; # } - if ($building->is_working && $building->work_ends->epoch <= $now_epoch) { - $todo{format_date($building->work_ends).$i} = { - object => $building, - type => 'building work complete', - }; - } - $i++; - } +# if ($building->is_working && $building->work_ends->epoch <= $now_epoch) { +# $todo{format_date($building->work_ends).$i} = { +# object => $building, +# type => 'building work complete', +# }; +# } +# $i++; +# } # get ship tasks my $ships = Lacuna->db->resultset('Lacuna::DB::Result::Ships')->search({ @@ -1373,10 +1376,10 @@ sub tick { $self->tick_to($object->date_available); $object->arrive; } - elsif ($job eq 'building work complete') { - $self->tick_to($object->work_ends); - $object->finish_work->update; - } +# elsif ($job eq 'building work complete') { +# $self->tick_to($object->work_ends); +# $object->finish_work->update; +# } # elsif ($job eq 'building upgraded') { # $self->tick_to($object->upgrade_ends); # $object->finish_upgrade; diff --git a/lib/Lacuna/DB/Result/Schedule.pm b/lib/Lacuna/DB/Result/Schedule.pm index 7a7fbabc..b2cb0110 100644 --- a/lib/Lacuna/DB/Result/Schedule.pm +++ b/lib/Lacuna/DB/Result/Schedule.pm @@ -32,6 +32,13 @@ after 'insert' => sub { return $self; }; +before 'delete' => sub { + my $self = shift; + + my $queue = Lacuna->queue; + # Delete the job off the queue + $queue->delete($self->job_id); +}; # Put this entry onto the beanstalk queue # @@ -44,7 +51,7 @@ sub queue_for_delivery { my $queue = Lacuna->queue || 'default'; my $priority = $self->priority || 1000; - +#print STDERR "DELAY: $delay table: ".$self->parent_table." task: ".$self->task."\n"; my $job = $queue->publish($self->queue, { id => $self->id, diff --git a/lib/Lacuna/Queue.pm b/lib/Lacuna/Queue.pm index 4e81549a..4bd71e2c 100644 --- a/lib/Lacuna/Queue.pm +++ b/lib/Lacuna/Queue.pm @@ -64,6 +64,7 @@ sub publish { my ($self, $queue, $payload, $options) = @_; my $beanstalk = $self->_beanstalk; + $queue = $queue || 'default'; $options = defined $options ? $options : {}, $beanstalk->use($queue); diff --git a/t/500_schedule.t b/t/500_schedule.t index de2e9feb..99744886 100644 --- a/t/500_schedule.t +++ b/t/500_schedule.t @@ -3,7 +3,7 @@ use lib '../lib'; use strict; use warnings; -use Test::More tests => 7; +use Test::More; use Test::Deep; use Test::Memory::Cycle; use Data::Dumper; @@ -13,55 +13,80 @@ use Lacuna; use TestHelper; -my $now = DateTime->now; -my $later = DateTime->now->add( seconds => 3); +my $tester = TestHelper->new->use_existing_test_empire; +my $session_id = $tester->session->id; +my $empire = $tester->empire; +my $home = $empire->home_planet; -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'}, +# Demolish all 'algae' test buildings +# +my $algae_rs = Lacuna->db->resultset('Building')->search({ + class => 'Lacuna::DB::Result::Building::Food::Algae', + body_id => $home->id, }); +while (my $algae_building = $algae_rs->next) { + diag("Demolishing algae ".$algae_building->id); + $algae_building->demolish; +}; -isa_ok($schedule, 'Lacuna::DB::Result::Schedule', 'Correct class'); -# Now test against beanstalk (it must be running) +# Test construction of a level 1 building # -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]"); +$tester->find_empty_plot; +my $building_1 = Lacuna->db->resultset('Lacuna::DB::Result::Building')->new({ + x => $tester->x, + y => $tester->y, + class => 'Lacuna::DB::Result::Building::Food::Algae', + level => 0, +}); +$home->build_building($building_1); +diag("Building ".$building_1->id." ends at ".$building_1->upgrade_ends); -# Delete this job, we no longer need it -$job->delete; +sleep 3; +$tester->find_empty_plot; +my $building_2 = Lacuna->db->resultset('Lacuna::DB::Result::Building')->new({ + x => $tester->x, + y => $tester->y, + class => 'Lacuna::DB::Result::Building::Food::Algae', + level => 0, +}); +$home->build_building($building_2); +diag("Building ".$building_2->id." ends at ".$building_2->upgrade_ends); + + +my $retry = 0; +my $building_1_complete; +my $building_2_complete; +TICK: +while () { + sleep 1; + $building_1->discard_changes; + $building_2->discard_changes; + diag("tick: ".++$retry); + + if (not $building_1_complete and not $building_1->is_upgrading) { + cmp_ok($retry, '>=', 12, "upgrade is not premature"); + cmp_ok($retry, '<', 15, "upgrade is not late"); + $building_1_complete = 1; + } + if (not $building_2_complete and not $building_2->is_upgrading) { + cmp_ok($retry, '>=', 27, "upgrade is not premature"); + cmp_ok($retry, '<', 30, "upgrade is not late"); + $building_2_complete = 1; + } + last TICK if $building_1_complete and $building_2_complete; + + if ($retry > 32) { + fail("Build did not terminate in time"); + last TICK; + } + +} + +sleep 1; + +done_testing; 1; diff --git a/t/510_bs_building.t b/t/510_bs_building.t new file mode 100644 index 00000000..13df2c11 --- /dev/null +++ b/t/510_bs_building.t @@ -0,0 +1,69 @@ +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 => 'default', + 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'); +exit; + + +# 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; + -- 2.51.2