lib/PublicInbox/LeiOverview.pm | 23 ++--------------------- lib/PublicInbox/LeiToMail.pm | 3 +-- lib/PublicInbox/LeiXSearch.pm | 16 ++++++++++++---- diff --git a/lib/PublicInbox/LeiOverview.pm b/lib/PublicInbox/LeiOverview.pm index 3125f015499e5e7833f3f0036129293623e19b31..d3df4faa72c2eeea1c8b3648a6f1a6bc1bac416f 100644 --- a/lib/PublicInbox/LeiOverview.pm +++ b/lib/PublicInbox/LeiOverview.pm @@ -147,17 +147,6 @@ } sub ovv_atexit_child { my ($self, $lei) = @_; - if (my $l2m = $lei->{l2m}) { - # wait for ->write_mail work we submitted to lei2mail - if (my $rd = delete $l2m->{each_smsg_done}) { - read($rd, my $buf, 1); # wait for EOF - } - } - # order matters, git->{-tmp}->DESTROY must not fire until - # {each_smsg_done} hits EOF above - if (my $git = delete $self->{git}) { - $git->async_wait_all; - } if (my $bref = delete $lei->{ovv_buf}) { my $lk = $self->lock_for_scope; $lei->out($$bref); @@ -213,19 +202,11 @@ my ($smsg, undef, $eml) = @_; # no mitem in $_[1] $wcb->(undef, $smsg, $eml); }; } elsif ($l2m && $l2m->{-wq_s1}) { - # $io->[0] becomes a notification pipe that triggers EOF - # in this wq worker when all outstanding ->write_mail - # calls are complete - my $io = []; - pipe($l2m->{each_smsg_done}, $io->[0]) or die "pipe: $!"; - fcntl($io->[0], 1031, 4096) if $^O eq 'linux'; # F_SETPIPE_SZ - my $git = $ibxish->git; # (LeiXSearch|Inbox|ExtSearch)->git - $self->{git} = $git; - my $git_dir = $git->{git_dir}; + my $git_dir = $ibxish->git->{git_dir}; sub { my ($smsg, $mitem) = @_; $smsg->{pct} = get_pct($mitem) if $mitem; - $l2m->wq_do('write_mail', $io, $git_dir, $smsg); + $l2m->wq_do('write_mail', [], $git_dir, $smsg); } } elsif ($self->{fmt} =~ /\A(concat)?json\z/ && $lei->{opt}->{pretty}) { my $EOR = ($1//'') eq 'concat' ? "\n}" : "\n},"; diff --git a/lib/PublicInbox/LeiToMail.pm b/lib/PublicInbox/LeiToMail.pm index 1f815e4079e9d7e9419194baa1f9b74482aa7408..4f84722188ad7c3c1fdb0d1261d903dc949f7791 100644 --- a/lib/PublicInbox/LeiToMail.pm +++ b/lib/PublicInbox/LeiToMail.pm @@ -490,10 +490,9 @@ } sub write_mail { # via ->wq_do my ($self, $git_dir, $smsg) = @_; - my $not_done = delete $self->{0} // die 'BUG: $not_done missing'; my $git = $self->{"$$\0$git_dir"} //= PublicInbox::Git->new($git_dir); git_async_cat($git, $smsg->{blob}, \&git_to_mail, - [$self->{wcb}, $smsg, $not_done]); + [$self->{wcb}, $smsg]); } sub wq_atexit_child { diff --git a/lib/PublicInbox/LeiXSearch.pm b/lib/PublicInbox/LeiXSearch.pm index e7f0ef6369e3d6fa400ad332a97314357c5e9ee9..2dc44414b29f12b17ffc4ccaaa8b6be7069bb4cb 100644 --- a/lib/PublicInbox/LeiXSearch.pm +++ b/lib/PublicInbox/LeiXSearch.pm @@ -287,12 +287,15 @@ undef $each_smsg; $lei->{ovv}->ovv_atexit_child($lei); } -sub git { +# called by LeiOverview::each_smsg_cb +sub git { $_[0]->{git_tmp} // die 'BUG: caller did not set {git_tmp}' } + +sub git_tmp ($) { my ($self) = @_; my (%seen, @dirs); - my $tmp = File::Temp->newdir('lei_xsrch_git-XXXXXXXX', TMPDIR => 1); - for my $ibx (@{$self->{shard2ibx} // []}) { - my $d = File::Spec->canonpath($ibx->git->{git_dir}); + my $tmp = File::Temp->newdir("lei_xsearch_git.$$-XXXX", TMPDIR => 1); + for my $ibxish (locals($self)) { + my $d = File::Spec->canonpath($ibxish->git->{git_dir}); $seen{$d} //= push @dirs, "$d/objects\n" } my $git_dir = $tmp->dirname; @@ -427,6 +430,11 @@ $lei->oldset, { lei => $lei }); pipe($lei->{startq}, $lei->{au_done}) or die "pipe: $!"; # 1031: F_SETPIPE_SZ fcntl($lei->{startq}, 1031, 4096) if $^O eq 'linux'; + } + if (!$lei->{opt}->{thread} && locals($self)) { # for query_mset + # lei->{git_tmp} is set for wq_wait_old so we don't + # delete until all lei2mail + lei_xsearch workers are reaped + $lei->{git_tmp} = $self->{git_tmp} = git_tmp($self); } $self->wq_workers_start('lei_xsearch', $self->{jobs}, $lei->oldset, { lei => $lei });