mirror of
https://github.com/mudler/LocalAI.git
synced 2026-10-05 04:24:39 -04:00
* feat(messaging): add shared subject rules Assisted-by: Claude:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io> * test(messaging): cover BroadcastRoots, ControlRoots and SubjectRoot Assisted-by: Claude:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io> * feat(messaging): add Broadcaster and enforce subject rules in every carrier Broadcaster is the fan-out half of MessagingClient. The NATS client and the in-memory FakeBus now refuse a subject outside the served roots and any wildcard other than a whole single token, and FakeBus shares MatchSubject instead of its own copy. FakeBus Unsubscribe now removes its own subscription instead of the first one with the same subject. A shared conformance suite in messagingtest runs against both carriers. The distributed e2e specs that used invented test.* subjects, and the one that subscribed with a > filter, now use subjects from subjects.go. Assisted-by: Claude:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io> * refactor: depend on Broadcaster where only publish and subscribe are used Narrowed to messaging.Broadcaster: nodes/staging_progress.go, nodes/install_progress_publisher.go, galleryop/operation.go, galleryop/service.go, agentpool/user_services.go, agentpool/agent_jobs.go, openresponses/store.go, openresponses/sync.go, syncstate/syncstate.go, finetune/service.go, quantization/service.go and failover/distsync/distsync.go. SubscribeJSON now takes a Broadcaster because it only calls Subscribe, which lets the narrowed consumers use it. Stayed wide: worker/supervisor.go, because its client field also serves the SubscribeReply handlers in worker/lifecycle.go. The request/reply, queue and wiring files (nodes/unloader.go, nodes/file_stager_s3.go, jobs/dispatcher.go, agents/dispatcher.go, agents/events.go, worker/file_staging.go, cli/agent_worker.go, http/app.go) are unchanged by design. Assisted-by: Claude:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io> * refactor(nodes): name the no-route condition and confine the carrier error Consumers matched nats.ErrNoResponders, which names an absence, to demote a node. They now match ErrNoRoute, the control path maps the carrier's failure onto it, and timeouts and worker refusals are pinned as not being no-route. Assisted-by: Claude:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io> * docs(nodes): state which FileStager implementations return ErrNoRoute Assisted-by: Claude:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io> * refactor(nodes): build backend clients through one node-aware seam Assisted-by: Claude:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io> * docs: describe the distributed transport seams Assisted-by: Claude:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io> * docs: correct comments that overclaim after the seams refactor Assisted-by: Claude:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io> * fix(agent-worker): refuse an unserved LOCALAI_AGENT_SUBJECT at startup The messaging client now refuses a subject whose root no carrier serves. An agent worker started with a custom LOCALAI_AGENT_SUBJECT such as tenant-a.agent.execute used to start and then wait on a subject the frontend never publishes to. After the subject rules landed it exited at subscribe time with an error that did not name the setting. Behaviour change: the worker now checks LOCALAI_AGENT_SUBJECT before it registers or connects, and exits with an error that names the variable and says to use a served subject under the agent root, for example agent.execute. The served roots are not widened: a custom root was never delivered by the frontend, and a wider set would reopen the drift the subject rules exist to close. The flag help and the agent worker docs state the constraint. Assisted-by: Claude:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io> * test(nodes): pin the reactions to ErrNoRoute Three callers react to ErrNoRoute and had no spec: the reconciler's upgrade drain falls back to the legacy forced install, the reconciler marks the node unhealthy when a pending op has no route, and the backend-op fan-out marks the node unhealthy. Each spec drives the real caller with a scripted no-responders reply and reads the result from the registry or the recorded requests. A fourth spec pins the other side: a pending op that times out leaves the node healthy and only counts the attempt, so mapping timeouts onto ErrNoRoute would fail here. Assisted-by: Claude:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io> * test(messaging): pin client subject checks and fail the carrier suite in CI Add specs that call Publish, Request, Subscribe, QueueSubscribe, SubscribeReply and QueueSubscribeReply on a client with no connection. Each call must return ErrUnservedSubject for bogus.thing and ErrUnsupportedWildcard for jobs.>. This proves that the subject check runs before the connection is used, and needs no server. The NATS conformance suite is the only check that runs the subject rules against a real carrier. Before this change it skipped without output when Docker was missing. Now it fails when CI is set, so a Linux runner without Docker cannot hide it. It still skips on local runs and on macOS CI, which has no Docker. Add SubjectNodeBackendInstallProgress to the list of constructors that must build served subjects, and ask contributors to extend the list. Assisted-by: Claude:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io> * docs: state what ErrNoRoute may change, and group the distributed guides The seams note said MarkUnhealthy was the only state change allowed on ErrNoRoute. A pending backend op still records the failed attempt, counts toward the reconciler's retry limit and is dead-lettered after the maximum attempts. The note now says that MarkUnhealthy is the only change to the node's own state, and that the per-op accounting is not a verdict about the node. The note also documents that the NATS conformance run fails under CI when Docker is missing. The distributed-seams row moves next to the distributed-state row in the topics table. The liveness ping spec header now says no route is a reason to skip the worker, not proof that the worker is gone. Assisted-by: Claude:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io> * refactor(nodes): give the backend client factory the node id Mechanical: the method gains a nodeID parameter and the eight test fakes are updated. No behaviour change. Assisted-by: Claude:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io> * refactor(nodes): drop the optional node-aware factory The node id is now in the main method, so the optional interface and its helper had no behaviour of their own. Assisted-by: Claude:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io> * refactor(nodes): dial backend probes through the client factory Assisted-by: Claude:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io> * refactor(nodes): dial workers' file servers through a per-node dialer Assisted-by: Claude:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io> * refactor(http): proxy backend logs through the per-node worker dialer The admin backend-logs proxy (list, lines and the WebSocket stream) now reaches a worker through the same per-node dialer as the HTTP file stager, so every frontend-to-worker dial goes through one seam. The shared direct dialer keeps alive for 15s where the proxy used 30s. Harmless for requests bounded at 15s. Assisted-by: Claude:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io> * fix(http): keep the backend-logs proxy independent of the admin connection The proxy request had no context before the dialer change and is bounded only by its 15s timeout. Keep it that way. Assisted-by: Claude:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io> * refactor: move the worker control payloads to workerctl Mechanical move of the request and reply structs, the install progress event and the file payloads out of messaging. The verbs no longer belong to one carrier. No alias is left behind. Assisted-by: Claude:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io> * refactor(worker): serve the lifecycle verbs through a controlServer The worker registers one handler per verb and a NATS server maps each verb to its subject. Registration errors now name the verb. node.stop is served with SubscribeReply, which is identical on the wire because the handler never replies. Assisted-by: Claude:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io> * refactor(worker): report install progress through the control sink Install and upgrade now emit download progress through the sink the control server hands them. The debounce and the terminal flush stay in the handler path, built over that sink by the new nodes.NewDebouncedInstallProgressSink, which replaces NewDebouncedInstallProgressPublisher. The subject and payload on the wire are unchanged. The supervisor no longer holds the bus, and installFn and upgradeFn let specs drive both verbs without a gallery. The malformed-request log lines are restored for install, upgrade, backend.delete, model.unload, model.stop and model.delete, with the reply bytes unchanged. The signal adapter is renamed noReply, which also lets worker.go import os/signal without an alias again. Assisted-by: Claude:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io> * refactor(worker): serve the file-staging verbs through a controlServer An empty list-dir answer is now {} rather than {"files":null}, because the typed reply omits an empty Files slice. The frontend decodes both to a nil slice in nodes/file_stager_s3.go. Assisted-by: Claude:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io> * feat(messaging): add WorkQueue and the NATS producer This is the producer side of the competing-consumer seam. The work kinds map one to one to today's subjects and queue groups: task to jobs.new and mcp-ci to jobs.mcp-ci.new (both in group workers), agent-run to agent.execute (group agent-workers). Enqueue publishes the payload as Publish does today, with one JSON marshal. FakeBus now records queue groups and keeps reply handlers so later specs can pin and drive them. Assisted-by: Claude:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io> * feat(messaging): add the NATS WorkConsumer An in-flight limit of one runs the handler inline on the delivery goroutine, as the MCP CI consumer does today. Any other limit spawns per delivery, as the agent consumer does. Queue groups are unchanged. Assisted-by: Claude:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io> * refactor: publish queued work through WorkQueue The job dispatcher, the agent pool and the agent scheduler enqueue through messaging.WorkQueue; the NATS implementation publishes to the same subjects as before. DistributedServices builds the queue next to the NATS client and hands it to the dispatcher and the agent pool, whose distributed mode switch now reads a non-nil WorkQueue. The unused AgentPoolService.SetNATSClient is removed. Assisted-by: Claude:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io> * refactor: consume queued work through WorkConsumer The agent dispatcher and the MCP CI consumer register through messaging.WorkConsumer. The NATS implementation keeps the inline one-at-a-time model for MCP CI and the per-delivery model for agent runs. handleMCPCIJob reports on the events publisher the carrier hands it instead of a captured client. Assisted-by: Claude:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io> * refactor: delete the consumers nothing in production reached jobs.new has a producer and no production consumer, and the agent dispatcher's Dispatch was only called from tests. Publishing jobs.new is unchanged. Assisted-by: Claude:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io> * refactor(mcp): send MCP requests to agent workers through AgentControl Timeouts still honour only the deadline, not cancellation, exactly as today. The NATS no-responders error maps to ErrNoRoute and a timeout does not. Assisted-by: Claude:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io> * refactor(agent-worker): serve MCP requests and backend.stop through agentRPCServer The agent worker's MCP tool and discovery reply subscriptions and its backend stop listener move behind an unexported agentRPCServer interface, served on NATS by nodes.NATSAgentRPCServer. The handlers become typed mcp.ToolHandler and mcp.DiscoveryHandler values that answer every failure with a reply carrying Error. Queue group (agent-workers), inline execution on the delivery goroutine, the background handler context, the unmarshal error reply texts and the reply-less backend stop subscription are unchanged. The backend stop handler takes the decoded backend name, so it can still close that backend's MCP sessions. Assisted-by: Claude:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io> * refactor(messaging): remove helpers that only tests used BroadcastRoots, ControlRoots and SubjectRoot had no production caller. The roots spec now asserts every served root through ValidateSubject instead. MatchSubject moves back into the test support package, the only place that used it, with its table. NATSAgentRPCServer drops the subscription list it stored and never read, and NewNATSAgentRPCServer gets a doc comment. Assisted-by: Claude:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io> * test(mcp): round trip the agent RPC server over a real NATS server One spec sends a tool request and a discovery request through NATSAgentControl to NATSAgentRPCServer and checks that the handlers see the decoded requests and the replies come back. It also puts an undecodable body on the tool subject and checks the server answers with an unmarshal error instead of leaving the requester to time out. Assisted-by: Claude:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io> * docs: describe the distributed transport seams The developer note now lists the final seams: fan-out, queues, both halves of the control verbs and of agent RPC, and the dial. It records the open items a second carrier has to handle. Assisted-by: Claude:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io> * test: pin the in-flight limit each queue consumer asks for The work queue specs pin what Consume does for a given limit, but nothing pinned which limit each production consumer passes. Changing the agent worker's MCP CI limit from 1 to 0 would have let MCP CI jobs run concurrently on each worker with every test green. Move the MCP CI Consume call into startMCPCIConsumer with the same wiring and pin that it asks for (WorkMCPCI, 1). Pin that NATSDispatcher.Start asks for (WorkAgentRun, maxConcurrent) for several limits. Assisted-by: Claude:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io> * refactor: remove helpers the branch left without a caller SubjectJobCancelWildcard lost its last subscriber when the frontend stopped listening on jobs.*.cancel; the NATS permissions and conformance suite spell the subject out, so nothing reads the constant. decodeBackendStopRequest returned a stopAll flag that production dropped and only a test read. decodeBackendStop is now the single decoder with the same semantics: an empty body is stop-all, an empty Backend is stop-all, malformed JSON is an error. stopBackends still derives stop-all from Backend, so no reply changes. Assisted-by: Claude:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io> * fix(messaging): keep an explicitly empty agent queue a plain subscription Before the work queue seam the agent worker passed LOCALAI_AGENT_QUEUE straight to QueueSubscribe, so an explicitly empty value made a plain subscription and every agent worker ran every agent run. WithAgentRunRoute replaced an empty queue with agent-workers, which silently changed that. Keep the queue as given once the option is applied. An empty subject still falls back to agent.execute, since it never had a meaning of its own. The flag default stays agent-workers, so only an explicitly empty value reaches this. Assisted-by: Claude:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io> * docs: correct comments and record the PR B notes Fix the recordingFactory comment (it also records the parallel flag), document that a negative maxInFlight is unbounded and that Unsubscribe from a handler deadlocks, and say a permanently undecodable payload returns nil. Record controlHandler's undecodable return as a kept exception, and add the second carrier notes to the developer note: the reconciler has no ClientFactory option, the logs proxy honours HTTP_PROXY, verbs one carrier serves need an opt-out, terminal replies come from the result event, and agent runs publish through the NATS-bound EventBridge, which is not an additive change. Assisted-by: Claude:claude-sonnet-5-5 Signed-off-by: Ettore Di Giacinto <mudler@localai.io> --------- Signed-off-by: Ettore Di Giacinto <mudler@localai.io> Co-authored-by: Ettore Di Giacinto <mudler@localai.io>
296 lines
9.1 KiB
Go
296 lines
9.1 KiB
Go
// Package syncstate provides SyncedMap, a reusable cross-replica in-memory map.
|
|
//
|
|
// LocalAI in distributed mode runs multiple frontend replicas behind a
|
|
// round-robin load balancer. Several features keep process-local in-memory state
|
|
// that is surfaced to the HTTP/UI API; without cross-replica sync a poll that
|
|
// lands on a replica which did not originate a change sees stale or missing data.
|
|
// SyncedMap collapses the three legs each feature otherwise hand-wires - an
|
|
// in-memory map, a NATS broadcast/apply path, and optional durable read-through -
|
|
// into one well-tested component so cross-replica consistency is a configuration
|
|
// choice rather than a bespoke re-implementation.
|
|
package syncstate
|
|
|
|
import (
|
|
"context"
|
|
"sync"
|
|
"time"
|
|
|
|
"github.com/mudler/LocalAI/core/services/messaging"
|
|
"github.com/mudler/xlog"
|
|
)
|
|
|
|
// Op values carried on the wire and passed to OnApply.
|
|
const (
|
|
opSet = "set"
|
|
opDelete = "delete"
|
|
)
|
|
|
|
// Store is optional durable backing for a SyncedMap. In distributed mode it is a
|
|
// single shared DB, so the apply path (a delta received from a peer) updates
|
|
// memory only and never re-writes the Store.
|
|
type Store[K comparable, V any] interface {
|
|
List(ctx context.Context) ([]V, error)
|
|
Upsert(ctx context.Context, v V) error
|
|
Delete(ctx context.Context, k K) error
|
|
}
|
|
|
|
// Config configures a SyncedMap.
|
|
type Config[K comparable, V any] struct {
|
|
Name string // subject namespace, e.g. "finetune.jobs"
|
|
Key func(V) K // extract the key from a value
|
|
Nats messaging.Broadcaster // nil => standalone: in-memory only, no broadcast/subscribe
|
|
Store Store[K, V] // optional read-through persistence
|
|
Loader func(ctx context.Context) ([]V, error) // source when there is no Store (e.g. disk reload)
|
|
OnApply func(op string, k K, v V) // optional hook after an applied change (e.g. ShutdownModel)
|
|
Reconcile time.Duration // optional periodic re-hydrate; 0 = off
|
|
}
|
|
|
|
// delta is the JSON wire envelope broadcast on every local mutation. Value is
|
|
// omitempty so a delete carries only op+key.
|
|
type delta[K comparable, V any] struct {
|
|
Op string `json:"op"`
|
|
Key K `json:"key"`
|
|
Value V `json:"value,omitempty"`
|
|
}
|
|
|
|
// SyncedMap is a cross-replica in-memory map. A local write (Set/Delete) updates
|
|
// the optional durable Store, then memory, then broadcasts a delta to peers. A peer's
|
|
// delta updates memory only and fires OnApply - it never re-broadcasts and never
|
|
// writes the Store. That structural split is the echo-loop guard (same pattern as
|
|
// galleryop.mergeStatus / OpCache.applyStart): receiving your own broadcast just
|
|
// re-applies an idempotent value to memory, so there is no storm and no
|
|
// double-write.
|
|
type SyncedMap[K comparable, V any] struct {
|
|
cfg Config[K, V]
|
|
|
|
mu sync.RWMutex
|
|
data map[K]V
|
|
|
|
sub Subscription
|
|
|
|
// lifeCtx outlives Start's argument: a reconnect callback or reconcile tick
|
|
// can fire long after Start returns, so they must not be tied to a ctx the
|
|
// caller may cancel. Close cancels it.
|
|
lifeCtx context.Context
|
|
cancel context.CancelFunc
|
|
wg sync.WaitGroup
|
|
}
|
|
|
|
// Subscription is the subset of messaging.Subscription the component holds onto.
|
|
type Subscription = messaging.Subscription
|
|
|
|
// New constructs a SyncedMap. Call Start to hydrate and begin syncing.
|
|
func New[K comparable, V any](cfg Config[K, V]) *SyncedMap[K, V] {
|
|
return &SyncedMap[K, V]{cfg: cfg, data: make(map[K]V)}
|
|
}
|
|
|
|
func (m *SyncedMap[K, V]) subject() string {
|
|
return messaging.SubjectSyncStateDelta(m.cfg.Name)
|
|
}
|
|
|
|
// Start hydrates from the source, subscribes for peer deltas, registers a
|
|
// reconnect re-hydrate (when the client supports it), and starts the optional
|
|
// reconcile ticker.
|
|
func (m *SyncedMap[K, V]) Start(ctx context.Context) error {
|
|
if err := m.hydrate(ctx); err != nil {
|
|
return err
|
|
}
|
|
|
|
// The cancel func is stored on the struct and invoked in Close (covered by
|
|
// tests); lifeCtx must outlive Start to drive the reconnect/reconcile
|
|
// goroutines, so it cannot be cancelled or deferred within this scope.
|
|
m.lifeCtx, m.cancel = context.WithCancel(context.Background()) // #nosec G118 -- cancel is invoked in Close()
|
|
|
|
if m.cfg.Nats != nil {
|
|
sub, err := messaging.SubscribeJSON(m.cfg.Nats, m.subject(), m.apply)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
m.sub = sub
|
|
|
|
// nats.go transparently resubscribes on reconnect, but it cannot know we
|
|
// kept derived in-memory state that may have drifted while the link was
|
|
// down, so re-hydrate from the durable source. Detected via an optional
|
|
// interface so Broadcaster itself stays minimal; standalone/test
|
|
// clients without the method simply fall back to the reconcile ticker.
|
|
if r, ok := m.cfg.Nats.(interface{ OnReconnect(func()) }); ok {
|
|
r.OnReconnect(func() {
|
|
if err := m.hydrate(m.lifeCtx); err != nil {
|
|
xlog.Warn("syncstate: reconnect re-hydrate failed", "name", m.cfg.Name, "error", err)
|
|
}
|
|
})
|
|
}
|
|
}
|
|
|
|
if m.cfg.Reconcile > 0 {
|
|
m.wg.Add(1)
|
|
go m.reconcileLoop()
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// Close unsubscribes and stops the reconcile ticker.
|
|
func (m *SyncedMap[K, V]) Close() error {
|
|
if m.cancel != nil {
|
|
m.cancel()
|
|
}
|
|
m.wg.Wait()
|
|
if m.sub != nil {
|
|
return m.sub.Unsubscribe()
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// Set writes through the Store, then updates the value locally, then
|
|
// broadcasts. The Store write comes first and happens under the lock so memory
|
|
// and durable state move together: when it fails, Set returns the error with
|
|
// memory and peers untouched. Keeping an unpersisted value in memory would let
|
|
// this replica serve it (and a caller that re-reads the map re-apply it) while
|
|
// the Store and every other replica disagree, until the next re-hydrate.
|
|
// The broadcast is best-effort after unlocking.
|
|
func (m *SyncedMap[K, V]) Set(ctx context.Context, v V) error {
|
|
k := m.cfg.Key(v)
|
|
m.mu.Lock()
|
|
if m.cfg.Store != nil {
|
|
if err := m.cfg.Store.Upsert(ctx, v); err != nil {
|
|
m.mu.Unlock()
|
|
return err
|
|
}
|
|
}
|
|
m.data[k] = v
|
|
m.mu.Unlock()
|
|
m.publish(opSet, k, v)
|
|
return nil
|
|
}
|
|
|
|
// Delete deletes the key from the Store, then removes it locally, then
|
|
// broadcasts. A failed Store delete leaves memory and peers untouched, for the
|
|
// same reason as Set.
|
|
func (m *SyncedMap[K, V]) Delete(ctx context.Context, k K) error {
|
|
m.mu.Lock()
|
|
if m.cfg.Store != nil {
|
|
if err := m.cfg.Store.Delete(ctx, k); err != nil {
|
|
m.mu.Unlock()
|
|
return err
|
|
}
|
|
}
|
|
delete(m.data, k)
|
|
m.mu.Unlock()
|
|
var zero V
|
|
m.publish(opDelete, k, zero)
|
|
return nil
|
|
}
|
|
|
|
// Get returns the value for k and whether it was present.
|
|
func (m *SyncedMap[K, V]) Get(k K) (V, bool) {
|
|
m.mu.RLock()
|
|
defer m.mu.RUnlock()
|
|
v, ok := m.data[k]
|
|
return v, ok
|
|
}
|
|
|
|
// List returns a snapshot slice of all values.
|
|
func (m *SyncedMap[K, V]) List() []V {
|
|
m.mu.RLock()
|
|
defer m.mu.RUnlock()
|
|
out := make([]V, 0, len(m.data))
|
|
for _, v := range m.data {
|
|
out = append(out, v)
|
|
}
|
|
return out
|
|
}
|
|
|
|
// Snapshot returns a copy of the underlying map.
|
|
func (m *SyncedMap[K, V]) Snapshot() map[K]V {
|
|
m.mu.RLock()
|
|
defer m.mu.RUnlock()
|
|
out := make(map[K]V, len(m.data))
|
|
for k, v := range m.data {
|
|
out[k] = v
|
|
}
|
|
return out
|
|
}
|
|
|
|
// publish broadcasts a delta. Standalone (nil Nats) is a strict no-op.
|
|
func (m *SyncedMap[K, V]) publish(op string, k K, v V) {
|
|
if m.cfg.Nats == nil {
|
|
return
|
|
}
|
|
if err := m.cfg.Nats.Publish(m.subject(), delta[K, V]{Op: op, Key: k, Value: v}); err != nil {
|
|
xlog.Warn("syncstate: failed to broadcast delta", "name", m.cfg.Name, "op", op, "error", err)
|
|
}
|
|
}
|
|
|
|
// apply handles a peer's delta: memory-only update plus OnApply. It deliberately
|
|
// never writes the Store nor re-publishes - that is the echo-loop guard.
|
|
func (m *SyncedMap[K, V]) apply(d delta[K, V]) {
|
|
switch d.Op {
|
|
case opSet:
|
|
m.mu.Lock()
|
|
m.data[d.Key] = d.Value
|
|
m.mu.Unlock()
|
|
case opDelete:
|
|
m.mu.Lock()
|
|
delete(m.data, d.Key)
|
|
m.mu.Unlock()
|
|
default:
|
|
xlog.Warn("syncstate: ignoring delta with unknown op", "name", m.cfg.Name, "op", d.Op)
|
|
return
|
|
}
|
|
if m.cfg.OnApply != nil {
|
|
m.cfg.OnApply(d.Op, d.Key, d.Value)
|
|
}
|
|
}
|
|
|
|
// hydrate replaces the whole map from the durable source: Store if present, else
|
|
// Loader. With neither, a late joiner starts empty and catches up via deltas
|
|
// (acceptable only for ephemeral state).
|
|
func (m *SyncedMap[K, V]) hydrate(ctx context.Context) error {
|
|
var (
|
|
vals []V
|
|
err error
|
|
)
|
|
switch {
|
|
case m.cfg.Store != nil:
|
|
vals, err = m.cfg.Store.List(ctx)
|
|
case m.cfg.Loader != nil:
|
|
vals, err = m.cfg.Loader(ctx)
|
|
default:
|
|
return nil
|
|
}
|
|
if err != nil {
|
|
return err
|
|
}
|
|
m.replaceAll(vals)
|
|
return nil
|
|
}
|
|
|
|
// replaceAll atomically swaps the map contents for the given values, keyed via
|
|
// cfg.Key.
|
|
func (m *SyncedMap[K, V]) replaceAll(vals []V) {
|
|
next := make(map[K]V, len(vals))
|
|
for _, v := range vals {
|
|
next[m.cfg.Key(v)] = v
|
|
}
|
|
m.mu.Lock()
|
|
m.data = next
|
|
m.mu.Unlock()
|
|
}
|
|
|
|
// reconcileLoop periodically re-hydrates to repair silent drift (missed deltas).
|
|
func (m *SyncedMap[K, V]) reconcileLoop() {
|
|
defer m.wg.Done()
|
|
t := time.NewTicker(m.cfg.Reconcile)
|
|
defer t.Stop()
|
|
for {
|
|
select {
|
|
case <-m.lifeCtx.Done():
|
|
return
|
|
case <-t.C:
|
|
if err := m.hydrate(m.lifeCtx); err != nil {
|
|
xlog.Warn("syncstate: reconcile re-hydrate failed", "name", m.cfg.Name, "error", err)
|
|
}
|
|
}
|
|
}
|
|
}
|