Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions Changes
Original file line number Diff line number Diff line change
@@ -1,3 +1,4 @@
- Allow log closing and deletion to be adjusted.
- Add documentation for the ElasticSearch document format

0.004
Expand Down
43 changes: 33 additions & 10 deletions lib/Message/Passing/Output/ElasticSearch.pm
Original file line number Diff line number Diff line change
Expand Up @@ -3,7 +3,7 @@ use Moose;
use ElasticSearch;
use AnyEvent;
use Scalar::Util qw/ weaken /;
use MooseX::Types::Moose qw/ ArrayRef Str Bool /;
use MooseX::Types::Moose qw/ ArrayRef Str Bool Int /;
use Scalar::Util qw/ weaken /;
use Try::Tiny qw/ try catch /;
use aliased 'DateTime' => 'DT';
Expand Down Expand Up @@ -65,6 +65,24 @@ has verbose => (
},
);

has housekeeping => (
isa => Bool,
is => 'ro',
default => 1,
);

has close_after_days => (
isa => Int,
is => 'ro',
default => 7,
);

has delete_after_days => (
isa => Int,
is => 'ro',
default => 30,
);

sub consume {
my ($self, $data) = @_;
return unless $data;
Expand Down Expand Up @@ -195,13 +213,18 @@ has _archive_timer => (
is => 'ro',
default => sub {
my $self = shift;
weaken($self);
my $time = 60 * 60 * 24; # Every day
AnyEvent->timer(
after => 60, # delay 1 hour to start first loop
interval => $time,
cb => sub { $self->_archive_index() },
);
if ($self->housekeeping) {
weaken($self);
my $time = 60 * 60 * 24; # Every day
return AnyEvent->timer(
after => 60, # delay 1 hour to start first loop
interval => $time,
cb => sub { $self->_archive_index() },
);
}
else {
return undef;
}
},
);

Expand All @@ -213,12 +236,12 @@ sub _archive_index {

my $dt = DT->from_epoch(epoch => time());

my $dt_to_close = $dt->clone->subtract(days => 7);
my $dt_to_close = $dt->clone->subtract(days => $self->close_after_days);
my $index_to_close = $self->_index_name_by_dt($dt_to_close);
$self->_es->close_index(index => $index_to_close)
->cb( sub { warn "Close index: $index_to_close \n" if $self->verbose; });

my $dt_to_delete = $dt->clone->subtract(days => 30);
my $dt_to_delete = $dt->clone->subtract(days => $self->delete_after_days);
my $index_to_delete = $self->_index_name_by_dt($dt_to_delete);
$self->_es->delete_index(
index => $index_to_delete,
Expand Down