Author: spadkins
Date: Fri Aug 3 12:34:40 2007
New Revision: 9819
Added:
p5ee/trunk/App-Context/lib/App/Context/POE/
p5ee/trunk/App-Context/lib/App/Context/POE.pm
p5ee/trunk/App-Context/lib/App/Context/POE/ClusterController.pm
p5ee/trunk/App-Context/lib/App/Context/POE/ClusterNode.pm
p5ee/trunk/App-Context/lib/App/Context/POE/Server.pm
Modified:
p5ee/trunk/App-Context/lib/App/Context/Server.pm
Log:
first version that shows promise for the POE-oriented Contexts
Added: p5ee/trunk/App-Context/lib/App/Context/POE.pm
==============================================================================
--- (empty file)
+++ p5ee/trunk/App-Context/lib/App/Context/POE.pm Fri Aug 3 12:34:40 2007
@@ -0,0 +1,17 @@
+#############################################################################
+## $Id: Conf.pm 3666 2006-03-11 20:34:10Z spadkins $
+#############################################################################
+
+package App::Context::POE;
+$VERSION = (q$Revision: 3666 $ =~ /(\d[\d\.]*)/)[0]; # VERSION numbers generated by svn
+
+use strict;
+use vars qw(@ISA);
+
+use App;
+use App::Context;
+
+@ISA = ( "App::Context" );
+
+1;
+
Added: p5ee/trunk/App-Context/lib/App/Context/POE/ClusterController.pm
==============================================================================
--- (empty file)
+++ p5ee/trunk/App-Context/lib/App/Context/POE/ClusterController.pm Fri Aug 3 12:34:40 2007
@@ -0,0 +1,470 @@
+
+#############################################################################
+## $Id: ClusterController.pm 6785 2006-08-11 23:13:19Z spadkins $
+#############################################################################
+
+package App::Context::POE::ClusterController;
+$VERSION = (q$Revision: 6785 $ =~ /(\d[\d\.]*)/)[0]; # VERSION numbers generated by svn
+
+use App;
+use App::Context::POE::Server;
+
+@ISA = ( "App::Context::POE::Server" );
+
+use Date::Format;
+use POE;
+
+use strict;
+
+=head1 NAME
+
+App::Context::POE::ClusterController - a runtime environment of a Cluster Controller served by many Cluster Nodes
+
+=head1 SYNOPSIS
+
+ # ... official way to get a Context object ...
+ use App;
+ $context = App->context();
+ $config = $context->config(); # get the configuration
+ $config->dispatch_events(); # dispatch events
+
+ # ... alternative way (used internally) ...
+ use App::Context::POE::ClusterController;
+ $context = App::Context::POE::ClusterController->new();
+
+=cut
+
+sub _init2a {
+ &App::sub_entry if ($App::trace);
+ my ($self, $options) = @_;
+ die "Controller must have a port defined (\$context->{options}{port})" if (!$self->{port});
+ $self->{is_controller} = 1;
+ $self->{num_async_events} = 0;
+ $self->{max_async_events_per_node} = $self->{options}{"app.context.max_async_events_per_node"} || 10;
+ $self->{max_async_events} = 0; # start with 0 because there are no nodes up
+
+ push(@{$self->{poe_states}},
+ "poe_remote_async_event_queued", "poe_set_node_status", "poe_run_event",
+ "poe_register_node", "poe_set_node_up", "poe_set_node_down");
+ push(@{$self->{poe_ikc_published_states}}, "poe_set_node_status");
+
+ $self->_init_poe($options);
+
+ &App::sub_exit() if ($App::trace);
+}
+
+sub _init2b {
+ &App::sub_entry if ($App::trace);
+ my ($self, $options) = @_;
+ $self->startup_nodes($options) if ($options->{startup});
+ &App::sub_exit() if ($App::trace);
+}
+
+sub dispatch_events_begin {
+ my ($self) = @_;
+ $self->log({level=>2},"Starting Cluster Controller on $self->{host}:$self->{port}\n");
+}
+
+sub dispatch_events_end {
+ my ($self) = @_;
+ $self->log({level=>2},"Stopping Cluster Controller\n");
+ # nothing special yet
+}
+
+sub send_async_event_now {
+ &App::sub_entry if ($App::trace);
+ my ($self, $event, $callback_event) = @_;
+
+ my $destination = $event->{destination};
+ if (! defined $destination) {
+ $self->log("ERROR: send_async_event_now() $event->{name}.$event->{method} : destination not assigned\n");
+ }
+ elsif ($event->{destination} eq "in_process") {
+ my $event_token = $self->send_async_event_in_process($event, $callback_event);
+ }
+ elsif ($destination =~ /^([^:]+):([0-9]+)$/) {
+ my $node_host = $1;
+ my $node_port = $2;
+ my $args = $event->{args};
+
+ my $remote_server_name = "poe_${node_host}_${node_port}";
+ my $remote_session_alias = $self->{poe_session_name}; # remote is same as local
+ my $remote_session_state = "poe_send_async_event";
+ my $local_callback_state = "poe_remote_async_event_queued";
+
+ $self->{num_async_events}++;
+ $self->{node}{$destination}{num_async_events}++;
+
+ my $kernel = $self->{poe_kernel};
+ $kernel->post("IKC", "call", "poe://$remote_server_name/$remote_session_alias/$remote_session_state",
+ [ $event, $callback_event ], "poe:$local_callback_state" );
+ }
+ else {
+ $self->SUPER::send_async_event_now($event, $callback_event);
+ }
+ &App::sub_exit() if ($App::trace);
+}
+
+sub ikc_register {
+ &App::sub_entry if ($App::trace);
+ my ($self, $kernel, $session_name) = @_[OBJECT, KERNEL, ARG1];
+ $self->log({level=>2},"POE: ikc_register ($session_name)\n");
+ if ($session_name =~ /^ikc_/) {
+ # do something
+ }
+ my ($retval);
+ &App::sub_exit($retval) if ($App::trace);
+ return($retval);
+}
+
+sub ikc_unregister {
+ &App::sub_entry if ($App::trace);
+ my ($self, $kernel, $remote_kernel_id) = @_[OBJECT, KERNEL, ARG0];
+ $self->log({level=>2},"POE: ikc_unregister ($remote_kernel_id)\n");
+ &App::sub_exit() if ($App::trace);
+}
+
+sub ikc_shutdown {
+ &App::sub_entry if ($App::trace);
+ my ($self, $kernel, $arg0, $arg1, $arg2, $arg3) = @_[OBJECT, KERNEL, ARG0, ARG1, ARG2, ARG3];
+ $self->log({level=>2},"POE: ikc_shutdown args=($arg0, $arg1, $arg2, $arg3)\n");
+ &App::sub_exit() if ($App::trace);
+ return;
+}
+
+sub poe_remote_async_event_queued {
+ &App::sub_entry if ($App::trace);
+ my ($self, $kernel, $runtime_event_token, $async_event) = @_[OBJECT, KERNEL, ARG0, ARG1];
+ $self->log({level=>2},"POE: poe_remote_async_event_queued ($async_event->[0]{name}.$async_event->[0]{method} => $runtime_event_token\n");
+ $self->{running_async_event}{$runtime_event_token} = $async_event;
+ &App::sub_exit() if ($App::trace);
+}
+
+# $runtime_event_tokens take the following forms:
+# $runtime_event_token = $pid; -- App::Context::Server::send_async_event_now() and ::finish_pid()
+# $runtime_event_token = "$host-$port-$serial"; -- i.e. a plain event token on the node
+sub _abort_running_async_event {
+ &App::sub_entry if ($App::trace);
+ my ($self, $runtime_event_token) = @_;
+ my $async_event = $self->{running_async_event}{$runtime_event_token};
+ my ($event, $callback_event) = @$async_event;
+ if ($runtime_event_token =~ /^[0-9]+$/) {
+ kill(9, $runtime_event_token);
+ }
+ elsif ($runtime_event_token =~ /^([^-]+)-([0-9]+)-/) {
+ my $node_host = $1;
+ my $node_port = $2;
+
+ my $remote_server_name = "poe_${node_host}_${node_port}";
+ my $remote_session_alias = $self->{poe_session_name}; # remote is same as local
+ my $remote_session_state = "poe_cancel_async_event";
+
+ my $kernel = $self->{poe_kernel};
+ $kernel->post("IKC", "post", "poe://$remote_server_name/$remote_session_alias/$remote_session_state",
+ [ $runtime_event_token ]);
+ }
+ else {
+ $self->log("ERROR: _abort_running_async_event() $event->{name}.$event->{method} : unparseable runtime event token [$runtime_event_token]\n");
+ }
+ &App::sub_exit() if ($App::trace);
+}
+
+sub assign_event_destination {
+ &App::sub_entry if ($App::trace);
+ my ($self, $event) = @_;
+ my $assigned = undef;
+ if ($self->{num_async_events} < $self->{max_async_events}) {
+ # SPA 2006-07-01: I just commented this out. I shouldn't need it.
+ # $event->{destination} = $self->{host};
+ my $main_service = $self->{main_service};
+ if ($main_service && $main_service->can("assign_event_destination")) {
+ $assigned = $main_service->assign_event_destination($event, $self->{nodes}, $self->{node});
+ }
+ else {
+ $assigned = $self->assign_event_destination_by_round_robin($event);
+ }
+ }
+ &App::sub_exit($assigned) if ($App::trace);
+ return($assigned);
+}
+
+sub assign_event_destination_by_round_robin {
+ &App::sub_entry if ($App::trace);
+ my ($self, $event) = @_;
+
+ my $assigned = undef;
+ my $nodes = $self->{nodes};
+ if ($#$nodes > -1) {
+ my $node_idx = $self->{node}{ALL}{last_node_idx};
+ $node_idx = (defined $node_idx) ? $node_idx + 1 : 0;
+ $node_idx = 0 if ($node_idx > $#$nodes);
+ $event->{destination} = $nodes->[$node_idx];
+ $self->{node}{ALL}{last_node_idx} = $node_idx;
+ $assigned = 1;
+ }
+
+ &App::sub_exit($assigned) if ($App::trace);
+ return($assigned);
+}
+
+sub poe_set_node_status {
+ &App::sub_entry if ($App::trace);
+ my ($self, $kernel, $args) = @_[OBJECT, KERNEL, ARG0];
+ $self->log("POE: poe_set_node_status args=(@$args)\n");
+ my ($retval);
+ &App::sub_exit($retval) if ($App::trace);
+ return($retval);
+}
+
+sub poe_run_event {
+ my ( $self, $kernel, $heap, $event ) = @_[ OBJECT, KERNEL, HEAP, ARG0 ];
+ &App::sub_entry if ($App::trace);
+ my ($event_str);
+ my $args = $event->{args} || [];
+ my $args_str = join(",", @$args);
+ if ($event->{name}) {
+ my $service_type = $event->{service_type} || "SessionObject";
+ $event_str = "$service_type($event->{name}).$event->{method}($args_str)";
+ }
+ else {
+ $event_str = "$event->{method}($args_str)";
+ }
+ $self->log({level=>2},"Run Event: $event_str\n");
+ $self->send_event($event);
+ &App::sub_exit() if ($App::trace);
+}
+
+sub state {
+ &App::sub_entry if ($App::trace);
+ my ($self) = @_;
+
+ my $datetime = time2str("%Y-%m-%d %H:%M:%S", time());
+ my $state = "Cluster Controller: $self->{host}:$self->{port} procs[$self->{num_procs}/$self->{max_procs}:max] async_events[$self->{num_async_events}/$self->{max_async_events}:max/$self->{max_async_events_per_node}:per]\n[$datetime]\n";
+ $state .= "\n";
+ $state .= $self->_state();
+
+ &App::sub_exit($state) if ($App::trace);
+ return($state);
+}
+
+sub _state {
+ &App::sub_entry if ($App::trace);
+ my ($self) = @_;
+
+ my $state = "";
+
+ my (@nodes);
+ @nodes = @{$self->{nodes}} if ($self->{nodes});
+ $state .= "Nodes: up [@nodes] last dispatched [$self->{node}{ALL}{last_node_idx}]\n";
+ my ($memfree, $memtotal, $swapfree, $swaptotal);
+ foreach my $node (sort keys %{$self->{node}}) {
+ next if ($node eq "ALL");
+ $state .= sprintf(" %-16s %4s : %3d/%3d max : [Load:%4.1f][Mem:%5.1f%%/%7d][Swap:%5.1f%%/%7d] : [%19s]\n", $node,
+ $self->{node}{$node}{up} ? "UP" : "down",
+ $self->{node}{$node}{num_async_events} || 0,
+ $self->{node}{$node}{max_async_events} || 0,
+ $self->{node}{$node}{load} || 0,
+ $self->{node}{$node}{memtotal} ? 100*($self->{node}{$node}{memtotal} - $self->{node}{$node}{memfree})/$self->{node}{$node}{memtotal} : 0,
+ $self->{node}{$node}{memtotal} || 0,
+ $self->{node}{$node}{swaptotal} ? 100*($self->{node}{$node}{swaptotal} - $self->{node}{$node}{swapfree})/$self->{node}{$node}{swaptotal} : 0,
+ $self->{node}{$node}{swaptotal} || 0,
+ $self->{node}{$node}{datetime});
+ }
+
+ $state .= $self->SUPER::_state();
+
+ &App::sub_exit($state) if ($App::trace);
+ return($state);
+}
+
+sub poe_register_node {
+ &App::sub_entry if ($App::trace);
+ my ($self, $kernel, $remote_kernel_id) = @_[OBJECT, KERNEL, ARG0];
+ $self->log({level=>2},"POE: poe_register_node ($remote_kernel_id)\n");
+
+ my ($node);
+ if ($remote_kernel_id =~ m!poe_([^_]+)_([0-9]+)!) {
+ $node = "$1:$2";
+ }
+ else {
+ $self->log("ERROR: poe_register_node: unparseable remote_kernel_id [$remote_kernel_id]\n");
+ }
+
+ if (!$self->{node}{$node}{up}) {
+ my $remote_server_name = "poe_${node}";
+ $remote_server_name =~ s/:/_/;
+ $kernel->post("IKC", "monitor", "poe://$remote_server_name",
+ {register => "poe_set_node_up",
+ unregister => "poe_set_node_down",
+ shutdown => "poe_set_node_down",
+ data => $node});
+ }
+
+ &App::sub_exit() if ($App::trace);
+}
+
+sub poe_set_node_up {
+ &App::sub_entry if ($App::trace);
+ my ($self, $kernel, $remote_kernel_id, $node) = @_[OBJECT, KERNEL, ARG0, ARG3];
+ $self->log({level=>2},"POE: poe_set_node_up ($remote_kernel_id; node=$node)\n");
+ my ($retval, $values);
+ if (!$self->{node}{$node}{up}) {
+ if ($node =~ /^([^:]+:\d+):(.*)/) {
+ $node = $1;
+ $values = $2;
+ if ($values) {
+ foreach my $value (split(/,/, $values)) {
+ if ($value =~ /^([^=]+)=(.*)/) {
+ $self->{node}{$node}{$1} = $2;
+ }
+ }
+ }
+ }
+ $self->{node}{$node}{datetime} = time2str("%Y-%m-%d %H:%M:%S", time());
+ if ($self->{node}{$node}{up}) {
+ $retval = "ok";
+ }
+ else {
+ $self->{node}{$node}{up} = 1;
+ $self->set_nodes();
+ $retval = "new";
+ }
+ }
+ &App::sub_exit($retval) if ($App::trace);
+ return($retval);
+}
+
+sub poe_set_node_down {
+ &App::sub_entry if ($App::trace);
+ my ($self, $kernel, $remote_kernel_id, $node) = @_[OBJECT, KERNEL, ARG0, ARG3];
+ $self->log({level=>2},"POE: poe_set_node_down ($remote_kernel_id; node=$node)\n");
+ my $runtime_event_token_prefix = $node;
+ $runtime_event_token_prefix =~ s/:/-/;
+ $self->reset_running_async_events($runtime_event_token_prefix);
+ $self->{node}{$node}{up} = 0;
+ $self->set_nodes();
+ &App::sub_exit() if ($App::trace);
+}
+
+sub set_nodes {
+ &App::sub_entry if ($App::trace);
+ my ($self) = @_;
+ my (@nodes);
+ foreach my $node (sort keys %{$self->{node}}) {
+ if ($self->{node}{$node}{up}) {
+ push(@nodes, $node);
+ }
+ }
+ $self->{nodes} = \@nodes;
+ $self->{max_async_events} = $self->{max_async_events_per_node} * ($#nodes + 1);
+ my $main_service = $self->{main_service};
+ if ($main_service && $main_service->can("capacity_change")) {
+ $main_service->capacity_change($self->{max_async_events}, \@nodes, $self->{node});
+ }
+ &App::sub_exit() if ($App::trace);
+}
+
+sub shutdown {
+ &App::sub_entry if ($App::trace);
+ my $self = shift;
+ $self->shutdown_nodes();
+ $self->write_node_file();
+ $self->SUPER::shutdown();
+ &App::sub_exit() if ($App::trace);
+}
+
+sub shutdown_nodes {
+ &App::sub_entry if ($App::trace);
+ my $self = shift;
+ foreach my $node (@{$self->{nodes}}) {
+ if ($node =~ /^([^:]+):([0-9]+)$/) {
+ my $remote_server_name = "poe_${1}_${2}";
+ my $remote_session_alias = $self->{poe_session_name}; # remote is same as local
+ my $remote_session_state = "poe_shutdown_node";
+ my $kernel = $self->{poe_kernel};
+ $kernel->post("IKC", "post", "poe://$remote_server_name/$remote_session_alias/$remote_session_state");
+ }
+ else {
+ $self->log("ERROR: shutdown_nodes: unparseable node [$node]\n");
+ }
+ }
+ &App::sub_exit() if ($App::trace);
+}
+
+sub startup_nodes {
+ &App::sub_entry if ($App::trace);
+ my ($self, $options) = @_;
+
+ my $startup = $options->{startup};
+
+ my ($node, $msg, $host, $port, $cmd);
+ if ($startup eq "1") {
+ $self->read_node_file();
+ }
+ else {
+ foreach $node (split(/,/,$startup)) {
+ $self->{node}{$node} = {};
+ }
+ }
+
+ my $cmd_fmt = $self->{options}{"app.context.node_start_cmd"} || "ssh -f {host} mvnode --port={port}";
+ foreach $node (keys %{$self->{node}}) {
+ if ($node =~ /^([^:]+):([0-9]+)$/) {
+ $host = $1;
+ $port = $2;
+ $cmd = $cmd_fmt;
+ $cmd =~ s/{host}/$host/g;
+ $cmd =~ s/{port}/$port/g;
+ $self->log("Starting Node [$node]: [$cmd]\n");
+ system("$cmd < /dev/null &");
+ }
+ else {
+ $self->log("ERROR: startup_nodes: unparseable node [$node]\n");
+ }
+ }
+ &App::sub_exit() if ($App::trace);
+}
+
+sub write_node_file {
+ &App::sub_entry if ($App::trace);
+ my $self = shift;
+ my $prefix = $self->{options}{prefix};
+ my $node_file = "$prefix/log/$self->{options}{app}-$self->{host}:$self->{port}.nodes";
+ if (open(FILE, "> $node_file")) {
+ foreach my $node (@{$self->{nodes}}) {
+ print App::Context::POE::ClusterController::FILE "$node\n";
+ }
+ close(App::Context::POE::ClusterController::FILE);
+ }
+ else {
+ $self->log("ERROR: Can't write node file [$node_file]: $!\n");
+ }
+ &App::sub_exit() if ($App::trace);
+}
+
+sub read_node_file {
+ &App::sub_entry if ($App::trace);
+ my $self = shift;
+ my $prefix = $self->{options}{prefix};
+ my $node_file = "$prefix/log/$self->{options}{app}-$self->{host}:$self->{port}.nodes";
+ my ($node);
+ if (open(FILE, "< $node_file")) {
+ while (<App::Context::POE::ClusterController::FILE>) {
+ chomp;
+ if (/^[^:]+:[0-9]+$/) {
+ $node = $_;
+ # just take note of its existence. we don't know yet if it is up.
+ $self->{node}{$node} = {} if (!defined $self->{node}{$node});
+ }
+ }
+ close(App::Context::POE::ClusterController::FILE);
+ }
+ else {
+ # This is not really a problem.
+ # $self->log("WARNING: Can't read node file [$node_file]: $!\n");
+ }
+ &App::sub_exit() if ($App::trace);
+}
+
+1;
+
Added: p5ee/trunk/App-Context/lib/App/Context/POE/ClusterNode.pm
==============================================================================
--- (empty file)
+++ p5ee/trunk/App-Context/lib/App/Context/POE/ClusterNode.pm Fri Aug 3 12:34:40 2007
@@ -0,0 +1,225 @@
+
+#############################################################################
+## $Id: ClusterNode.pm 3666 2006-03-11 20:34:10Z spadkins $
+#############################################################################
+
+package App::Context::POE::ClusterNode;
+$VERSION = (q$Revision: 3666 $ =~ /(\d[\d\.]*)/)[0]; # VERSION numbers generated by svn
+
+use App;
+use App::Context::POE::Server;
+
+@ISA = ( "App::Context::POE::Server" );
+
+use strict;
+
+use Date::Format;
+
+use POE;
+use POE::Component::IKC::Client;
+use POE::Component::IKC::Responder;
+use POE::Component::Server::SimpleHTTP;
+use HTTP::Status qw/RC_OK/;
+use Socket qw(INADDR_ANY);
+
+=head1 NAME
+
+App::Context::ClusterNode - a runtime environment for a Cluster Node that serves a Cluster Controller
+
+=head1 SYNOPSIS
+
+ # ... official way to get a Context object ...
+ use App;
+ $context = App->context();
+ $config = $context->config(); # get the configuration
+ $config->dispatch_events(); # dispatch events
+
+ # ... alternative way (used internally) ...
+ use App::Context::ClusterNode;
+ $context = App::Context::ClusterNode->new();
+
+=cut
+
+sub _init2a {
+ &App::sub_entry if ($App::trace);
+ my ($self, $options) = @_;
+ $self->{controller_host} = $options->{controller_host};
+ $self->{controller_port} = $options->{controller_port};
+ $self->{disable_event_loop_extensions} = 1;
+ die "Node must have a controller host and port defined (\$context->{options}{controller_host} and {controller_port})"
+ if (!$self->{controller_host} || !$self->{controller_port});
+
+ #push(@{$self->{poe_states}}, "foo", "bar");
+ #push(@{$self->{poe_ikc_published_states}}, "more", "states");
+
+ $self->_init_poe($options);
+
+ &App::sub_exit() if ($App::trace);
+}
+
+sub _init_poe {
+ &App::sub_entry if ($App::trace);
+ my ($self, $options) = @_;
+
+ my $ikc_name = "ikc_$self->{host}_$self->{port}";
+ ### Set up a server
+ POE::Component::IKC::Responder->spawn();
+ POE::Component::IKC::Client->spawn(
+ ip => $self->{controller_host},
+ port => $self->{controller_port},
+ name => $ikc_name,
+ timeout => 60,
+ );
+ $self->log({level=>2},"Listening for Inter-Kernel Communications on $self->{host}:$self->{port}\n");
+
+ my $session_name = $self->{poe_session_name};
+ POE::Component::Server::SimpleHTTP->new(
+ 'ALIAS' => $self->{poe_kernel_httpd_name},
+ 'ADDRESS' => INADDR_ANY,
+ 'PORT' => $self->{options}{http_port},
+ 'HANDLERS' => [
+ { 'DIR' => '/testrun', 'SESSION' => $session_name, 'EVENT' => 'poe_http_test_run', },
+ { 'DIR' => '.*', 'SESSION' => $session_name, 'EVENT' => 'poe_http_server_state', },
+ ],
+ );
+ $self->log({level=>2},"Listening for HTTP Requests on $self->{host}:$self->{options}{http_port}\n");
+
+ &App::sub_exit() if ($App::trace);
+}
+
+sub _init2b {
+ &App::sub_entry if ($App::trace);
+ my ($self, $options) = @_;
+
+ &App::sub_exit() if ($App::trace);
+}
+
+sub _start {
+ &App::sub_entry if ($App::trace);
+ my ( $self, $kernel, $heap ) = @_[ OBJECT, KERNEL, HEAP ];
+ $self->log({level=>2},"POE: _start\n");
+
+ my $name = $self->{poe_session_name};
+ $kernel->alias_set($name);
+
+ $kernel->sig(CHLD => "poe_sigchld");
+ $kernel->sig(HUP => "poe_sigignore");
+ $kernel->sig(INT => "poe_sigterm");
+ $kernel->sig(QUIT => "poe_sigterm");
+ $kernel->sig(USR1 => "poe_sigignore");
+ $kernel->sig(USR2 => "poe_sigignore");
+ $kernel->sig(TERM => "poe_sigterm");
+
+ $kernel->call( IKC => publish => $name, $self->{poe_ikc_published_states} );
+
+ my $remote_server_name = "poe_$self->{controller_host}_$self->{controller_port}";
+ my $node = "$self->{host}:$self->{port}";
+
+ $kernel->post("IKC", "monitor", "poe://$remote_server_name",
+ {register => "ikc_register",
+ unregister => "ikc_unregister",
+ shutdown => "ikc_shutdown",
+ data => $node});
+
+ # don't start kicking off async events until we give the nodes a chance to register themselves
+ $kernel->delay_set("poe_event_loop_extension", 5) if (!$self->{disable_event_loop_extensions});
+ $kernel->delay_set("poe_alarm", 5);
+
+ &App::sub_exit() if ($App::trace);
+}
+
+sub ikc_register {
+ &App::sub_entry if ($App::trace);
+ my ($self, $kernel, $remote_kernel_id, $node) = @_[OBJECT, KERNEL, ARG0, ARG3];
+ $self->log({level=>2},"POE: ikc_register ($remote_kernel_id; node=$node)\n");
+ $self->{controller_up} = 1;
+ my ($retval);
+ &App::sub_exit($retval) if ($App::trace);
+ return($retval);
+}
+
+sub ikc_unregister {
+ &App::sub_entry if ($App::trace);
+ my ($self, $kernel, $remote_kernel_id, $node) = @_[OBJECT, KERNEL, ARG0, ARG3];
+ $self->log({level=>2},"POE: ikc_unregister ($remote_kernel_id; node=$node)\n");
+ $self->{controller_up} = 0;
+ &App::sub_exit() if ($App::trace);
+}
+
+sub ikc_shutdown {
+ &App::sub_entry if ($App::trace);
+ my ( $self, $kernel, $session, $heap ) = @_[ OBJECT, KERNEL, SESSION, HEAP ];
+ $self->log({level=>2},"POE: ikc_shutdown\n");
+ &App::sub_exit() if ($App::trace);
+ return;
+}
+
+sub dispatch_events_begin {
+ &App::sub_entry if ($App::trace);
+ my ($self) = @_;
+ $self->log({level=>2},"Starting Cluster Node on $self->{host}:$self->{port}\n");
+ my $node_heartbeat = $self->{options}{node_heartbeat} || 60;
+ $self->schedule_event(
+ method => "send_node_status",
+ time => time(), # immediately ...
+ interval => $node_heartbeat, # and every X seconds hereafter
+ );
+ &App::sub_exit() if ($App::trace);
+}
+
+sub dispatch_events_end {
+ &App::sub_entry if ($App::trace);
+ my ($self) = @_;
+ $self->log({level=>2},"Stopping Cluster Node\n");
+ &App::sub_exit() if ($App::trace);
+}
+
+sub send_node_status {
+ &App::sub_entry if ($App::trace);
+ my ($self) = @_;
+ my $controller_host = $self->{controller_host};
+ my $controller_port = $self->{controller_port};
+ my $node_host = $self->{host};
+ my $node_port = $self->{port};
+
+ my $remote_server_name = "poe_${controller_host}_${controller_port}";
+ my $remote_session_alias = $self->{poe_session_name}; # remote is same as local
+ my $remote_session_state = "poe_set_node_status";
+ my $sys_info = $self->get_sys_info();
+
+ if ($self->{controller_up}) {
+ my $kernel = $self->{poe_kernel};
+ $kernel->post("IKC", "post", "poe://$remote_server_name/$remote_session_alias/$remote_session_state",
+ [ $sys_info ]);
+ }
+
+ &App::sub_exit() if ($App::trace);
+}
+
+sub state {
+ &App::sub_entry if ($App::trace);
+ my ($self) = @_;
+
+ my $datetime = time2str("%Y-%m-%d %H:%M:%S", time());
+ my $state = "Cluster Node: $self->{host}:$self->{port} procs[$self->{num_procs}/$self->{max_procs}:max] async_events[$self->{num_async_events}/$self->{max_async_events}:max]\n[$datetime]\n";
+ $state .= "\n";
+ $state .= $self->_state();
+
+ &App::sub_exit($state) if ($App::trace);
+ return($state);
+}
+
+sub _state {
+ &App::sub_entry if ($App::trace);
+ my ($self) = @_;
+
+ my $state = "";
+
+ $state .= $self->SUPER::_state();
+
+ &App::sub_exit($state) if ($App::trace);
+ return($state);
+}
+
+1;
+
Added: p5ee/trunk/App-Context/lib/App/Context/POE/Server.pm
==============================================================================
--- (empty file)
+++ p5ee/trunk/App-Context/lib/App/Context/POE/Server.pm Fri Aug 3 12:34:40 2007
@@ -0,0 +1,941 @@
+#############################################################################
+## $Id: Server.pm 6786 2006-08-11 23:22:48Z zroberts $
+#############################################################################
+
+package App::Context::POE::Server;
+$VERSION = (q$Revision: 6786 $ =~ /(\d[\d\.]*)/)[0]; # VERSION numbers generated by svn
+
+use strict;
+use vars qw(@ISA);
+use warnings;
+
+use App::Context::POE;
+
+@ISA = ( "App::Context::POE" );
+
+use POSIX ":sys_wait_h";
+use Sys::Hostname;
+use Date::Format;
+use Date::Parse;
+
+use POE;
+use POE::Component::Server::SimpleHTTP;
+use POE::Component::IKC::Server;
+use HTTP::Status qw/RC_OK/;
+use Socket qw(INADDR_ANY);
+
+sub _init {
+ &App::sub_entry if ($App::trace);
+ my ($self, $options) = @_;
+ $options = {} if (!defined $options);
+
+ $self->SUPER::_init($options);
+
+ App->mkdir($options->{prefix}, "data", "app", "Context");
+
+ $| = 1; # autoflush STDOUT (not sure this is required)
+
+ ### Configuration stuff
+ my $host = hostname;
+ $self->{hostname} = $host;
+ $host =~ s/\..*//; # get rid of fully qualified domain name
+ $self->{host} = $host;
+ $options->{port} ||= 8080;
+ $self->{port} = $options->{port};
+ $options->{http_port} ||= $options->{port}+1;
+ $self->{poe_kernel_name} = "poe_$self->{host}_$self->{port}";
+ $self->{poe_kernel_httpd_name} = $self->{poe_kernel_name} . "_httpd";
+ $self->{poe_session_name} = "poe_session";
+ $self->{poe_kernel} = $poe_kernel;
+
+ $self->{num_procs} = 0;
+ $self->{max_procs} = $self->{options}{"app.context.max_procs"} || 10;
+ $self->{max_async_events} = $self->{options}{"app.context.max_async_events"}
+ if (defined $self->{options}{"app.context.max_async_events"});
+ $self->{max_async_events} ||= 10;
+ $self->{num_async_events} = 0;
+ $self->{async_event_count} = 0;
+ $self->{pending_async_events} = [];
+ $self->{running_async_event} = {};
+
+ $self->{verbose} = $options->{verbose};
+
+ $self->{poe_states} = [qw(
+ _start _stop _default poe_sigchld poe_sigterm poe_sigignore poe_shutdown poe_alarm
+ ikc_register ikc_unregister ikc_shutdown
+ poe_run_event poe_event_loop_extension poe_dispatch_pending_async_events
+ poe_server_state poe_http_server_state poe_http_test_run
+ )];
+ $self->{poe_ikc_published_states} = ["poe_server_state"];
+
+ ### Does nothing by default, used by ClusterController, maybe other subclasses?
+ $self->_init2a($options);
+
+ ### Do log rotation
+ ### TODO: this should be refactored out
+ if ($self->{options}{log_rotate}) {
+ my $rotate_sec = $self->{options}{log_rotate};
+ $rotate_sec = $rotate_sec*(24*3600) if ($rotate_sec <= 31); # interpret as days
+ my $time = time();
+ my $base_time = str2time(time2str("%Y-%m-%d 00:00:00", $time)); # I need a base which is midnight local time
+ my $next_rotate_time = ((int(($time - $base_time)/$rotate_sec)+1)*$rotate_sec) + $base_time;
+ $self->schedule_event(
+ tag => "context-log-rotation",
+ method => "log_file_open",
+ args => [0], # don't overwrite
+ time => $next_rotate_time,
+ interval => $rotate_sec, # and every X seconds hereafter
+ );
+ }
+
+ ### Does nothing by default
+ $self->_init2b($options);
+
+ &App::sub_exit() if ($App::trace);
+}
+
+### Used by subclasses
+sub _init2a {
+ &App::sub_entry if ($App::trace);
+ my ($self, $options) = @_;
+ $self->_init_poe($options);
+ &App::sub_exit() if ($App::trace);
+}
+
+sub _init_poe {
+ &App::sub_entry if ($App::trace);
+ my ($self, $options) = @_;
+
+ ### Set up a server
+ POE::Component::IKC::Server->spawn(
+ port => $self->{port},
+ name => $self->{poe_kernel_name},
+ );
+ $self->log({level=>2},"Listening for Inter-Kernel Communications on $self->{host}:$self->{port}\n");
+ POE::Component::IKC::Responder->spawn();
+
+ my $session_name = $self->{poe_session_name};
+ POE::Component::Server::SimpleHTTP->new(
+ 'ALIAS' => $self->{poe_kernel_httpd_name},
+ 'ADDRESS' => INADDR_ANY,
+ 'PORT' => $self->{options}{http_port},
+ 'HANDLERS' => [
+ { 'DIR' => '/testrun', 'SESSION' => $session_name, 'EVENT' => 'poe_http_test_run', },
+ { 'DIR' => '.*', 'SESSION' => $session_name, 'EVENT' => 'poe_http_server_state', },
+ ],
+ );
+ $self->log({level=>2},"Listening for HTTP Requests on $self->{host}:$self->{options}{http_port}\n");
+
+ &App::sub_exit() if ($App::trace);
+}
+
+### Used by subclasses
+sub _init2b {
+ &App::sub_entry if ($App::trace);
+ my ($self, $options) = @_;
+ &App::sub_exit() if ($App::trace);
+}
+
+sub shutdown {
+ &App::sub_entry if ($App::trace);
+ my $self = shift;
+
+ ### Shut down servers
+ ### TODO
+
+ ### Shut down children
+ $self->shutdown_child_processes();
+
+ ### Call SUPER shutdown
+ $self->SUPER::shutdown();
+
+ &App::sub_exit() if ($App::trace);
+}
+
+sub shutdown_child_processes {
+ &App::sub_entry if ($App::trace);
+ my $self = shift;
+ if ($self->{proc}) {
+ foreach my $pid (keys %{$self->{proc}}) {
+ kill(15, $pid);
+ }
+ }
+ &App::sub_exit() if ($App::trace);
+}
+
+sub dispatch_events {
+ &App::sub_entry if ($App::trace);
+ my ($self, $max_events_occurred) = @_;
+
+ my $verbose = $self->{verbose};
+ my $options = $self->{options};
+
+ my ($role, $port, $startup, $shutdown);
+ $self->dispatch_events_begin();
+
+ ### Set up init_objects, untouched and snagged from App::Context::POE::Server
+ my $objects = $options->{init_objects};
+ my ($service_type, $name, $service);
+ foreach my $object (split(/ *[;,]+ */, $objects)) {
+ if ($object) {
+ if ($object =~ /^([A-Z][A-Za-z0-9]+)\.([A-Za-z0-9_-]+)$/) {
+ $service_type = $1;
+ $name = $2;
+ }
+ else {
+ $service_type = "SessionObject";
+ $name = $object;
+ }
+ $service = $self->service($service_type, $name); # instantiate it. that's all.
+ $self->log({level=>3},"$service_type $name instantiated [$service]\n");
+ $self->{main_service} = $service if (!$self->{main_service});
+ }
+ }
+
+ eval {
+ ### POE Server begins here
+ POE::Session->create( object_states => [ $self => $self->{poe_states} ] );
+ $poe_kernel->run();
+ };
+ if ($@) {
+ $self->log($@);
+ }
+
+ $self->dispatch_events_end();
+ $self->shutdown();
+ &App::sub_exit() if ($App::trace);
+}
+
+sub dispatch_events_begin {
+ &App::sub_entry if ($App::trace);
+ my ($self) = @_;
+ my $verbose = $self->{verbose};
+ $self->log({level=>2},"Starting Dispatching Events on Server on $self->{host}:$self->{port}\n");
+ &App::sub_exit() if ($App::trace);
+}
+
+sub dispatch_events_end {
+ my ($self) = @_;
+ my $verbose = $self->{verbose};
+ $self->log({level=>2},"Stopping Dispatching Events on Server.\n");
+}
+
+sub process_msg {
+ my ($self, $msg) = @_;
+ $self->log({level=>3},"process_msg: [$msg]\n");
+ my $verbose = $self->{verbose};
+ my $return_value = $self->process_custom_msg($msg);
+ if (!$return_value) {
+ $return_value = "unknown [$msg]\n";
+ }
+ &App::sub_exit($return_value) if ($App::trace);
+ return($return_value);
+}
+
+sub process_custom_msg {
+ &App::sub_entry if ($App::trace);
+ my ($self, $msg) = @_;
+ my $return_value = "";
+ &App::sub_exit($return_value) if ($App::trace);
+ return($return_value);
+}
+
+sub state {
+ &App::sub_entry if ($App::trace);
+ my ($self) = @_;
+
+ my $datetime = time2str("%Y-%m-%d %H:%M:%S", time());
+ my $state = "Server: $self->{host}:$self->{port} procs[$self->{num_procs}/$self->{max_procs}:max] async_events[$self->{num_async_events}/$self->{max_async_events}:max]\n[$datetime]\n";
+ $state .= "\n";
+ $state .= $self->_state();
+
+ &App::sub_exit($state) if ($App::trace);
+ return($state);
+}
+
+sub _state {
+ &App::sub_entry if ($App::trace);
+ my ($self) = @_;
+
+ my $state = "";
+
+ my $options = $self->{options};
+ my $objects = $options->{init_objects};
+ my ($service_type, $name, $service);
+ foreach my $object (split(/ *[;,]+ */, $objects)) {
+ if ($object) {
+ if ($object =~ /^([A-Z][A-Za-z0-9]+)\.([A-Za-z0-9_-]+)$/) {
+ $service_type = $1;
+ $name = $2;
+ }
+ else {
+ $service_type = "SessionObject";
+ $name = $object;
+ }
+ $service = $self->service($service_type, $name); # instantiate it. that's all.
+ if ($service->can("state")) {
+ $state .= "\n";
+ $state .= $service->state();
+ }
+ }
+ }
+
+ my $main_service = $self->{main_service};
+
+ $state .= "\n";
+ $state .= "Running Async Events:\n";
+ my ($async_event, $event, $callback_event, @args, $args_str, $event_token, $runtime_event_token, $str);
+ foreach $runtime_event_token (sort keys %{$self->{running_async_event}}) {
+ $async_event = $self->{running_async_event}{$runtime_event_token};
+ ($event, $callback_event) = @$async_event;
+ $str = "";
+ if ($main_service && $main_service->can("format_async_event")) {
+ $str = $main_service->format_async_event($event, $callback_event, $runtime_event_token);
+ }
+ if ($str) {
+ $state .= " ";
+ $state .= $main_service->format_async_event($event, $callback_event, $runtime_event_token);
+ $state .= "\n";
+ }
+ else {
+ @args = ();
+ @args = @{$event->{args}} if ($event->{args});
+ $args_str = join(",",@args);
+ $state .= sprintf(" %-20s %-20s %-24s", $event->{event_token}, $runtime_event_token, "$event->{name}.$event->{method}($args_str)");
+ if ($callback_event) {
+ @args = ();
+ @args = @{$callback_event->{args}} if ($callback_event->{args});
+ $args_str = join(",",@args);
+ $state .= "$callback_event->{name}.$callback_event->{method}($args_str)";
+ }
+ $state .= "\n";
+ }
+ }
+
+ $state .= "\n";
+ $state .= "Pending Async Events: count [$self->{async_event_count}]\n";
+ foreach $async_event (@{$self->{pending_async_events}}) {
+ ($event, $callback_event) = @$async_event;
+ $str = "";
+ if ($main_service && $main_service->can("format_async_event")) {
+ $str = $main_service->format_async_event($event, $callback_event);
+ }
+ if ($str) {
+ $state .= " ";
+ $state .= $main_service->format_async_event($event, $callback_event);
+ $state .= "\n";
+ }
+ else {
+ @args = ();
+ @args = @{$event->{args}} if ($event->{args});
+ $args_str = join(",",@args);
+ $state .= sprintf(" %-20s %-40s", $event->{event_token}, "$event->{name}.$event->{method}($args_str)");
+ if ($callback_event) {
+ @args = ();
+ @args = @{$callback_event->{args}} if ($callback_event->{args});
+ $args_str = join(",",@args);
+ $state .= " => $callback_event->{name}.$callback_event->{method}($args_str)";
+ }
+ $state .= "\n";
+ }
+ }
+
+ $state .= "\n";
+
+ $state .= $self->SUPER::_state();
+
+ &App::sub_exit($state) if ($App::trace);
+ return($state);
+}
+
+# TODO: Implement this as a fork() or a context-level message to a node to fork().
+# i.e. messages such as "EVENT:" and "EVENT-OK:"
+# Save the callback_event according to an event_token.
+# Then implement cleanup_pid to send the callback_event.
+
+sub send_async_event {
+ &App::sub_entry if ($App::trace);
+ my ($self, $event, $callback_event) = @_;
+ my $event_token = $self->new_event_token();
+ $event->{event_token} = $event_token;
+ $callback_event->{event_token} = $event_token if ($callback_event);
+ push(@{$self->{pending_async_events}}, [ $event, $callback_event ]);
+ &App::sub_exit($event_token) if ($App::trace);
+ return($event_token);
+}
+
+sub new_event_token {
+ &App::sub_entry if ($App::trace);
+ my ($self) = @_;
+ $self->{async_event_count} ++;
+ my $event_token = "$self->{host}-$self->{port}-$self->{async_event_count}";
+ &App::sub_exit($event_token) if ($App::trace);
+ return($event_token);
+}
+
+sub dispatch_pending_async_events {
+ &App::sub_entry if ($App::trace);
+ my ($self, $max_events) = @_;
+ $max_events ||= 9999;
+ my $pending_async_events = $self->{pending_async_events};
+ my ($async_event, $assigned, $event, $in_process);
+ my $events_occurred = 0;
+ my $i = 0;
+ my $event_capacity_exists = 1;
+ my $max_i = $#$pending_async_events;
+ while ($i <= $max_i && $events_occurred < $max_events) {
+ $async_event = $pending_async_events->[$i];
+ $event = $async_event->[0];
+ if ($event->{destination}) {
+ $self->send_async_event_now(@$async_event);
+ $events_occurred ++;
+ splice(@$pending_async_events, $i, 1); # remove $pending_async_events->[$i]
+ $max_i--;
+ }
+ elsif ($event_capacity_exists) {
+ $assigned = $self->assign_event_destination($event);
+ if ($assigned) {
+ $self->send_async_event_now(@$async_event);
+ $events_occurred ++;
+ # keep $i the same
+ splice(@$pending_async_events, $i, 1); # remove $pending_async_events->[$i]
+ $max_i--;
+ }
+ else { # [undef] no servers are eligible for assignment
+ $event_capacity_exists = 0; # there's no sense looking at the other pending async events
+ $i++; # look at the next one
+ }
+ }
+ else { # [0] this async_event is not eligible to run
+ $i++; # look at the next one
+ }
+ }
+ &App::sub_exit($events_occurred) if ($App::trace);
+ return($events_occurred);
+}
+
+sub assign_event_destination {
+ &App::sub_entry if ($App::trace);
+ my ($self, $event) = @_;
+ my $assigned = undef;
+ if ($self->{num_procs} < $self->{max_procs} &&
+ (!defined $self->{max_async_events} || $self->{num_async_events} < $self->{max_async_events})) {
+ $event->{destination} = $self->{host};
+ $assigned = 1;
+ }
+ &App::sub_exit($assigned) if ($App::trace);
+ return($assigned);
+}
+
+sub send_async_event_now {
+ &App::sub_entry if ($App::trace);
+ my ($self, $event, $callback_event) = @_;
+ if ($event->{destination} eq "in_process") {
+ my $event_token = $self->send_async_event_in_process($event, $callback_event);
+ }
+ else {
+ my $pid = $self->fork();
+ if (!$pid) { # running in child
+ my $exitval = 0;
+ my (@results);
+ eval {
+ @results = $self->send_event($event);
+ };
+ if ($@) {
+ @results = ($@);
+ }
+ if ($#results > -1 && defined $results[0] && $results[0] ne "") {
+ my $textfile = $self->{options}{prefix} . "/data/app/Context/$$";
+ if (open(FILE, "> $textfile")) {
+ print App::Context::POE::Server::FILE @results;
+ close(App::Context::POE::Server::FILE);
+ }
+ else {
+ $exitval = 1;
+ }
+ }
+ $self->shutdown();
+ $self->exit($exitval);
+ }
+ my $destination = $event->{destination} || "local";
+ $self->{num_async_events}++;
+ $self->{node}{$destination}{num_async_events}++;
+ my $runtime_event_token = $pid;
+ $self->{running_async_event}{$runtime_event_token} = [ $event, $callback_event ];
+ }
+ &App::sub_exit() if ($App::trace);
+}
+=head2 wait_for_event()
+
+ * Signature: $self->wait_for_event($event_token)
+ * Param: $event_token string
+ * Return: void
+ * Throws: App::Exception
+ * Since: 0.01
+
+ Sample Usage:
+
+ $self->wait_for_event($event_token);
+
+The wait_for_event() method is called when an asynchronous event has been
+sent and no more processing can be completed before it is done.
+
+=cut
+
+sub wait_for_event {
+ &App::sub_entry if ($App::trace);
+ my ($self, $event_token) = @_;
+ &App::sub_exit() if ($App::trace);
+}
+
+sub fork {
+ &App::sub_entry if ($App::trace);
+ my ($self) = @_;
+ my $pid = $self->SUPER::fork();
+ if ($pid) { # the parent process has a new child process
+ $self->{num_procs}++;
+ $self->{proc}{$pid} = {};
+ }
+ else { # the new child process has no sub-processes
+ $self->{num_procs} = 0;
+ $self->{proc} = {};
+ $SIG{INT} = sub { $self->log({level=>2},"Caught Signal: @_ (quitting)\n"); $self->exit(102); }; # SIG 2
+ $SIG{QUIT} = sub { $self->log({level=>2},"Caught Signal: @_ (quitting)\n"); $self->exit(103); }; # SIG 3
+ $SIG{TERM} = sub { $self->log({level=>2},"Caught Signal: @_ (quitting)\n"); $self->exit(115); }; # SIG 15
+ }
+ &App::sub_exit($pid) if ($App::trace);
+ return($pid);
+}
+
+sub finish_pid {
+ &App::sub_entry if ($App::trace);
+ my ($self, $pid, $exitval, $sig) = @_;
+
+ $self->{num_procs}--;
+ delete $self->{proc}{$pid};
+
+ my $runtime_event_token = $pid;
+ my $async_event = $self->{running_async_event}{$runtime_event_token};
+ if ($async_event) {
+ my ($event, $callback_event) = @$async_event;
+ my $returnval = "";
+ my $returnvalfile = $self->{options}{prefix} . "/data/app/Context/$pid";
+ if (open(FILE, $returnvalfile)) {
+ if ($callback_event) {
+ $returnval = join("",<App::Context::POE::Server::FILE>);
+ }
+ close(App::Context::POE::Server::FILE);
+ unlink($returnvalfile);
+ }
+
+ my $destination = $event->{destination} || "local";
+ $self->{num_async_events}--;
+ $self->{node}{$destination}{num_async_events}--;
+ delete $self->{running_async_event}{$runtime_event_token};
+
+ if ($callback_event) {
+ $callback_event->{args} = [] if (! $callback_event->{args});
+ my $errmsg = ($exitval || $sig) ? "Exit $exitval on $pid [sig=$sig]" : "";
+ push(@{$callback_event->{args}},
+ {event_token => $callback_event->{event_token}, returnval => $returnval, errnum => $exitval, errmsg => $errmsg});
+ $self->send_event($callback_event);
+ }
+ elsif ($sig == 9) { # killed without a chance to finish its work
+ $self->finish_killed_async_event($event);
+ }
+ }
+
+ &App::sub_exit() if ($App::trace);
+}
+
+sub finish_killed_async_event {
+ &App::sub_entry if ($App::trace);
+ my ($self, $event) = @_;
+ &App::sub_exit() if ($App::trace);
+}
+
+sub find_runtime_event_token {
+ &App::sub_entry if ($App::trace);
+ my ($self, $event_token) = @_;
+ my $running_async_event = $self->{running_async_event};
+ my ($runtime_event_token_found, $async_event);
+ foreach my $runtime_event_token (keys %$running_async_event) {
+ $async_event = $running_async_event->{$runtime_event_token};
+ if ($async_event->[0]{event_token} eq $event_token) {
+ $runtime_event_token_found = $runtime_event_token;
+ last;
+ }
+ }
+ &App::sub_exit($runtime_event_token_found) if ($App::trace);
+ return($runtime_event_token_found);
+}
+
+sub reset_running_async_events {
+ &App::sub_entry if ($App::trace);
+ my ($self, $runtime_event_token_prefix) = @_;
+ $runtime_event_token_prefix =~ s/:/-/; # in case they send "localhost:8080" instead of "localhost-8080"
+ my $running_async_event = $self->{running_async_event};
+ my ($runtime_event_token, $async_event);
+ foreach $runtime_event_token (keys %$running_async_event) {
+ $async_event = $running_async_event->{$runtime_event_token};
+ if ($async_event && $runtime_event_token =~ /^$runtime_event_token_prefix\b/) {
+ $self->reset_running_async_event($runtime_event_token);
+ }
+ }
+ &App::sub_exit() if ($App::trace);
+}
+
+sub reset_running_async_event {
+ &App::sub_entry if ($App::trace);
+ my ($self, $runtime_event_token) = @_;
+ my $async_event = $self->abort_running_async_event($runtime_event_token);
+ if ($async_event) {
+ my $pending_async_events = $self->{pending_async_events};
+ unshift(@$pending_async_events, $async_event);
+ }
+ &App::sub_exit($async_event) if ($App::trace);
+ return($async_event);
+}
+
+sub abort_async_event {
+ &App::sub_entry if ($App::trace);
+ my ($self, $event_token) = @_;
+ my $pending_async_events = $self->{pending_async_events};
+ my ($async_event);
+ my $aborted = 0;
+ # first look for it in the pending list
+ for (my $i = 0; $i <= $#$pending_async_events; $i++) {
+ $async_event = $pending_async_events->[$i];
+ if ($async_event->[0]{event_token} eq $event_token) {
+ splice(@$pending_async_events, $i, 1);
+ $aborted = 1;
+ last;
+ }
+ }
+ # then look for it in the running list
+ if (!$aborted) {
+ my $runtime_event_token = $self->find_runtime_event_token($event_token);
+ if ($runtime_event_token) {
+ $async_event = $self->abort_running_async_event($runtime_event_token);
+ }
+ }
+ &App::sub_exit($async_event) if ($App::trace);
+ return($async_event);
+}
+
+sub abort_running_async_event {
+ &App::sub_entry if ($App::trace);
+ my ($self, $runtime_event_token) = @_;
+ my $running_async_event = $self->{running_async_event};
+ my $pending_async_events = $self->{pending_async_events};
+ my $async_event = $running_async_event->{$runtime_event_token};
+ if ($async_event) {
+ $self->{num_async_events}--;
+ delete $self->{running_async_event}{$runtime_event_token};
+ unshift(@$pending_async_events, $async_event);
+ $self->_abort_running_async_event($runtime_event_token, @$async_event);
+ }
+ &App::sub_exit($async_event) if ($App::trace);
+ return($async_event);
+}
+
+# $runtime_event_tokens take the following forms:
+# $runtime_event_token = $pid; -- App::Context::POE::Server::send_async_event_now() and ::finish_pid()
+sub _abort_running_async_event {
+ &App::sub_entry if ($App::trace);
+ my ($self, $runtime_event_token, $event, $callback_event) = @_;
+ if ($runtime_event_token =~ /^[0-9]+$/) {
+ kill(15, $runtime_event_token);
+ }
+ else {
+ $self->log("Unable to abort running async event [$runtime_event_token]\n");
+ }
+ &App::sub_exit() if ($App::trace);
+}
+
+#############################################################################
+# user()
+#############################################################################
+
+=head2 user()
+
+The user() method returns the username of the authenticated user.
+The special name, "guest", refers to the unauthenticated (anonymous) user.
+
+ * Signature: $username = $context->user();
+ * Param: void
+ * Return: string
+ * Throws: <none>
+ * Since: 0.01
+
+ Sample Usage:
+
+ $username = $context->user();
+
+=cut
+
+sub user {
+ &App::sub_entry if ($App::trace);
+ my $self = shift;
+ my $user = $self->{user} || getlogin || (getpwuid($<))[0] || "guest";
+ &App::sub_exit($user) if ($App::trace);
+ $user;
+}
+
+#############################################################################
+### POE state routines
+#############################################################################
+
+sub _default {
+ &App::sub_entry if ($App::trace);
+ my ( $self, $kernel, $heap, $state, $args ) = @_[ OBJECT, KERNEL, HEAP, ARG0, ARG1 ];
+ my (@args);
+ @args = @$args if (ref($args) eq "ARRAY");
+ @args = ($args) if (!ref($args));
+ $self->log({level=>2},"POE: _default - WARNING: Entered an unhandled state ($state) with args (@args)\n");
+ &App::sub_exit() if ($App::trace);
+}
+
+sub _start {
+ &App::sub_entry if ($App::trace);
+ my ( $self, $kernel, $heap ) = @_[ OBJECT, KERNEL, HEAP ];
+ $self->log({level=>2},"POE: _start\n");
+
+ my $name = $self->{poe_session_name};
+ $kernel->alias_set($name);
+
+ $kernel->sig(CHLD => "poe_sigchld");
+ $kernel->sig(HUP => "poe_sigignore");
+ $kernel->sig(INT => "poe_sigterm");
+ $kernel->sig(QUIT => "poe_sigterm");
+ $kernel->sig(USR1 => "poe_sigignore");
+ $kernel->sig(USR2 => "poe_sigignore");
+ $kernel->sig(TERM => "poe_sigterm");
+
+ $kernel->call( IKC => publish => $name, $self->{poe_ikc_published_states} );
+
+ $kernel->post("IKC", "monitor", "*",
+ {register => "ikc_register",
+ unregister => "ikc_unregister",
+ shutdown => "ikc_shutdown"});
+
+ # don't start kicking off async events until we give the nodes a chance to register themselves
+ $kernel->delay_set("poe_event_loop_extension", 5) if (!$self->{disable_event_loop_extensions});
+ $kernel->delay_set("poe_alarm", 5);
+
+ &App::sub_exit() if ($App::trace);
+}
+
+sub _stop {
+ &App::sub_entry if ($App::trace);
+ my ( $self, $kernel, $heap, $state, $args ) = @_[ OBJECT, KERNEL, HEAP, ARG0, ARG1 ];
+ $self->log({level=>2},"POE: _start\n");
+ #sleep(1); # take a second to let child processes to die (perhaps not necessary, perhaps necessary when using POE::Wheel::Run)
+ &App::sub_exit() if ($App::trace);
+}
+
+sub ikc_register {
+ &App::sub_entry if ($App::trace);
+ my ($self, $kernel, $session_name) = @_[OBJECT, KERNEL, ARG1];
+ $self->log({level=>2},"POE: ikc_register ($session_name)\n");
+ my ($retval);
+ &App::sub_exit($retval) if ($App::trace);
+ return($retval);
+}
+
+sub ikc_unregister {
+ &App::sub_entry if ($App::trace);
+ my ($self, $kernel, $remote_kernel_id) = @_[OBJECT, KERNEL, ARG0];
+ $self->log({level=>2},"POE: ikc_unregister ($remote_kernel_id)\n");
+ &App::sub_exit() if ($App::trace);
+}
+
+sub ikc_shutdown {
+ &App::sub_entry if ($App::trace);
+ my ($self, $kernel, $arg0, $arg1, $arg2, $arg3) = @_[OBJECT, KERNEL, ARG0, ARG1, ARG2, ARG3];
+ $self->log({level=>2},"POE: ikc_shutdown args=($arg0, $arg1, $arg2, $arg3)\n");
+ &App::sub_exit() if ($App::trace);
+ return;
+}
+
+sub poe_sigterm {
+ &App::sub_entry if ($App::trace);
+ my ( $self, $kernel, $heap, $signame ) = @_[ OBJECT, KERNEL, HEAP, ARG0 ];
+ $self->log({level=>2},"POE: poe_sigterm (Caught signal $signame. quitting.)\n");
+
+ # How do I shut down the POE kernel now and exit?
+ # I think I need to shut down the last session and the kernel will exit.
+ # As per http://poe.perl.org/?POE_FAQ/How_do_I_force_a_session_to_shut_down
+ # $kernel->yield("poe_shutdown");
+ # However the signals which bring me here seem to do the shutdown for me, so it's unnecessary
+
+ &App::sub_exit() if ($App::trace);
+}
+
+sub poe_sigignore {
+ &App::sub_entry if ($App::trace);
+ my ( $self, $kernel, $heap, $signame ) = @_[ OBJECT, KERNEL, HEAP, ARG0 ];
+ $self->log({level=>2},"POE: poe_sigignore (Caught signal $signame. quitting.)\n");
+ &App::sub_exit() if ($App::trace);
+}
+
+sub poe_sigchld {
+ &App::sub_entry if ($App::trace);
+ my ( $self, $kernel, $heap, $pid, $status ) = @_[ OBJECT, KERNEL, HEAP, ARG1, ARG2 ];
+ #print STDERR "NOTICE: STATE (poe_sigchld) invoked with ($pid, $status) args\n";
+ my $exitval = $status >> 8;
+ my $sig = $status & 255;
+ $self->log({level=>2},"POE: poe_sigchld (Child $pid finished [exitval=$exitval,sig=$sig])\n");
+ $self->finish_pid($pid, $exitval, $sig);
+ $kernel->yield("poe_dispatch_pending_async_events");
+ &App::sub_exit() if ($App::trace);
+}
+
+sub poe_alarm {
+ &App::sub_entry if ($App::trace);
+ my ( $self, $kernel, $heap ) = @_[ OBJECT, KERNEL, HEAP ];
+ $self->log({level=>2},"POE: poe_alarm (Dispatching pending events and queueing scheduled events)\n");
+ $kernel->yield("poe_dispatch_pending_async_events");
+ my $time = time();
+ my (@events);
+ my $events_occurred = 0;
+ my $time_of_next_event = 0;
+ while ($time_of_next_event <= $time) {
+ $time_of_next_event = $self->get_current_events(\@events, $time);
+ if ($#events > -1) {
+ foreach my $event (@events) {
+ $kernel->yield("poe_run_event", $event); # put on the POE run queue
+ $events_occurred++;
+ }
+ $time = time();
+ }
+ }
+ # every time we process an alarm, we need to set the next one
+ my $sec_until_next_event = $time_of_next_event - $time;
+ $self->{alarm_id} = $kernel->delay_set("poe_alarm", $sec_until_next_event);
+ &App::sub_exit() if ($App::trace);
+}
+
+# NOTE: see http://poe.perl.org/?POE_FAQ/How_do_I_force_a_session_to_shut_down
+sub poe_shutdown {
+ &App::sub_entry if ($App::trace);
+ my ( $self, $kernel, $session, $heap ) = @_[ OBJECT, KERNEL, SESSION, HEAP ];
+ $self->log({level=>2},"POE: poe_shutdown\n");
+
+ # delete all wheels.
+ delete $heap->{wheel};
+
+ # clear your alias
+ $kernel->alias_remove( $heap->{alias} );
+
+ # clear all alarms you might have set
+ $kernel->alarm_remove_all();
+
+ # get rid of external ref count
+ $kernel->refcount_decrement( $session, $self->{poe_session_name} );
+
+ # propagate the message to children
+ $kernel->post( $heap->{child_session}, 'poe_shutdown' );
+ &App::sub_exit() if ($App::trace);
+ return;
+}
+
+sub poe_dispatch_pending_async_events {
+ &App::sub_entry if ($App::trace);
+ my ( $self, $kernel, $heap ) = @_[ OBJECT, KERNEL, HEAP ];
+ $self->log({level=>2},"POE: poe_dispatch_pending_async_events\n");
+ my $events_occurred = $self->dispatch_pending_async_events();
+ $kernel->yield("poe_dispatch_pending_async_events") if ($events_occurred > 0);
+ &App::sub_exit() if ($App::trace);
+}
+
+sub poe_event_loop_extension {
+ &App::sub_entry if ($App::trace);
+ my ( $self, $kernel, $heap ) = @_[ OBJECT, KERNEL, HEAP ];
+ $self->log({level=>2},"POE: poe_event_loop_extension\n");
+ my $event_loop_extensions = $self->{event_loop_extensions};
+ #$self->log({level=>2},"Event Loop extension ($event_loop_extensions: #=" . ($#$event_loop_extensions+1) . ").\n");
+ if ($event_loop_extensions && $#$event_loop_extensions > -1) {
+ my ($extension, $obj, $method, $args, $event_executed);
+ for (my $i = 0; $i <= $#$event_loop_extensions; $i++) {
+ $extension = $event_loop_extensions->[$i];
+ ($obj, $method, $args) = @$extension;
+ $event_executed = $obj->$method(@$args); # execute extension
+ #if ($event_executed) {
+ # $self->log({level=>2},"Event Loop extension: ${obj}->${method}(@$args) = $event_executed\n");
+ #}
+ }
+ }
+ $kernel->delay_set("poe_event_loop_extension", 1);
+ &App::sub_exit() if ($App::trace);
+}
+
+sub poe_run_event {
+ &App::sub_entry if ($App::trace);
+ my ( $self, $kernel, $heap, $event ) = @_[ OBJECT, KERNEL, HEAP, ARG0 ];
+ my ($event_str);
+ my $args = $event->{args} || [];
+ my $args_str = join(",", @$args);
+ if ($event->{name}) {
+ my $service_type = $event->{service_type} || "SessionObject";
+ $event_str = "$service_type($event->{name}).$event->{method}($args_str)";
+ }
+ else {
+ $event_str = "$event->{method}($args_str)";
+ }
+ $self->log({level=>2},"POE: poe_run_event ($event_str)\n");
+ $self->send_event($event);
+ &App::sub_exit() if ($App::trace);
+}
+
+sub poe_server_state {
+ &App::sub_entry if ($App::trace);
+ my ( $self, $kernel, $heap ) = @_[ OBJECT, KERNEL, HEAP ];
+ $self->log({level=>2},"POE: poe_server_state\n");
+ my $server_state = $self->state();
+ &App::sub_exit($server_state) if ($App::trace);
+ return $server_state;
+}
+
+sub poe_http_server_state {
+ &App::sub_entry if ($App::trace);
+ my ( $self, $kernel, $heap, $request, $response ) = @_[ OBJECT, KERNEL, HEAP, ARG0, ARG1 ];
+ $self->log({level=>2},"POE: poe_http_server_state\n");
+ my $server_state = $kernel->call( $self->{poe_session_name}, 'poe_server_state' );
+
+ # Build the response.
+ $response->code(RC_OK);
+ $response->push_header( "Content-Type", "text/plain" );
+ $response->content($server_state);
+
+ # Signal that the request was handled okay.
+ $kernel->post( $self->{poe_kernel_httpd_name}, 'DONE', $response );
+ &App::sub_exit(RC_OK) if ($App::trace);
+ return RC_OK;
+}
+
+sub poe_http_test_run {
+ &App::sub_entry if ($App::trace);
+ my ( $self, $kernel, $heap, $request, $response ) = @_[ OBJECT, KERNEL, HEAP, ARG0, ARG1 ];
+ $self->log({level=>2},"POE: poe_http_test_run\n");
+ my $event = {
+ service_type => "SessionObject",
+ name => "mvworkd",
+ method => "sleep2",
+ args => [ 3 ],
+ };
+ $self->send_async_event_now($event);
+
+ # Build the response.
+ $response->code(RC_OK);
+ $response->push_header( "Content-Type", "text/plain" );
+ $response->content("SessionObject(mvworkd).sleep(30)");
+
+ # Signal that the request was handled okay.
+ $kernel->post( $self->{poe_kernel_httpd_name}, 'DONE', $response );
+ &App::sub_exit(RC_OK) if ($App::trace);
+ return RC_OK;
+}
+
+1;
+
Modified: p5ee/trunk/App-Context/lib/App/Context/Server.pm
==============================================================================
--- p5ee/trunk/App-Context/lib/App/Context/Server.pm (original)
+++ p5ee/trunk/App-Context/lib/App/Context/Server.pm Fri Aug 3 12:34:40 2007
@@ -600,20 +600,33 @@
my ($self, $max_events) = @_;
$max_events ||= 9999;
my $pending_async_events = $self->{pending_async_events};
- my ($async_event, $assigned);
+ my ($async_event, $assigned, $event, $in_process);
my $events_occurred = 0;
my $i = 0;
- while ($i <= $#$pending_async_events && $events_occurred < $max_events) {
+ my $event_capacity_exists = 1;
+ my $max_i = $#$pending_async_events;
+ while ($i <= $max_i && $events_occurred < $max_events) {
$async_event = $pending_async_events->[$i];
- $assigned = $self->assign_event_destination($async_event->[0]);
- if ($assigned) {
+ $event = $async_event->[0];
+ if ($event->{destination}) {
$self->send_async_event_now(@$async_event);
$events_occurred ++;
splice(@$pending_async_events, $i, 1); # remove $pending_async_events->[$i]
- # keep $i the same
+ $max_i--;
}
- elsif (! defined $assigned) { # [undef] no servers are eligible for assignment
- last; # there's no sense looking at the other pending async events
+ elsif ($event_capacity_exists) {
+ $assigned = $self->assign_event_destination($event);
+ if ($assigned) {
+ $self->send_async_event_now(@$async_event);
+ $events_occurred ++;
+ # keep $i the same
+ splice(@$pending_async_events, $i, 1); # remove $pending_async_events->[$i]
+ $max_i--;
+ }
+ else { # [undef] no servers are eligible for assignment
+ $event_capacity_exists = 0; # there's no sense looking at the other pending async events
+ $i++; # look at the next one
+ }
}
else { # [0] this async_event is not eligible to run
$i++; # look at the next one
@@ -639,34 +652,39 @@
sub send_async_event_now {
&App::sub_entry if ($App::trace);
my ($self, $event, $callback_event) = @_;
- my $pid = $self->fork();
- if (!$pid) { # running in child
- my $exitval = 0;
- my (@results);
- eval {
- @results = $self->send_event($event);
- };
- if ($@) {
- @results = ($@);
- }
- if ($#results > -1 && defined $results[0] && $results[0] ne "") {
- my $textfile = $self->{options}{prefix} . "/data/app/Context/$$";
- if (open(FILE, "> $textfile")) {
- print App::Context::Server::FILE @results;
- close(App::Context::Server::FILE);
+ if ($event->{destination} eq "in_process") {
+ my $event_token = $self->send_async_event_in_process($event, $callback_event);
+ }
+ else {
+ my $pid = $self->fork();
+ if (!$pid) { # running in child
+ my $exitval = 0;
+ my (@results);
+ eval {
+ @results = $self->send_event($event);
+ };
+ if ($@) {
+ @results = ($@);
}
- else {
- $exitval = 1;
+ if ($#results > -1 && defined $results[0] && $results[0] ne "") {
+ my $textfile = $self->{options}{prefix} . "/data/app/Context/$$";
+ if (open(FILE, "> $textfile")) {
+ print App::Context::Server::FILE @results;
+ close(App::Context::Server::FILE);
+ }
+ else {
+ $exitval = 1;
+ }
}
+ $self->shutdown();
+ $self->exit($exitval);
}
- $self->shutdown();
- $self->exit($exitval);
+ my $destination = $event->{destination} || "local";
+ $self->{num_async_events}++;
+ $self->{node}{$destination}{num_async_events}++;
+ my $runtime_event_token = $pid;
+ $self->{running_async_event}{$runtime_event_token} = [ $event, $callback_event ];
}
- my $destination = $event->{destination} || "local";
- $self->{num_async_events}++;
- $self->{node}{$destination}{num_async_events}++;
- my $runtime_event_token = $pid;
- $self->{running_async_event}{$runtime_event_token} = [ $event, $callback_event ];
&App::sub_exit() if ($App::trace);
}
lmpx.com only provides a reader for public news (NNTP) servers. It is not
affiliated with the servers or forums shown here and is not responsible for
the content of articles, which is written by their respective authors.