Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
18 changes: 14 additions & 4 deletions docs/architecture.md
Original file line number Diff line number Diff line change
Expand Up @@ -502,7 +502,14 @@ index (stock-tooling plaintext recovery)** — a deliberate trade, not settled h
The **catalog** is different — many tiny MST/manifest blocks share a CAR — so catalog blocks resolve
via the indexer's index-claim / sharded-dag-index path (block CID → byte range in its shard). That path
is retained for the catalog regardless. *(In the R0/R1 appliance topology, catalog-block lookup is
served from the local Postgres location table rather than the indexing-service.)*
served from the local Postgres mirror of that contract: `shard_inclusions` (block digest → shard
digest + byte range, recorded by the flush path before a segment is marked shipped) joined to the
shard's `blob_locations` row. A whole-blob location table alone cannot serve catalog blocks — they
are interior slices of shipped CARs, not stored blobs — which is exactly what breaks
retention-retired catalog reads without the inclusion table. Because the local mirror is what
ingot's own reads depend on, the network index publication (`/index/add` at ship time) is
best-effort: its failure is logged, not a ship failure — a wedged indexer must not wedge
retention. A retry queue for failed publications is a TODO.)*

Consuming a bare location commitment is a capability Ingot's locator must gain: today it surfaces a
location only via the inclusion → shard → commitment path and never returns a stored bare
Expand Down Expand Up @@ -780,9 +787,12 @@ These are intentional simplifications of the target topology, not bugs:
Regime-A/B compaction, and subroots ([§6](#6-the-forgechain-layer)) are Piri-side concerns. From Ingot's side a delete is
just `remove(digest)`; Piri decides the on-chain regime. `max_blob_size` is the only size knob Ingot
carries.
- **Local location table instead of the indexer (R0/R1 appliance reduction).** Body-blob locations
are recorded in a local Postgres `blob_locations` table behind a `Locator` seam, rather than read
back through the indexing-service ([§5](#5-the-data-layer), [§8](#8-retrieval-addressing-when-bodies-need-a-sharded-dag-index)). The indexer-backed `Locator` is an `indexer-ready` swap-in.
- **Local location + inclusion tables instead of the indexer (R0/R1 appliance reduction).** Body-blob
and shipped-shard locations are recorded in a local Postgres `blob_locations` table, and each
shipped catalog shard's inner-block byte ranges in `shard_inclusions`, behind a `Locator` seam,
rather than read back through the indexing-service ([§5](#5-the-data-layer), [§8](#8-retrieval-addressing-when-bodies-need-a-sharded-dag-index)). The two tables mirror the
indexing-service contract (location commitments + inclusions), so the indexer-backed `Locator`
remains an `indexer-ready` swap-in.

### Object lifecycle (not implemented)

Expand Down
6 changes: 3 additions & 3 deletions go.mod
Original file line number Diff line number Diff line change
Expand Up @@ -7,11 +7,11 @@ require (
github.com/aws/aws-sdk-go-v2/config v1.32.26
github.com/aws/aws-sdk-go-v2/credentials v1.19.25
github.com/aws/aws-sdk-go-v2/service/s3 v1.104.1
github.com/fil-forge/hilt v0.0.1-0.20260716084626-7ddddf09ecc0
github.com/fil-forge/hilt v0.0.1-0.20260724134448-ba71f843f6a4
github.com/fil-forge/indexing-service v1.13.5-0.20260619142411-efe3f5fab717
github.com/fil-forge/libforge v0.0.0-20260713100115-a3aa293b990c
github.com/fil-forge/libforge v0.0.0-20260724113901-7fc3b2cec1ef
github.com/fil-forge/smelt v0.0.0-20260720130429-63116166a06c
github.com/fil-forge/ucantone v0.0.0-20260706102443-79141c5cc52e
github.com/fil-forge/ucantone v0.0.0-20260727203046-ccb77059de44
github.com/fil-forge/versitygw v0.0.0-20260716095011-7a65883d595a
github.com/fxamacker/cbor/v2 v2.9.2
github.com/go-jose/go-jose/v4 v4.1.4
Expand Down
12 changes: 6 additions & 6 deletions go.sum
Original file line number Diff line number Diff line change
Expand Up @@ -246,16 +246,16 @@ github.com/fil-forge/go-ipni-tools v0.0.0-20260519194815-545b9421aec0 h1:HAfXUPv
github.com/fil-forge/go-ipni-tools v0.0.0-20260519194815-545b9421aec0/go.mod h1:3NRV/7wc4/0uzzrGdI7NoN/yeF1UvqKRwMyjBqGc5s0=
github.com/fil-forge/go-ucanto v0.0.0-20260507172450-5cb5d073f8ab h1:2J2cDThqTKP6/0k3SfdlSxfyPa3aLqjTYnmvbEcryfg=
github.com/fil-forge/go-ucanto v0.0.0-20260507172450-5cb5d073f8ab/go.mod h1:lZF3UXZ2hGLKYmXdquG50JqI9pRlUrV6lubGtgOYfwc=
github.com/fil-forge/hilt v0.0.1-0.20260716084626-7ddddf09ecc0 h1:dma5d9PcBPyvMqZTOmoGUeZBKvnZztdko4hWqBpX7ds=
github.com/fil-forge/hilt v0.0.1-0.20260716084626-7ddddf09ecc0/go.mod h1:zzPrQQ/VhgShDIi9IsfeFdiNXJftACRXxUaJI6opW9s=
github.com/fil-forge/hilt v0.0.1-0.20260724134448-ba71f843f6a4 h1:pDDN87a4dMuH8mqqEZZOnrCZ+Myc1ArEjJxg0ZNpQzA=
github.com/fil-forge/hilt v0.0.1-0.20260724134448-ba71f843f6a4/go.mod h1:AO/+NYsz//BoqdHjVfik0fbxjiq+H+86qzGBt2/2uxQ=
github.com/fil-forge/indexing-service v1.13.5-0.20260619142411-efe3f5fab717 h1:Wke8qgaDgy7DGaIS28VpHip+YHSdQP1hGNEZTrXXzb4=
github.com/fil-forge/indexing-service v1.13.5-0.20260619142411-efe3f5fab717/go.mod h1:wFcakLohOqpMRkJzWRdGFGHRFpGB+PpEMCm9wkt2cqU=
github.com/fil-forge/libforge v0.0.0-20260713100115-a3aa293b990c h1:ElgAxd+08QifEpnhO1vH58A/CD+i/Hdic6B+k3W/QPY=
github.com/fil-forge/libforge v0.0.0-20260713100115-a3aa293b990c/go.mod h1:0kXihIQ4L2uZ00nR5XrZ/Y8Db7Ht/qQNuiWslwMJ95M=
github.com/fil-forge/libforge v0.0.0-20260724113901-7fc3b2cec1ef h1:xgkciShyWdCQ2Pyl2qEwP7JMbsdbHdECj/v6E5+fj6U=
github.com/fil-forge/libforge v0.0.0-20260724113901-7fc3b2cec1ef/go.mod h1:0kXihIQ4L2uZ00nR5XrZ/Y8Db7Ht/qQNuiWslwMJ95M=
github.com/fil-forge/smelt v0.0.0-20260720130429-63116166a06c h1:WHvsleEU6ZiNYDFgLx6KtorXulD+IuLiorMgmp4Th8s=
github.com/fil-forge/smelt v0.0.0-20260720130429-63116166a06c/go.mod h1:NM/mk/XiP1Kzsy9HWGQeLsosPUvVY3bIrkUzqswThpU=
github.com/fil-forge/ucantone v0.0.0-20260706102443-79141c5cc52e h1:di/SseJVEO6bznSH/UEnAE2nSwwp/gqgnPLgE9/x2Zg=
github.com/fil-forge/ucantone v0.0.0-20260706102443-79141c5cc52e/go.mod h1:oFY5BfD0bDeodGlbBHh3/nK99MAS93rGXjoQz7s5qgE=
github.com/fil-forge/ucantone v0.0.0-20260727203046-ccb77059de44 h1:ofvb2Qq7++VPRelGsLbtnd1ZMKVT4n4QGoa79BhJ6VQ=
github.com/fil-forge/ucantone v0.0.0-20260727203046-ccb77059de44/go.mod h1:oFY5BfD0bDeodGlbBHh3/nK99MAS93rGXjoQz7s5qgE=
github.com/fil-forge/versitygw v0.0.0-20260716095011-7a65883d595a h1:lDwnNmF4LNbevx/YYNCC0CazMiL4eIe38MvfqbRR3RA=
github.com/fil-forge/versitygw v0.0.0-20260716095011-7a65883d595a/go.mod h1:t73Wa2xqpT0NdY+SzRZDiO+tQD13b3UEZRoe2wYFiIE=
github.com/filecoin-project/go-data-segment v0.0.1 h1:1wmDxOG4ubWQm3ZC1XI5nCon5qgSq7Ra3Rb6Dbu10Gs=
Expand Down
34 changes: 18 additions & 16 deletions inmem/store.go
Original file line number Diff line number Diff line change
Expand Up @@ -54,12 +54,13 @@ type MemStore struct {
// The architecture's relational surface (docs/architecture.md §5–§7),
// mirroring the Postgres tables so the in-process suite exercises the
// same code paths. See stores.go for the methods over these.
blobRefs map[claimKey]registry.BlobClaim
intents map[string]registry.UploadIntent // keyed by string(digest)
locations map[locKey]registry.BlobLocation // keyed by (space, digest)
sessions map[string]registry.MultipartSession // keyed by uploadID
parts map[string]map[int]registry.MultipartPart // uploadID -> partNumber -> part
gcCands map[string]struct{} // keyed by string(cid)
blobRefs map[claimKey]registry.BlobClaim
intents map[string]registry.UploadIntent // keyed by string(digest)
locations map[locKey]registry.BlobLocation // keyed by (space, digest)
inclusions map[locKey]registry.BlobInclusion // keyed by (space, digest)
sessions map[string]registry.MultipartSession // keyed by uploadID
parts map[string]map[int]registry.MultipartPart // uploadID -> partNumber -> part
gcCands map[string]struct{} // keyed by string(cid)
}

// claimKey / locKey are the composite map keys for the blob_refs and
Expand All @@ -76,14 +77,15 @@ type locKey struct {
// NewMemStore returns an empty MemStore.
func NewMemStore() *MemStore {
return &MemStore{
buckets: map[string]*registry.State{},
segments: map[uint64]*logstore.SegmentMeta{},
blobRefs: map[claimKey]registry.BlobClaim{},
intents: map[string]registry.UploadIntent{},
locations: map[locKey]registry.BlobLocation{},
sessions: map[string]registry.MultipartSession{},
parts: map[string]map[int]registry.MultipartPart{},
gcCands: map[string]struct{}{},
buckets: map[string]*registry.State{},
segments: map[uint64]*logstore.SegmentMeta{},
blobRefs: map[claimKey]registry.BlobClaim{},
intents: map[string]registry.UploadIntent{},
locations: map[locKey]registry.BlobLocation{},
inclusions: map[locKey]registry.BlobInclusion{},
sessions: map[string]registry.MultipartSession{},
parts: map[string]map[int]registry.MultipartPart{},
gcCands: map[string]struct{}{},
}
}

Expand Down Expand Up @@ -339,8 +341,8 @@ func (NopBaseReader) OpenBlob(_ context.Context, _ did.DID, _ multihash.Multihas
// network, so the spool's local copy serves all reads.
type NopUploader struct{}

func (NopUploader) SubmitShard(_ context.Context, _ blockstore.Plane, _ did.DID, _ uploader.CARShard) error {
return nil
func (NopUploader) SubmitShard(_ context.Context, _ blockstore.Plane, _ did.DID, _ uploader.CARShard) (uploader.BlobLocation, error) {
return uploader.BlobLocation{}, nil
}

func (NopUploader) UploadBlob(_ context.Context, _ did.DID, _ multihash.Multihash, size int64, _ string) (uploader.BlobLocation, error) {
Expand Down
28 changes: 28 additions & 0 deletions inmem/stores.go
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,7 @@ var (
_ registry.BlobRefStore = (*MemStore)(nil)
_ registry.IntentStore = (*MemStore)(nil)
_ registry.LocationStore = (*MemStore)(nil)
_ registry.InclusionStore = (*MemStore)(nil)
_ registry.MultipartStore = (*MemStore)(nil)
_ registry.GCStore = (*MemStore)(nil)
)
Expand Down Expand Up @@ -147,6 +148,33 @@ func (m *MemStore) DeleteLocation(_ context.Context, space did.DID, digest []byt
return nil
}

// InclusionStore =============================================================

func (m *MemStore) PutInclusions(_ context.Context, incs []registry.BlobInclusion) error {
m.mu.Lock()
defer m.mu.Unlock()
for _, inc := range incs {
cp := inc
cp.Digest = cloneBytes(inc.Digest)
cp.ShardDigest = cloneBytes(inc.ShardDigest)
m.inclusions[locKey{inc.Space, string(inc.Digest)}] = cp
}
return nil
}

func (m *MemStore) GetInclusion(_ context.Context, space did.DID, digest []byte) (*registry.BlobInclusion, error) {
m.mu.Lock()
defer m.mu.Unlock()
inc, ok := m.inclusions[locKey{space, string(digest)}]
if !ok {
return nil, registry.ErrNotFound
}
cp := inc
cp.Digest = cloneBytes(inc.Digest)
cp.ShardDigest = cloneBytes(inc.ShardDigest)
return &cp, nil
}

// MultipartStore =============================================================

func (m *MemStore) CreateSession(_ context.Context, s registry.MultipartSession) error {
Expand Down
141 changes: 141 additions & 0 deletions itest/forge_retention_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,141 @@
//go:build itest

package itest

import (
"bytes"
"fmt"
"strings"
"testing"
"time"

ingottest "github.com/fil-forge/ingot/testing"
"github.com/fil-forge/smelt/pkg/stack"
)

// TestForgeReadAfterCatalogRetention proves catalog retention behaves as a
// cache, not an availability cliff: after a shipped catalog segment is retired
// off local disk, its blocks (object manifests, MST nodes) must still resolve
// through the fallthrough read tier — the shard_inclusions row (block → shard
// CAR + byte range, recorded at ship) joined to the shard's blob_locations
// row, retrieved as a ranged /content/retrieve against piri.
//
// The config seals the catalog every 1s, retains ONE shipped segment, and
// disables the read cache. An early object is written, then later writes roll
// the catalog past the retain window until the early object's segment is
// retired; the early object must then still GET (its manifest is only in the
// retired segment) and the bucket must still list without a delimiter (the
// walk fetches every leaf's manifest).
//
// go test -tags itest ./itest -run TestForgeReadAfterCatalogRetention -v -timeout 900s
func TestForgeReadAfterCatalogRetention(t *testing.T) {
ctx := t.Context()

s, ingotEndpoint := forgeStack(t, stack.WithServiceConfig("ingot", "testdata/config-retention.yaml"))
accessKey, secretKey := hiltProvisionTenant(t, ctx, s, "retention")
cfg := forgeConfig(ingotEndpoint, accessKey, secretKey)

const bucket = "retention-bucket"
if err := ingottest.CreateBucket(ctx, cfg, bucket); err != nil {
t.Fatalf("create bucket: %v", err)
}

// The early object: its manifest lands in the first catalog segment(s),
// which the later writes will push out of the retain window.
early := patternBytes(64 << 10)
if err := ingottest.PutBytes(ctx, cfg, bucket, "early-obj", early); err != nil {
t.Fatalf("put early object: %v", err)
}

// Count the segment CARs currently on disk; the early manifest lives in
// one of these. Retirement is proven when every one of them is gone.
initialSegs := catalogSegments(t, s, bucket)
if len(initialSegs) == 0 {
// The first segment may not have sealed yet; wait for it so we have
// a concrete set to watch retire.
waitFor(t, time.Minute, "first catalog segment to seal", func() bool {
initialSegs = catalogSegments(t, s, bucket)
return len(initialSegs) > 0
})
}
t.Logf("early object's manifest is in segment(s): %v", initialSegs)

// Roll the catalog: spaced writes each seal (1s seal_age) and ship a new
// segment; retain=1 retires everything older. Keep writing until every
// initial segment file is gone from the container.
waitFor(t, 5*time.Minute, "initial catalog segments to retire", func() bool {
key := fmt.Sprintf("filler/obj-%d", time.Now().UnixNano())
if err := ingottest.PutBytes(ctx, cfg, bucket, key, patternBytes(4<<10)); err != nil {
t.Fatalf("put filler object: %v", err)
}
// Slower than seal_age so each segment seals and ships with headroom
// (no flush-queue overflow → no re-ships through the dedup path).
time.Sleep(3 * time.Second)
remaining := catalogSegments(t, s, bucket)
for _, seg := range initialSegs {
for _, r := range remaining {
if seg == r {
return false
}
}
}
return true
})
t.Logf("initial segments retired; segments now on disk: %v", catalogSegments(t, s, bucket))

// GET the early object: its manifest exists only in a retired segment, so
// this read MUST come through inclusion → shard → ranged piri retrieval.
got, err := ingottest.GetBytes(ctx, cfg, bucket, "early-obj")
if err != nil {
t.Fatalf("get early object after its catalog segment retired: %v", err)
}
if !bytes.Equal(got, early) {
t.Fatalf("early object mismatch after retention: got %d bytes, want %d", len(got), len(early))
}

// Undelimited list walks every leaf and fetches every manifest — the
// regression that motivated this test (root listing failed with
// "blockstore: not found" once a manifest's segment was retired).
keys, err := ingottest.ListKeys(ctx, cfg, bucket)
if err != nil {
t.Fatalf("list bucket after retention: %v", err)
}
found := false
for _, k := range keys {
if k == "early-obj" {
found = true
break
}
}
if !found {
t.Fatalf("early-obj missing from listing: %v", keys)
}
t.Logf("read-after-retention OK: %d bytes via ranged shard retrieval; %d keys listed", len(got), len(keys))
}

// catalogSegments lists the catalog-plane CAR files currently on the ingot
// container's disk for bucket (segments live under
// /data/segments/<bucket>/catalog/).
func catalogSegments(t *testing.T, s *stack.Stack, bucket string) []string {
t.Helper()
out, _, err := s.Exec(t.Context(), "ingot", "sh", "-c",
"ls /data/segments/"+bucket+"/catalog/ 2>/dev/null | grep '\\.car$' || true")
if err != nil {
t.Fatalf("list catalog segments: %v", err)
}
fields := strings.Fields(strings.TrimSpace(out))
return fields
}

// waitFor polls cond (which may do work per attempt) until true or fatal after
// timeout.
func waitFor(t *testing.T, timeout time.Duration, what string, cond func() bool) {
t.Helper()
deadline := time.Now().Add(timeout)
for time.Now().Before(deadline) {
if cond() {
return
}
}
t.Fatalf("timed out after %s waiting for %s", timeout, what)
}
5 changes: 3 additions & 2 deletions itest/scenarios_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -291,8 +291,9 @@ func TestForgeScenarios(t *testing.T) {
}

resp := preflight(t, origin)
if resp.StatusCode != http.StatusNoContent {
t.Errorf("preflight status = %d, want 204", resp.StatusCode)
// CORS preflight spec allows any "ok status" (200-299).
if resp.StatusCode < 200 || resp.StatusCode >= 300 {
t.Errorf("preflight status = %d, want 2xx", resp.StatusCode)
}
if got := resp.Header.Get("Access-Control-Allow-Origin"); got != origin {
t.Errorf("preflight Allow-Origin = %q, want the request origin echoed", got)
Expand Down
11 changes: 10 additions & 1 deletion itest/stack_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -104,7 +104,9 @@ func forgeStack(t *testing.T, extra ...stack.Option) (*stack.Stack, string) {
t.Helper()
t.Logf("booting the smelt Forge stack (~1-2 min; first run also compiles ingot and pulls images)")
opts := []stack.Option{
stack.WithPiriNodes(stack.PiriNodeConfig{}),
// Postgres-backed piri: piri:main's curio PDP pipeline refuses
// sqlite ("curio PDP pipeline requires Postgres") as of 2026-07-24.
stack.WithPiriNodes(stack.PiriNodeConfig{Postgres: true}),
stack.WithServiceBinary("ingot", localIngotBinary(t)),
}
// Local-dev escape hatches: run against upload-service (sprue) / piri
Expand All @@ -122,6 +124,13 @@ func forgeStack(t *testing.T, extra ...stack.Option) (*stack.Stack, string) {
t.Logf("using piri image override: %s", img)
opts = append(opts, stack.WithPiriImage(img))
}
// Same idea one step earlier in the pipeline: mount a locally-built piri
// binary (linux, static) over the image's /usr/bin/piri — validates an
// unreleased piri/ucantone change with no image build at all.
if bin := os.Getenv("INGOT_ITEST_PIRI_BINARY"); bin != "" {
t.Logf("using piri binary override: %s", bin)
opts = append(opts, stack.WithPiriBinary(bin))
}
opts = append(opts, extra...)
s := stack.MustNewStack(t, opts...)
endpoint := s.IngotEndpoint()
Expand Down
40 changes: 40 additions & 0 deletions itest/testdata/config-retention.yaml
Original file line number Diff line number Diff line change
@@ -0,0 +1,40 @@
# Ingot daemon config for the catalog-retention fallthrough itest — smelt's
# default forge-mode config (systems/ingot/config/config.yaml) with the
# catalog plane tuned to seal fast and retain only one shipped segment, and
# the read cache disabled so every fallthrough read exercises the locator +
# ranged /content/retrieve path rather than a warm cache.
# Mounted over /etc/ingot/config.yaml via stack.WithServiceConfig.

addr: "0.0.0.0:9000"
region: us-west-1
data_dir: /data

root_access: ingot
root_secret: ingotsecret

identity:
key_file: /keys/ingot.pem

postgres_dsn: "postgres://ingot:ingot@ingot-postgres:5432/ingot?sslmode=disable"

upload_service_url: "http://upload:80"
upload_service_did: "did:web:upload"
upload_receipts_url: "http://upload:80/receipt"

indexer_endpoint: "http://indexer:80"
indexer_did: "did:web:indexer"

# S3 authorization service (hilt); the DID is used verbatim (no resolution)
auth_service_url: "http://hilt:80"
auth_service_did: "did:web:hilt"
# hilt → ingot delegations for /s3/request/authorize and /s3/bucket/*
auth_service_proofs: "/proofs/hilt-ingot-s3-proof.txt"

# The point of this config: roll catalog segments fast and retire aggressively,
# with no read cache masking the fallthrough. seal_age is fast but leaves the
# flusher headroom per ship (a 1s seal outruns the ship pipeline, overflowing
# the flush queue and forcing re-ships through sprue's dedup path).
read_cache_bytes: -1
catalog_plane:
seal_age: 2s
retain: 1
Loading
Loading