v0.80.0: PBS DR tier slice 2 — the apply-bridge (pbs_dr consumer, felhom-pbs-apply set-only wrapper, verify-pin-before-consume, adoption-first, loud consumed-failed, escrow seed)
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01PSK5g6qYLknKj8u3QAFEr6
This commit is contained in:
@@ -0,0 +1,82 @@
|
||||
package pbsdr
|
||||
|
||||
import (
|
||||
"context"
|
||||
"log/slog"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"gitea.dooplex.hu/admin/felhom-agent/internal/hub"
|
||||
)
|
||||
|
||||
// Loop drives the Manager on its own cadence AND consumes fetched desired-state via the
|
||||
// desired.Syncer raw-consumer seam (the wgtunnel Loop shape). fetched=false ("no desired data
|
||||
// seen yet") is never a signal; an absent pbs_dr block on a PRESENT desired-state is a plain
|
||||
// no-op this slice (no teardown — deprovision is a deliberate future op).
|
||||
type Loop struct {
|
||||
mgr *Manager
|
||||
interval time.Duration
|
||||
logger *slog.Logger
|
||||
|
||||
mu sync.Mutex
|
||||
fetched bool
|
||||
block *hub.WirePBSDR
|
||||
|
||||
nudge chan struct{}
|
||||
}
|
||||
|
||||
// NewLoop builds the loop. interval defaults to 60s.
|
||||
func NewLoop(mgr *Manager, interval time.Duration, logger *slog.Logger) *Loop {
|
||||
if interval <= 0 {
|
||||
interval = 60 * time.Second
|
||||
}
|
||||
if logger == nil {
|
||||
logger = slog.Default()
|
||||
}
|
||||
return &Loop{mgr: mgr, interval: interval, logger: logger, nudge: make(chan struct{}, 1)}
|
||||
}
|
||||
|
||||
// OnDesiredState implements desired.RawConsumer: store the latest pbs_dr block (or its absence)
|
||||
// and nudge the loop. Non-blocking and panic-free by construction.
|
||||
func (l *Loop) OnDesiredState(_ context.Context, resp *hub.DesiredStateResponse) {
|
||||
if resp == nil {
|
||||
return
|
||||
}
|
||||
l.mu.Lock()
|
||||
l.fetched = true
|
||||
l.block = resp.DesiredState.PBSDR
|
||||
l.mu.Unlock()
|
||||
select {
|
||||
case l.nudge <- struct{}{}:
|
||||
default:
|
||||
}
|
||||
}
|
||||
|
||||
func (l *Loop) snapshot() (bool, *hub.WirePBSDR) {
|
||||
l.mu.Lock()
|
||||
defer l.mu.Unlock()
|
||||
return l.fetched, l.block
|
||||
}
|
||||
|
||||
// Run applies immediately, then on every tick or desired-state nudge, until ctx is cancelled.
|
||||
func (l *Loop) Run(ctx context.Context) error {
|
||||
fetched, block := l.snapshot()
|
||||
l.mgr.Apply(ctx, fetched, block)
|
||||
t := time.NewTicker(l.interval)
|
||||
defer t.Stop()
|
||||
for {
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
return ctx.Err()
|
||||
case <-t.C:
|
||||
case <-l.nudge:
|
||||
}
|
||||
fetched, block = l.snapshot()
|
||||
l.mgr.Apply(ctx, fetched, block)
|
||||
}
|
||||
}
|
||||
|
||||
// PBSDRStatus implements the hub collector's PBSDRReporter seam.
|
||||
func (l *Loop) PBSDRStatus(_ context.Context) *hub.PBSDRStatus {
|
||||
return l.mgr.Status()
|
||||
}
|
||||
Reference in New Issue
Block a user