Skip to content

Instantly share code, notes, and snippets.

@Mic92
Last active September 30, 2026 08:08
Show Gist options
  • Select an option

  • Save Mic92/5a4c5f6f6988355c250218ef74d51a18 to your computer and use it in GitHub Desktop.

Select an option

Save Mic92/5a4c5f6f6988355c250218ef74d51a18 to your computer and use it in GitHub Desktop.
SeaweedFS filer patch: answer cache misses from a remote key index (niks3/nix-grpc-store binary cache)

SeaweedFS filer patch plus the NixOS deployment of the gateway, for a niks3 binary cache on an S3 bucket.

Problem

The bucket is mounted into the SeaweedFS filer as a remote mount, and niks3 talks S3 to the gateway instead of the provider. Every cache miss under the mount then makes the filer HEAD the object in the bucket. Most misses are names that do not exist at all — clients probing narinfos they might already have, paths GC collected — and providers rate limit those requests, so a busy cache returns errors instead of misses. As all writes go through the filer, that HEAD can only confirm the miss.

Patch

seaweedfs-remote-key-index.patch lists each remote mount once at start-up, keeps its keys in memory, and answers a miss for a name outside that set with not-found instead of asking the remote. Upstream behaviour is kept until a listing has loaded, or if it failed. remote_key_index_dir also writes the listing to disk, so a restart gates lookups from the first request with the last snapshot instead of waiting for a fresh LIST. Tests in weed/filer/filer_lazy_remote_test.go.

WEED_FILER_OPTIONS_REMOTE_KEY_INDEX=true
WEED_FILER_OPTIONS_REMOTE_KEY_INDEX_DIR=/var/lib/seaweedfs/remote-key-index

Objects written to the bucket behind the filer's back stay hidden until a restart or remote.meta.sync.

Architecture

flowchart LR
  subgraph builders["remote builders"]
    runner["CI job"] --> daemon["build node"]
  end
  subgraph gateway["gateway host"]
    nginx["nginx"] --> niks3["niks3: presign, sign, GC"]
    weed["weed mini: master, volume, filer, S3 gateway"] --- ki["remote key index"]
    sync["filer.remote.sync"] --> weed
  end
  b2[("S3 bucket")]
  cdn["CDN"]
  client["nix substitutes"]
  daemon -->|build output over mTLS| nginx
  niks3 -->|presigned PUT| nginx
  nginx -->|S3| weed
  weed -->|read-through, write-back| b2
  ki -.->|initial LIST| b2
  b2 --> cdn --> client
  nginx -->|read host for fresh paths| client
Loading

The patched path is the one every miss takes: nginx → filer → key index → bucket only for names the listing knows. Reads normally go through the CDN and never reach this host; the read host exists so a path is readable the moment it is uploaded, before the CDN has synced.

Deploy

setup.nix is a NixOS module for the gateway host:

  • systemd.services.seaweedfs — weed mini (master, volume, filer, S3 gateway) bound to localhost, TLS on the API host's ACME certificate, the key index enabled, and its -s3.config identity file composed by ExecStartPre so secret keys stay out of the store
  • systemd.services.seaweedfs-mount — remote.configure plus an idempotent remote.mount of the bucket under /buckets/<bucket>
  • systemd.services.seaweedfs-remote-sync — filer.remote.sync, the write-back half
  • nginx: the API vhost that proxies /<bucket>/ to the gateway with the Host header preserved (it is part of the presigned signature), and the read host, which re-adds Content-Encoding: zstd that objects lose when they arrive through the remote mount
  • the niks3 settings that point its S3 client at the gateway, in comments
services.seaweedfsCache = {
  enable = true;
  bucket = "nix-cache";
  apiHost = "cache-api.example.com";
  readHost = "cache-local.example.com";
  publicCacheUrl = "https://cache.example.com";
  region = "eu-central-003";
  endpoint = "s3.eu-central-003.backblazeb2.com";
  secretsDir = "/run/secrets"; # gateway-access-key, gateway-secret-key, b2-access-key, b2-secret-key
};

Build

nix build git+https://gist.github.com/Mic92/5a4c5f6f6988355c250218ef74d51a18.git
./result/bin/weed --version

The gist is a git repo, so the flake works directly; flake.lock pins nixpkgs nixos-unstable, and .#default / .#seaweedfs exist for x86_64/aarch64 Linux and macOS — both are nixpkgs' seaweedfs with nothing changed but the extra patch. No new Go dependencies, so vendorHash is unchanged; if a bump breaks the context, the patchPhase failure is the only breakage to fix. Not submitted upstream yet.

{
"nodes": {
"nixpkgs": {
"locked": {
"lastModified": 1790578696,
"narHash": "sha256-ZoxIApko70jCdbH3l20HWXOBaT2HZd87orzd2yJ9dVE=",
"owner": "NixOS",
"repo": "nixpkgs",
"rev": "7a0f122f5090cf4c2ade2a13a0e229d4e19ba71f",
"type": "github"
},
"original": {
"owner": "NixOS",
"ref": "nixos-unstable",
"repo": "nixpkgs",
"type": "github"
}
},
"root": {
"inputs": {
"nixpkgs": "nixpkgs"
}
}
},
"root": "root",
"version": 7
}
{
description = "SeaweedFS with the filer remote key index patch applied";
inputs.nixpkgs.url = "github:NixOS/nixpkgs/nixos-unstable";
outputs =
{ nixpkgs, ... }:
let
systems = [
"x86_64-linux"
"aarch64-linux"
"aarch64-darwin"
"x86_64-darwin"
];
forAllSystems = f: nixpkgs.lib.genAttrs systems f;
pkgsFor = system: import nixpkgs { inherit system; };
in
{
# The gateway deployment: three systemd units plus two nginx vhosts, and the
# niks3 side spelled out in comments. It builds the patched `weed` itself with
# `pkgs.callPackage ./seaweedfs.nix`, so importing it is all a host needs:
#
# nixosConfigurations.host = nixpkgs.lib.nixosSystem {
# modules = [ inputs.seaweedfs-patches.nixosModules.default myCacheConfig ];
# };
nixosModules.default = ./setup.nix;
packages = forAllSystems (
system:
let
seaweedfs = (pkgsFor system).callPackage ./seaweedfs.nix { };
in
{
inherit seaweedfs;
default = seaweedfs;
}
);
devShells = forAllSystems (
system:
let
pkgs = pkgsFor system;
in
{
default = pkgs.mkShell {
packages = [
pkgs.callPackage
./seaweedfs.nix
{ }
pkgs.go
];
};
}
);
};
}
From 3e9633eaf0fbc52d93b40948e8c289b88d1b1a32 Mon Sep 17 00:00:00 2001
From: Mic92 <joerg@thalheim.io>
Date: Fri, 25 Sep 2026 16:19:15 +0200
Subject: [PATCH] filer: skip the remote stat for names missing from a remote
listing
MIME-Version: 1.0
Content-Type: text/plain; charset=UTF-8
Content-Transfer-Encoding: 8bit
A local miss under a remote mount makes the filer stat the object in
the remote bucket. Most misses are for objects that do not exist, and
each one costs a request that some providers rate limit.
With filer.options.remote_key_index enabled, list every mount once in
the background after the mounts load. Once a mount's listing is stored,
a miss for a name outside it returns not-found without a remote
request. Before the listing finishes, or if it fails, lookups behave as
before. Objects written to the remote afterwards stay hidden until a
restart or remote.meta.sync.
filer.options.remote_key_index_dir additionally keeps each listing on
disk. A restart then gates lookups from the start with the last
snapshot, and the fresh listing replaces it once it completes.
Signed-off-by: Jörg Thalheim <joerg@thalheim.io>
---
weed/filer/filer.go | 22 ++--
weed/filer/filer_lazy_remote.go | 6 ++
weed/filer/filer_lazy_remote_test.go | 140 +++++++++++++++++++++++-
weed/filer/filer_on_meta_event.go | 1 +
weed/filer/filer_remote_key_index.go | 156 +++++++++++++++++++++++++++
weed/server/filer_server.go | 2 +
6 files changed, 320 insertions(+), 7 deletions(-)
create mode 100644 weed/filer/filer_remote_key_index.go
diff --git a/weed/filer/filer.go b/weed/filer/filer.go
index ef9b18f..45d8a1d 100644
--- a/weed/filer/filer.go
+++ b/weed/filer/filer.go
@@ -7,6 +7,7 @@ import (
"os"
"sort"
"strings"
+ "sync"
"time"
"github.com/seaweedfs/seaweedfs/weed/remote_storage"
@@ -67,12 +68,21 @@ type Filer struct {
lazyListGroup singleflight.Group
Dlm *lock_manager.DistributedLockManager
MaxFilenameLength uint32
- deletionQuit chan struct{}
- DeletionRetryQueue *DeletionRetryQueue
- EmptyFolderCleaner *empty_folder_cleanup.EmptyFolderCleaner
- EmptyFolderCleanupDelay time.Duration
- persistedLogCache *persistedLogCache
- metaLogInflight metaLogInflight
+ // EnableRemoteKeyIndex lists each remote mount once and skips the remote
+ // stat for names missing from that listing. Objects written to the remote
+ // behind the filer stay hidden until a restart or remote.meta.sync.
+ EnableRemoteKeyIndex bool
+ // RemoteKeyIndexDir keeps each listing across restarts, so lookups are
+ // gated from the start instead of after the first listing.
+ RemoteKeyIndexDir string
+ remoteKeys sync.Map // mount dir -> *remoteKeySet
+ remoteKeysLoading sync.Map // mount dir -> struct{}
+ deletionQuit chan struct{}
+ DeletionRetryQueue *DeletionRetryQueue
+ EmptyFolderCleaner *empty_folder_cleanup.EmptyFolderCleaner
+ EmptyFolderCleanupDelay time.Duration
+ persistedLogCache *persistedLogCache
+ metaLogInflight metaLogInflight
}
func NewFiler(masters pb.ServerDiscovery, grpcDialOption grpc.DialOption, filerHost pb.ServerAddress, filerGroup string, collection string, replication string, dataCenter string, maxFilenameLength uint32, notifyFn func()) *Filer {
diff --git a/weed/filer/filer_lazy_remote.go b/weed/filer/filer_lazy_remote.go
index 1b6983b..7d332b5 100644
--- a/weed/filer/filer_lazy_remote.go
+++ b/weed/filer/filer_lazy_remote.go
@@ -48,6 +48,12 @@ func (f *Filer) maybeLazyFetchFromRemote(ctx context.Context, p util.FullPath) (
return nil, nil
}
+ if keys, ok := f.remoteKeys.Load(string(mountDir)); ok {
+ if _, known := keys.(*remoteKeySet).keys[p]; !known {
+ return nil, nil
+ }
+ }
+
relPath := strings.TrimPrefix(string(p), string(mountDir))
if relPath != "" && !strings.HasPrefix(relPath, "/") {
relPath = "/" + relPath
diff --git a/weed/filer/filer_lazy_remote_test.go b/weed/filer/filer_lazy_remote_test.go
index f9b3442..cc00037 100644
--- a/weed/filer/filer_lazy_remote_test.go
+++ b/weed/filer/filer_lazy_remote_test.go
@@ -201,6 +201,7 @@ type stubRemoteClient struct {
deleteCalls []*remote_pb.RemoteStorageLocation
removeCalls []*remote_pb.RemoteStorageLocation
+ traverseFn func(loc *remote_pb.RemoteStorageLocation, visitFn remote_storage.VisitFunc) error
listDirFn func(loc *remote_pb.RemoteStorageLocation, visitFn remote_storage.VisitFunc) error
listDirCalls int
}
@@ -208,7 +209,10 @@ type stubRemoteClient struct {
func (c *stubRemoteClient) StatFile(*remote_pb.RemoteStorageLocation) (*filer_pb.RemoteEntry, error) {
return c.statResult, c.statErr
}
-func (c *stubRemoteClient) Traverse(*remote_pb.RemoteStorageLocation, remote_storage.VisitFunc) error {
+func (c *stubRemoteClient) Traverse(loc *remote_pb.RemoteStorageLocation, visitFn remote_storage.VisitFunc) error {
+ if c.traverseFn != nil {
+ return c.traverseFn(loc, visitFn)
+ }
return nil
}
func (c *stubRemoteClient) ReadFile(*remote_pb.RemoteStorageLocation, int64, int64) ([]byte, error) {
@@ -1300,3 +1304,137 @@ func TestMaybeLazyListFromRemote_ContextGuardPreventsRecursion(t *testing.T) {
f.maybeLazyListFromRemote(fetchCtx, util.FullPath("/buckets/mybucket"))
assert.Equal(t, 0, stub.listDirCalls)
}
+
+func newKeyIndexFiler(t *testing.T, storageType string, keys []string) (*Filer, *countingRemoteClient, func()) {
+ t.Helper()
+ stub := &countingRemoteClient{
+ stubRemoteClient: stubRemoteClient{
+ statResult: &filer_pb.RemoteEntry{RemoteMtime: 1, RemoteSize: 1},
+ traverseFn: func(_ *remote_pb.RemoteStorageLocation, visit remote_storage.VisitFunc) error {
+ for _, k := range keys {
+ dir, name := util.FullPath("/" + k).DirAndName()
+ if err := visit(dir, name, false, &filer_pb.RemoteEntry{}); err != nil {
+ return err
+ }
+ }
+ return nil
+ },
+ },
+ }
+ restore := registerStubMaker(t, storageType, stub)
+ conf := &remote_pb.RemoteConf{Name: storageType, Type: storageType}
+ rs := NewFilerRemoteStorage()
+ rs.storageNameToConf[conf.Name] = conf
+ rs.mapDirectoryToRemoteStorage("/buckets/mybucket", &remote_pb.RemoteStorageLocation{
+ Name: storageType,
+ Bucket: "mybucket",
+ Path: "/",
+ })
+ f := newTestFiler(t, newStubFilerStore(), rs)
+ f.EnableRemoteKeyIndex = true
+ return f, stub, restore
+}
+
+func TestRemoteKeyIndex_UnknownNameSkipsRemoteStat(t *testing.T) {
+ f, stub, restore := newKeyIndexFiler(t, "stub_keyidx_unknown", []string{"known.narinfo", "log/a"})
+ defer restore()
+ require.NoError(t, f.buildRemoteKeyIndexes())
+
+ entry, err := f.maybeLazyFetchFromRemote(context.Background(), "/buckets/mybucket/other.narinfo")
+ require.NoError(t, err)
+ assert.Nil(t, entry)
+ assert.Equal(t, 0, stub.statCalls, "a name absent from the listing must not reach the remote")
+}
+
+func TestRemoteKeyIndex_KnownNameIsFetched(t *testing.T) {
+ f, stub, restore := newKeyIndexFiler(t, "stub_keyidx_known", []string{"known.narinfo", "log/a"})
+ defer restore()
+ require.NoError(t, f.buildRemoteKeyIndexes())
+
+ for _, p := range []util.FullPath{"/buckets/mybucket/known.narinfo", "/buckets/mybucket/log/a"} {
+ entry, err := f.maybeLazyFetchFromRemote(context.Background(), p)
+ require.NoError(t, err)
+ assert.NotNil(t, entry, p)
+ }
+ assert.Equal(t, 2, stub.statCalls)
+}
+
+func TestRemoteKeyIndex_BeforeListingLooksUpRemote(t *testing.T) {
+ f, stub, restore := newKeyIndexFiler(t, "stub_keyidx_early", nil)
+ defer restore()
+
+ entry, err := f.maybeLazyFetchFromRemote(context.Background(), "/buckets/mybucket/x.narinfo")
+ require.NoError(t, err)
+ assert.NotNil(t, entry)
+ assert.Equal(t, 1, stub.statCalls)
+}
+
+func TestRemoteKeyIndex_FailedListingKeepsLookingUpRemote(t *testing.T) {
+ f, stub, restore := newKeyIndexFiler(t, "stub_keyidx_fail", nil)
+ defer restore()
+ stub.traverseFn = func(*remote_pb.RemoteStorageLocation, remote_storage.VisitFunc) error {
+ return errors.New("list failed")
+ }
+ require.Error(t, f.buildRemoteKeyIndexes())
+
+ entry, err := f.maybeLazyFetchFromRemote(context.Background(), "/buckets/mybucket/x.narinfo")
+ require.NoError(t, err)
+ assert.NotNil(t, entry)
+}
+
+func TestRemoteKeyIndex_DisabledByDefault(t *testing.T) {
+ f, stub, restore := newKeyIndexFiler(t, "stub_keyidx_off", []string{"known.narinfo"})
+ defer restore()
+ f.EnableRemoteKeyIndex = false
+ require.NoError(t, f.buildRemoteKeyIndexes())
+
+ entry, err := f.maybeLazyFetchFromRemote(context.Background(), "/buckets/mybucket/other.narinfo")
+ require.NoError(t, err)
+ assert.NotNil(t, entry)
+ assert.Equal(t, 1, stub.statCalls)
+}
+
+func TestRemoteKeyIndex_SnapshotGatesLookupsBeforeListing(t *testing.T) {
+ dir := t.TempDir()
+ f, _, restore := newKeyIndexFiler(t, "stub_keyidx_snap1", []string{"known.narinfo"})
+ defer restore()
+ f.RemoteKeyIndexDir = dir
+ require.NoError(t, f.buildRemoteKeyIndexes())
+
+ f2, stub2, restore2 := newKeyIndexFiler(t, "stub_keyidx_snap1b", nil)
+ defer restore2()
+ f2.RemoteKeyIndexDir = dir
+ stub2.traverseFn = func(*remote_pb.RemoteStorageLocation, remote_storage.VisitFunc) error {
+ return errors.New("b2 unreachable")
+ }
+ require.Error(t, f2.buildRemoteKeyIndexes())
+
+ entry, err := f2.maybeLazyFetchFromRemote(context.Background(), "/buckets/mybucket/other.narinfo")
+ require.NoError(t, err)
+ assert.Nil(t, entry)
+ entry, err = f2.maybeLazyFetchFromRemote(context.Background(), "/buckets/mybucket/known.narinfo")
+ require.NoError(t, err)
+ assert.NotNil(t, entry)
+ assert.Equal(t, 1, stub2.statCalls, "only the snapshotted name reaches the remote")
+}
+
+func TestRemoteKeyIndex_ListingReplacesSnapshot(t *testing.T) {
+ dir := t.TempDir()
+ f, _, restore := newKeyIndexFiler(t, "stub_keyidx_snap2", []string{"old.narinfo"})
+ defer restore()
+ f.RemoteKeyIndexDir = dir
+ require.NoError(t, f.buildRemoteKeyIndexes())
+
+ f2, stub2, restore2 := newKeyIndexFiler(t, "stub_keyidx_snap2b", []string{"new.narinfo"})
+ defer restore2()
+ f2.RemoteKeyIndexDir = dir
+ require.NoError(t, f2.buildRemoteKeyIndexes())
+
+ entry, err := f2.maybeLazyFetchFromRemote(context.Background(), "/buckets/mybucket/old.narinfo")
+ require.NoError(t, err)
+ assert.Nil(t, entry)
+ entry, err = f2.maybeLazyFetchFromRemote(context.Background(), "/buckets/mybucket/new.narinfo")
+ require.NoError(t, err)
+ assert.NotNil(t, entry)
+ assert.Equal(t, 1, stub2.statCalls)
+}
diff --git a/weed/filer/filer_on_meta_event.go b/weed/filer/filer_on_meta_event.go
index d57810f..30cc0a0 100644
--- a/weed/filer/filer_on_meta_event.go
+++ b/weed/filer/filer_on_meta_event.go
@@ -134,6 +134,7 @@ func (f *Filer) LoadRemoteStorageConfAndMapping() {
glog.Errorf("read remote conf and mapping: %v", err)
return
}
+ f.startRemoteKeyIndexes()
}
func (f *Filer) maybeReloadRemoteStorageConfigurationAndMapping(event *filer_pb.SubscribeMetadataResponse) {
if !filer_pb.MetadataEventTouchesDirectory(event, DirectoryEtcRemote) {
diff --git a/weed/filer/filer_remote_key_index.go b/weed/filer/filer_remote_key_index.go
new file mode 100644
index 0000000..4e87d2f
--- /dev/null
+++ b/weed/filer/filer_remote_key_index.go
@@ -0,0 +1,156 @@
+package filer
+
+import (
+ "bufio"
+ "context"
+ "fmt"
+ "net/url"
+ "os"
+ "path/filepath"
+
+ "github.com/seaweedfs/seaweedfs/weed/glog"
+ "github.com/seaweedfs/seaweedfs/weed/pb/filer_pb"
+ "github.com/seaweedfs/seaweedfs/weed/pb/remote_pb"
+ "github.com/seaweedfs/seaweedfs/weed/util"
+)
+
+// remoteKeySet holds the object paths of one mount. A set loaded from a
+// snapshot is not listed yet and is replaced by the first listing.
+type remoteKeySet struct {
+ keys map[util.FullPath]struct{}
+ listed bool
+}
+
+func (rs *FilerRemoteStorage) mounts() map[util.FullPath]*remote_pb.RemoteStorageLocation {
+ rs.mu.RLock()
+ defer rs.mu.RUnlock()
+ mounts := make(map[util.FullPath]*remote_pb.RemoteStorageLocation)
+ rs.rules.Walk(func(key []byte, loc *remote_pb.RemoteStorageLocation) bool {
+ mounts[util.FullPath(key[:len(key)-1])] = loc
+ return true
+ })
+ return mounts
+}
+
+// startRemoteKeyIndexes lists every mount that has no listing yet in the
+// background. Until a listing or snapshot is stored, lookups behave as
+// without an index.
+func (f *Filer) startRemoteKeyIndexes() {
+ if !f.EnableRemoteKeyIndex || f.RemoteStorage == nil {
+ return
+ }
+ go func() {
+ if err := f.buildRemoteKeyIndexes(); err != nil {
+ glog.Warningf("remote key index: %v", err)
+ }
+ }()
+}
+
+func (f *Filer) buildRemoteKeyIndexes() error {
+ if !f.EnableRemoteKeyIndex || f.RemoteStorage == nil {
+ return nil
+ }
+ var firstErr error
+ for mountDir, loc := range f.RemoteStorage.mounts() {
+ if set, ok := f.remoteKeys.Load(string(mountDir)); ok && set.(*remoteKeySet).listed {
+ continue
+ }
+ if _, busy := f.remoteKeysLoading.LoadOrStore(string(mountDir), struct{}{}); busy {
+ continue
+ }
+ err := f.buildRemoteKeyIndex(mountDir, loc)
+ f.remoteKeysLoading.Delete(string(mountDir))
+ if err != nil && firstErr == nil {
+ firstErr = fmt.Errorf("list %s: %w", mountDir, err)
+ }
+ }
+ return firstErr
+}
+
+func (f *Filer) buildRemoteKeyIndex(mountDir util.FullPath, loc *remote_pb.RemoteStorageLocation) error {
+ if _, ok := f.remoteKeys.Load(string(mountDir)); !ok {
+ if keys, err := f.loadRemoteKeySnapshot(mountDir); err != nil {
+ glog.Warningf("remote key index: snapshot for %s: %v", mountDir, err)
+ } else if keys != nil {
+ f.remoteKeys.Store(string(mountDir), &remoteKeySet{keys: keys})
+ glog.V(0).Infof("remote key index: %d objects under %s from snapshot", len(keys), mountDir)
+ }
+ }
+
+ conf, found := f.RemoteStorage.GetRemoteStorageConf(loc.Name)
+ if !found {
+ return fmt.Errorf("no remote storage %q", loc.Name)
+ }
+ client, err := f.buildRemoteStorageClient(context.Background(), conf)
+ if err != nil {
+ return err
+ }
+ keys := make(map[util.FullPath]struct{})
+ err = client.Traverse(loc, func(dir string, name string, _ bool, _ *filer_pb.RemoteEntry) error {
+ remotePath := util.FullPath(dir).Child(name)
+ keys[MapRemoteStorageLocationPathToFullPath(mountDir, loc, string(remotePath))] = struct{}{}
+ return nil
+ })
+ if err != nil {
+ return err
+ }
+ f.remoteKeys.Store(string(mountDir), &remoteKeySet{keys: keys, listed: true})
+ glog.V(0).Infof("remote key index: %d objects under %s", len(keys), mountDir)
+ if err := f.saveRemoteKeySnapshot(mountDir, keys); err != nil {
+ glog.Warningf("remote key index: save snapshot for %s: %v", mountDir, err)
+ }
+ return nil
+}
+
+func (f *Filer) remoteKeySnapshotPath(mountDir util.FullPath) string {
+ return filepath.Join(f.RemoteKeyIndexDir, url.PathEscape(string(mountDir)))
+}
+
+func (f *Filer) loadRemoteKeySnapshot(mountDir util.FullPath) (map[util.FullPath]struct{}, error) {
+ if f.RemoteKeyIndexDir == "" {
+ return nil, nil
+ }
+ file, err := os.Open(f.remoteKeySnapshotPath(mountDir))
+ if os.IsNotExist(err) {
+ return nil, nil
+ }
+ if err != nil {
+ return nil, err
+ }
+ defer file.Close()
+ keys := make(map[util.FullPath]struct{})
+ scanner := bufio.NewScanner(file)
+ for scanner.Scan() {
+ keys[util.FullPath(scanner.Text())] = struct{}{}
+ }
+ return keys, scanner.Err()
+}
+
+// saveRemoteKeySnapshot renames a finished file into place so a crash never
+// leaves a truncated snapshot that would hide existing objects.
+func (f *Filer) saveRemoteKeySnapshot(mountDir util.FullPath, keys map[util.FullPath]struct{}) error {
+ if f.RemoteKeyIndexDir == "" {
+ return nil
+ }
+ if err := os.MkdirAll(f.RemoteKeyIndexDir, 0o755); err != nil {
+ return err
+ }
+ tmp, err := os.CreateTemp(f.RemoteKeyIndexDir, ".snapshot.")
+ if err != nil {
+ return err
+ }
+ defer os.Remove(tmp.Name())
+ w := bufio.NewWriter(tmp)
+ for k := range keys {
+ w.WriteString(string(k))
+ w.WriteByte('\n')
+ }
+ if err := w.Flush(); err != nil {
+ tmp.Close()
+ return err
+ }
+ if err := tmp.Close(); err != nil {
+ return err
+ }
+ return os.Rename(tmp.Name(), f.remoteKeySnapshotPath(mountDir))
+}
diff --git a/weed/server/filer_server.go b/weed/server/filer_server.go
index 7d998ad..2f2107e 100644
--- a/weed/server/filer_server.go
+++ b/weed/server/filer_server.go
@@ -228,6 +228,8 @@ func NewFilerServer(defaultMux, readonlyMux *http.ServeMux, option *FilerOption)
glog.V(0).Infof("max_file_name_length %d", maxFilenameLength)
fs.filer = filer.NewFiler(*option.Masters, fs.grpcDialOption, option.Host, option.FilerGroup, option.Collection, option.DefaultReplication, option.DataCenter, maxFilenameLength, nil)
fs.filer.Cipher = option.Cipher
+ fs.filer.EnableRemoteKeyIndex = v.GetBool("filer.options.remote_key_index")
+ fs.filer.RemoteKeyIndexDir = v.GetString("filer.options.remote_key_index_dir")
fs.filer.DefaultDiskType = option.DiskType
fs.filer.BuildGuardedRemoteClient = BuildGuardedRemoteStorageClient
fs.filer.AllowUntrustedRemoteEndpoints = option.AllowUntrustedRemoteEndpoints
--
2.55.0
# The nixpkgs seaweedfs derivation plus the filer patch. Called with
# `pkgs.callPackage`, so the version and vendorHash come from nixpkgs and only
# the patch is ours.
{ seaweedfs }:
seaweedfs.overrideAttrs (old: {
patches = (old.patches or [ ]) ++ [ ./seaweedfs-remote-key-index.patch ];
})
/*
A SeaweedFS S3 gateway in front of an object store, used as the S3 endpoint of
a niks3 binary cache.
Layout: one `weed mini` node (master, volume, filer and S3 gateway, all bound to
localhost). The bucket is mounted into the filer as a remote mount, so existing
objects are readable without being downloaded and writes sync back in the
background. niks3 talks S3 to the gateway instead of the provider, and the
filer answers misses from the key index of ./seaweedfs-remote-key-index.patch.
Secrets are read from `services.seaweedfsCache.secretsDir`, a root-only
directory containing:
gateway-access-key credential for the gateway's own S3 identity
gateway-secret-key
b2-access-key credential for the remote bucket, scoped to the bucket
b2-secret-key
plus, for the niks3 side: signing-key (nix cache secret key) and api-token.
*/
{
config,
lib,
pkgs,
...
}:
let
cfg = config.services.seaweedfsCache;
weed = lib.getExe cfg.package;
s3Port = 8333;
s3HttpsPort = 8334;
filerPort = 8888;
masterPort = 9333;
stateDir = "/var/lib/seaweedfs";
remoteDir = "/buckets/${cfg.bucket}";
certDir = config.security.acme.certs.${cfg.apiHost}.directory;
in
{
options.services.seaweedfsCache = {
enable = lib.mkEnableOption "SeaweedFS S3 gateway for a niks3 binary cache";
package = lib.mkOption {
type = lib.types.package;
default = pkgs.callPackage ./seaweedfs.nix { };
defaultText = "pkgs.callPackage ./seaweedfs.nix { } (nixpkgs' seaweedfs plus the patch)";
description = "The `weed` binary. Must carry the remote key index patch.";
};
bucket = lib.mkOption {
type = lib.types.str;
example = "nix-cache";
description = "Name of the remote bucket, also the gateway's bucket name.";
};
apiHost = lib.mkOption {
type = lib.types.str;
example = "cache-api.example.com";
description = ''
Host of the niks3 API and of the presigned URLs. The S3 gateway serves TLS
with this host's ACME certificate, so niks3 can use https without trusting
a self-signed cert, and a signature that names this host verifies when
nginx forwards to the gateway with the Host header preserved.
'';
};
readHost = lib.mkOption {
type = lib.types.nullOr lib.types.str;
default = null;
example = "cache-local.example.com";
description = ''
Host that serves the bucket straight from this gateway, so a path is
readable as soon as it is uploaded instead of after the CDN caught up.
Anonymous GET and HEAD only.
'';
};
publicCacheUrl = lib.mkOption {
type = lib.types.str;
default = "https://${cfg.apiHost}";
example = "https://cache.example.com";
description = ''
Public URL of the cache, as trusted by clients: usually a CDN in front of
the remote bucket, so reads never reach this machine.
'';
};
region = lib.mkOption {
type = lib.types.str;
example = "eu-central-003";
description = ''
Region of the remote bucket. Names its S3 endpoint; set it explicitly,
since SigV4 signing against some providers fails when it is inferred from
the endpoint.
'';
};
endpoint = lib.mkOption {
type = lib.types.str;
example = "s3.eu-central-003.backblazeb2.com";
description = "S3 endpoint of the remote bucket.";
};
secretsDir = lib.mkOption {
type = lib.types.str;
example = "/run/secrets";
description = ''
Root-only directory with the credential files listed above. A string, not a
path, so that a runtime location is not copied into the store.
'';
};
enableTagging = lib.mkOption {
type = lib.types.bool;
default = false;
description = ''
Whether the remote bucket supports S3 object tagging. Some providers answer
tagging requests with an error, so the gateway must not forward them.
'';
};
};
config = lib.mkIf cfg.enable {
users.users.seaweedfs = {
isSystemUser = true;
group = "seaweedfs";
# Reads the ACME certificate for the gateway's TLS listener.
extraGroups = [ "nginx" ];
};
users.groups.seaweedfs = { };
# niks3 reaches the gateway by its API hostname, so the certificate matches.
networking.hosts."127.0.0.1" = [ cfg.apiHost ];
# The gateway reads the certificate when it starts.
security.acme.certs.${cfg.apiHost}.reloadServices = [ "seaweedfs.service" ];
systemd.services.seaweedfs = {
wantedBy = [ "multi-user.target" ];
serviceConfig = {
User = "seaweedfs";
Group = "seaweedfs";
StateDirectory = "seaweedfs";
Restart = "always";
# Names missing from the remote listing skip the remote stat; the
# snapshot on disk gates lookups right after a restart.
Environment = [
"WEED_FILER_OPTIONS_REMOTE_KEY_INDEX=true"
"WEED_FILER_OPTIONS_REMOTE_KEY_INDEX_DIR=${stateDir}/remote-key-index"
];
# Composed here so the secret keys stay out of the world-readable Nix store.
ExecStartPre = "${pkgs.writeShellScript "seaweedfs-s3-config" ''
set -euo pipefail
printf '{"identities":[{"name":"s3gateway","credentials":[{"accessKey":"%s","secretKey":"%s"}],"actions":["Admin","Read","Write","List","Tagging"]},{"name":"anonymous","actions":["Read:${cfg.bucket}"]}]}' \
"$(cat ${cfg.secretsDir}/gateway-access-key)" \
"$(cat ${cfg.secretsDir}/gateway-secret-key)" \
> ${stateDir}/s3.json
chmod 0600 ${stateDir}/s3.json
''}";
ExecStart = lib.escapeShellArgs [
weed
"mini"
"-dir=${stateDir}"
"-ip=127.0.0.1"
"-ip.bind=127.0.0.1"
"-master.port=${toString masterPort}"
"-filer.port=${toString filerPort}"
"-s3.port=${toString s3Port}"
"-s3.port.https=${toString s3HttpsPort}"
"-s3.cert.file=${certDir}/fullchain.pem"
"-s3.key.file=${certDir}/key.pem"
"-s3.config=${stateDir}/s3.json"
"-s3.port.iceberg=0"
"-s3.port.lance=0"
"-master.telemetry=false"
"-admin.ui=false"
];
};
};
# Mounts the bucket under the gateway: existing objects appear without being
# downloaded, deletes and new objects sync back.
systemd.services.seaweedfs-mount = {
wantedBy = [ "multi-user.target" ];
requires = [ "seaweedfs.service" ];
after = [ "seaweedfs.service" ];
path = [ pkgs.coreutils ];
serviceConfig = {
Type = "oneshot";
RemainAfterExit = true;
Restart = "on-failure";
RestartSec = 5;
};
script = ''
until (exec 3<>/dev/tcp/127.0.0.1/${toString masterPort}) 2>/dev/null; do sleep 1; done
shell() { ${weed} shell -master=127.0.0.1:${toString masterPort}; }
printf 'remote.configure -name=bucket -type=s3 -s3.access_key=%s -s3.secret_key=%s -s3.region=${cfg.region} -s3.support_tagging=${lib.boolToString cfg.enableTagging} -s3.endpoint=https://${cfg.endpoint}\nexit\n' \
"$(cat ${cfg.secretsDir}/b2-access-key)" \
"$(cat ${cfg.secretsDir}/b2-secret-key)" | shell
# remote.mount fails on a directory that is already mounted.
if ! printf 'remote.mount\nexit\n' | shell | grep -q '"${remoteDir}"'; then
printf 'remote.mount -dir=${remoteDir} -remote=bucket/${cfg.bucket}\nexit\n' | shell
fi
'';
};
systemd.services.seaweedfs-remote-sync = {
wantedBy = [ "multi-user.target" ];
requires = [ "seaweedfs-mount.service" ];
after = [ "seaweedfs-mount.service" ];
serviceConfig = {
User = "seaweedfs";
Group = "seaweedfs";
Restart = "always";
ExecStart = "${weed} filer.remote.sync -filer=127.0.0.1:${toString filerPort} -dir=${remoteDir}";
};
};
services.nginx.enable = true;
# Uploads and presigned downloads from niks3. The Host header is part of the
# presigned signature, so it has to survive the hop to the gateway. Issuing the
# ACME certificate here is what the gateway's TLS listener reads.
services.nginx.virtualHosts = {
${cfg.apiHost} = {
enableACME = true;
forceSSL = true;
locations."/${cfg.bucket}/" = {
proxyPass = "http://127.0.0.1:${toString s3Port}";
extraConfig = ''
proxy_set_header Host $host;
proxy_request_buffering off;
proxy_buffering off;
client_max_body_size 0;
'';
};
};
}
// lib.optionalAttrs (cfg.readHost != null) {
${cfg.readHost} = {
enableACME = true;
forceSSL = true;
locations."/" = {
proxyPass = "http://127.0.0.1:${toString s3Port}/${cfg.bucket}/";
extraConfig = ''
limit_except GET HEAD { deny all; }
proxy_buffering off;
'';
};
# niks3 stores these zstd-compressed, and objects that arrive through the
# remote mount lose Content-Encoding, so Nix would read them as corrupt.
locations."~ ^/(?<zstdobj>[a-z0-9]+\\.(narinfo|ls)|log/[A-Za-z0-9._+-]+|realisations/[A-Za-z0-9._+!-]+\\.doi)$" =
{
proxyPass = "http://127.0.0.1:${toString s3Port}/${cfg.bucket}/$zstdobj";
extraConfig = ''
limit_except GET HEAD { deny all; }
proxy_buffering off;
proxy_hide_header Content-Encoding;
add_header Content-Encoding zstd;
'';
};
};
};
# The niks3 side of the setup: its S3 client talks to the gateway over TLS, and
# presigned URLs name the public API host that nginx forwards from.
#
# services.niks3.s3 = {
# endpoint = "${cfg.apiHost}:${toString s3HttpsPort}";
# bucket = cfg.bucket;
# region = "us-east-1";
# useSSL = true;
# bucketLookup = "path";
# publicUrl = "https://${cfg.apiHost}";
# accessKeyFile = "${cfg.secretsDir}/gateway-access-key";
# secretKeyFile = "${cfg.secretsDir}/gateway-secret-key";
# };
# services.niks3 = {
# enable = true;
# httpAddr = "127.0.0.1:5751";
# cacheUrl = cfg.publicCacheUrl;
# serverUrl = "https://${cfg.apiHost}";
# signKeyFiles = [ "${cfg.secretsDir}/signing-key" ];
# apiTokenFile = "${cfg.secretsDir}/api-token";
# };
# # Uploads come from machines that have no cloud credential of their own, so
# # authenticate them with a client certificate instead of a shared token.
# services.niks3.nginx = {
# enable = true;
# domain = cfg.apiHost;
# mtls = { enable = true; clientCAFile = "${cfg.secretsDir}/client-ca.crt"; };
# };
# # niks3 must not start before the bucket is mounted, or GC would see an
# # empty listing.
# systemd.services.niks3.after = [ "seaweedfs-mount.service" ];
# systemd.services.niks3.requires = [ "seaweedfs-mount.service" ];
};
}
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment