Files
felhom-agent/internal/desired/syncer.go
T
admin 8ecf8929fb slice 10A: activate the control envelope (Down channel) + hub-backed desired provider (v0.15.0)
The control envelope becomes live: the agent caches the hub's desired-state +
generation and re-fetches GET /hosts/{id}/desired-state only when the
generation advances. A new internal/desired Syncer maps the wire shape into a
reconcile.CachingProvider feeding the engine; benign deltas reconcile, an
explicit guest decommission is gated pending_signature (exec is 10B). Adds the
DesiredStateResponse/WireDesiredState wire types + Client.FetchDesiredState +
the loop EnvelopeObserver seam. Cross-repo golden (envelope + desired-state)
byte-identical with the hub.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-06-10 19:02:59 +02:00

98 lines
4.0 KiB
Go

// Package desired bridges the hub's "Down" channel (the control-envelope generation signal +
// the desired-state fetch) to the reconcile engine's provider (slice 10A). It implements
// hub.EnvelopeObserver: on each heartbeat it inspects the envelope's DesiredGeneration and, only
// when it has ADVANCED past the cached one, fetches the full desired-state and updates the
// engine's CachingProvider. So the heartbeat stays light; the heavy state moves on change.
//
// It lives in its own package because it imports BOTH hub (the wire client + types) and reconcile
// (the domain DesiredState + CachingProvider). hub does not import it (the loop sees only the
// hub.EnvelopeObserver seam) and reconcile does not import it — so there is no import cycle.
package desired
import (
"context"
"log/slog"
"gitea.dooplex.hu/admin/felhom-agent/internal/hub"
"gitea.dooplex.hu/admin/felhom-agent/internal/reconcile"
)
// Fetcher fetches this host's desired-state from the hub. Satisfied by *hub.Client.
type Fetcher interface {
FetchDesiredState(ctx context.Context) (*hub.DesiredStateResponse, error)
}
// Syncer keeps the engine's CachingProvider in step with the hub's authoritative desired-state.
type Syncer struct {
fetcher Fetcher
provider *reconcile.CachingProvider
logger *slog.Logger
}
// NewSyncer builds a Syncer over the hub fetcher and the engine's provider.
func NewSyncer(fetcher Fetcher, provider *reconcile.CachingProvider, logger *slog.Logger) *Syncer {
if logger == nil {
logger = slog.Default()
}
return &Syncer{fetcher: fetcher, provider: provider, logger: logger}
}
// OnEnvelope implements hub.EnvelopeObserver. It fetches + caches the desired-state ONLY when the
// envelope's generation advances past the provider's cached generation — otherwise it is a no-op
// (the cached state is already current). A fetch failure keeps the last-known state (the engine
// keeps reconciling toward it) and is retried on the next advance signal.
func (s *Syncer) OnEnvelope(ctx context.Context, env *hub.ControlEnvelope) {
if env == nil || s.provider == nil {
return
}
have := s.provider.Generation()
if env.DesiredGeneration <= have {
return // cached: the heavy desired-state moves only on a generation advance
}
resp, err := s.fetcher.FetchDesiredState(ctx)
if err != nil {
s.logger.Warn("desired: fetch failed; keeping cached desired-state",
"have_generation", have, "envelope_generation", env.DesiredGeneration, "err", err)
return
}
state := mapWire(resp.DesiredState, s.logger)
// Cache against the FETCHED generation (not the envelope's) — robust to a generation that
// advanced again between the heartbeat and this fetch (we won't re-fetch the same state).
s.provider.Update(resp.Generation, state)
s.logger.Info("desired: updated from hub",
"generation", resp.Generation, "guests", len(state.Guests))
if env.HasSignedOps {
// 10A only notes the flag; fetching + verifying + executing signed ops is slice 10B.
s.logger.Info("desired: hub reports pending signed ops (fetch/execute is slice 10B)")
}
}
// mapWire maps the hub wire desired-state to the reconcile domain. 10A acts only on guests; the
// forward-compat fields (restore_directive — 10D — etc.) are carried on the wire and logged, but
// not translated into actions here.
func mapWire(w hub.WireDesiredState, logger *slog.Logger) reconcile.DesiredState {
guests := make(map[int]reconcile.DesiredGuest, len(w.Guests))
for _, g := range w.Guests {
dg := reconcile.DesiredGuest{
VMID: g.VMID,
Spec: g.Spec,
Description: g.Description,
Decommission: g.Decommission,
}
switch g.Run {
case "running":
dg.Run = reconcile.RunRunning
case "stopped":
dg.Run = reconcile.RunStopped
default:
dg.Run = reconcile.RunUnspecified // unknown/empty → unmanaged (planner leaves run alone)
}
guests[g.VMID] = dg
}
if w.RestoreDirective != nil {
logger.Info("desired: restore_directive present (consumed in slice 10D — ignored in 10A)",
"mode", w.RestoreDirective.Mode)
}
return reconcile.DesiredState{Guests: guests}
}