#-*- perl -*-
#
# Copyright (C) 2001,2002,2003,2004,2005,2006 Ken'ichi Fukamachi
# All rights reserved. This program is free software; you can
# redistribute it and/or modify it under the same terms as Perl itself.
#
# $FML: QueueManager.pm,v 1.43 2006/04/22 08:55:32 fukachan Exp $
#
package FML::Process::QueueManager;
use strict;
use Carp;
=head1 NAME
FML::Process::QueueManager - provide queue manipulation functions.
=head1 SYNOPSIS
To flush all entries in the queue,
use FML::Process::QueueManager;
my $qmgr_args = { directory => $queue_dir };
my $queue = new FML::Process::QueueManager $curproc, $qmgr_args;
$queue->send($curproc);
or if you send specific queue C<$queue_id>, use
$queue->send($curproc, $queue_id);
where C<$queue_id> is queue id such as 1000390413.14775.1,
not file path.
=head1 DESCRIPTION
queue flush!
=head1 METHODS
=head2 new($qmgr_args)
constructor.
=cut
# Descriptions: constructor.
# Arguments: OBJ($self) OBJ($curproc) HASH_REF($qmgr_args)
# Side Effects: none
# Return Value: OBJ
sub new
{
my ($self, $curproc, $qmgr_args) = @_;
my ($type) = ref($self) || $self;
my $me = {};
$me->{ _curproc } = $curproc;
$me->{ _directory } = $qmgr_args->{ directory };
return bless $me, $type;
}
=head2 send( [ $id ] )
try to send all messages in the queue. If the queue id C<$id> is
specified, send only the queue corresponding to C<$id>.
=cut
# Descriptions: send message(s) in queue directory sequentially.
# Arguments: OBJ($self) STR($id)
# Side Effects: queue flush-ed
# Return Value: none
sub send
{
my ($self, $id) = @_;
my $curproc = $self->{ _curproc };
my $queue_dir = $self->{ _directory };
my $max_count = 100;
my $count = 0;
my $count_ok = 0;
my $count_err = 0;
my $channel = 'qmgr_reschedule';
my $fp_log = sub { $curproc->log(@_);};
my $fp_logerror = sub { $curproc->logerror(@_);};
my $fp_logdebug = sub { $curproc->logdebug(@_);};
use Mail::Delivery::Queue;
my $queue = new Mail::Delivery::Queue { directory => $queue_dir };
$queue->set_log_function($fp_log);
$queue->set_log_error_function($fp_logerror);
$queue->set_log_debug_function($fp_logdebug);
# XXX-TODO: more readable variable name: $ra -> $msg_queue_list ?
my $ra = [];
if (defined $id) {
$ra = [ $id ];
}
else {
# XXX-TODO: customizable
$queue->set_policy("fair-queue");
$ra = $queue->list();
unless (@$ra) {
$curproc->logdebug("qmgr: empty active queue. re-schedule");
$queue->reschedule();
$ra = $queue->list();
}
}
QUEUE:
for my $qid (@$ra) {
last QUEUE if $curproc->is_process_time_limit();
my $q = new Mail::Delivery::Queue {
id => $qid,
directory => $queue_dir,
};
$q->set_log_info_function($fp_log);
$q->set_log_error_function($fp_logerror);
$q->set_log_debug_function($fp_logdebug);
# check before try lock (XXX not enough check)
unless ($q->is_valid_active_queue()) {
next QUEUE;
}
if ( $q->lock( { wait => 10 } ) && $q->is_valid_active_queue() ) {
$curproc->logdebug("qmgr: got lock qid=$qid");
my $is_locked = 1;
if ( $q->is_valid_active_queue() ) {
my $r = $self->_send($q);
if ($r) {
$q->remove();
$count_ok++;
}
else {
$curproc->logdebug("qmgr: qid=$qid try later.");
$q->unlock();
$q->sleep_queue();
$is_locked = 0;
$count_err++;
}
$count++;
}
else {
# XXX-TODO: $q->remove() if invalid queue ?
$curproc->logerror("qmgr: qid=$qid is invalid");
}
$q->unlock() if $is_locked;
}
else {
$curproc->logdebug("qmgr: qid=$qid is locked or invalid. retry");
}
# upper limit of processing done on one process.
last QUEUE if $count >= $max_count;
}
if ($count) {
$curproc->logdebug("qmgr: total=$count ok=$count_ok error=$count_err");
}
if ($curproc->is_event_timeout($channel)) {
if (defined $queue) {
$curproc->logdebug("qmgr: re-schedule");
$queue->reschedule();
}
# XXX-TODO: customizable
$curproc->event_set_timeout($channel, time + 300);
}
}
# Descriptions: send message object $q.
# Arguments: OBJ($self) OBJ($q)
# Side Effects: queue flush-ed
# Return Value: STR
sub _send
{
my ($self, $q) = @_;
my $curproc = $self->{ _curproc };
my $qid = $q->id();
my $qf_act = $q->active_file_path($qid);
my $recipient_map = $q->recipients_file_path($qid);
my $sender = $q->get_sender($qid);
use Mail::Message;
my $msg = Mail::Message->parse( { file => $qf_act } );
# XXX lock for recipient maps is NOT needed since already a copy.
# XXX queue is already locked and need no lock for recipient maps here.
use FML::Mailer;
my $mailwrapper = new FML::Mailer $curproc;
my $r = $mailwrapper->send($q, {
sender => $sender,
recipient_maps => $recipient_map,
message => $msg,
});
if ($r) {
my $delay = '?';
if ($qid =~ /^(\d+)\./) { $delay = time - $1;}
$curproc->log("qmgr: status=sent qid=$qid delay=$delay");
}
else {
$curproc->log("qmgr: status=deferred qid=$qid");
}
return $r;
}
=head2 cleanup( [ $id ] )
clean up queue directory.
=cut
# Descriptions: clean up directory.
# Arguments: OBJ($self)
# Side Effects: queue flush-ed
# Return Value: none
sub cleanup
{
my ($self) = @_;
my $curproc = $self->{ _curproc };
my $queue_dir = $self->{ _directory };
my $fp_log = sub { $curproc->log(@_);};
my $fp_logerror = sub { $curproc->logerror(@_);};
my $fp_logdebug = sub { $curproc->logdebug(@_);};
use Mail::Delivery::Queue;
my $queue = new Mail::Delivery::Queue { directory => $queue_dir };
$queue->set_log_function($fp_log);
$queue->set_log_error_function($fp_logerror);
$queue->set_log_debug_function($fp_logdebug);
# XXX-TODO: customizable. $mail_queue_max_lifetime = 5d ?
my $list = $queue->list_all() || [];
my $limit = 5 * 24 * 3600; # 5 days.
my $now = time;
for my $qid (@$list) {
my $q = new Mail::Delivery::Queue {
id => $qid,
directory => $queue_dir,
};
$q->set_log_function($fp_log);
$q->set_log_error_function($fp_logerror);
$q->set_log_debug_function($fp_logdebug);
unless ( $q->is_valid_active_queue() ) {
my $mtime = $q->last_modified_time();
# enough old.
if ($mtime < $now - $limit) {
$curproc->logdebug("qmgr: remove too old queue qid=$qid");
$q->remove();
}
}
}
}
=head1 CODING STYLE
See C on fml coding style guide.
=head1 AUTHOR
Ken'ichi Fukamachi
=head1 COPYRIGHT
Copyright (C) 2001,2002,2003,2004,2005,2006 Ken'ichi Fukamachi
All rights reserved. This program is free software; you can
redistribute it and/or modify it under the same terms as Perl itself.
=head1 HISTORY
FML::Process::QueueManager first appeared in fml8 mailing list driver package.
See C for more details.
=cut
1;