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
59 changes: 43 additions & 16 deletions server/filestore.go
Original file line number Diff line number Diff line change
Expand Up @@ -1520,8 +1520,8 @@ func (mb *msgBlock) rebuildStateLocked() (*LostStreamData, []uint64, error) {
// For tombstones that we find and collect.
var (
tombstones []uint64
minTombstoneSeq uint64
minTombstoneTs int64
maxTombstoneSeq uint64
maxTombstoneTs int64
)

// To detect gaps from compaction, and to ensure the sequence keeps moving up.
Expand Down Expand Up @@ -1580,8 +1580,8 @@ func (mb *msgBlock) rebuildStateLocked() (*LostStreamData, []uint64, error) {
seq = seq &^ tbit
// Need to process this here and make sure we have accounted for this properly.
tombstones = append(tombstones, seq)
if minTombstoneSeq == 0 || seq < minTombstoneSeq {
minTombstoneSeq, minTombstoneTs = seq, ts
if maxTombstoneSeq == 0 || seq > maxTombstoneSeq {
maxTombstoneSeq, maxTombstoneTs = seq, ts
}
index += rl
continue
Expand Down Expand Up @@ -1664,12 +1664,12 @@ func (mb *msgBlock) rebuildStateLocked() (*LostStreamData, []uint64, error) {
fseq := atomic.LoadUint64(&mb.first.seq)
if fseq > 0 {
atomic.StoreUint64(&mb.last.seq, fseq-1)
} else if fseq == 0 && minTombstoneSeq > 0 {
atomic.StoreUint64(&mb.first.seq, minTombstoneSeq+1)
} else if fseq == 0 && maxTombstoneSeq > 0 {
atomic.StoreUint64(&mb.first.seq, maxTombstoneSeq+1)
mb.first.ts = 0
if mb.last.seq == 0 {
atomic.StoreUint64(&mb.last.seq, minTombstoneSeq)
mb.last.ts = minTombstoneTs
atomic.StoreUint64(&mb.last.seq, maxTombstoneSeq)
mb.last.ts = maxTombstoneTs
}
}
}
Expand Down Expand Up @@ -2285,6 +2285,12 @@ func (fs *fileStore) recoverMsgs() error {
fs.removeMsgBlockFromList(mb)
continue
}
// If the stream is empty, reset the first/last sequences so these can
// properly move up based purely on tombstones spread over multiple blocks.
if fs.state.Msgs == 0 {
fs.state.FirstSeq, fs.state.LastSeq = 0, 0
fs.state.FirstTime, fs.state.LastTime = time.Time{}, time.Time{}
}
fseq := atomic.LoadUint64(&mb.first.seq)
if fs.state.FirstSeq == 0 || (fseq < fs.state.FirstSeq && mb.first.ts != 0) {
fs.state.FirstSeq = fseq
Expand Down Expand Up @@ -5064,9 +5070,9 @@ func (fs *fileStore) removeMsg(seq uint64, secure, viaLimits, needFSLock bool) (

// If erase but block is empty, we can simply remove the block later.
if secure && !isEmpty {
// Grab record info.
ri, rl, _, _ := mb.slotInfo(int(seq - mb.cache.fseq))
if err := mb.eraseMsg(seq, int(ri), int(rl), isLastBlock); err != nil {
// Grab record info, but use the pre-computed record length.
ri, _, _, _ := mb.slotInfo(int(seq - mb.cache.fseq))
if err := mb.eraseMsg(seq, int(ri), int(msz), isLastBlock); err != nil {
mb.finishedWithCache()
return false, err
}
Expand Down Expand Up @@ -5665,9 +5671,21 @@ func (mb *msgBlock) selectNextFirst() {
}

// Select the next FirstSeq
// Also cleans up empty blocks at the start only containing tombstones.
// Lock should be held.
func (fs *fileStore) selectNextFirst() {
if len(fs.blks) > 0 {
for len(fs.blks) > 1 {
mb := fs.blks[0]
mb.mu.Lock()
empty := mb.msgs == 0
if !empty {
mb.mu.Unlock()
break
}
fs.forceRemoveMsgBlock(mb)
mb.mu.Unlock()
}
mb := fs.blks[0]
mb.mu.RLock()
fs.state.FirstSeq = atomic.LoadUint64(&mb.first.seq)
Expand Down Expand Up @@ -9138,7 +9156,7 @@ func (fs *fileStore) Truncate(seq uint64) error {
}
mb.mu.Lock()
}
fs.removeMsgBlock(mb)
fs.forceRemoveMsgBlock(mb)
mb.mu.Unlock()
}

Expand All @@ -9160,7 +9178,7 @@ func (fs *fileStore) Truncate(seq uint64) error {
}
smb.mu.Lock()
}
fs.removeMsgBlock(smb)
fs.forceRemoveMsgBlock(smb)
smb.mu.Unlock()
goto SKIP
}
Expand Down Expand Up @@ -9202,7 +9220,7 @@ SKIP:
if !hasWrittenTombstones {
fs.lmb = smb
tmb.mu.Lock()
fs.removeMsgBlock(tmb)
fs.forceRemoveMsgBlock(tmb)
tmb.mu.Unlock()
}

Expand Down Expand Up @@ -9282,17 +9300,26 @@ func (fs *fileStore) removeMsgBlockFromList(mb *msgBlock) {
// Both locks should be held.
func (fs *fileStore) removeMsgBlock(mb *msgBlock) {
// Check for us being last message block
lseq, lts := atomic.LoadUint64(&mb.last.seq), mb.last.ts
if mb == fs.lmb {
lseq, lts := atomic.LoadUint64(&mb.last.seq), mb.last.ts
// Creating a new message write block requires that the lmb lock is not held.
mb.mu.Unlock()
// Write the tombstone to remember since this was last block.
if lmb, _ := fs.newMsgBlockForWrite(); lmb != nil {
fs.writeTombstone(lseq, lts)
}
mb.mu.Lock()
} else if lseq == fs.state.LastSeq {
// Need to write a tombstone for the last sequence if we're removing the block containing it.
fs.writeTombstone(lseq, lts)
}
// Only delete message block after (potentially) writing a new lmb.
// Only delete message block after (potentially) writing a tombstone.
fs.forceRemoveMsgBlock(mb)
}

// Removes the msgBlock, without writing tombstones to ensure the last sequence is preserved.
// Both locks should be held.
func (fs *fileStore) forceRemoveMsgBlock(mb *msgBlock) {
mb.dirtyCloseWithRemove(true)
fs.removeMsgBlockFromList(mb)
}
Expand Down
233 changes: 233 additions & 0 deletions server/filestore_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -10613,3 +10613,236 @@ func TestFileStoreCacheLookupOnEmptyBlock(t *testing.T) {
require_True(t, lmb.cache == nil)
})
}

func TestFileStoreEraseMsgDoesNotLoseTombstones(t *testing.T) {
testFileStoreAllPermutations(t, func(t *testing.T, fcfg FileStoreConfig) {
cfg := StreamConfig{Name: "zzz", Subjects: []string{"foo"}, Storage: FileStorage}
created := time.Now()
fs, err := newFileStoreWithCreated(fcfg, cfg, created, prf(&fcfg), nil)
require_NoError(t, err)
defer fs.Stop()

secret := []byte("secret!")
// The first message will remain throughout.
_, _, err = fs.StoreMsg("foo", nil, nil, 0)
require_NoError(t, err)
// The second message wil be removed, so a tombstone will be placed.
_, _, err = fs.StoreMsg("foo", nil, nil, 0)
require_NoError(t, err)
// The third message is secret and will be erased.
_, _, err = fs.StoreMsg("foo", nil, secret, 0)
require_NoError(t, err)

// Removing the second message places a tombstone.
_, err = fs.RemoveMsg(2)
require_NoError(t, err)

// A fourth message gets placed after the tombstone.
_, _, err = fs.StoreMsg("foo", nil, nil, 0)
require_NoError(t, err)

// Now we erase the third message.
// This erases this message and should not lose the tombstone that comes after it.
_, err = fs.EraseMsg(3)
require_NoError(t, err)

before := fs.State()
require_Equal(t, before.Msgs, 2)
require_Equal(t, before.FirstSeq, 1)
require_Equal(t, before.LastSeq, 4)
require_True(t, slices.Equal(before.Deleted, []uint64{2, 3}))

_, err = fs.LoadMsg(2, nil)
require_Error(t, err, errDeletedMsg)
_, err = fs.LoadMsg(3, nil)
require_Error(t, err, errDeletedMsg)

// Make sure we can recover properly with no index.db present.
fs.Stop()
os.Remove(filepath.Join(fs.fcfg.StoreDir, msgDir, streamStreamStateFile))

fs, err = newFileStoreWithCreated(fcfg, cfg, created, prf(&fcfg), nil)
require_NoError(t, err)
defer fs.Stop()

if state := fs.State(); !reflect.DeepEqual(state, before) {
t.Fatalf("Expected state\n of %+v, \ngot %+v without index.db state", before, state)
}

_, err = fs.LoadMsg(2, nil)
require_Error(t, err, errDeletedMsg)
_, err = fs.LoadMsg(3, nil)
require_Error(t, err, errDeletedMsg)
})
}

func TestFileStoreTombstonesNoFirstSeqRollback(t *testing.T) {
testFileStoreAllPermutations(t, func(t *testing.T, fcfg FileStoreConfig) {
fcfg.BlockSize = 10 * 33 // 10 messages per block.
cfg := StreamConfig{Name: "zzz", Subjects: []string{"foo"}, Storage: FileStorage}
created := time.Now()
fs, err := newFileStoreWithCreated(fcfg, cfg, created, prf(&fcfg), nil)
require_NoError(t, err)
defer fs.Stop()

for i := 0; i < 20; i++ {
_, _, err = fs.StoreMsg("foo", nil, nil, 0)
require_NoError(t, err)
}

before := fs.State()
require_Equal(t, before.Msgs, 20)
require_Equal(t, before.FirstSeq, 1)
require_Equal(t, before.LastSeq, 20)

// Expect 2 blocks with messages.
fs.mu.RLock()
lblks := len(fs.blks)
fs.mu.RUnlock()
require_Equal(t, lblks, 2)

// Write some tombstones for all messages, these will be in multiple blocks.
for seq := uint64(1); seq <= 20; seq++ {
_, err = fs.RemoveMsg(seq)
require_NoError(t, err)
}

before = fs.State()
require_Equal(t, before.Msgs, 0)
require_Equal(t, before.FirstSeq, 21)
require_Equal(t, before.LastSeq, 20)

// Expect 1 block purely with tombstones.
fs.mu.RLock()
lblks = len(fs.blks)
fs.mu.RUnlock()
require_Equal(t, lblks, 1)

// Make sure we can recover properly with no index.db present.
fs.Stop()
os.Remove(filepath.Join(fs.fcfg.StoreDir, msgDir, streamStreamStateFile))

fs, err = newFileStoreWithCreated(fcfg, cfg, created, prf(&fcfg), nil)
require_NoError(t, err)
defer fs.Stop()

if state := fs.State(); !reflect.DeepEqual(state, before) {
t.Fatalf("Expected state\n of %+v, \ngot %+v without index.db state", before, state)
}
})
}

func TestFileStoreTombstonesSelectNextFirstCleanup(t *testing.T) {
testFileStoreAllPermutations(t, func(t *testing.T, fcfg FileStoreConfig) {
fcfg.BlockSize = 10 * 33 // 10 messages per block.
cfg := StreamConfig{Name: "zzz", Subjects: []string{"foo"}, Storage: FileStorage}
created := time.Now()
fs, err := newFileStoreWithCreated(fcfg, cfg, created, prf(&fcfg), nil)
require_NoError(t, err)
defer fs.Stop()

// Write a bunch of messages in multiple blocks.
for i := 0; i < 50; i++ {
_, _, err = fs.StoreMsg("foo", nil, nil, 0)
require_NoError(t, err)
}

for seq := uint64(2); seq <= 49; seq++ {
_, err = fs.RemoveMsg(seq)
require_NoError(t, err)
}

_, err = fs.newMsgBlockForWrite()
require_NoError(t, err)
for i := 0; i < 50; i++ {
_, _, err = fs.StoreMsg("foo", nil, nil, 0)
require_NoError(t, err)
}

for seq := uint64(50); seq <= 100; seq++ {
_, err = fs.RemoveMsg(seq)
require_NoError(t, err)
}

before := fs.State()
require_Equal(t, before.Msgs, 1)
require_Equal(t, before.FirstSeq, 1)
require_Equal(t, before.LastSeq, 100)

_, err = fs.RemoveMsg(1)
require_NoError(t, err)

before = fs.State()
require_Equal(t, before.Msgs, 0)
require_Equal(t, before.FirstSeq, 101)
require_Equal(t, before.LastSeq, 100)

// Make sure we can recover properly with no index.db present.
fs.Stop()
os.Remove(filepath.Join(fs.fcfg.StoreDir, msgDir, streamStreamStateFile))

fs, err = newFileStoreWithCreated(fcfg, cfg, created, prf(&fcfg), nil)
require_NoError(t, err)
defer fs.Stop()

if state := fs.State(); !reflect.DeepEqual(state, before) {
t.Fatalf("Expected state\n of %+v, \ngot %+v without index.db state", before, state)
}
})
}

func TestFileStoreTombstonesSelectNextFirstCleanupOnRecovery(t *testing.T) {
testFileStoreAllPermutations(t, func(t *testing.T, fcfg FileStoreConfig) {
fcfg.BlockSize = 10 * 33 // 10 messages per block.
cfg := StreamConfig{Name: "zzz", Subjects: []string{"foo"}, Storage: FileStorage}
created := time.Now()
fs, err := newFileStoreWithCreated(fcfg, cfg, created, prf(&fcfg), nil)
require_NoError(t, err)
defer fs.Stop()

// Write a bunch of messages in multiple blocks.
for i := 0; i < 50; i++ {
_, _, err = fs.StoreMsg("foo", nil, nil, 0)
require_NoError(t, err)
}

for seq := uint64(2); seq <= 49; seq++ {
_, err = fs.RemoveMsg(seq)
require_NoError(t, err)
}

_, err = fs.newMsgBlockForWrite()
require_NoError(t, err)
for i := 0; i < 50; i++ {
_, _, err = fs.StoreMsg("foo", nil, nil, 0)
require_NoError(t, err)
}

for seq := uint64(50); seq <= 100; seq++ {
_, err = fs.RemoveMsg(seq)
require_NoError(t, err)
}

before := fs.State()
require_Equal(t, before.Msgs, 1)
require_Equal(t, before.FirstSeq, 1)
require_Equal(t, before.LastSeq, 100)

// Explicitly write tombstone instead of calling fs.RemoveMsg,
// so we need to recover from a hard kill.
require_NoError(t, fs.writeTombstone(1, 0))
before = StreamState{FirstSeq: 101, FirstTime: time.Time{}, LastSeq: 100, LastTime: before.LastTime}

// Make sure we can recover properly with no index.db present.
fs.Stop()
os.Remove(filepath.Join(fs.fcfg.StoreDir, msgDir, streamStreamStateFile))

fs, err = newFileStoreWithCreated(fcfg, cfg, created, prf(&fcfg), nil)
require_NoError(t, err)
defer fs.Stop()

if state := fs.State(); !reflect.DeepEqual(state, before) {
t.Fatalf("Expected state\n of %+v, \ngot %+v without index.db state", before, state)
}
})
}
Loading