Skip to content

Commit a21f9b0

Browse files
committed
fix: update SOC chunk only with newer ts
1 parent a5c800f commit a21f9b0

2 files changed

Lines changed: 210 additions & 8 deletions

File tree

pkg/storer/internal/reserve/reserve.go

Lines changed: 41 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -102,6 +102,10 @@ func New(
102102
// 4. Two different chunk addresses that share the same batch stamp index and timestamp are settled by a tie-break:
103103
// the lexicographically lower chunk address wins. The loser is rejected; the winner replaces the stored chunk
104104
// through the usual remove-and-store path (including a fresh bin ID for pullsync).
105+
// 5. Two single owner chunks that share an address under different stamps (any batch or stamp
106+
// index) settle on one shared payload: a strictly higher stamp timestamp replaces it; equal
107+
// timestamps are settled by the lexicographically lower stamp hash. An older stamp is
108+
// rejected. Same-stamp divergence remains handled by resolveDivergence above.
105109
func (r *Reserve) Put(ctx context.Context, chunk swarm.Chunk) error {
106110
socReplaced, err := r.putChunk(ctx, chunk)
107111
if err != nil {
@@ -174,6 +178,12 @@ func (r *Reserve) putChunk(ctx context.Context, chunk swarm.Chunk) (socReplaced
174178
var shouldIncReserveSize bool
175179

176180
err = r.st.Run(ctx, func(s transaction.Store) error {
181+
if chunkType == swarm.ChunkTypeSingleOwner {
182+
if err := checkSOCStampOverwrite(ctx, s, chunk, stampHash); err != nil {
183+
return err
184+
}
185+
}
186+
177187
oldStampIndex, loadedStampIndex, err := stampindex.LoadOrStore(s.IndexStore(), reserveScope, chunk)
178188
if err != nil {
179189
return fmt.Errorf("load or store stamp index for chunk %v has fail: %w", chunk, err)
@@ -381,6 +391,37 @@ func (r *Reserve) putChunk(ctx context.Context, chunk swarm.Chunk) (socReplaced
381391
return socReplaced, nil
382392
}
383393

394+
// checkSOCStampOverwrite rejects an incoming single owner chunk when the
395+
// address already holds a payload under a stamp that should keep winning:
396+
// a strictly higher timestamp, or an equal timestamp with a lower or equal
397+
// stamp hash. Must run before LoadOrStore, which writes immediately.
398+
func checkSOCStampOverwrite(ctx context.Context, s transaction.Store, chunk swarm.Chunk, stampHash []byte) error {
399+
hasPayload, err := s.ChunkStore().Has(ctx, chunk.Address())
400+
if err != nil || !hasPayload {
401+
return err
402+
}
403+
404+
curr := binary.BigEndian.Uint64(chunk.Stamp().Timestamp())
405+
return chunkstamp.IterateAll(s.IndexStore(), reserveScope, chunk.Address(), func(st swarm.Stamp) (bool, error) {
406+
prev := binary.BigEndian.Uint64(st.Timestamp())
407+
if prev > curr {
408+
return true, fmt.Errorf("overwrite same chunk. prev %d cur %d batch %s: %w",
409+
prev, curr, hex.EncodeToString(chunk.Stamp().BatchID()), storage.ErrOverwriteNewerChunk)
410+
}
411+
if prev == curr {
412+
prevHash, err := st.Hash()
413+
if err != nil {
414+
return true, err
415+
}
416+
if bytes.Compare(prevHash, stampHash) <= 0 {
417+
return true, fmt.Errorf("overwrite same chunk. prev %d cur %d batch %s: %w",
418+
prev, curr, hex.EncodeToString(chunk.Stamp().BatchID()), storage.ErrOverwriteNewerChunk)
419+
}
420+
}
421+
return false, nil
422+
})
423+
}
424+
384425
// refreshSiblingSums recomputes the divergence checksum of every reserve entry
385426
// at the given address after its shared payload was replaced. Without the
386427
// refresh, entries under other stamps keep advertising content the node no

pkg/storer/internal/reserve/reserve_test.go

Lines changed: 169 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -201,14 +201,12 @@ func TestSameChunkAddress(t *testing.T) {
201201
bin := swarm.Proximity(baseAddr.Bytes(), ch1.Address().Bytes())
202202
binBinIDs[bin] += 1
203203
err = r.Put(ctx, ch2)
204-
if err != nil {
205-
t.Fatal(err)
204+
if !errors.Is(err, storage.ErrOverwriteNewerChunk) {
205+
t.Fatal("expected error")
206206
}
207-
bin2 := swarm.Proximity(baseAddr.Bytes(), ch2.Address().Bytes())
208-
binBinIDs[bin2] += 1
209207
size2 := r.Size()
210-
if size2-size1 != 2 {
211-
t.Fatalf("expected reserve size to increase by 2, got %d", size2-size1)
208+
if size2-size1 != 1 {
209+
t.Fatalf("expected reserve size to increase by 1, got %d", size2-size1)
212210
}
213211
})
214212

@@ -1241,7 +1239,7 @@ func TestSOCSiblingSumRefresh(t *testing.T) {
12411239
s3 := soctesting.GenerateMockSocWithSigner(t, []byte("v3"), signer)
12421240

12431241
stampA := postagetesting.MustNewFields(batchA.ID, 0, 1)
1244-
stampB := postagetesting.MustNewFields(batchB.ID, 0, 1)
1242+
stampB := postagetesting.MustNewFields(batchB.ID, 0, 2)
12451243
chA := s1.Chunk().WithStamp(stampA)
12461244
chB := s2.Chunk().WithStamp(stampB)
12471245

@@ -1308,7 +1306,7 @@ func TestSOCSiblingSumRefresh(t *testing.T) {
13081306
// the same-batch replacement path (higher stamp timestamp) must refresh
13091307
// batch B's entry the same way.
13101308
staleSumB := sumOf(batchB.ID, stampHashB)
1311-
chA2 := s3.Chunk().WithStamp(postagetesting.MustNewFields(batchA.ID, 0, 2))
1309+
chA2 := s3.Chunk().WithStamp(postagetesting.MustNewFields(batchA.ID, 0, 3))
13121310
if err := r.Put(ctx, chA2); err != nil {
13131311
t.Fatal(err)
13141312
}
@@ -1447,6 +1445,169 @@ func TestChunkSumIndexRandomOps(t *testing.T) {
14471445
checkInvariant(200)
14481446
}
14491447

1448+
// TestSOCCrossBatchTimestamp covers two single owner chunks that share an
1449+
// address but are stamped under different batches. The shared chunkstore
1450+
// payload is replaced when the incoming stamp timestamp is strictly higher,
1451+
// or when timestamps are equal and the incoming stamp hash is lower. An older
1452+
// stamp is rejected so neighborhoods converge.
1453+
func TestSOCCrossBatchTimestamp(t *testing.T) {
1454+
t.Parallel()
1455+
1456+
ctx := context.Background()
1457+
signer := getSigner(t)
1458+
batchA := postagetesting.MustNewBatch()
1459+
batchB := postagetesting.MustNewBatch()
1460+
1461+
sOlder := soctesting.GenerateMockSocWithSigner(t, []byte("older"), signer)
1462+
sNewer := soctesting.GenerateMockSocWithSigner(t, []byte("newer"), signer)
1463+
if !sOlder.Chunk().Address().Equal(sNewer.Chunk().Address()) {
1464+
t.Fatal("expected shared SOC address")
1465+
}
1466+
1467+
t.Run("higher timestamp replaces", func(t *testing.T) {
1468+
t.Parallel()
1469+
1470+
baseAddr := swarm.RandAddress(t)
1471+
ts := internal.NewInmemStorage()
1472+
r, err := reserve.New(baseAddr, ts, 0, kademlia.NewTopologyDriver(), log.Noop)
1473+
if err != nil {
1474+
t.Fatal(err)
1475+
}
1476+
1477+
older := sOlder.Chunk().WithStamp(postagetesting.MustNewFields(batchA.ID, 0, 1))
1478+
newer := sNewer.Chunk().WithStamp(postagetesting.MustNewFields(batchB.ID, 0, 2))
1479+
1480+
if err := r.Put(ctx, older); err != nil {
1481+
t.Fatal(err)
1482+
}
1483+
if err := r.Put(ctx, newer); err != nil {
1484+
t.Fatal(err)
1485+
}
1486+
1487+
got, err := ts.ChunkStore().Get(ctx, newer.Address())
1488+
if err != nil {
1489+
t.Fatal(err)
1490+
}
1491+
if !bytes.Equal(got.Data(), newer.Data()) {
1492+
t.Fatal("expected payload from the higher-timestamp stamp")
1493+
}
1494+
})
1495+
1496+
t.Run("equal timestamp stamp hash tie-break", func(t *testing.T) {
1497+
t.Parallel()
1498+
1499+
chA := sOlder.Chunk().WithStamp(postagetesting.MustNewFields(batchA.ID, 0, 5))
1500+
chB := sNewer.Chunk().WithStamp(postagetesting.MustNewFields(batchB.ID, 0, 5))
1501+
hashA, err := chA.Stamp().Hash()
1502+
if err != nil {
1503+
t.Fatal(err)
1504+
}
1505+
hashB, err := chB.Stamp().Hash()
1506+
if err != nil {
1507+
t.Fatal(err)
1508+
}
1509+
var winner, loser swarm.Chunk
1510+
if bytes.Compare(hashA, hashB) < 0 {
1511+
winner, loser = chA, chB
1512+
} else {
1513+
winner, loser = chB, chA
1514+
}
1515+
1516+
for _, order := range [][]swarm.Chunk{{winner, loser}, {loser, winner}} {
1517+
baseAddr := swarm.RandAddress(t)
1518+
ts := internal.NewInmemStorage()
1519+
r, err := reserve.New(baseAddr, ts, 0, kademlia.NewTopologyDriver(), log.Noop)
1520+
if err != nil {
1521+
t.Fatal(err)
1522+
}
1523+
if err := r.Put(ctx, order[0]); err != nil {
1524+
t.Fatal(err)
1525+
}
1526+
_ = r.Put(ctx, order[1]) // may reject when winner is already stored
1527+
1528+
got, err := ts.ChunkStore().Get(ctx, winner.Address())
1529+
if err != nil {
1530+
t.Fatal(err)
1531+
}
1532+
if !bytes.Equal(got.Data(), winner.Data()) {
1533+
t.Fatal("expected payload from the lower stamp-hash claim")
1534+
}
1535+
}
1536+
})
1537+
1538+
t.Run("lower timestamp rejected", func(t *testing.T) {
1539+
t.Parallel()
1540+
1541+
baseAddr := swarm.RandAddress(t)
1542+
ts := internal.NewInmemStorage()
1543+
r, err := reserve.New(baseAddr, ts, 0, kademlia.NewTopologyDriver(), log.Noop)
1544+
if err != nil {
1545+
t.Fatal(err)
1546+
}
1547+
1548+
newer := sNewer.Chunk().WithStamp(postagetesting.MustNewFields(batchA.ID, 0, 9))
1549+
older := sOlder.Chunk().WithStamp(postagetesting.MustNewFields(batchB.ID, 0, 3))
1550+
1551+
if err := r.Put(ctx, newer); err != nil {
1552+
t.Fatal(err)
1553+
}
1554+
err = r.Put(ctx, older)
1555+
if !errors.Is(err, storage.ErrOverwriteNewerChunk) {
1556+
t.Fatalf("expected ErrOverwriteNewerChunk, got %v", err)
1557+
}
1558+
1559+
got, err := ts.ChunkStore().Get(ctx, newer.Address())
1560+
if err != nil {
1561+
t.Fatal(err)
1562+
}
1563+
if !bytes.Equal(got.Data(), newer.Data()) {
1564+
t.Fatal("expected newer payload to remain")
1565+
}
1566+
})
1567+
1568+
t.Run("same batch different stamp index", func(t *testing.T) {
1569+
t.Parallel()
1570+
1571+
chLow := sOlder.Chunk().WithStamp(postagetesting.MustNewFields(batchA.ID, 0, 5))
1572+
chHigh := sNewer.Chunk().WithStamp(postagetesting.MustNewFields(batchA.ID, 1, 5))
1573+
hashLow, err := chLow.Stamp().Hash()
1574+
if err != nil {
1575+
t.Fatal(err)
1576+
}
1577+
hashHigh, err := chHigh.Stamp().Hash()
1578+
if err != nil {
1579+
t.Fatal(err)
1580+
}
1581+
var winner, loser swarm.Chunk
1582+
if bytes.Compare(hashLow, hashHigh) < 0 {
1583+
winner, loser = chLow, chHigh
1584+
} else {
1585+
winner, loser = chHigh, chLow
1586+
}
1587+
1588+
for _, order := range [][]swarm.Chunk{{winner, loser}, {loser, winner}} {
1589+
baseAddr := swarm.RandAddress(t)
1590+
ts := internal.NewInmemStorage()
1591+
r, err := reserve.New(baseAddr, ts, 0, kademlia.NewTopologyDriver(), log.Noop)
1592+
if err != nil {
1593+
t.Fatal(err)
1594+
}
1595+
if err := r.Put(ctx, order[0]); err != nil {
1596+
t.Fatal(err)
1597+
}
1598+
_ = r.Put(ctx, order[1])
1599+
1600+
got, err := ts.ChunkStore().Get(ctx, winner.Address())
1601+
if err != nil {
1602+
t.Fatal(err)
1603+
}
1604+
if !bytes.Equal(got.Data(), winner.Data()) {
1605+
t.Fatal("expected payload from the lower stamp-hash claim")
1606+
}
1607+
}
1608+
})
1609+
}
1610+
14501611
// TestSOCDivergence covers two single owner chunks that share an address, batch
14511612
// and stamp while wrapping different content. Both are valid, so the storage
14521613
// layer settles which one the neighborhood keeps, and it must reach the same

0 commit comments

Comments
 (0)