From 872d48ebad7f135588973366b3114e222e5ddacb Mon Sep 17 00:00:00 2001 From: Eelco Dolstra Date: Wed, 30 Sep 2026 10:55:21 +0200 Subject: [PATCH 1/3] git: Don't require the git-hashing experimental feature in low-level helpers The helpers for parsing and dumping Git blobs and trees are not user-facing. The git-hashing feature is still checked where the Git content-address method is selected. Assisted-by: Claude Fable 5.1 --- src/libutil-tests/git.cc | 78 ++++++++++------------------- src/libutil/git.cc | 47 +++++------------ src/libutil/include/nix/util/git.hh | 31 +++--------- 3 files changed, 46 insertions(+), 110 deletions(-) diff --git a/src/libutil-tests/git.cc b/src/libutil-tests/git.cc index 81e34a065cac..f4fd4c0d9893 100644 --- a/src/libutil-tests/git.cc +++ b/src/libutil-tests/git.cc @@ -17,19 +17,6 @@ class GitTest : public CharacterizationTest { return unitTestData / std::string(testStem); } - - /** - * We set these in tests rather than the regular globals so we don't have - * to worry about race conditions if the tests run concurrently. - */ - ExperimentalFeatureSettings mockXpSettings; - -private: - - void SetUp() override - { - mockXpSettings.set("experimental-features", "git-hashing"); - } }; TEST(GitMode, gitMode_directory) @@ -75,8 +62,8 @@ TEST_F(GitTest, blob_read) StringSource in{encoded}; StringSink out; RegularFileSink out2{out}; - ASSERT_EQ(parseObjectType(in, mockXpSettings), ObjectType::Blob); - parseBlob(out2, CanonPath::root, in, BlobMode::Regular, mockXpSettings); + ASSERT_EQ(parseObjectType(in), ObjectType::Blob); + parseBlob(out2, CanonPath::root, in, BlobMode::Regular); auto expected = readFile(goldenMaster("hello-world.bin")); @@ -90,7 +77,7 @@ TEST_F(GitTest, blob_write) writeTest("hello-world-blob.bin", [&]() { auto decoded = readFile(goldenMaster("hello-world.bin")); StringSink s; - dumpBlobPrefix(decoded.size(), s, mockXpSettings); + dumpBlobPrefix(decoded.size(), s); s(decoded); return s.s; }); @@ -176,26 +163,20 @@ const static git::Tree treeSha256 = { }, }; -static auto mkTreeReadTest(HashAlgorithm hashAlgo, git::Tree tree, const ExperimentalFeatureSettings & mockXpSettings) +static auto mkTreeReadTest(HashAlgorithm hashAlgo, git::Tree tree) { using namespace git; - return [hashAlgo, tree, mockXpSettings](const auto & encoded) { + return [hashAlgo, tree](const auto & encoded) { StringSource in{encoded}; NullFileSystemObjectSink out; Tree got; - ASSERT_EQ(parseObjectType(in, mockXpSettings), ObjectType::Tree); - parseTree( - out, - CanonPath::root, - in, - hashAlgo, - [&](auto & name, auto entry) { - auto name2 = std::string{name.rel()}; - if (entry.mode == Mode::Directory) - name2 += '/'; - got.insert_or_assign(name2, std::move(entry)); - }, - mockXpSettings); + ASSERT_EQ(parseObjectType(in), ObjectType::Tree); + parseTree(out, CanonPath::root, in, hashAlgo, [&](auto & name, auto entry) { + auto name2 = std::string{name.rel()}; + if (entry.mode == Mode::Directory) + name2 += '/'; + got.insert_or_assign(name2, std::move(entry)); + }); ASSERT_EQ(got, tree); }; @@ -203,12 +184,12 @@ static auto mkTreeReadTest(HashAlgorithm hashAlgo, git::Tree tree, const Experim TEST_F(GitTest, tree_sha1_read) { - readTest("tree-sha1.bin", mkTreeReadTest(HashAlgorithm::SHA1, treeSha1, mockXpSettings)); + readTest("tree-sha1.bin", mkTreeReadTest(HashAlgorithm::SHA1, treeSha1)); } TEST_F(GitTest, tree_sha256_read) { - readTest("tree-sha256.bin", mkTreeReadTest(HashAlgorithm::SHA256, treeSha256, mockXpSettings)); + readTest("tree-sha256.bin", mkTreeReadTest(HashAlgorithm::SHA256, treeSha256)); } TEST_F(GitTest, tree_sha1_write) @@ -216,7 +197,7 @@ TEST_F(GitTest, tree_sha1_write) using namespace git; writeTest("tree-sha1.bin", [&]() { StringSink s; - dumpTree(treeSha1, s, mockXpSettings); + dumpTree(treeSha1, s); return s.s; }); } @@ -226,7 +207,7 @@ TEST_F(GitTest, tree_sha256_write) using namespace git; writeTest("tree-sha256.bin", [&]() { StringSink s; - dumpTree(treeSha256, s, mockXpSettings); + dumpTree(treeSha256, s); return s.s; }); } @@ -249,7 +230,7 @@ TEST_F(GitTest, both_roundrip) StringSink s; HashSink hashSink{hashAlgo}; TeeSink s2{s, hashSink}; - auto mode = dump(path, s2, dumpHook, defaultPathFilter, mockXpSettings); + auto mode = dump(path, s2, dumpHook, defaultPathFilter); auto hash = hashSink.finish().hash; cas.insert_or_assign(hash, std::move(s.s)); return TreeEntry{ @@ -267,22 +248,15 @@ TEST_F(GitTest, both_roundrip) std::function mkSinkHook; mkSinkHook = [&](auto prefix, auto & hash, auto blobMode) { StringSource in{cas[hash]}; - parse( - sinkFiles2, - prefix, - in, - blobMode, - hashAlgo, - [&](const CanonPath & name, const auto & entry) { - mkSinkHook( - prefix / name, - entry.hash, - // N.B. this cast would not be acceptable in real - // code, because it would make an assert reachable, - // but it should harmless in this test. - static_cast(entry.mode)); - }, - mockXpSettings); + parse(sinkFiles2, prefix, in, blobMode, hashAlgo, [&](const CanonPath & name, const auto & entry) { + mkSinkHook( + prefix / name, + entry.hash, + // N.B. this cast would not be acceptable in real + // code, because it would make an assert reachable, + // but it should harmless in this test. + static_cast(entry.mode)); + }); }; mkSinkHook(CanonPath::root, root.hash, BlobMode::Regular); diff --git a/src/libutil/git.cc b/src/libutil/git.cc index 96c6dd28791d..d29092d5770b 100644 --- a/src/libutil/git.cc +++ b/src/libutil/git.cc @@ -48,15 +48,8 @@ static std::string getString(Source & source, int n) return v; } -void parseBlob( - FileSystemObjectSink & sink, - const CanonPath & sinkPath, - Source & source, - BlobMode blobMode, - const ExperimentalFeatureSettings & xpSettings) +void parseBlob(FileSystemObjectSink & sink, const CanonPath & sinkPath, Source & source, BlobMode blobMode) { - xpSettings.require(Xp::GitHashing); - const unsigned long long size = std::stoi(getStringUntil(source, 0)); auto doRegularFile = [&](bool executable) { @@ -103,8 +96,7 @@ void parseTree( const CanonPath & sinkPath, Source & source, HashAlgorithm hashAlgo, - fun hook, - const ExperimentalFeatureSettings & xpSettings) + fun hook) { const unsigned long long size = std::stoi(getStringUntil(source, 0)); unsigned long long left = size; @@ -146,10 +138,8 @@ void parseTree( } } -ObjectType parseObjectType(Source & source, const ExperimentalFeatureSettings & xpSettings) +ObjectType parseObjectType(Source & source) { - xpSettings.require(Xp::GitHashing); - auto type = getString(source, 5); if (type == "blob ") { @@ -166,19 +156,16 @@ void parse( Source & source, BlobMode rootModeIfBlob, HashAlgorithm hashAlgo, - fun hook, - const ExperimentalFeatureSettings & xpSettings) + fun hook) { - xpSettings.require(Xp::GitHashing); - - auto type = parseObjectType(source, xpSettings); + auto type = parseObjectType(source); switch (type) { case ObjectType::Blob: - parseBlob(sink, sinkPath, source, rootModeIfBlob, xpSettings); + parseBlob(sink, sinkPath, source, rootModeIfBlob); break; case ObjectType::Tree: - parseTree(sink, sinkPath, source, hashAlgo, hook, xpSettings); + parseTree(sink, sinkPath, source, hashAlgo, hook); break; default: assert(false); @@ -228,17 +215,14 @@ void restore(FileSystemObjectSink & sink, Source & source, HashAlgorithm hashAlg }); } -void dumpBlobPrefix(uint64_t size, Sink & sink, const ExperimentalFeatureSettings & xpSettings) +void dumpBlobPrefix(uint64_t size, Sink & sink) { - xpSettings.require(Xp::GitHashing); auto s = fmt("blob %d\0"s, std::to_string(size)); sink(s); } -void dumpTree(const Tree & entries, Sink & sink, const ExperimentalFeatureSettings & xpSettings) +void dumpTree(const Tree & entries, Sink & sink) { - xpSettings.require(Xp::GitHashing); - std::string v1; for (auto & [name, entry] : entries) { @@ -260,18 +244,13 @@ void dumpTree(const Tree & entries, Sink & sink, const ExperimentalFeatureSettin sink(v1); } -Mode dump( - const SourcePath & path, - Sink & sink, - fun hook, - PathFilter & filter, - const ExperimentalFeatureSettings & xpSettings) +Mode dump(const SourcePath & path, Sink & sink, fun hook, PathFilter & filter) { auto st = path.lstat(); switch (st.type) { case SourceAccessor::tRegular: { - path.readFile(sink, [&](uint64_t size) { dumpBlobPrefix(size, sink, xpSettings); }); + path.readFile(sink, [&](uint64_t size) { dumpBlobPrefix(size, sink); }); return st.isExecutable ? Mode::Executable : Mode::Regular; } @@ -290,13 +269,13 @@ Mode dump( entries.insert_or_assign(std::move(name2), std::move(entry)); } - dumpTree(entries, sink, xpSettings); + dumpTree(entries, sink); return Mode::Directory; } case SourceAccessor::tSymlink: { auto target = path.readLink(); - dumpBlobPrefix(target.size(), sink, xpSettings); + dumpBlobPrefix(target.size(), sink); sink(target); return Mode::Symlink; } diff --git a/src/libutil/include/nix/util/git.hh b/src/libutil/include/nix/util/git.hh index 01f4c4b88762..a81f31e979af 100644 --- a/src/libutil/include/nix/util/git.hh +++ b/src/libutil/include/nix/util/git.hh @@ -72,8 +72,7 @@ using SinkHook = void(const CanonPath & name, TreeEntry entry); * * @throws if prefix not recognized */ -ObjectType -parseObjectType(Source & source, const ExperimentalFeatureSettings & xpSettings = experimentalFeatureSettings); +ObjectType parseObjectType(Source & source); /** * These 3 modes are represented by blob objects. @@ -87,12 +86,7 @@ enum struct BlobMode : RawMode { Symlink = static_cast(Mode::Symlink), }; -void parseBlob( - FileSystemObjectSink & sink, - const CanonPath & sinkPath, - Source & source, - BlobMode blobMode, - const ExperimentalFeatureSettings & xpSettings = experimentalFeatureSettings); +void parseBlob(FileSystemObjectSink & sink, const CanonPath & sinkPath, Source & source, BlobMode blobMode); /** * @param hashAlgo must be `HashAlgo::SHA1` or `HashAlgo::SHA256` for now. @@ -102,8 +96,7 @@ void parseTree( const CanonPath & sinkPath, Source & source, HashAlgorithm hashAlgo, - fun hook, - const ExperimentalFeatureSettings & xpSettings = experimentalFeatureSettings); + fun hook); /** * Helper putting the previous three `parse*` functions together. @@ -120,8 +113,7 @@ void parse( Source & source, BlobMode rootModeIfBlob, HashAlgorithm hashAlgo, - fun hook, - const ExperimentalFeatureSettings & xpSettings = experimentalFeatureSettings); + fun hook); /** * Assists with writing a `SinkHook` step (2). @@ -145,17 +137,13 @@ void restore(FileSystemObjectSink & sink, Source & source, HashAlgorithm hashAlg /** * Dumps a single file to a sink - * - * @param xpSettings for testing purposes */ -void dumpBlobPrefix( - uint64_t size, Sink & sink, const ExperimentalFeatureSettings & xpSettings = experimentalFeatureSettings); +void dumpBlobPrefix(uint64_t size, Sink & sink); /** * Dumps a representation of a git tree to a sink */ -void dumpTree( - const Tree & entries, Sink & sink, const ExperimentalFeatureSettings & xpSettings = experimentalFeatureSettings); +void dumpTree(const Tree & entries, Sink & sink); /** * Callback for processing a child with `dump` @@ -168,12 +156,7 @@ void dumpTree( */ using DumpHook = TreeEntry(const SourcePath & path); -Mode dump( - const SourcePath & path, - Sink & sink, - fun hook, - PathFilter & filter = defaultPathFilter, - const ExperimentalFeatureSettings & xpSettings = experimentalFeatureSettings); +Mode dump(const SourcePath & path, Sink & sink, fun hook, PathFilter & filter = defaultPathFilter); /** * Recursively dumps path, hashing as we go. From f436270fc0466824d7983e861973d11a150abd3b Mon Sep 17 00:00:00 2001 From: Eelco Dolstra Date: Wed, 30 Sep 2026 10:55:21 +0200 Subject: [PATCH 2/3] Replace the Git-based tarball cache by a SQLite database Tarballs are now unpacked into ~/.cache/nix/tarball-cache-v3.sqlite instead of a bare Git repository written via libgit2. Files and directories are still identified by their Git blob and tree hashes, so tree hashes and accessor fingerprints are unchanged. Blobs are compressed individually using zstd. This avoids the accumulation of packfiles (one per import), which made object lookups slow, and libgit2's concurrency issues. There is no garbage collection yet, and every file is stored as a single blob. Assisted-by: Claude Fable 5.1 --- src/libfetchers/git-utils.cc | 26 - src/libfetchers/github.cc | 7 +- .../include/nix/fetchers/fetch-settings.hh | 6 +- .../include/nix/fetchers/meson.build | 1 + .../include/nix/fetchers/tarball-cache.hh | 59 ++ src/libfetchers/meson.build | 4 + src/libfetchers/package.nix | 2 + src/libfetchers/tarball-cache.cc | 780 ++++++++++++++++++ src/libfetchers/tarball.cc | 8 +- src/libstore/include/nix/store/sqlite.hh | 7 + src/libstore/sqlite.cc | 7 + 11 files changed, 872 insertions(+), 35 deletions(-) create mode 100644 src/libfetchers/include/nix/fetchers/tarball-cache.hh create mode 100644 src/libfetchers/tarball-cache.cc diff --git a/src/libfetchers/git-utils.cc b/src/libfetchers/git-utils.cc index 8fb02c79fcec..226daf3021d1 100644 --- a/src/libfetchers/git-utils.cc +++ b/src/libfetchers/git-utils.cc @@ -1632,32 +1632,6 @@ std::vector> GitRepoImpl::getSubmodules return result; } -namespace fetchers { - -ref Settings::getTarballCache() const -{ - /* v1: Had either only loose objects or thin packfiles referring to loose objects - * v2: Must have only packfiles with no loose objects. Should get repacked periodically - * for optimal packfiles. - */ - static auto repoDir = std::filesystem::path(getCacheDir()) / "tarball-cache-v2"; - auto tarballCache(_tarballCache.lock()); - if (!*tarballCache) - *tarballCache = GitRepo::openRepo( - repoDir, - { - .create = true, - .bare = true, - .packfilesOnly = true, - /* Tarball unpacking is not expected to benefit from deltas much, - compared to how much CPU times it takes to find. */ - .dontFindDeltas = true, - }); - return ref(*tarballCache); -} - -} // namespace fetchers - static Sync> workdirInfoCache_; GitRepo::WorkdirInfo GitRepo::getCachedWorkdirInfo(const std::filesystem::path & path) diff --git a/src/libfetchers/github.cc b/src/libfetchers/github.cc index 930aeb050b5e..94ed500cfe84 100644 --- a/src/libfetchers/github.cc +++ b/src/libfetchers/github.cc @@ -9,6 +9,7 @@ #include "nix/fetchers/tarball.hh" #include "nix/util/tarfile.hh" #include "nix/fetchers/git-utils.hh" +#include "nix/fetchers/tarball-cache.hh" #include #include @@ -287,7 +288,7 @@ struct GitArchiveInputScheme : InputScheme if (auto lastModifiedAttrs = cache->lookup(lastModifiedKey)) { auto treeHash = getRevAttr(*treeHashAttrs, "treeHash"); auto lastModified = getIntAttr(*lastModifiedAttrs, "lastModified"); - if (settings.getTarballCache()->hasObject(treeHash)) + if (settings.getTarballCache()->hasTree(treeHash)) return { {std::move(input), TarballInfo{.treeHash = treeHash, .lastModified = (time_t) lastModified}}}; else @@ -308,7 +309,7 @@ struct GitArchiveInputScheme : InputScheme }); auto act = std::make_unique( - *logger, lvlInfo, actUnknown, fmt("unpacking '%s' into the Git cache", input.to_string())); + *logger, lvlInfo, actUnknown, fmt("unpacking '%s' into the tarball cache", input.to_string())); TarArchive archive{*source}; auto tarballCache = settings.getTarballCache(); @@ -351,7 +352,7 @@ struct GitArchiveInputScheme : InputScheme input.attrs.insert_or_assign("lastModified", uint64_t(tarballInfo.lastModified)); auto accessor = - settings.getTarballCache()->getAccessor(tarballInfo.treeHash, {}, "«" + input.to_string(true) + "»"); + settings.getTarballCache()->getAccessor(tarballInfo.treeHash, "«" + input.to_string(true) + "»"); if (!settings.trustTarballsFromGitForges) // FIXME: computing the NAR hash here is wasteful if diff --git a/src/libfetchers/include/nix/fetchers/fetch-settings.hh b/src/libfetchers/include/nix/fetchers/fetch-settings.hh index b68e2c0316f6..7413e77bc427 100644 --- a/src/libfetchers/include/nix/fetchers/fetch-settings.hh +++ b/src/libfetchers/include/nix/fetchers/fetch-settings.hh @@ -13,7 +13,6 @@ namespace nix { -struct GitRepo; struct SrcToStore; } // namespace nix @@ -21,6 +20,7 @@ struct SrcToStore; namespace nix::fetchers { struct Cache; +struct TarballCache; struct Settings : public Config { @@ -155,7 +155,7 @@ struct Settings : public Config ref getCache() const; - ref getTarballCache() const; + ref getTarballCache() const; /** * In-memory cache for calls to fetchToStore(); maps source paths to their store @@ -171,7 +171,7 @@ private: mutable Sync> _cache; - mutable Sync> _tarballCache; + mutable Sync> _tarballCache; }; } // namespace nix::fetchers diff --git a/src/libfetchers/include/nix/fetchers/meson.build b/src/libfetchers/include/nix/fetchers/meson.build index f3bb80942a28..a616e213e684 100644 --- a/src/libfetchers/include/nix/fetchers/meson.build +++ b/src/libfetchers/include/nix/fetchers/meson.build @@ -12,5 +12,6 @@ headers = files( 'input-cache.hh', 'provenance.hh', 'registry.hh', + 'tarball-cache.hh', 'tarball.hh', ) diff --git a/src/libfetchers/include/nix/fetchers/tarball-cache.hh b/src/libfetchers/include/nix/fetchers/tarball-cache.hh new file mode 100644 index 000000000000..13aabb0d33c8 --- /dev/null +++ b/src/libfetchers/include/nix/fetchers/tarball-cache.hh @@ -0,0 +1,59 @@ +#pragma once +///@file + +#include "nix/fetchers/git-utils.hh" +#include "nix/util/hash.hh" +#include "nix/util/ref.hh" +#include "nix/util/source-accessor.hh" + +#include + +namespace nix::fetchers { + +struct Settings; + +/** + * A content-addressed cache of unpacked tarballs, stored in a SQLite + * database. Files and directories are identified by their Git (SHA-1) + * blob and tree hashes. + */ +struct TarballCache +{ + virtual ~TarballCache() = default; + + /** + * Open (and create if necessary) the tarball cache stored in the + * SQLite database `dbPath`. + */ + static ref open(const std::filesystem::path & dbPath); + + /** + * Whether the cache contains the tree `treeHash`. If so, then + * everything reachable from that tree is present as well. + */ + virtual bool hasTree(const Hash & treeHash) = 0; + + virtual ref getAccessor(const Hash & treeHash, std::string displayPrefix) = 0; + + /** + * Return a sink that imports a file system tree into the cache. + * Calling `flush()` on the sink returns the hash of the root + * tree. + */ + virtual ref getFileSystemObjectSink() = 0; + + /** + * Given a Git tree hash, compute the hash of its NAR + * serialisation. This is memoised on-disk. + */ + virtual Hash treeHashToNarHash(const Settings & settings, const Hash & treeHash) = 0; + + /** + * If the specified tree is a directory with a single entry that + * is a directory, return the hash of that entry. Otherwise return + * the passed hash unchanged. + */ + virtual Hash dereferenceSingletonDirectory(const Hash & treeHash) = 0; +}; + +} // namespace nix::fetchers diff --git a/src/libfetchers/meson.build b/src/libfetchers/meson.build index 134fd496a7a1..c175cd6a8a4b 100644 --- a/src/libfetchers/meson.build +++ b/src/libfetchers/meson.build @@ -31,6 +31,9 @@ deps_public += nlohmann_json libgit2 = dependency('libgit2', version : '>= 1.9') deps_private += libgit2 +zstd = dependency('libzstd', version : '>= 1.4.0') +deps_private += zstd + subdir('nix-meson-build-support/common') sources = files( @@ -51,6 +54,7 @@ sources = files( 'path.cc', 'provenance.cc', 'registry.cc', + 'tarball-cache.cc', 'tarball.cc', ) diff --git a/src/libfetchers/package.nix b/src/libfetchers/package.nix index 1a30ac293018..f52b7e9668e0 100644 --- a/src/libfetchers/package.nix +++ b/src/libfetchers/package.nix @@ -6,6 +6,7 @@ nix-store, nlohmann_json, libgit2, + zstd, # Configuration Options @@ -35,6 +36,7 @@ mkMesonLibrary (finalAttrs: { buildInputs = [ libgit2 + zstd ]; propagatedBuildInputs = [ diff --git a/src/libfetchers/tarball-cache.cc b/src/libfetchers/tarball-cache.cc new file mode 100644 index 000000000000..229277d36b5e --- /dev/null +++ b/src/libfetchers/tarball-cache.cc @@ -0,0 +1,780 @@ +#include "nix/fetchers/tarball-cache.hh" +#include "nix/fetchers/cache.hh" +#include "nix/fetchers/fetch-settings.hh" +#include "nix/store/globals.hh" +#include "nix/store/sqlite.hh" +#include "nix/util/file-system.hh" +#include "nix/util/finally.hh" +#include "nix/util/git.hh" +#include "nix/util/pool.hh" +#include "nix/util/signals.hh" +#include "nix/util/sync.hh" +#include "nix/util/thread-pool.hh" +#include "nix/util/users.hh" + +#include +#include + +#include + +#include +#include +#include +#include + +namespace nix::fetchers { + +namespace { + +const char * schema = R"sql( + +create table if not exists Blobs ( + oid blob primary key not null, + size integer not null, + compression integer not null, + data blob not null +); + +create table if not exists Trees ( + oid blob primary key not null +) without rowid; + +create table if not exists TreeEntries ( + tree blob not null, + name text not null, + mode integer not null, + child blob not null, + primary key (tree, name) +) without rowid; + +)sql"; + +/** + * Values of the `Blobs.compression` column. + */ +enum struct Compression : int64_t { + None = 0, + Zstd = 1, +}; + +Hash toOid(std::string_view s) +{ + Hash oid(HashAlgorithm::SHA1); + if (s.size() != oid.hashSize) + throw Error("tarball cache contains an object ID of %d bytes", s.size()); + memcpy(oid.hash, s.data(), oid.hashSize); + return oid; +} + +Hash hashBlob(std::string_view contents) +{ + HashSink sink(HashAlgorithm::SHA1); + git::dumpBlobPrefix(contents.size(), sink); + sink(contents); + return sink.finish().hash; +} + +/* Note: we call libzstd directly rather than using `compress()` / `makeDecompressionSink()` from + `nix/util/compression.hh`. The latter set up a new zstd context (or for decompression, a libarchive reader) for every + call, which is expensive when (de)compressing a large number of small blobs: it made importing and reading Nixpkgs + about 60% slower. Here we reuse one zstd context per thread instead. + + The zstd contexts are per-thread, but they're only used for the duration of a single (de)compression call, so + they're never shared between fibers running on the same thread. */ + +/** + * Compress `contents` using zstd. Return `std::nullopt` if that doesn't make it smaller. + */ +std::optional compress(std::string_view contents) +{ + if (contents.empty()) + return std::nullopt; + + thread_local std::unique_ptr cctx{ZSTD_createCCtx(), ZSTD_freeCCtx}; + if (!cctx) + throw Error("unable to create a zstd compression context"); + + std::string res; + res.resize(ZSTD_compressBound(contents.size())); + auto n = + ZSTD_compressCCtx(cctx.get(), res.data(), res.size(), contents.data(), contents.size(), ZSTD_CLEVEL_DEFAULT); + if (ZSTD_isError(n)) + throw Error("zstd compression failed: %s", ZSTD_getErrorName(n)); + if (n >= contents.size()) + return std::nullopt; + res.resize(n); + return res; +} + +std::string decompress(std::string_view data, uint64_t size) +{ + thread_local std::unique_ptr dctx{ZSTD_createDCtx(), ZSTD_freeDCtx}; + if (!dctx) + throw Error("unable to create a zstd decompression context"); + + std::string res; + res.resize(size); + auto n = ZSTD_decompressDCtx(dctx.get(), res.data(), res.size(), data.data(), data.size()); + if (ZSTD_isError(n)) + throw Error("zstd decompression failed: %s", ZSTD_getErrorName(n)); + if (n != size) + throw Error("blob in tarball cache has size %d, expected %d", n, size); + return res; +} + +struct Entry +{ + git::Mode mode; + Hash oid; +}; + +using Dir = std::map; + +struct Connection +{ + SQLite db; + SQLiteStmt hasTree, queryEntries, hasBlob, queryBlob, insertBlob, insertTree, insertEntry; + + Connection(const std::filesystem::path & dbPath) + { + db = SQLite(dbPath, {.useWAL = settings.useSQLiteWAL}); + db.exec("pragma synchronous = normal"); + + hasTree.create(db, "select 1 from Trees where oid = ?"); + queryEntries.create(db, "select name, mode, child from TreeEntries where tree = ?"); + hasBlob.create(db, "select 1 from Blobs where oid = ?"); + queryBlob.create(db, "select size, compression, data from Blobs where oid = ?"); + insertBlob.create(db, "insert or ignore into Blobs(oid, size, compression, data) values (?, ?, ?, ?)"); + insertTree.create(db, "insert or ignore into Trees(oid) values (?)"); + insertEntry.create(db, "insert or ignore into TreeEntries(tree, name, mode, child) values (?, ?, ?, ?)"); + } +}; + +/** + * A blob that is ready to be written to the database. + */ +struct PendingBlob +{ + Hash oid; + uint64_t size; + Compression compression; + std::string data; +}; + +struct TarballCacheImpl : TarballCache, std::enable_shared_from_this +{ + std::filesystem::path dbPath; + + /** + * SQLite connections. Every thread that accesses the database + * takes its own connection from this pool. + */ + Pool pool; + + /** + * Mutex to serialize write transactions within this process, since SQLite only allows one writer at a time anyway. + * (Writers in other processes are handled by SQLite's busy handler.) + */ + std::mutex writeMutex; + + TarballCacheImpl(std::filesystem::path _dbPath) + : dbPath(std::move(_dbPath)) + , pool(std::numeric_limits::max(), [this]() { return make_ref(dbPath); }) + { + createDirs(dbPath.parent_path()); + + SQLite db(dbPath, {.useWAL = settings.useSQLiteWAL}); + if (settings.useSQLiteWAL) + db.exec("pragma main.journal_mode = wal"); + db.exec(schema); + } + + bool hasTree(const Hash & treeHash) override + { + auto conn(pool.get()); + return conn->hasTree.use().apply(treeHash.hash, treeHash.hashSize).next(); + } + + ref readTree(const Hash & treeHash) + { + auto dir = make_ref(); + auto conn(pool.get()); + auto stmt(conn->queryEntries.use().apply(treeHash.hash, treeHash.hashSize)); + while (stmt.next()) { + auto rawMode = (git::RawMode) stmt.getInt(1); + auto mode = git::decodeMode(rawMode); + if (!mode) + throw Error("tarball cache contains an entry with unknown mode %o", rawMode); + dir->emplace(stmt.getStr(0), Entry{.mode = *mode, .oid = toOid(stmt.getBlob(2))}); + } + return dir; + } + + void readBlob(const Hash & oid, Sink & sink, fun sizeCallback) + { + auto conn(pool.get()); + auto stmt(conn->queryBlob.use().apply(oid.hash, oid.hashSize)); + if (!stmt.next()) + throw Error("blob '%s' is missing from the tarball cache", oid.gitRev()); + + uint64_t size = stmt.getInt(0); + auto compression = (Compression) stmt.getInt(1); + auto data = stmt.getBlob(2); + + switch (compression) { + case Compression::None: + sizeCallback(size); + sink(data); + break; + case Compression::Zstd: { + auto contents = decompress(data, size); + sizeCallback(size); + sink(contents); + break; + } + default: + throw Error( + "blob '%s' in the tarball cache has unknown compression type %d", oid.gitRev(), (int64_t) compression); + } + } + + void writeBlobs(const std::vector & blobs) + { + if (blobs.empty()) + return; + std::lock_guard lock(writeMutex); + auto conn(pool.get()); + retrySQLite([&]() { + SQLiteTxn txn(conn->db); + for (auto & blob : blobs) + conn->insertBlob.use() + .apply(blob.oid.hash, blob.oid.hashSize) + .apply((int64_t) blob.size) + .apply((int64_t) blob.compression) + .apply((const unsigned char *) blob.data.data(), blob.data.size()) + .exec(); + txn.commit(); + }); + } + + ref getAccessor(const Hash & treeHash, std::string displayPrefix) override; + + ref getFileSystemObjectSink() override; + + Hash treeHashToNarHash(const Settings & settings, const Hash & treeHash) override + { + auto accessor = getAccessor(treeHash, ""); + + Cache::Key cacheKey{"treeHashToNarHash", {{"treeHash", treeHash.gitRev()}}}; + + if (auto res = settings.getCache()->lookup(cacheKey)) + return Hash::parseAny(getStrAttr(*res, "narHash"), HashAlgorithm::SHA256); + + auto narHash = accessor->hashPath(CanonPath::root); + + settings.getCache()->upsert(cacheKey, Attrs({{"narHash", narHash.to_string(HashFormat::SRI, true)}})); + + return narHash; + } + + Hash dereferenceSingletonDirectory(const Hash & treeHash) override + { + auto dir = readTree(treeHash); + if (dir->size() == 1 && dir->begin()->second.mode == git::Mode::Directory) + return dir->begin()->second.oid; + return treeHash; + } +}; + +struct TarballCacheAccessor : SourceAccessor +{ + ref cache; + + Hash root; + + /** + * Cache of directory listings. A null value denotes a path that is not a directory. + */ + SharedSync>> dirCache; + + TarballCacheAccessor(ref cache, const Hash & root) + : cache(cache) + , root(root) + { + if (!cache->hasTree(root)) + throw Error("tree '%s' is missing from the tarball cache", root.gitRev()); + fingerprint = GitAccessorOptions{}.makeFingerprint(root); + } + + void anchor() override {} + + /** + * Return the contents of the directory `path`, or null if `path` doesn't exist or is not a directory. + */ + std::shared_ptr getDir(const CanonPath & path) + { + { + auto dirCache_(dirCache.readLock()); + if (auto i = dirCache_->find(path); i != dirCache_->end()) + return i->second; + } + + std::shared_ptr dir; + if (auto entry = lookup(path); entry && entry->mode == git::Mode::Directory) + dir = cache->readTree(entry->oid).get_ptr(); + + dirCache.lock()->emplace(path, dir); + + return dir; + } + + std::optional lookup(const CanonPath & path) + { + if (path.isRoot()) + return Entry{.mode = git::Mode::Directory, .oid = root}; + + auto dir = getDir(*path.parent()); + if (!dir) + return std::nullopt; + + auto i = dir->find(std::string(*path.baseName())); + if (i == dir->end()) + return std::nullopt; + + return i->second; + } + + Entry need(const CanonPath & path) + { + auto entry = lookup(path); + if (!entry) + throw FileNotFound("path '%s' does not exist", showPath(path)); + return *entry; + } + + void readFile(const CanonPath & path, Sink & sink, fun sizeCallback) override + { + auto entry = need(path); + if (entry.mode != git::Mode::Regular && entry.mode != git::Mode::Executable) + throw Error("'%s' is not a regular file", showPath(path)); + cache->readBlob(entry.oid, sink, sizeCallback); + } + + bool pathExists(const CanonPath & path) override + { + return (bool) lookup(path); + } + + std::optional maybeLstat(const CanonPath & path) override + { + auto entry = lookup(path); + if (!entry) + return std::nullopt; + + switch (entry->mode) { + case git::Mode::Directory: + return Stat{.type = tDirectory}; + case git::Mode::Regular: + return Stat{.type = tRegular}; + case git::Mode::Executable: + return Stat{.type = tRegular, .isExecutable = true}; + case git::Mode::Symlink: + return Stat{.type = tSymlink}; + } + unreachable(); + } + + DirEntries readDirectory(const CanonPath & path) override + { + auto entry = need(path); + if (entry.mode != git::Mode::Directory) + throw Error("'%s' is not a directory", showPath(path)); + + auto dir = getDir(path); + assert(dir); + + DirEntries res; + for (auto & [name, entry] : *dir) + res.emplace( + name, + entry.mode == git::Mode::Directory ? tDirectory + : entry.mode == git::Mode::Symlink ? tSymlink + : tRegular); + return res; + } + + std::string readLink(const CanonPath & path) override + { + auto entry = need(path); + if (entry.mode != git::Mode::Symlink) + throw Error("'%s' is not a symlink", showPath(path)); + StringSink sink; + cache->readBlob(entry.oid, sink, [&](uint64_t size) { sink.s.reserve(size); }); + return std::move(sink.s); + } +}; + +ref TarballCacheImpl::getAccessor(const Hash & treeHash, std::string displayPrefix) +{ + auto accessor = make_ref(ref(shared_from_this()), treeHash); + accessor->setPathDisplay(std::move(displayPrefix)); + return accessor; +} + +struct TarballCacheSink : GitFileSystemObjectSink +{ + ref cache; + + unsigned int concurrency = std::min(std::thread::hardware_concurrency(), 10U); + + ThreadPool workers{concurrency}; + + /** Total file contents in flight. */ + std::atomic totalBufSize{0}; + + /** If the file contents in flight exceed this threshold, files are processed synchronously. */ + static constexpr size_t maxBufSize = 64 * 1024 * 1024; + + /** The amount of (compressed) blob data to write to the database in a single transaction. */ + static constexpr size_t maxBatchSize = 16 * 1024 * 1024; + + struct Child; + + /// A directory to be written as a tree. + struct Directory + { + std::map children; + std::optional oid; + + Child & lookup(const CanonPath & path) + { + assert(!path.isRoot()); + auto parent = path.parent(); + auto cur = this; + for (auto & name : *parent) { + auto i = cur->children.find(std::string(name)); + if (i == cur->children.end()) + throw Error("path '%s' does not exist", path); + auto dir = std::get_if(&i->second.file); + if (!dir) + throw Error("path '%s' has a non-directory parent", path); + cur = dir; + } + + auto i = cur->children.find(std::string(*path.baseName())); + if (i == cur->children.end()) + throw Error("path '%s' does not exist", path); + return i->second; + } + }; + + size_t nextId = 0; // for Child.id + + struct Child + { + git::Mode mode; + std::variant file; + + /// Sequential numbering of the file in the tarball. This is + /// used to make sure we only import the latest version of a + /// path. + size_t id{0}; + }; + + struct State + { + Directory root; + }; + + Sync _state; + + struct Pending + { + std::vector blobs; + size_t size = 0; + }; + + /** Blobs waiting to be written to the database. */ + Sync _pending; + + /** The blobs that have already been seen during this import. */ + Sync> _seen; + + TarballCacheSink(ref cache) + : cache(cache) + { + } + + ~TarballCacheSink() + { + // Make sure the worker threads are destroyed before any state + // they're referring to. + workers.shutdown(); + } + + void addNode(State & state, const CanonPath & path, Child && child) + { + assert(!path.isRoot()); + auto parent = path.parent(); + + Directory * cur = &state.root; + + for (auto & i : *parent) { + auto child = std::get_if( + &cur->children.emplace(std::string(i), Child{git::Mode::Directory, {Directory()}}).first->second.file); + if (!child) + throw Error("path '%s' has a non-directory parent", path); + cur = child; + } + + std::string name(*path.baseName()); + + if (auto prev = cur->children.find(name); prev == cur->children.end() || prev->second.id < child.id) + cur->children.insert_or_assign(name, std::move(child)); + } + + /** + * Add a blob to the cache, unless it's already there. Returns its Git hash. + */ + Hash addBlob(std::string_view contents) + { + auto oid = hashBlob(contents); + + if (!_seen.lock()->insert(oid).second) + return oid; + + { + auto conn(cache->pool.get()); + if (conn->hasBlob.use().apply(oid.hash, oid.hashSize).next()) + return oid; + } + + PendingBlob blob{.oid = oid, .size = contents.size()}; + if (auto compressed = compress(contents)) { + blob.compression = Compression::Zstd; + blob.data = std::move(*compressed); + } else { + blob.compression = Compression::None; + blob.data = contents; + } + + std::vector batch; + + { + auto pending(_pending.lock()); + pending->size += blob.data.size(); + pending->blobs.push_back(std::move(blob)); + if (pending->size < maxBatchSize) + return oid; + batch = std::move(pending->blobs); + pending->blobs.clear(); + pending->size = 0; + } + + cache->writeBlobs(batch); + + return oid; + } + + /** + * Asynchronously add a blob to the cache and register it as `path` in the tree. + */ + void addBlobNode(const CanonPath & path, git::Mode mode, std::string contents_, size_t id) + { + auto contents = std::make_shared(std::move(contents_)); + + totalBufSize += contents->size(); + + auto work = [this, path, mode, contents, id]() { + Finally releaseBuf([&]() { totalBufSize -= contents->size(); }); + auto oid = addBlob(*contents); + addNode(*_state.lock(), path, Child{mode, oid, id}); + }; + + /* To avoid unbounded memory usage, process the file synchronously if there is too much data in flight. */ + if (totalBufSize > maxBufSize) + work(); + else + workers.enqueue(std::move(work)); + } + + void createRegularFile(const CanonPath & path, fun func) override + { + checkInterrupt(); + + struct CRF : CreateRegularFileSink + { + std::string contents; + bool executable = false; + + void operator()(std::string_view data) override + { + contents.append(data); + } + + void isExecutable() override + { + executable = true; + } + + void preallocateContents(uint64_t size) override + { + contents.reserve(size); + } + }; + + CRF crf; + + func(crf); + + addBlobNode( + path, crf.executable ? git::Mode::Executable : git::Mode::Regular, std::move(crf.contents), nextId++); + } + + void createDirectory(const CanonPath & path) override + { + if (path.isRoot()) + return; + auto state(_state.lock()); + addNode(*state, path, {git::Mode::Directory, Directory()}); + } + + void createSymlink(const CanonPath & path, const std::string & target) override + { + addBlobNode(path, git::Mode::Symlink, target, 0); + } + + std::map hardLinks; + + void createHardlink(const CanonPath & path, const CanonPath & target) override + { + hardLinks.insert_or_assign(path, target); + } + + /** + * Compute the Git tree hashes of `dir` and its subdirectories. + */ + static Hash hashTree(Directory & dir) + { + git::Tree tree; + + for (auto & [name, child] : dir.children) { + if (auto subdir = std::get_if(&child.file)) + // `git::Tree` expects directory names to have a trailing slash. + tree.emplace(name + "/", git::TreeEntry{.mode = git::Mode::Directory, .hash = hashTree(*subdir)}); + else + tree.emplace(name, git::TreeEntry{.mode = child.mode, .hash = std::get(child.file)}); + } + + HashSink sink(HashAlgorithm::SHA1); + git::dumpTree(tree, sink); + dir.oid = sink.finish().hash; + return *dir.oid; + } + + Hash flush() override + { + workers.process(); + + // Write the remaining blobs. + { + auto pending(_pending.lock()); + cache->writeBlobs(pending->blobs); + pending->blobs.clear(); + pending->size = 0; + } + + auto state(_state.lock()); + + /* Create hard links. */ + for (auto & [path, target] : hardLinks) { + if (target.isRoot()) + continue; + try { + auto child = state->root.lookup(target); + auto oid = std::get_if(&child.file); + if (!oid) + throw Error("cannot create a hard link to a directory"); + addNode(*state, path, {child.mode, *oid}); + } catch (Error & e) { + e.addTrace(nullptr, "while creating a hard link from '%s' to '%s'", path, target); + throw; + } + } + + auto rootOid = hashTree(state->root); + + /* Figure out which trees are not in the database yet. Note that if a tree is already in the database, then so + are all its children. */ + std::vector missing; + { + auto conn(cache->pool.get()); + boost::unordered_flat_set visited; + + [&](this const auto & visit, const Directory & dir) -> void { + checkInterrupt(); + + auto & oid = dir.oid.value(); + if (!visited.insert(oid).second) + return; + if (conn->hasTree.use().apply(oid.hash, oid.hashSize).next()) + return; + + for (auto & child : dir.children) + if (auto subdir = std::get_if(&child.second.file)) + visit(*subdir); + + missing.push_back(&dir); + }(state->root); + } + + /* Write the missing trees. This must be done after writing the blobs, since the existence of a tree implies + that all its children exist. */ + if (!missing.empty()) { + std::lock_guard lock(cache->writeMutex); + auto conn(cache->pool.get()); + retrySQLite([&]() { + SQLiteTxn txn(conn->db); + for (auto dir : missing) { + auto & oid = dir->oid.value(); + for (auto & [name, child] : dir->children) { + auto subdir = std::get_if(&child.file); + auto & childOid = subdir ? subdir->oid.value() : std::get(child.file); + conn->insertEntry.use() + .apply(oid.hash, oid.hashSize) + .apply(name) + .apply((int64_t) child.mode) + .apply(childOid.hash, childOid.hashSize) + .exec(); + } + conn->insertTree.use().apply(oid.hash, oid.hashSize).exec(); + } + txn.commit(); + }); + } + + return rootOid; + } +}; + +ref TarballCacheImpl::getFileSystemObjectSink() +{ + return make_ref(ref(shared_from_this())); +} + +} // namespace + +ref TarballCache::open(const std::filesystem::path & dbPath) +{ + return make_ref(dbPath); +} + +ref Settings::getTarballCache() const +{ + auto tarballCache(_tarballCache.lock()); + if (!*tarballCache) + *tarballCache = TarballCache::open(getCacheDir() / "tarball-cache-v3.sqlite").get_ptr(); + return ref(*tarballCache); +} + +} // namespace nix::fetchers diff --git a/src/libfetchers/tarball.cc b/src/libfetchers/tarball.cc index 5586229e56cf..cc98be8ae3bf 100644 --- a/src/libfetchers/tarball.cc +++ b/src/libfetchers/tarball.cc @@ -8,6 +8,7 @@ #include "nix/util/types.hh" #include "nix/store/store-api.hh" #include "nix/fetchers/git-utils.hh" +#include "nix/fetchers/tarball-cache.hh" #include "nix/fetchers/fetch-settings.hh" #include "nix/fetchers/provenance.hh" @@ -152,11 +153,11 @@ static std::optional downloadTarball_( .treeHash = treeHash, .lastModified = (time_t) getIntAttr(infoAttrs, "lastModified"), .immutableUrl = maybeGetStrAttr(infoAttrs, "immutableUrl"), - .accessor = settings.getTarballCache()->getAccessor(treeHash, {}, displayPrefix), + .accessor = settings.getTarballCache()->getAccessor(treeHash, displayPrefix), }; }; - if (cached && !settings.getTarballCache()->hasObject(getRevAttr(cached->value, "treeHash"))) + if (cached && !settings.getTarballCache()->hasTree(getRevAttr(cached->value, "treeHash"))) cached.reset(); if (cached && !cached->expired) @@ -177,7 +178,8 @@ static std::optional downloadTarball_( // TODO: fall back to cached value if download fails. - auto act = std::make_unique(*logger, lvlInfo, actUnknown, fmt("unpacking '%s' into the Git cache", url)); + auto act = + std::make_unique(*logger, lvlInfo, actUnknown, fmt("unpacking '%s' into the tarball cache", url)); AutoDelete cleanupTemp; diff --git a/src/libstore/include/nix/store/sqlite.hh b/src/libstore/include/nix/store/sqlite.hh index e416c44bd04b..f4aad0f330f0 100644 --- a/src/libstore/include/nix/store/sqlite.hh +++ b/src/libstore/include/nix/store/sqlite.hh @@ -157,6 +157,13 @@ struct SQLiteStmt bool next(); std::string getStr(int col); + + /** + * Return the contents of a blob (or text) column. The result is + * only valid until the next call to `next()` or until this + * `Use` is destroyed. + */ + std::string_view getBlob(int col); int64_t getInt(int col); bool isNull(int col); }; diff --git a/src/libstore/sqlite.cc b/src/libstore/sqlite.cc index 57ea4639dea0..efa6f9afcca8 100644 --- a/src/libstore/sqlite.cc +++ b/src/libstore/sqlite.cc @@ -270,6 +270,13 @@ std::string SQLiteStmt::Use::getStr(int col) return s; } +std::string_view SQLiteStmt::Use::getBlob(int col) +{ + auto data = (const char *) sqlite3_column_blob(stmt, col); + /* Note: `sqlite3_column_blob()` returns null for zero-length blobs. */ + return data ? std::string_view(data, sqlite3_column_bytes(stmt, col)) : std::string_view(); +} + int64_t SQLiteStmt::Use::getInt(int col) { // FIXME: detect nulls? From 81c78da4209bd509c0e1dcad708b7935b985f99a Mon Sep 17 00:00:00 2001 From: Eelco Dolstra Date: Tue, 6 Oct 2026 11:57:10 +0200 Subject: [PATCH 3/3] TarballCache: Store blobs in chunks Blobs are now stored as a sequence of chunks of at most 1 MiB, each compressed independently. Files larger than a chunk are spilled to a temporary file during import. So memory use no longer scales with file size, and blobs larger than SQLite's blob size limit (1 GB) can be stored. Importing a tarball with a 300 MB file now takes 184 MB of RSS, compared to 1.2 GB with the Git-based cache. The cost for Nixpkgs is about 5% on cold imports and 4 MB of database size, since every blob now has two rows. Assisted-by: Claude Fable 5.1 --- src/libfetchers/tarball-cache.cc | 338 ++++++++++++++++++++++++------- 1 file changed, 263 insertions(+), 75 deletions(-) diff --git a/src/libfetchers/tarball-cache.cc b/src/libfetchers/tarball-cache.cc index 229277d36b5e..13431823152e 100644 --- a/src/libfetchers/tarball-cache.cc +++ b/src/libfetchers/tarball-cache.cc @@ -3,10 +3,12 @@ #include "nix/fetchers/fetch-settings.hh" #include "nix/store/globals.hh" #include "nix/store/sqlite.hh" +#include "nix/util/file-descriptor.hh" #include "nix/util/file-system.hh" #include "nix/util/finally.hh" #include "nix/util/git.hh" #include "nix/util/pool.hh" +#include "nix/util/serialise.hh" #include "nix/util/signals.hh" #include "nix/util/sync.hh" #include "nix/util/thread-pool.hh" @@ -22,17 +24,28 @@ #include #include +#include + namespace nix::fetchers { namespace { const char * schema = R"sql( +-- A blob is stored as a sequence of chunks of at most `chunkSize` bytes, each compressed independently. The `Blobs` +-- row is written after all of its chunks, so its existence implies that the blob is complete. create table if not exists Blobs ( - oid blob primary key not null, + oid blob primary key not null, + size integer not null +); + +create table if not exists BlobChunks ( + oid blob not null, + seq integer not null, size integer not null, compression integer not null, - data blob not null + data blob not null, + primary key (oid, seq) ); create table if not exists Trees ( @@ -50,13 +63,19 @@ create table if not exists TreeEntries ( )sql"; /** - * Values of the `Blobs.compression` column. + * Values of the `BlobChunks.compression` column. */ enum struct Compression : int64_t { None = 0, Zstd = 1, }; +/** + * The maximum (uncompressed) size of a blob chunk. Files larger than this are spilled to a temporary file during + * import, so this bounds the memory used per file. + */ +constexpr size_t chunkSize = 1024 * 1024; + Hash toOid(std::string_view s) { Hash oid(HashAlgorithm::SHA1); @@ -133,7 +152,7 @@ using Dir = std::map; struct Connection { SQLite db; - SQLiteStmt hasTree, queryEntries, hasBlob, queryBlob, insertBlob, insertTree, insertEntry; + SQLiteStmt hasTree, queryEntries, hasBlob, queryBlob, insertBlob, insertChunk, insertTree, insertEntry; Connection(const std::filesystem::path & dbPath) { @@ -143,24 +162,81 @@ struct Connection hasTree.create(db, "select 1 from Trees where oid = ?"); queryEntries.create(db, "select name, mode, child from TreeEntries where tree = ?"); hasBlob.create(db, "select 1 from Blobs where oid = ?"); - queryBlob.create(db, "select size, compression, data from Blobs where oid = ?"); - insertBlob.create(db, "insert or ignore into Blobs(oid, size, compression, data) values (?, ?, ?, ?)"); + queryBlob.create( + db, + "select b.size, c.seq, c.size, c.compression, c.data from Blobs b join BlobChunks c on c.oid = b.oid " + "where b.oid = ? order by c.seq"); + insertBlob.create(db, "insert or ignore into Blobs(oid, size) values (?, ?)"); + insertChunk.create( + db, "insert or ignore into BlobChunks(oid, seq, size, compression, data) values (?, ?, ?, ?, ?)"); insertTree.create(db, "insert or ignore into Trees(oid) values (?)"); insertEntry.create(db, "insert or ignore into TreeEntries(tree, name, mode, child) values (?, ?, ?, ?)"); } }; /** - * A blob that is ready to be written to the database. + * A blob chunk that is ready to be written to the database. */ -struct PendingBlob +struct PendingChunk { Hash oid; + uint64_t seq; uint64_t size; Compression compression; std::string data; }; +/** + * A blob whose chunks have all been queued for writing. + */ +struct PendingBlob +{ + Hash oid; + uint64_t size; +}; + +/** + * Rows waiting to be written to the database in a single transaction. Chunks are written before blobs, so a blob + * must be added after its chunks. + */ +struct Pending +{ + std::vector chunks; + std::vector blobs; + + /** Total size of `chunks[*].data`. */ + size_t size = 0; + + void addChunk(PendingChunk && chunk) + { + size += chunk.data.size(); + chunks.push_back(std::move(chunk)); + } + + void clear() + { + chunks.clear(); + blobs.clear(); + size = 0; + } +}; + +/** + * Compress `contents` into a chunk of a blob. + */ +PendingChunk makeChunk(const Hash & oid, uint64_t seq, std::string_view contents) +{ + PendingChunk chunk{.oid = oid, .seq = seq, .size = contents.size()}; + if (auto compressed = compress(contents)) { + chunk.compression = Compression::Zstd; + chunk.data = std::move(*compressed); + } else { + chunk.compression = Compression::None; + chunk.data = contents; + } + return chunk; +} + struct TarballCacheImpl : TarballCache, std::enable_shared_from_this { std::filesystem::path dbPath; @@ -214,45 +290,71 @@ struct TarballCacheImpl : TarballCache, std::enable_shared_from_thisqueryBlob.use().apply(oid.hash, oid.hashSize)); - if (!stmt.next()) - throw Error("blob '%s' is missing from the tarball cache", oid.gitRev()); - uint64_t size = stmt.getInt(0); - auto compression = (Compression) stmt.getInt(1); - auto data = stmt.getBlob(2); - - switch (compression) { - case Compression::None: - sizeCallback(size); - sink(data); - break; - case Compression::Zstd: { - auto contents = decompress(data, size); - sizeCallback(size); - sink(contents); - break; - } - default: - throw Error( - "blob '%s' in the tarball cache has unknown compression type %d", oid.gitRev(), (int64_t) compression); + auto corrupt = [&](std::string_view what) { + throw Error("blob '%s' in the tarball cache is corrupt: %s", oid.gitRev(), what); + }; + + uint64_t size = 0, seq = 0, bytesRead = 0; + + while (stmt.next()) { + if (seq == 0) { + size = stmt.getInt(0); + sizeCallback(size); + } + + if ((uint64_t) stmt.getInt(1) != seq) + corrupt(fmt("chunk %d is missing", seq)); + + uint64_t chunkSize = stmt.getInt(2); + auto compression = (Compression) stmt.getInt(3); + auto data = stmt.getBlob(4); + + switch (compression) { + case Compression::None: + if (data.size() != chunkSize) + corrupt(fmt("chunk %d has size %d, expected %d", seq, data.size(), chunkSize)); + sink(data); + break; + case Compression::Zstd: + sink(decompress(data, chunkSize)); + break; + default: + corrupt(fmt("unknown compression type %d", (int64_t) compression)); + } + + bytesRead += chunkSize; + seq++; } + + if (seq == 0) + throw Error("blob '%s' is missing from the tarball cache", oid.gitRev()); + + if (bytesRead != size) + corrupt(fmt("expected %d bytes, got %d", size, bytesRead)); } - void writeBlobs(const std::vector & blobs) + /** + * Write the chunks and blobs in `pending` to the database in a single transaction. + */ + void writePending(const Pending & pending) { - if (blobs.empty()) + if (pending.chunks.empty() && pending.blobs.empty()) return; std::lock_guard lock(writeMutex); auto conn(pool.get()); retrySQLite([&]() { SQLiteTxn txn(conn->db); - for (auto & blob : blobs) - conn->insertBlob.use() - .apply(blob.oid.hash, blob.oid.hashSize) - .apply((int64_t) blob.size) - .apply((int64_t) blob.compression) - .apply((const unsigned char *) blob.data.data(), blob.data.size()) + for (auto & chunk : pending.chunks) + conn->insertChunk.use() + .apply(chunk.oid.hash, chunk.oid.hashSize) + .apply((int64_t) chunk.seq) + .apply((int64_t) chunk.size) + .apply((int64_t) chunk.compression) + .apply((const unsigned char *) chunk.data.data(), chunk.data.size()) .exec(); + for (auto & blob : pending.blobs) + conn->insertBlob.use().apply(blob.oid.hash, blob.oid.hashSize).apply((int64_t) blob.size).exec(); txn.commit(); }); } @@ -488,13 +590,7 @@ struct TarballCacheSink : GitFileSystemObjectSink Sync _state; - struct Pending - { - std::vector blobs; - size_t size = 0; - }; - - /** Blobs waiting to be written to the database. */ + /** Rows waiting to be written to the database. */ Sync _pending; /** The blobs that have already been seen during this import. */ @@ -533,45 +629,105 @@ struct TarballCacheSink : GitFileSystemObjectSink cur->children.insert_or_assign(name, std::move(child)); } + /** + * Queue a chunk for writing, and write the queued rows if there are enough of them. + */ + void queueChunk(PendingChunk && chunk) + { + Pending batch; + + { + auto pending(_pending.lock()); + pending->addChunk(std::move(chunk)); + if (pending->size < maxBatchSize) + return; + batch = std::move(*pending); + pending->clear(); + } + + cache->writePending(batch); + } + + /** + * Return whether the blob `oid` still needs to be written, i.e. it's not in the cache and hasn't been queued + * during this import. + */ + bool needBlob(const Hash & oid) + { + if (!_seen.lock()->insert(oid).second) + return false; + + auto conn(cache->pool.get()); + return !conn->hasBlob.use().apply(oid.hash, oid.hashSize).next(); + } + /** * Add a blob to the cache, unless it's already there. Returns its Git hash. */ Hash addBlob(std::string_view contents) { + assert(contents.size() <= chunkSize); + auto oid = hashBlob(contents); - if (!_seen.lock()->insert(oid).second) + if (!needBlob(oid)) return oid; - { - auto conn(cache->pool.get()); - if (conn->hasBlob.use().apply(oid.hash, oid.hashSize).next()) - return oid; - } - - PendingBlob blob{.oid = oid, .size = contents.size()}; - if (auto compressed = compress(contents)) { - blob.compression = Compression::Zstd; - blob.data = std::move(*compressed); - } else { - blob.compression = Compression::None; - blob.data = contents; - } + auto chunk = makeChunk(oid, 0, contents); - std::vector batch; + Pending batch; { auto pending(_pending.lock()); - pending->size += blob.data.size(); - pending->blobs.push_back(std::move(blob)); + pending->addChunk(std::move(chunk)); + pending->blobs.push_back({.oid = oid, .size = contents.size()}); if (pending->size < maxBatchSize) return oid; - batch = std::move(pending->blobs); - pending->blobs.clear(); - pending->size = 0; + batch = std::move(*pending); + pending->clear(); } - cache->writeBlobs(batch); + cache->writePending(batch); + + return oid; + } + + /** + * Add a blob that was spilled to a temporary file to the cache, unless it's already there. Returns its Git hash. + */ + Hash addBlobFromFile(Descriptor fd, uint64_t size) + { + auto rewind = [&]() { + if (lseek(fd, 0, SEEK_SET) == -1) + throw SysError("seeking in temporary file"); + }; + + /* First pass: compute the Git hash. */ + auto oid = [&]() { + rewind(); + HashSink hashSink(HashAlgorithm::SHA1); + git::dumpBlobPrefix(size, hashSink); + FdSource source(fd); + source.drainInto(hashSink); + return hashSink.finish().hash; + }(); + + if (!needBlob(oid)) + return oid; + + /* Second pass: compress and queue the chunks. */ + rewind(); + FdSource source(fd); + std::string buf; + uint64_t seq = 0, left = size; + do { + buf.resize(std::min(chunkSize, left)); + source(buf.data(), buf.size()); + left -= buf.size(); + queueChunk(makeChunk(oid, seq++, buf)); + } while (left); + + _pending.lock()->blobs.push_back({.oid = oid, .size = size}); return oid; } @@ -602,14 +758,36 @@ struct TarballCacheSink : GitFileSystemObjectSink { checkInterrupt(); + /* The contents are buffered in memory up to `chunkSize` bytes. Beyond that, they're spilled to a temporary + file, which is processed after the whole file has been received. */ struct CRF : CreateRegularFileSink { std::string contents; bool executable = false; + AutoCloseFD tempFd; + AutoDelete tempDel; + uint64_t size = 0; + void operator()(std::string_view data) override { - contents.append(data); + size += data.size(); + + if (!tempFd) { + if (contents.size() + data.size() <= chunkSize) { + contents.append(data); + return; + } + + auto [fd, path] = createTempFile("nix-tarball-cache"); + tempFd = std::move(fd); + tempDel = AutoDelete(path, /*recursive=*/false); + writeFull(tempFd.get(), contents); + contents.clear(); + contents.shrink_to_fit(); + } + + writeFull(tempFd.get(), data); } void isExecutable() override @@ -619,16 +797,27 @@ struct TarballCacheSink : GitFileSystemObjectSink void preallocateContents(uint64_t size) override { - contents.reserve(size); + if (size <= chunkSize) + contents.reserve(size); } }; - CRF crf; + auto crf = std::make_shared(); - func(crf); + func(*crf); - addBlobNode( - path, crf.executable ? git::Mode::Executable : git::Mode::Regular, std::move(crf.contents), nextId++); + auto mode = crf->executable ? git::Mode::Executable : git::Mode::Regular; + auto id = nextId++; + + if (!crf->tempFd) { + addBlobNode(path, mode, std::move(crf->contents), id); + return; + } + + workers.enqueue([this, path, mode, crf, id]() { + auto oid = addBlobFromFile(crf->tempFd.get(), crf->size); + addNode(*_state.lock(), path, Child{mode, oid, id}); + }); } void createDirectory(const CanonPath & path) override @@ -676,12 +865,11 @@ struct TarballCacheSink : GitFileSystemObjectSink { workers.process(); - // Write the remaining blobs. + // Write the remaining chunks and blobs. { auto pending(_pending.lock()); - cache->writeBlobs(pending->blobs); - pending->blobs.clear(); - pending->size = 0; + cache->writePending(*pending); + pending->clear(); } auto state(_state.lock());