Files
worker/internal/workercluster/store_test.go
Gleb Tv 2c7a0236da feat: publish standalone worker
Separate worker packaging and service lifecycle from the control plane.
2026-07-13 17:55:14 +03:00

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)
}