use v5.10.1;
use parent qw(PublicInbox::Lock);
use PublicInbox::SearchIdxShard;
+use PublicInbox::IPC;
use PublicInbox::Eml;
use PublicInbox::Git;
use PublicInbox::Import;
use PublicInbox::OverIdx;
use PublicInbox::Msgmap;
use PublicInbox::Spawn qw(spawn popen_rd run_die);
-use PublicInbox::SearchIdx qw(log2stack crlf_adjust is_ancestor check_size
- is_bad_blob);
+use PublicInbox::Search;
+use PublicInbox::SearchIdx qw(log2stack is_ancestor check_size is_bad_blob);
use IO::Handle; # ->autoflush
use File::Temp ();
+use POSIX ();
my $OID = qr/[a-f0-9]{40,}/;
# an estimate of the post-packed size to the raw uncompressed size
# to increase Xapian shards
our $NPROC_MAX_DEFAULT = 4;
-sub detect_nproc () {
- # getconf(1) is POSIX, but *NPROCESSORS* vars are not
- for (qw(_NPROCESSORS_ONLN NPROCESSORS_ONLN)) {
- `getconf $_ 2>/dev/null` =~ /^(\d+)$/ and return $1;
- }
- for my $nproc (qw(nproc gnproc)) { # GNU coreutils nproc
- `$nproc 2>/dev/null` =~ /^(\d+)$/ and return $1;
- }
-
- # should we bother with `sysctl hw.ncpu`? Those only give
- # us total processor count, not online processor count.
- undef
-}
-
sub nproc_shards ($) {
my ($creat_opt) = @_;
my $n = $creat_opt->{nproc} if ref($creat_opt) eq 'HASH';
$n //= $ENV{NPROC};
if (!$n) {
# assume 2 cores if not detectable or zero
- state $NPROC_DETECTED = detect_nproc() || 2;
+ state $NPROC_DETECTED = PublicInbox::IPC::detect_nproc() || 2;
$n = $NPROC_DETECTED;
$n = $NPROC_MAX_DEFAULT if $n > $NPROC_MAX_DEFAULT;
}
}
# indexes a message, returns true if checkpointing is needed
-sub do_idx ($$$$) {
- my ($self, $msgref, $mime, $smsg) = @_;
- $smsg->{bytes} = $smsg->{raw_bytes} + crlf_adjust($$msgref);
- $self->{oidx}->add_overview($mime, $smsg);
- my $idx = idx_shard($self, $smsg->{num});
- $idx->index_raw($msgref, $mime, $smsg);
- my $n = $self->{transact_bytes} += $smsg->{raw_bytes};
+sub do_idx ($$$) {
+ my ($self, $eml, $smsg) = @_;
+ $self->{oidx}->add_overview($eml, $smsg);
+ if ($self->{-need_xapian}) {
+ my $idx = idx_shard($self, $smsg->{num});
+ $idx->index_eml($eml, $smsg);
+ }
+ my $n = $self->{transact_bytes} += $smsg->{bytes};
$n >= $self->{batch_bytes};
}
$cmt = $im->get_mark($cmt);
$self->{last_commit}->[$self->{epoch_max}] = $cmt;
- my $msgref = delete $smsg->{-raw_email};
- if (do_idx($self, $msgref, $mime, $smsg)) {
+ if (do_idx($self, $mime, $smsg)) {
$self->checkpoint;
}
my $max = $self->{shards} - 1;
my $idx = $self->{idx_shards} = [];
push @$idx, PublicInbox::SearchIdxShard->new($self, $_) for (0..$max);
+ $self->{-need_xapian} = $idx->[0]->need_xapian;
# SearchIdxShard may do their own flushing, so don't scale
# until after forking
for my $smsg (@$need_reindex) {
my $new_smsg = bless {
blob => $blob,
- raw_bytes => $bytes,
num => $smsg->{num},
mid => $smsg->{mid},
}, 'PublicInbox::Smsg';
my $sync = { autime => $smsg->{ds}, cotime => $smsg->{ts} };
$new_smsg->populate($new_mime, $sync);
- do_idx($self, \$raw, $new_mime, $new_smsg);
+ $new_smsg->set_bytes($raw, $bytes);
+ do_idx($self, $new_mime, $new_smsg);
}
$rewritten->{rewrites};
}
}
my $midx = $self->{midx}; # misc index
- $midx->commit_txn if $midx;
+ if ($midx) {
+ $midx->commit_txn;
+ $PublicInbox::Search::X{CLOEXEC_UNSET} and
+ $self->git->cleanup;
+ }
# last_commit is special, don't commit these until
# Xapian shards are done:
$dbh->commit;
$dbh->begin_work;
}
- $midx->begin_txn if $midx;
}
$self->{total_bytes} += $self->{transact_bytes};
$self->{transact_bytes} = 0;
sub write_alternates ($$$) {
my ($info_dir, $mode, $out) = @_;
- my $fh = File::Temp->new(TEMPLATE => 'alt-XXXXXXXX', DIR => $info_dir);
+ my $fh = File::Temp->new(TEMPLATE => 'alt-XXXX', DIR => $info_dir);
my $tmp = $fh->filename;
print $fh @$out or die "print $tmp: $!\n";
chmod($mode, $fh) or die "fchmod $tmp: $!\n";
sub diff ($$$) {
my ($mid, $cur, $new) = @_;
- my $ah = File::Temp->new(TEMPLATE => 'email-cur-XXXXXXXX', TMPDIR => 1);
+ my $ah = File::Temp->new(TEMPLATE => 'email-cur-XXXX', TMPDIR => 1);
print $ah $cur->as_string or die "print: $!";
$ah->flush or die "flush: $!";
PublicInbox::Import::drop_unwanted_headers($new);
- my $bh = File::Temp->new(TEMPLATE => 'email-new-XXXXXXXX', TMPDIR => 1);
+ my $bh = File::Temp->new(TEMPLATE => 'email-new-XXXX', TMPDIR => 1);
print $bh $new->as_string or die "print: $!";
$bh->flush or die "flush: $!";
my $cmd = [ qw(diff -u), $ah->filename, $bh->filename ];
sub atfork_child {
my ($self) = @_;
if (my $older_siblings = $self->{idx_shards}) {
- $_->shard_atfork_child for @$older_siblings;
+ $_->ipc_sibling_atfork_child for @$older_siblings;
}
if (my $im = $self->{im}) {
$im->atfork_child;
}
# {unindexed} is unlikely
- if ((my $unindexed = $arg->{unindexed}) && scalar(@$mids) == 1) {
- $num = delete($unindexed->{$mids->[0]});
+ if (my $unindexed = $arg->{unindexed}) {
+ my $oidbin = pack('H*', $oid);
+ my $u = $unindexed->{$oidbin};
+ ($num, $mid0) = splice(@$u, 0, 2) if $u;
if (defined $num) {
- $mid0 = $mids->[0];
$self->{mm}->mid_set($num, $mid0);
- delete($arg->{unindexed}) if !keys(%$unindexed);
+ if (scalar(@$u) == 0) { # done with current OID
+ delete $unindexed->{$oidbin};
+ delete($arg->{unindexed}) if !keys(%$unindexed);
+ }
}
}
if (!defined($num)) { # reuse if reindexing (or duplicates)
}
++${$arg->{nr}};
my $smsg = bless {
- raw_bytes => $size,
num => $num,
blob => $oid,
mid => $mid0,
}, 'PublicInbox::Smsg';
$smsg->populate($eml, $arg);
- if (do_idx($self, $bref, $eml, $smsg)) {
+ $smsg->set_bytes($$bref, $size);
+ if (do_idx($self, $eml, $smsg)) {
${$arg->{need_checkpoint}} = 1;
}
index_finalize($arg, 1);
sub unindex_oid_aux ($$$) {
my ($self, $oid, $mid) = @_;
my @removed = $self->{oidx}->remove_oid($oid, $mid);
+ return unless $self->{-need_xapian};
for my $num (@removed) {
- my $idx = idx_shard($self, $num);
- $idx->shard_remove($num);
+ idx_shard($self, $num)->ipc_do('xdb_remove', $num);
}
}
warn "BUG: multiple articles linked to $oid\n",
join(',',sort keys %gone), "\n";
}
- foreach my $num (keys %gone) {
+ # reuse (num => mid) mapping in ascending numeric order
+ for my $num (sort { $a <=> $b } keys %gone) {
+ $num += 0;
if ($unindexed) {
my $mid0 = $mm->mid_for($num);
- $unindexed->{$mid0} = $num;
+ my $oidbin = pack('H*', $oid);
+ push @{$unindexed->{$oidbin}}, $num, $mid0;
}
$mm->num_delete($num);
}
sub unindex_todo ($$$) {
my ($self, $sync, $unit) = @_;
my $unindex_range = delete($unit->{unindex_range}) // return;
- my $unindexed = $sync->{unindexed} //= {}; # $mid0 => $num
+ my $unindexed = $sync->{unindexed} //= {}; # $oidbin => [$num, $mid0]
my $before = scalar keys %$unindexed;
# order does not matter, here:
my $fh = $unit->{git}->popen(qw(log --raw -r --no-notes --no-color
sub index_xap_only { # git->cat_async callback
my ($bref, $oid, $type, $size, $smsg) = @_;
- my $self = $smsg->{self};
+ my $self = delete $smsg->{self};
my $idx = idx_shard($self, $smsg->{num});
- $smsg->{raw_bytes} = $size;
- $idx->index_raw($bref, undef, $smsg);
- $self->{transact_bytes} += $size;
+ $idx->index_eml(PublicInbox::Eml->new($bref), $smsg);
+ $self->{transact_bytes} += $smsg->{bytes};
}
sub index_xap_step ($$$;$) {
}
}
$self->git->cat_async_wait;
+ $self->{ibx}->cleanup;
$self->done;
}