AC-MrGamoo

 view release on metacpan or  search on metacpan

lib/AC/MrGamoo/Job.pm  view on Meta::CPAN


our $MAXFILE = `sh -c "ulimit -n"`;
$MAXFILE = 255 if $^O eq 'solaris' && $MAXFILE > 255;

################################################################

# schedule periodic "cronjob"
AC::DC::Sched->new(
    info	=> "job periodic",
    freq	=> 2,
    func	=> \&periodic,
   );

################################################################

sub new {
    my $class = shift;
    # %{ APCMRMJobCreate }

    my $me = bless {
        request		=> { @_ },
        phase_no	=> -1,
        file_info	=> {},
        tmp_file	=> [],
        server_info	=> {},
        task_running	=> {},
        task_pending	=> {},
        xfer_running	=> {},
        xfer_pending	=> {},
        request_running	=> {},
        request_pending	=> {},
        statistics	=> { job_start => time() },
    }, $class;

    if( $REGISTRY{ $me->{request}{jobid} } ){
        verbose("ignoring duplicate request job $me->{request}{jobid}");
        # will cause a 200 OK, so the requestor will not retry
        return $REGISTRY{ $me->{request}{jobid} };
    }

    verbose("new job: $me->{request}{jobid} ($me->{request}{traceinfo})");

    my $cf = $me->{options} = decode_json( $me->{request}{options} ) if $me->{request}{options};

    # open connection  to eu-console
    $me->{euconsole} = AC::MrGamoo::EUConsole->new( $me->{request}{jobid}, $me->{request}{console} );

    # partially compile
    eval {
        $me->{mr} = AC::MrGamoo::Submit::Compile->new( text => $me->{request}{jobsrc} );
    };
    if(my $e = $@){
        problem("cannot compile job: $e");
        return;
    }

    # RSN - get_file_list + Plan may take too long - do in sub-process

    # get file list
    my $files = get_file_list( $cf );
    #print STDERR "files: ", dumper($files), "\n";

    for my $f (@$files){
        $me->{file_info}{ $f->{filename} } = $f;
    }

    # get server list
    my $servers = get_peer_list( $cf );
    #print STDERR "servers: ", dumper($servers), "\n";

    # plan job
    my $plan = AC::MrGamoo::Job::Plan->new( $me, $servers, $files );
    #print STDERR "plan: ", dumper($plan), "\n";

    $me->{plan} = $plan;

    $me->{maxfail} = 5 * ( (keys %{$plan->{taskidx}}) + @{$plan->{copying}});

    $me->{server_info}{$_->{id}} = {} for @$servers;

    $me->_preload_file_copies();
    $REGISTRY{ $me->{request}{jobid} } = $me;
    return $me;
}

sub start {
    my $me = shift;

    debug("start job");
    $me->{euconsole}->send_msg('debug', 'starting job');
    $me->_try_to_do_something();
    1;
}

################################################################

# record status rcvd from task
sub task_status {
    my $me = _find(shift, @_);
    my %p  = @_;
    my $taskid = $p{taskid};

    return unless $me;
    my $t = $me->{task_running}{$taskid};
    return unless $t;

    $t->update_status( $me, $p{phase}, $p{progress} );

    1;
}

# record status rcvd from file xfer
sub xfer_status {
    my $me = _find(shift, @_);
    my %p  = @_;
    my $copyid = $p{copyid};

    return unless $me;
    my $c = $me->{xfer_running}{$copyid};
    return unless $c;

    $c->update_status( $me, $p{status_code} );

    1;
}

################################################################

sub periodic {
    # debug("periodic check");

    $_trying = 0;
    for my $job (values %REGISTRY){

lib/AC/MrGamoo/Job.pm  view on Meta::CPAN

    for my $r (@rp){
        last if $startreqs <= 0;
        next unless $me->_maybe_start_request( $r );
        $startreqs --;
    }

    unless( $me->{aborted} ){

        # are there tasks that can start
        my $started = 0;
        my @tp = sort { $a->{created} <=> $b->{created} } values %{$me->{task_pending}};
        for my $t (@tp){
            $started += $me->_maybe_start_task( $t );
            last if $started >= $me->{plan}{nserver} / 4;	# keep them from getting in to lockstep
        }

        # are there copies that can start
        my @cp = sort { $a->{created} <=> $b->{created} } values %{$me->{xfer_pending}};
        for my $c (@cp){
            $me->_maybe_start_xfer( $c );
        }
    }

    # should we speculatively copy some files

    # should we speculatively retry a task


    $_trying --;
}

################################################################

sub something_failed {
    my $me = shift;

    return if ++$me->{total_fails} < $me->{maxfail};
    $me->abort();
    return 1;
}

################################################################

sub report {

    my $txt;

    for my $j (values %REGISTRY){
        my $ph;
        $ph = 'start'     if $j->{phase_no} < 0;
        $ph ||= 'cleanup' if $j->{phase_no} >= @{$j->{plan}{phases}};
        $ph ||= $j->{plan}{phases}[ $j->{phase_no} ];

        my $tr = keys %{$j->{task_running}};
        my $tp = keys %{$j->{task_pending}};
        my $cr = keys %{$j->{copy_running}};
        my $cp = keys %{$j->{copy_pending}};
        my $rr = keys %{$j->{request_running}};
        my $rp = keys %{$j->{request_pending}};

        $txt .= sprintf("%s %8s %4d %4d %4d %4d %4d %4d\n", $j->{request}{jobid}, $ph, $tr, $tp, $cr, $cp, $rr, $rp);
        # ...
    }

    return $txt;
}


1;



( run in 0.568 second using v1.01-cache-2.11-cpan-bbc515a03b3 )