174 строки
4.9 KiB
Go
174 строки
4.9 KiB
Go
package workercluster
|
|
|
|
import (
|
|
"testing"
|
|
"time"
|
|
|
|
"github.com/hashicorp/raft"
|
|
"github.com/stretchr/testify/assert"
|
|
"github.com/stretchr/testify/require"
|
|
)
|
|
|
|
// TestStoreRoundTrip_Log writes a single log entry, reads it back, and
|
|
// verifies every field matches. This is the contract the raft library
|
|
// relies on.
|
|
func TestStoreRoundTrip_Log(t *testing.T) {
|
|
s, err := NewBoltStore(t.TempDir())
|
|
require.NoError(t, err)
|
|
defer s.Close() //nolint:errcheck
|
|
|
|
store := s.LogStore()
|
|
|
|
in := &raft.Log{
|
|
Index: 5,
|
|
Term: 3,
|
|
Type: raft.LogCommand,
|
|
Data: []byte(`{"type":"config.adopt","adopted":{"version":1}}`),
|
|
AppendedAt: nowRFC3339(),
|
|
}
|
|
require.NoError(t, store.StoreLog(in))
|
|
|
|
got := &raft.Log{}
|
|
require.NoError(t, store.GetLog(5, got))
|
|
assert.Equal(t, in.Index, got.Index)
|
|
assert.Equal(t, in.Term, got.Term)
|
|
assert.Equal(t, in.Type, got.Type)
|
|
assert.Equal(t, in.Data, got.Data)
|
|
assert.True(t, in.AppendedAt.Equal(got.AppendedAt))
|
|
}
|
|
|
|
// TestStoreRoundTrip_FirstLastIndex exercises the index range after
|
|
// StoreLog and DeleteRange.
|
|
func TestStoreRoundTrip_FirstLastIndex(t *testing.T) {
|
|
s, err := NewBoltStore(t.TempDir())
|
|
require.NoError(t, err)
|
|
defer s.Close() //nolint:errcheck
|
|
|
|
store := s.LogStore()
|
|
|
|
first, err := store.FirstIndex()
|
|
require.NoError(t, err)
|
|
assert.EqualValues(t, 0, first)
|
|
|
|
last, err := store.LastIndex()
|
|
require.NoError(t, err)
|
|
assert.EqualValues(t, 0, last)
|
|
|
|
for i := uint64(10); i <= 20; i++ {
|
|
require.NoError(t, store.StoreLog(&raft.Log{Index: i, Term: 1, Type: raft.LogCommand, Data: []byte{byte(i)}}))
|
|
}
|
|
|
|
first, err = store.FirstIndex()
|
|
require.NoError(t, err)
|
|
assert.EqualValues(t, 10, first)
|
|
last, err = store.LastIndex()
|
|
require.NoError(t, err)
|
|
assert.EqualValues(t, 20, last)
|
|
|
|
// Delete a middle range.
|
|
require.NoError(t, store.DeleteRange(12, 15))
|
|
|
|
first, err = store.FirstIndex()
|
|
require.NoError(t, err)
|
|
assert.EqualValues(t, 10, first)
|
|
last, err = store.LastIndex()
|
|
require.NoError(t, err)
|
|
assert.EqualValues(t, 20, last, "LastIndex should still be the highest stored")
|
|
|
|
// The deleted slot should now report ErrLogNotFound.
|
|
err = store.GetLog(13, &raft.Log{})
|
|
assert.ErrorIs(t, err, raft.ErrLogNotFound)
|
|
|
|
// Surviving slot still readable.
|
|
got := &raft.Log{}
|
|
require.NoError(t, store.GetLog(11, got))
|
|
assert.Equal(t, byte(11), got.Data[0])
|
|
}
|
|
|
|
// TestStoreRoundTrip_StableStore checks the StableStore contract: Set /
|
|
// Get / SetUint64 / GetUint64.
|
|
func TestStoreRoundTrip_StableStore(t *testing.T) {
|
|
s, err := NewBoltStore(t.TempDir())
|
|
require.NoError(t, err)
|
|
defer s.Close() //nolint:errcheck
|
|
|
|
ss := s.StableStore()
|
|
|
|
require.NoError(t, ss.SetUint64([]byte("term"), 7))
|
|
require.NoError(t, ss.Set([]byte("voted_for"), []byte("w-1")))
|
|
|
|
term, err := ss.GetUint64([]byte("term"))
|
|
require.NoError(t, err)
|
|
assert.EqualValues(t, 7, term)
|
|
|
|
voted, err := ss.Get([]byte("voted_for"))
|
|
require.NoError(t, err)
|
|
assert.Equal(t, []byte("w-1"), voted)
|
|
|
|
// Unknown key returns zero value, not error.
|
|
missing, err := ss.GetUint64([]byte("nope"))
|
|
require.NoError(t, err)
|
|
assert.EqualValues(t, 0, missing)
|
|
}
|
|
|
|
// TestStoreRoundTrip_SnapshotStore_CreateOpenList exercises the
|
|
// SnapshotStore lifecycle: Create, write, Close, then List + Open.
|
|
func TestStoreRoundTrip_SnapshotStore_CreateOpenList(t *testing.T) {
|
|
s, err := NewBoltStore(t.TempDir())
|
|
require.NoError(t, err)
|
|
defer s.Close() //nolint:errcheck
|
|
|
|
ss := s.SnapshotStore()
|
|
|
|
cfg := raft.Configuration{Servers: []raft.Server{{ID: "w-1", Address: "127.0.0.1:1", Suffrage: raft.Voter}}}
|
|
sink, err := ss.Create(raft.SnapshotVersionMax, 100, 4, cfg, 50, nil)
|
|
require.NoError(t, err)
|
|
|
|
body := []byte(`{"hello":"world","config_count":3}`)
|
|
_, err = sink.Write(body)
|
|
require.NoError(t, err)
|
|
require.NoError(t, sink.Close())
|
|
|
|
list, err := ss.List()
|
|
require.NoError(t, err)
|
|
require.Len(t, list, 1)
|
|
assert.EqualValues(t, 100, list[0].Index)
|
|
assert.EqualValues(t, 4, list[0].Term)
|
|
|
|
meta, r, err := ss.Open(list[0].ID)
|
|
require.NoError(t, err)
|
|
defer r.Close() //nolint:errcheck
|
|
|
|
buf := make([]byte, len(body))
|
|
n, err := r.Read(buf)
|
|
require.NoError(t, err)
|
|
assert.Equal(t, len(body), n)
|
|
assert.Equal(t, body, buf[:n])
|
|
|
|
assert.EqualValues(t, 100, meta.Index)
|
|
}
|
|
|
|
// TestStoreRoundTrip_SnapshotStore_DeleteRangeAfterSnapshot ensures the
|
|
// log store + snapshot store coexist on disk under the same bbolt file.
|
|
func TestStoreRoundTrip_SnapshotStore_DeleteRangeAfterSnapshot(t *testing.T) {
|
|
s, err := NewBoltStore(t.TempDir())
|
|
require.NoError(t, err)
|
|
defer s.Close() //nolint:errcheck
|
|
|
|
ls := s.LogStore()
|
|
for i := uint64(1); i <= 50; i++ {
|
|
require.NoError(t, ls.StoreLog(&raft.Log{Index: i, Term: 1, Type: raft.LogCommand}))
|
|
}
|
|
require.NoError(t, ls.DeleteRange(1, 25))
|
|
|
|
first, err := ls.FirstIndex()
|
|
require.NoError(t, err)
|
|
assert.EqualValues(t, 26, first)
|
|
}
|
|
|
|
// nowRFC3339 returns a fixed-format RFC3339 timestamp for tests so
|
|
// append log entries have a deterministic AppendedAt.
|
|
func nowRFC3339() time.Time {
|
|
return time.Date(2026, 6, 27, 12, 0, 0, 0, time.UTC)
|
|
}
|