Files
mesh-controller/internal/conditions/memory.go
T
jochen fe857ba081
mesh/merge-gate pass: builds build-agent, mesh-controller, route-proxy → ace, g14, novox, shanks; no bus step; every machine composes with the change as it…
mesh/repo-check pass: its merge-check.sh passed
mesh/delivery delivered
Switch the memory store's failure under its lock, so a test cannot race the keeper
The keeper keeps each transition from its own goroutine, which reads the memory
store's Fail field under the store's lock; tests assigned the exported field
bare, so TestAnUnreadableStoreClearsNothing failed under -race whenever the
goroutine appended in that window. The field is now set only through SetFail
(and Told's likewise), and a test makes the race certain rather than rare.
Test-only: the controller runs the bus store, never InMemory.
2026-10-09 01:25:04 +02:00

165 lines
3.7 KiB
Go

package conditions
import (
"context"
"encoding/json"
"errors"
"sort"
"sync"
"time"
)
// InMemory is a store and a history held in this process: for tests, and for nothing else — a
// condition kept here is forgotten by a restart, which is the fault the store exists to remove.
type InMemory struct {
mu sync.Mutex
values map[string]Entry
revision uint64
events []Event
// fail, when set, is what every read and write answers: a store that is away. It is set only
// through SetFail, under the lock, because the keeper's telling goroutine reads it while a test
// takes the store away.
fail error
}
// NewInMemory is an empty store.
func NewInMemory() *InMemory { return &InMemory{values: map[string]Entry{}} }
// SetFail makes every read and write answer err from now on, or none when err is nil.
func (m *InMemory) SetFail(err error) {
m.mu.Lock()
defer m.mu.Unlock()
m.fail = err
}
func (m *InMemory) Get(_ context.Context, key string) (Entry, bool, error) {
m.mu.Lock()
defer m.mu.Unlock()
if m.fail != nil {
return Entry{}, false, m.fail
}
e, ok := m.values[key]
return e, ok, nil
}
func (m *InMemory) Create(_ context.Context, key string, value []byte) error {
m.mu.Lock()
defer m.mu.Unlock()
if m.fail != nil {
return m.fail
}
if _, ok := m.values[key]; ok {
return ErrMoved
}
m.revision++
m.values[key] = Entry{Value: value, Revision: m.revision}
return nil
}
func (m *InMemory) Update(_ context.Context, key string, value []byte, revision uint64) error {
m.mu.Lock()
defer m.mu.Unlock()
if m.fail != nil {
return m.fail
}
if e, ok := m.values[key]; !ok || e.Revision != revision {
return ErrMoved
}
m.revision++
m.values[key] = Entry{Value: value, Revision: m.revision}
return nil
}
func (m *InMemory) Delete(_ context.Context, key string, revision uint64) error {
m.mu.Lock()
defer m.mu.Unlock()
if m.fail != nil {
return m.fail
}
if e, ok := m.values[key]; !ok || e.Revision != revision {
return ErrMoved
}
delete(m.values, key)
return nil
}
func (m *InMemory) All(context.Context) (map[string]Entry, error) {
m.mu.Lock()
defer m.mu.Unlock()
if m.fail != nil {
return nil, m.fail
}
out := make(map[string]Entry, len(m.values))
for k, v := range m.values {
out[k] = v
}
return out, nil
}
func (m *InMemory) Append(_ context.Context, e Event) error {
m.mu.Lock()
defer m.mu.Unlock()
if m.fail != nil {
return m.fail
}
m.events = append(m.events, e)
return nil
}
func (m *InMemory) Since(_ context.Context, since time.Time) ([]Event, error) {
m.mu.Lock()
defer m.mu.Unlock()
if m.fail != nil {
return nil, m.fail
}
var out []Event
for _, e := range m.events {
if !e.At.Before(since) {
out = append(out, e)
}
}
sort.SliceStable(out, func(i, j int) bool { return out[i].At.Before(out[j].At) })
return out, nil
}
// Told is a teller that remembers what it was told, for tests.
type Told struct {
mu sync.Mutex
Events []Event
Names []string
// fail, when set, is what every publish answers; set only through SetFail, under the lock.
fail error
}
// SetFail makes every publish answer err from now on, or none when err is nil.
func (t *Told) SetFail(err error) {
t.mu.Lock()
defer t.mu.Unlock()
t.fail = err
}
func (t *Told) PublishSeatEvent(_ context.Context, seat, event string, body []byte) error {
t.mu.Lock()
defer t.mu.Unlock()
if t.fail != nil {
return t.fail
}
if seat != Seat {
return errors.New("told under the wrong seat: " + seat)
}
var e Event
if err := json.Unmarshal(body, &e); err != nil {
return err
}
t.Events = append(t.Events, e)
t.Names = append(t.Names, event)
return nil
}
// Said is a copy of what was told so far.
func (t *Told) Said() []Event {
t.mu.Lock()
defer t.mu.Unlock()
return append([]Event(nil), t.Events...)
}