Files
felhom-controller/controller/internal/report/waiter.go
T
admin 2dd05670ae feat(report): Direction-2 immediate-sync wait channel client (v0.140.0)
report.Waiter holds a hanging GET against hub /api/v1/wait?gen=N (same hub
URL+key as the pusher). On a generation change it fires the v0.139.0
report.Trigger — nothing else; the report ACK delivers everything through the
unchanged machinery. No overall client timeout (held GET); per-request ctx
bounds a dead connection. First-observation records-not-fires; same-gen
timeout fires nothing; errors (incl. 404 from a pre-v0.58.0 hub) back off
5s->5min while the 15-min cycle reconciles. Wired beside the trigger under the
same hubPusher!=nil && Hub.Enabled gate.

Red-proof: disable the baseline branch -> first observation fires
(TestWaiter_FirstObservationRecordsNoFire), run-fail-reverted.

Copy soften: backups_remote/escrow "néhány másodperc, legfeljebb 15 perc".
Pairs with hub v0.58.0. Grounding:
felhom.eu/documentation/audits/SPIKE-immediate-sync-transport-2026-07-16.md
2026-07-16 20:55:30 +02:00

215 lines
6.9 KiB
Go

// waiter.go — the Direction-2 immediate-sync long-poll client (v0.140.0).
//
// The box holds a hanging authenticated GET against the hub's /api/v1/wait, carrying its last-seen
// operator-intent generation. The hub completes the hold the instant that generation moves (an
// operator saved config, re-issued offsite, changed a floor, …) or after its hold window. On a
// CHANGE the Waiter simply fires the v0.139.0 report.Trigger — the immediate report's ACK then
// delivers config/escrow/claim/floor through the UNCHANGED machinery (this file adds NO delivery
// logic). On a timeout (same generation) it fires nothing and reconnects. Every failure degrades to
// the ~15-min scheduled cycle, which stays the reconciliation backbone.
//
// The wait response is CONTENTLESS — the Waiter reads only a generation number from it (a
// {"gen":N} line preceded by heartbeat newlines) and never interprets anything else. Grounding:
// felhom.eu/documentation/audits/SPIKE-immediate-sync-transport-2026-07-16.md (option b).
package report
import (
"bufio"
"context"
"encoding/json"
"fmt"
"log"
"math/rand"
"net"
"net/http"
"strings"
"time"
"gitea.dooplex.hu/admin/felhom-controller/internal/logx"
)
const (
// waiterHoldCeiling bounds ONE held request from the client side. The hub holds ~240 s; this
// margin lets a healthy completion arrive while capping a black-holed connection (hub died
// mid-hold, NAT dropped the mapping) so the loop can back off instead of hanging forever.
waiterHoldCeiling = 300 * time.Second
waiterBackoffMin = 5 * time.Second
waiterBackoffMax = 5 * time.Minute
// Reconnect jitter after a clean completion — spreads fleet reconnects so they don't align.
waiterJitterMin = 1 * time.Second
waiterJitterMax = 3 * time.Second
)
// Waiter long-polls the hub wait channel and fires onWake when the operator-intent generation
// advances. Construct with NewWaiter; Run owns all state (single goroutine).
type Waiter struct {
url string // "{hub}/api/v1/wait"
apiKey string
client *http.Client
onWake func()
logger *log.Logger
// Tunables — production defaults set by NewWaiter; tests shrink them.
holdCeiling time.Duration
backoffMin time.Duration
backoffMax time.Duration
jitterMin time.Duration
jitterMax time.Duration
// Run-goroutine-only state.
lastSeen uint64
seen bool
failing bool
}
// NewWaiter builds a Waiter against the hub. The HTTP client deliberately has NO overall Timeout (a
// held GET must be able to stay open for the hub's full hold) — only sane connect/TLS/header
// deadlines, so a dead hub is detected quickly while a healthy hold is never cut.
func NewWaiter(hubURL, apiKey string, onWake func(), logger *log.Logger) *Waiter {
client := &http.Client{
Timeout: 0, // NO total timeout — the whole point is to hold
Transport: &http.Transport{
DialContext: (&net.Dialer{Timeout: 10 * time.Second}).DialContext,
TLSHandshakeTimeout: 10 * time.Second,
ResponseHeaderTimeout: 30 * time.Second, // the hub flushes headers immediately, before holding
ExpectContinueTimeout: 1 * time.Second,
IdleConnTimeout: 90 * time.Second,
},
}
return &Waiter{
url: strings.TrimRight(hubURL, "/") + "/api/v1/wait",
apiKey: apiKey,
client: client,
onWake: onWake,
logger: logger,
holdCeiling: waiterHoldCeiling,
backoffMin: waiterBackoffMin,
backoffMax: waiterBackoffMax,
jitterMin: waiterJitterMin,
jitterMax: waiterJitterMax,
}
}
// Run is the loop; main.go starts it under the process context. Exits promptly on ctx cancel.
func (w *Waiter) Run(ctx context.Context) {
logx.Infof(w.logger, "[report] hub wait channel active (hold ≤240s)")
backoff := w.backoffMin
for {
if ctx.Err() != nil {
return
}
gen, err := w.pollOnce(ctx)
if err != nil {
if ctx.Err() != nil {
return // shutdown, not a real failure
}
if !w.failing {
w.failing = true
logx.Warnf(w.logger, "[report] wait channel error: %v — backing off (the 15-min cycle still reconciles)", err)
}
if !w.sleep(ctx, backoff) {
return
}
backoff *= 2
if backoff > w.backoffMax {
backoff = w.backoffMax
}
continue
}
// Success: reset failure state + backoff.
if w.failing {
w.failing = false
logx.Infof(w.logger, "[report] wait channel recovered")
}
backoff = w.backoffMin
switch {
case !w.seen:
// First observation records the baseline WITHOUT firing — the startup report already
// covered current state; firing here would echo one extra report on every process start.
w.seen = true
w.lastSeen = gen
logx.Debugf(w.logger, "[report] wait baseline generation=%d", gen)
case gen != w.lastSeen:
w.lastSeen = gen
logx.Debugf(w.logger, "[report] wait woke: generation=%d — firing out-of-cycle report", gen)
if w.onWake != nil {
w.onWake()
}
default:
// Same generation → the hub's hold simply timed out; nothing to do (Scenario B).
w.lastSeen = gen
}
if !w.sleep(ctx, w.jitter()) {
return
}
}
}
// pollOnce performs one held GET and returns the generation the hub reported. A non-2xx status
// (incl. 404 from a hub that predates the endpoint), a transport error, or an unparsable body all
// return an error so the caller backs off.
func (w *Waiter) pollOnce(ctx context.Context) (uint64, error) {
reqCtx, cancel := context.WithTimeout(ctx, w.holdCeiling)
defer cancel()
req, err := http.NewRequestWithContext(reqCtx, http.MethodGet, fmt.Sprintf("%s?gen=%d", w.url, w.lastSeen), nil)
if err != nil {
return 0, err
}
if w.apiKey != "" {
req.Header.Set("Authorization", "Bearer "+w.apiKey)
}
resp, err := w.client.Do(req)
if err != nil {
return 0, err
}
defer resp.Body.Close()
if resp.StatusCode != http.StatusOK {
return 0, fmt.Errorf("wait status %d", resp.StatusCode)
}
// The hold streams heartbeat newlines then a single {"gen":N} line. Read line-wise, ignore
// blanks, parse the first non-blank line as the completion.
sc := bufio.NewScanner(resp.Body)
for sc.Scan() {
line := strings.TrimSpace(sc.Text())
if line == "" {
continue // heartbeat
}
var r struct {
Gen uint64 `json:"gen"`
}
if err := json.Unmarshal([]byte(line), &r); err != nil {
return 0, fmt.Errorf("malformed wait completion %q: %w", line, err)
}
return r.Gen, nil
}
if err := sc.Err(); err != nil {
return 0, err
}
return 0, fmt.Errorf("wait closed without a completion line")
}
// jitter returns a random reconnect delay in [jitterMin, jitterMax].
func (w *Waiter) jitter() time.Duration {
span := w.jitterMax - w.jitterMin
if span <= 0 {
return w.jitterMin
}
return w.jitterMin + time.Duration(rand.Int63n(int64(span)))
}
// sleep waits d or until ctx is cancelled; false = cancelled (caller returns promptly).
func (w *Waiter) sleep(ctx context.Context, d time.Duration) bool {
timer := time.NewTimer(d)
defer timer.Stop()
select {
case <-ctx.Done():
return false
case <-timer.C:
return true
}
}