Files
mu4e/lib/mu-store.cc
Daniel Colascione 31835200bd New incremental cleanup strategy: 63%-88% faster
This change adds a new cleanup mode that avoids cleanup having
re-traverse the directories the index pass just looked at.
Additionally, we efficiently query the Xapian database by walking the
term list instead of doing multiple point-wise path lookups.

I'd noticed that most of my time in mu's cleanup pass consisted of
B-tree lookups in Xapian (one 8KB pread64 at a time).  The point
lookups forced Xapian to traverse from the root of the B-tree to the
leaf for every single message.  Additionally, in order to join on the
message path, we had to do *another* B-tree traversal after locating
each message term.  Now we just walk the terms in order, which is much
more efficient, as we touch each B-tree node only once.

On my system, with 1371861 total messages, the total time of mu
index (no lazy check):

--nocleanup: 3.6s
incremental cleanup: 4.2s (0.6s in cleanup)
legacy cleanup: 5.2s (1.6s in cleanup)

With the new mode, we save 1.0s of the 1.6s cleanup, so we're
~63% faster.

But the incremental cleanup works even better with lazy checking.
If I enable --lazy-check, dirty only my INBOX (360778 messages), and
run index, I get:

--nocleanup: 0.9s
incremental cleanup: 1.1s (0.2s in cleanup)
legacy cleanup: 2.5s (1.6s in cleanup)

We save 1.4s out of 1.6s for ~88% speedup.

This change also fixes a timestamp bug: we should be storing
the *start* time of the index pass in metadata, not the end time, so
that on the next index pass, we notice messages that arrived between
the two times.

All tests pass.  You can set the environment variable
MU_NO_INCREMENTAL_CLEANUP to use the legacy cleanup path instead.
2026-08-20 15:05:28 -07:00

836 lines
21 KiB
C++

/*
** Copyright (C) 2021-2025 Dirk-Jan C. Binnema <djcb@djcbsoftware.nl>
**
** This program is free software; you can redistribute it and/or modify it
** under the terms of the GNU General Public License as published by the
** Free Software Foundation; either version 3, or (at your option) any
** later version.
**
** This program is distributed in the hope that it will be useful,
** but WITHOUT ANY WARRANTY; without even the implied warranty of
** MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
** GNU General Public License for more details.
**
** You should have received a copy of the GNU General Public License
** along with this program; if not, write to the Free Software Foundation,
** Inc., 51 Franklin Street, Fifth Floor, Boston, MA 02110-1301, USA.
**
*/
#include "config.h"
#include "mu-store.hh"
#include <chrono>
#include <mutex>
#include <algorithm>
#include <array>
#include <cstdlib>
#include <stdexcept>
#include <unordered_map>
#include <atomic>
#include <type_traits>
#include <iostream>
#include <cstring>
#include "mu-maildir.hh"
#include "mu-query.hh"
#include "mu-xapian-db.hh"
#include "mu-scanner.hh"
#include "mu-store-labels.hh"
#include "utils/mu-error.hh"
#include "utils/mu-utils.hh"
#include <utils/mu-utils-file.hh>
using namespace Mu;
static_assert(std::is_same<Store::Id, Xapian::docid>::value, "wrong type for Store::Id");
// Properties
constexpr auto ExpectedSchemaVersion = MU_STORE_SCHEMA_VERSION;
static std::string
remove_slash(const std::string& str)
{
auto clean{str};
while (!clean.empty() && clean[clean.length() - 1] == '/')
clean.pop_back();
return clean;
}
struct Store::Private {
Private(const std::string& path, bool readonly):
xapian_db_{XapianDb(path, readonly ? XapianDb::Flavor::ReadOnly
: XapianDb::Flavor::Open)},
config_{xapian_db_},
contacts_cache_{config_},
labels_cache_{config_},
root_maildir_{remove_slash(config_.get<Config::Id::RootMaildir>())},
message_opts_{make_message_options(config_)}
{}
Private(const std::string& path, const std::string& root_maildir,
Option<const Config&> conf):
xapian_db_{XapianDb(path, XapianDb::Flavor::CreateOverwrite)},
config_{make_config(xapian_db_, root_maildir, conf)},
contacts_cache_{config_},
labels_cache_{config_},
root_maildir_{remove_slash(config_.get<Config::Id::RootMaildir>())},
message_opts_{make_message_options(config_)} {
// so tell xapian-db to update its internal cached values from
// config. In practice: batch-size.
xapian_db_.reinit();
}
~Private() try {
mu_debug("closing store @ {}", xapian_db_.path());
} catch (...) {
mu_critical("caught exception in store dtor");
}
Config make_config(XapianDb& xapian_db, const std::string& root_maildir,
Option<const Config&> conf) {
if (!g_path_is_absolute(root_maildir.c_str()))
throw Error{Error::Code::File,
"root maildir path is not absolute ({})",
root_maildir};
Config config{xapian_db};
if (conf)
config.import_configurable(*conf);
config.set<Config::Id::RootMaildir>(remove_slash(root_maildir));
config.set<Config::Id::SchemaVersion>(ExpectedSchemaVersion);
return config;
}
Message::Options make_message_options(const Config& conf) {
if (conf.get<Config::Id::SupportNgrams>())
return Message::Options::SupportNgrams;
else
return Message::Options::None;
}
Result<void> serialize();
Option<Message> find_message_unlocked(Store::Id docid) const;
using IdDocument = std::pair<Store::Id, Xapian::Document>;
Option<IdDocument> find_document_unlocked(const std::string& path) const;
Store::IdVec find_duplicates_unlocked(const Store& store,
const std::string& message_id) const;
using PathMessage = std::pair<std::string, Message>;
Result<PathMessage> move_message_unlocked(Message&& msg,
Option<const std::string&> target_mdir,
Option<Flags> new_flags,
MoveOptions opts);
bool remove_message_by_id_unlocked(Store::Id id);
XapianDb xapian_db_;
Config config_;
ContactsCache contacts_cache_;
LabelsCache labels_cache_;
std::unique_ptr<Indexer> indexer_;
const std::string root_maildir_;
const Message::Options message_opts_;
std::mutex lock_;
};
Result<void>
Store::Private::serialize()
{
if (xapian_db_.read_only())
return Err(Error::Code::Store, "store must be writable");
mu_debug("serialize data into store @ {}", xapian_db_.path());
contacts_cache_.serialize(); // does the 'dirty' check internally.
labels_cache_.serialize();
return Ok();
}
Option<Message>
Store::Private::find_message_unlocked(Store::Id docid) const
{
if (auto&& doc{xapian_db_.document(docid)}; !doc)
return Nothing;
else if (auto&& msg{Message::make_from_document(std::move(*doc))}; !msg)
return Nothing;
else
return Some(std::move(*msg));
}
Option<Store::Private::IdDocument>
Store::Private::find_document_unlocked(const std::string& path) const
{
constexpr auto path_field{field_from_id(Field::Id::Path)};
auto enq{xapian_db_.enquire()};
enq.set_query(Xapian::Query{path_field.xapian_term(path)});
if (auto mset{enq.get_mset(0, 1)}; mset.empty())
return Nothing; // message not found
else
return Some(IdDocument{*mset.begin(), mset.begin().get_document()});
}
Store::IdVec
Store::Private::find_duplicates_unlocked(const Store& store,
const std::string& message_id) const
{
if (message_id.empty() || message_id.size() > MaxTermLength) {
mu_warning("invalid message-id '{}'", message_id);
return {};
}
const auto expr{mu_format("{}:{}",
field_from_id(Field::Id::MessageId).shortcut,
message_id)};
if (auto&& res{store.run_query(expr)}; !res) {
mu_warning("error finding message-ids: {}", res.error().what());
return {};
} else {
Store::IdVec ids;
ids.reserve(res->size());
for (auto&& mi: *res)
ids.emplace_back(mi.doc_id());
return ids;
}
}
static void
reinit_export_labels(const Store& store)
{
// slightly hacky way to get the cache-path...
const auto cache_path{canonicalize_filename(
join_paths(store.path(), "..")) + "/"};
if (const auto res =
Mu::export_labels(store, "", cache_path); !res)
throw Mu::Error(Error::Code::CannotReinit,
"cannot re-init; "
"failed to export labels: {}",
res.error().what())
.add_hint("see mu-init(1) for details");
else {
mu_info("exported labels to: {}", *res);
mu_println("exported labels to: {}", *res);
}
}
Store::Store(const std::string& path, Store::Options opts)
: priv_{std::make_unique<Private>(
path, none_of(opts & Store::Options::Writable))}
{
if (none_of(opts & Store::Options::Writable) &&
any_of(opts & Store::Options::ReInit))
throw Mu::Error(Error::Code::InvalidArgument,
"Options::ReInit requires Options::Writable");
const auto s_version{config().get<Config::Id::SchemaVersion>()};
if (any_of(opts & Store::Options::ReInit)) {
/* export labels, if necessary (throws if there's an error) */
if (!label_map().empty())
reinit_export_labels(*this);
/* don't try to recover from version with an incompatible
* scheme */
if (s_version < 500)
throw Mu::Error(
Error::Code::CannotReinit,
"old schema ({}) is too old to re-initialize from",
s_version).add_hint("Invoke 'mu init' without '--reinit'; "
"see mu-init(1) for details");
const auto old_root_maildir{root_maildir()};
MemDb mem_db;
Config old_config(mem_db);
old_config.import_configurable(config());
this->priv_.reset();
/* and create a new one "in place" */
Store new_store(path, old_root_maildir, old_config);
this->priv_ = std::move(new_store.priv_);
}
/* otherwise, the schema version should match. */
if (s_version != ExpectedSchemaVersion)
throw Mu::Error(Error::Code::SchemaMismatch,
"expected schema-version {}, but got {}",
ExpectedSchemaVersion, s_version).
add_hint("Please (re)initialize with 'mu init'; see mu-init(1) for details");
}
Store::Store(const std::string& path,
const std::string& root_maildir,
Option<const Config&> conf):
priv_{std::make_unique<Private>(path, root_maildir, conf)}
{}
Store::Store(Store&& other)
{
priv_ = std::move(other.priv_);
priv_->indexer_.reset();
}
Store::~Store() = default;
Store::Statistics
Store::statistics() const
{
Statistics stats{};
stats.size = size();
stats.last_change = config().get<Config::Id::LastChange>();
stats.last_index = config().get<Config::Id::LastIndex>();
return stats;
}
const XapianDb&
Store::xapian_db() const
{
return priv_->xapian_db_;
}
XapianDb&
Store::xapian_db()
{
return priv_->xapian_db_;
}
const Config&
Store::config() const
{
return priv_->config_;
}
Config&
Store::config()
{
return priv_->config_;
}
const std::string&
Store::root_maildir() const {
return priv_->root_maildir_;
}
const ContactsCache&
Store::contacts_cache() const
{
return priv_->contacts_cache_;
}
Result<void>
Store::serialize()
{
std::lock_guard guard{priv_->lock_};
return priv_->serialize();
}
Indexer&
Store::indexer()
{
std::lock_guard guard{priv_->lock_};
if (xapian_db().read_only())
throw Error{Error::Code::Store, "no indexer for read-only store"};
else if (!priv_->indexer_)
priv_->indexer_ = std::make_unique<Indexer>(*this);
return *priv_->indexer_.get();
}
Result<Store::Id>
Store::add_message(Message& msg, bool is_new)
{
const auto mdir{maildir_from_path(msg.path(), root_maildir())};
if (!mdir)
return Err(mdir.error());
if (auto&& res = msg.set_maildir(mdir.value()); !res)
return Err(res.error());
// we shouldn't mix ngrams/non-ngrams messages.
if (any_of(msg.options() & Message::Options::SupportNgrams) !=
any_of(message_options() & Message::Options::SupportNgrams))
return Err(Error::Code::InvalidArgument, "incompatible message options");
/* add contacts from this message to cache; this cache
* also determines whether those contacts are _personal_, i.e. match
* our personal addresses.
*
* if a message has any personal contacts, mark it as personal; do
* this by updating the message flags.
*/
bool is_personal{};
priv_->contacts_cache_.add(msg.all_contacts(), is_personal);
if (is_personal)
msg.set_flags(msg.flags() | Flags::Personal);
std::lock_guard guard{priv_->lock_};
auto&& res = is_new ?
xapian_db().add_document(msg.document().xapian_document()) :
xapian_db().replace_document(msg.xapian_term(), msg.document().xapian_document());
if (!res)
return Err(res.error());
mu_debug("added {}{}message @ {}; docid = {}",
is_new ? "new " : "", is_personal ? "personal " : "", msg.path(), *res);
return res;
}
Result<Store::Id>
Store::add_message(const std::string& path, bool is_new)
{
if (auto msg{Message::make_from_path(path, priv_->message_opts_)}; !msg)
return Err(msg.error());
else
return add_message(msg.value(), is_new);
}
bool
Store::remove_message(const std::string& path)
{
const auto term{field_from_id(Field::Id::Path).xapian_term(path)};
std::lock_guard guard{priv_->lock_};
if (const auto& id_doc = priv_->find_document_unlocked(path); !id_doc)
return false;
else {
const Document doc{id_doc->second};
for (auto&& label : doc.string_vec_value(Field::Id::Labels))
priv_->labels_cache_.decrease(label);
xapian_db().delete_document(term);
mu_debug("deleted message @ {} from store", path);
return true;
}
}
void
Store::remove_messages(const std::vector<Store::Id>& ids)
{
std::lock_guard guard{priv_->lock_};
xapian_db().request_transaction();
for (auto&& id : ids)
priv_->remove_message_by_id_unlocked(id);
xapian_db().request_commit(true/*force*/);
}
bool
Store::Private::remove_message_by_id_unlocked(Store::Id id)
{
if (auto msg = labels_cache_.empty() ? Nothing : find_message_unlocked(id); msg)
for (auto&& label: msg->labels())
labels_cache_.decrease(label);
if (!xapian_db_.delete_document(id)) {
mu_warning("cannot find document {} for deletion", id);
return false;
}
return true;
}
size_t
Store::remove_messages_by_term(std::span<const std::string> terms,
std::function<void (size_t)> progress_fn)
{
std::unique_lock lock{priv_->lock_};
size_t nr_removed = 0;
std::vector<Xapian::Query> qvec;
std::vector<Store::Id> ids_to_remove;
xapian_db().request_transaction();
while (!terms.empty()) {
auto chunk = terms.subspan(0, std::min<size_t>(terms.size(), 1024));
terms = terms.subspan(chunk.size());
qvec.clear();
qvec.reserve(chunk.size());
std::ranges::copy(chunk, std::back_inserter(qvec));
auto enq = xapian_db().enquire();
enq.set_weighting_scheme(Xapian::BoolWeight()); // No score
enq.set_docid_order(Xapian::Enquire::ASCENDING);
enq.set_query(Xapian::Query{Xapian::Query::OP_OR, qvec.begin(), qvec.end()});
auto mset = enq.get_mset(0, xapian_db().size());
ids_to_remove.insert(ids_to_remove.end(), mset.begin(), mset.end());
}
// Sort the IDs to remove to make Xapian tree traversal easier
std::ranges::sort(ids_to_remove);
for (Id id : ids_to_remove) {
nr_removed += priv_->remove_message_by_id_unlocked(id);
if (nr_removed % 500 == 0)
progress_fn(nr_removed);
}
xapian_db().request_commit(true/*force*/);
return nr_removed;
}
Option<Message>
Store::find_message(Store::Id docid) const
{
std::lock_guard guard{priv_->lock_};
return priv_->find_message_unlocked(docid);
}
Option<Store::Id>
Store::find_message_id(const std::string& path) const
{
std::lock_guard guard{priv_->lock_};
if (const auto id_doc = priv_->find_document_unlocked(path); !id_doc)
return Nothing; // message not found
else
return Some(std::move(id_doc->first));
}
Store::IdMessageVec
Store::find_messages(IdVec ids) const
{
std::lock_guard guard{priv_->lock_};
IdMessageVec id_msgs;
for (auto&& id: ids) {
if (auto&& msg{priv_->find_message_unlocked(id)}; msg)
id_msgs.emplace_back(std::make_pair(id, std::move(*msg)));
}
return id_msgs;
}
/**
* Move a message in store and filesystem; with DryRun, only calculate the target name.
*
* Lock is assumed taken already
*
* @param id message id
* @param target_mdir target_mdir (or Nothing for current)
* @param new_flags new flags (or Nothing)
* @param opts move_options
*
* @return the Message after the moving, or an Error
*/
Result<Store::Private::PathMessage>
Store::Private::move_message_unlocked(Message&& msg,
Option<const std::string&> target_mdir,
Option<Flags> new_flags,
MoveOptions opts)
{
const auto old_path = msg.path();
const auto target_flags = new_flags.value_or(msg.flags());
const auto target_maildir = target_mdir.value_or(msg.maildir());
/* 1. first determine the file system path of the target */
const auto target_path =
maildir_determine_target(msg.path(), root_maildir_,
target_maildir, target_flags,
any_of(opts & MoveOptions::ChangeName));
if (!target_path)
return Err(target_path.error());
// in dry-run mode, we only determine the target-path
if (none_of(opts & MoveOptions::DryRun)) {
/* 2. let's move it */
if (const auto res = maildir_move_message(msg.path(), target_path.value()); !res)
return Err(res.error());
/* 3. file move worked, now update the message with the new info.*/
if (auto&& res = msg.update_after_move(
target_path.value(), target_maildir, target_flags); !res)
return Err(res.error());
/* 4. update message worked; re-store it */
if (auto&& res = xapian_db_.replace_document(
field_from_id(Field::Id::Path).xapian_term(old_path),
msg.document().xapian_document()); !res)
return Err(res.error());
}
/* 6. Profit! */
return Ok(PathMessage{std::move(*target_path), std::move(msg)});
}
Store::IdVec
Store::find_duplicates(const std::string& message_id) const
{
std::lock_guard guard{priv_->lock_};
return priv_->find_duplicates_unlocked(*this, message_id);
}
Result<Store::IdPathVec>
Store::move_message(Store::Id id,
Option<const std::string&> target_mdir,
Option<Flags> new_flags,
MoveOptions opts)
{
auto filter_dup_flags=[](Flags old_flags, Flags new_flags) -> Flags {
new_flags = flags_keep_unmutable(old_flags, new_flags, Flags::Draft);
new_flags = flags_keep_unmutable(old_flags, new_flags, Flags::Flagged);
new_flags = flags_keep_unmutable(old_flags, new_flags, Flags::Trashed);
return new_flags;
};
std::lock_guard guard{priv_->lock_};
auto msg{priv_->find_message_unlocked(id)};
if (!msg)
return Err(Error::Code::Store, "cannot find message <{}>", id);
const auto message_id{msg->message_id()};
auto res{priv_->move_message_unlocked(std::move(*msg), target_mdir, new_flags, opts)};
if (!res)
return Err(res.error());
IdPathVec id_paths{{id, res->first}};
if (none_of(opts & Store::MoveOptions::DupFlags) || message_id.empty() || !new_flags)
return Ok(std::move(id_paths));
/* handle the dup-flags case; i.e. apply (a subset of) the flags to
* all messages with the same message-id as well */
auto dups{priv_->find_duplicates_unlocked(*this, message_id)};
for (auto&& dupid: dups) {
if (dupid == id)
continue; // already
auto dup_msg{priv_->find_message_unlocked(dupid)};
if (!dup_msg)
continue; // no such message
/* For now, don't change Draft/Flagged/Trashed */
const auto dup_flags{filter_dup_flags(dup_msg->flags(), *new_flags)};
/* use the updated new_flags and MoveOptions without DupFlags (so we don't
* recurse) */
opts = opts & ~MoveOptions::DupFlags;
if (auto dup_res = priv_->move_message_unlocked(
std::move(*dup_msg), Nothing, dup_flags, opts); !dup_res)
mu_warning("failed to move dup: {}", dup_res.error().what());
else
id_paths.emplace_back(dupid, dup_res->first);
}
// sort the dup paths by name;
std::sort(id_paths.begin() + 1, id_paths.end(),
[](const auto& idp1, const auto& idp2) { return idp1.second < idp2.second; });
return Ok(std::move(id_paths));
}
Store::IdVec
Store::id_vec(const IdPathVec& ips)
{
IdVec idv;
for (auto&& ip: ips)
idv.emplace_back(ip.first);
return idv;
}
time_t
Store::dirstamp(const std::string& path) const
{
std::string ts;
{
std::unique_lock lock{priv_->lock_};
ts = xapian_db().metadata(path);
}
return ts.empty() ? 0 /*epoch*/ : ::strtoll(ts.c_str(), {}, 16);
}
void
Store::set_dirstamp(const std::string& path, time_t tstamp)
{
std::unique_lock lock{priv_->lock_};
xapian_db().set_metadata(path, mu_format("{:x}", tstamp));
}
bool
Store::contains_message(const std::string& path) const
{
std::unique_lock lock{priv_->lock_};
return xapian_db().term_exists(field_from_id(Field::Id::Path).xapian_term(path));
}
Result<Labels::DeltaLabelVec>
Store::update_labels(Message& msg, const Labels::DeltaLabelVec& labels_delta)
{
if (msg.docid() == Store::InvalidId)
return Err(Error::Code::Store, "invalid doc-id");
std::unique_lock lock{priv_->lock_};
// i.e. the set of effective labels. and the set up updates, the "diff"
auto updates{updated_labels(msg.labels(), labels_delta)};
if (updates.second.empty())
return Ok(std::move(updates.second)); // nothing to do
msg.set_labels(updates.first);
const auto res = xapian_db().replace_document(msg.xapian_term(),
msg.document().xapian_document());
if (!res)
return Err(res.error());
priv_->labels_cache_.update(updates.second);
return Ok(std::move(updates.second));
}
Result<Labels::DeltaLabelVec>
Store::clear_labels(Message& msg)
{
using namespace Labels;
DeltaLabelVec updates;
{
std::unique_lock lock{priv_->lock_};
for (auto& label: msg.labels())
updates.emplace_back(DeltaLabel{Delta::Remove, std::move(label)});
}
return update_labels(msg, updates);
}
Result<void>
Store::restore_label_map()
{
if (xapian_db().read_only())
throw Error{Error::Code::Store, "store is read-only"};
std::unique_lock lock{priv_->lock_};
return priv_->labels_cache_.restore(*this);
}
LabelsCache::Map
Store::label_map() const
{
std::unique_lock lock{priv_->lock_};
return priv_->labels_cache_.label_map();
}
std::size_t
Store::for_each_message_path(Store::ForEachMessageFunc msg_func) const
{
size_t n{};
xapian_try([&] {
std::lock_guard guard{priv_->lock_};
auto enq{xapian_db().enquire()};
enq.set_query(Xapian::Query::MatchAll);
enq.set_cutoff(0, 0);
Xapian::MSet matches(enq.get_mset(0, xapian_db().size()));
constexpr auto path_no{field_from_id(Field::Id::Path).value_no()};
for (auto&& it = matches.begin(); it != matches.end(); ++it, ++n)
if (!msg_func(*it, it.get_document().get_value(path_no)))
break;
});
return n;
}
std::size_t
Store::for_each_term(Field::Id field_id, Store::ForEachTermFunc func) const
{
return xapian_db().all_terms(field_from_id(field_id).xapian_term(), func);
}
std::mutex&
Store::lock() const
{
return priv_->lock_;
}
Result<QueryResults>
Store::run_query(const std::string& expr,
Field::Id sortfield_id,
QueryFlags flags, size_t maxnum) const
{
return Query{*this}.run(expr, sortfield_id, flags, maxnum);
}
size_t
Store::count_query(const std::string& expr) const
{
return xapian_try([&] {
std::lock_guard guard{priv_->lock_};
Query q{*this};
return q.count(expr); }, 0);
}
std::string
Store::parse_query(const std::string& expr, bool xapian) const
{
return xapian_try([&] {
std::lock_guard guard{priv_->lock_};
Query q{*this};
return q.parse(expr, xapian);
}, std::string{});
}
std::vector<std::string>
Store::maildirs() const
{
std::vector<std::string> mdirs;
const auto prefix_size{root_maildir().size()};
Scanner::Handler handler = [&](const std::string& path, auto&& _1, auto&& _2) {
auto md{path.substr(prefix_size)};
mdirs.emplace_back(md.empty() ? "/" : std::move(md));
return true;
};
Scanner scanner{root_maildir(), handler, Scanner::Mode::MaildirsOnly};
scanner.start();
std::sort(mdirs.begin(), mdirs.end());
return mdirs;
}
Message::Options
Store::message_options() const
{
return priv_->message_opts_;
}