From 19f564ddd94d3a9889275f9963195fa69a680cde Mon Sep 17 00:00:00 2001 From: "pulse-triage[bot]" <249995291+pulse-triage[bot]@users.noreply.github.com> Date: Tue, 8 Sep 2026 17:46:16 +0100 Subject: [PATCH] fix(truenas): maintain idle sessions and bound successful poll cadence Persistent authenticated WebSockets sent nothing between polls, allowing idle-timeout appliances to discard them before the next refresh. Send bounded control pings for each authenticated session and stop the sender on disposal without treating pongs as fresh inventory. Successful refreshes also added a full interval after completion. Target start-to-start cadence while preserving up to five seconds of minimum idle time, unchanged failure backoff and honest completion timestamps. Authenticated idle-timeout and cadence regressions failed before the repairs; focused transport and poller race tests pass. Addresses the reproduced paths in #1893, not native appliance or reporter confirmation. Change-source: pulse-maintainer --- .../v6/internal/subsystems/monitoring.md | 27 +++++ internal/monitoring/truenas_poller.go | 16 ++- internal/monitoring/truenas_poller_test.go | 43 +++++++- internal/truenas/client.go | 6 +- internal/truenas/transport.go | 44 +++++++- internal/truenas/transport_test.go | 102 ++++++++++++++++++ 6 files changed, 231 insertions(+), 7 deletions(-) diff --git a/docs/release-control/v6/internal/subsystems/monitoring.md b/docs/release-control/v6/internal/subsystems/monitoring.md index 232042dbf..841f35a0e 100644 --- a/docs/release-control/v6/internal/subsystems/monitoring.md +++ b/docs/release-control/v6/internal/subsystems/monitoring.md @@ -17,6 +17,33 @@ ## Purpose +### TrueNAS persistent-session liveness and successful poll cadence + +Authenticated JSON-RPC WebSocket sessions send transport-only PING controls +every 25 seconds with a five-second write deadline. The sender belongs to the +session, not its opening request context, and disposal stops and joins it. +A failed control write closes the socket so existing transport handling owns +retry; actions retain their no-replay boundary. Existing serialized RPC and +stream readers consume pongs. Neither ping nor pong updates inventory freshness +or appliance health, and sending keepalives is not proactive pong-timeout detection. + +Successful TrueNAS refreshes target start-to-start cadence, bounded below by +completion plus min(five seconds, configured interval), preventing back-to-back +load on slow appliances. Failed refreshes retain a full completion-based retry +interval. Manual connection tests retain completion-based scheduling. Observed +last-attempt and last-success timestamps remain completion timestamps; resource +freshness thresholds must not be extended to hide genuinely stale devices. + +`TestAuthenticatedRPCSurvivesIdleTransportTimeout` and +`TestRPCSessionKeepaliveConcurrentCallsAndShutdown` in +`internal/truenas/transport_test.go` verify an authenticated synthetic idle-timeout +server, opening-context cancellation, concurrent controls/RPCs and sender disposal. +`TestTrueNASSuccessfulPollCadenceIncludesBoundedIdleGap` in +`internal/monitoring/truenas_poller_test.go` verifies due boundaries, short/slow +cycles, completion timestamps and unchanged failure backoff. These are synthetic +runtime proofs, not native firmware timeout or reporter-resolution evidence. + + **Availability backfill preserves concurrent discovery changes (7 September 2026)** The backfill List snapshot is a work list, not an authoritative record to save. diff --git a/internal/monitoring/truenas_poller.go b/internal/monitoring/truenas_poller.go index 80b59b2c2..b44e26049 100644 --- a/internal/monitoring/truenas_poller.go +++ b/internal/monitoring/truenas_poller.go @@ -517,7 +517,7 @@ func (p *TrueNASPoller) pollAll(ctx context.Context) { snapshot := entry.provider.Snapshot() p.mu.Lock() - p.recordConnectionSuccessLocked(entry.orgID, entry.id, entry.config, end, snapshot) + p.recordConnectionSuccessLocked(entry.orgID, entry.id, entry.config, start, end, snapshot) p.mu.Unlock() refreshedOrgs[entry.orgID] = struct{}{} p.ingestRecoveryPoints(ctx, entry.orgID, entry.id, entry.provider) @@ -811,6 +811,7 @@ func (p *TrueNASPoller) recordConnectionSuccessLocked( orgID string, connID string, instance config.TrueNASInstance, + startedAt time.Time, at time.Time, snapshot *truenas.FixtureSnapshot, ) { @@ -819,7 +820,16 @@ func (p *TrueNASPoller) recordConnectionSuccessLocked( status.lastSuccessAt = at status.lastError = nil status.consecutiveFailures = 0 - status.nextPollAt = at.Add(p.effectiveRuntimePollInterval(instance)) + // Successful cycles target start-to-start cadence, but never immediately + // hammer a slow appliance with another refresh. Keep at least five seconds + // idle (or the entire interval for short-interval configurations). Failures + // continue to use a full completion-based retry interval below. + interval := p.effectiveRuntimePollInterval(instance) + idleGap := min(5*time.Second, interval) + status.nextPollAt = startedAt.Add(interval) + if earliest := at.Add(idleGap); status.nextPollAt.Before(earliest) { + status.nextPollAt = earliest + } if snapshot != nil { status.observed = buildTrueNASObservedSummary(snapshot) } @@ -869,7 +879,7 @@ func (p *TrueNASPoller) RecordConnectionTestSuccess( p.mu.Lock() defer p.mu.Unlock() - p.recordConnectionSuccessLocked(orgID, connID, instance, at, nil) + p.recordConnectionSuccessLocked(orgID, connID, instance, at, at, nil) } // RecordConnectionTestFailure updates one saved TrueNAS connection summary after diff --git a/internal/monitoring/truenas_poller_test.go b/internal/monitoring/truenas_poller_test.go index 0c6548fa8..caa605a7e 100644 --- a/internal/monitoring/truenas_poller_test.go +++ b/internal/monitoring/truenas_poller_test.go @@ -414,7 +414,7 @@ func TestTrueNASPollerManualConnectionTestsUpdateSummariesWithoutClearingObserve manualSuccessAt := failureAt.Add(2 * time.Minute) poller.mu.Lock() - poller.recordConnectionSuccessLocked("default", connection.ID, connection, firstSuccess, snapshot) + poller.recordConnectionSuccessLocked("default", connection.ID, connection, firstSuccess, firstSuccess, snapshot) poller.recordConnectionFailureLocked("default", connection.ID, connection, errors.New("manual auth failed"), failureAt) poller.mu.Unlock() @@ -2164,3 +2164,44 @@ func TestTrueNASPollerKeysSystemsByConnection(t *testing.T) { return firstOK && secondOK && len(ids) == 2 }, "expected one connection-scoped system source ID per configured connection") } + +func TestTrueNASSuccessfulPollCadenceIncludesBoundedIdleGap(t *testing.T) { + for _, tc := range []struct { + name string + interval, duration, next time.Duration + }{ + {"fast", time.Minute, 2 * time.Second, time.Minute}, + {"nearly_due", time.Minute, 58 * time.Second, 63 * time.Second}, + {"slow", time.Minute, 90 * time.Second, 95 * time.Second}, + {"genuinely_stale", time.Minute, 122 * time.Second, 127 * time.Second}, + {"short_interval", time.Second, 2 * time.Second, 3 * time.Second}, + } { + t.Run(tc.name, func(t *testing.T) { + poller := NewTrueNASPoller(nil, 0, nil) + instance := config.TrueNASInstance{ID: "cadence", PollIntervalSecs: int(tc.interval / time.Second)} + start := time.Date(2026, 9, 8, 0, 0, 0, 0, time.UTC) + end := start.Add(tc.duration) + poller.recordConnectionSuccessLocked("default", instance.ID, instance, start, end, nil) + status := poller.ensureConnectionRuntimeStatusLocked("default", instance.ID) + want := start.Add(tc.next) + if !status.nextPollAt.Equal(want) { + t.Errorf("next poll = %v, want %v", status.nextPollAt, want) + } + if !status.lastSuccessAt.Equal(end) || !status.lastAttemptAt.Equal(end) { + t.Error("scheduling altered observed completion timestamps") + } + if poller.connectionPollDueLocked("default", instance.ID, instance, want.Add(-time.Nanosecond)) { + t.Error("poll due before bounded idle gap elapsed") + } + if !poller.connectionPollDueLocked("default", instance.ID, instance, want) { + t.Error("poll not due at scheduled time") + } + // Failed refreshes retain the full interval after completion; this change + // must not accelerate retries against an unavailable appliance. + poller.recordConnectionFailureLocked("default", instance.ID, instance, errors.New("offline"), end) + if !status.nextPollAt.Equal(end.Add(tc.interval)) { + t.Error("failure backoff changed") + } + }) + } +} diff --git a/internal/truenas/client.go b/internal/truenas/client.go index ff15206a8..d67d2e3e2 100644 --- a/internal/truenas/client.go +++ b/internal/truenas/client.go @@ -2061,8 +2061,10 @@ func appendDiskTemperature(out map[string]int, diskName string, value any) { } type trueNASRPCClient struct { - conn *websocket.Conn - nextID int64 + conn *websocket.Conn + nextID int64 + keepaliveStop chan struct{} + keepaliveDone chan struct{} } func (c *trueNASRPCClient) subscribe(ctx context.Context, event string) (string, error) { diff --git a/internal/truenas/transport.go b/internal/truenas/transport.go index 9c79cfc06..68ef382ce 100644 --- a/internal/truenas/transport.go +++ b/internal/truenas/transport.go @@ -9,6 +9,8 @@ import ( "strconv" "strings" "time" + + "github.com/gorilla/websocket" ) var errRPCStreamSessionConsumed = errors.New("truenas rpc stream session cannot be reused") @@ -276,6 +278,7 @@ func (c *Client) openAuthenticatedRPC(ctx context.Context) (*trueNASRPCClient, s _ = conn.Close() return nil, "", err } + rpc.startKeepalive(25 * time.Second) return rpc, authMechanism, nil } @@ -444,7 +447,7 @@ func (c *Client) waitReconnectBackoff(ctx context.Context) error { func (c *Client) closeRPCLocked() { if c.rpc != nil && c.rpc.conn != nil { - _ = c.rpc.conn.Close() + c.rpc.close() } c.rpc = nil c.updateTransportStatus(func(status *TransportStatus) { @@ -452,6 +455,45 @@ func (c *Client) closeRPCLocked() { }) } +// startKeepalive belongs to the authenticated session, not the context of the +// call which opened it. WriteControl is safe alongside the serialized RPC +// reader/writer; it must not change their deadlines or acquire rpcMu. +// Pongs are consumed by the existing RPC/stream readers. A ping is only idle +// transport maintenance, never evidence of fresh inventory or appliance health. +func (c *trueNASRPCClient) startKeepalive(interval time.Duration) { + c.keepaliveStop = make(chan struct{}) + c.keepaliveDone = make(chan struct{}) + go func() { + defer close(c.keepaliveDone) + ticker := time.NewTicker(interval) + defer ticker.Stop() + for { + select { + case <-c.keepaliveStop: + return + case <-ticker.C: + if err := c.conn.WriteControl(websocket.PingMessage, nil, time.Now().Add(5*time.Second)); err != nil { + // Unblock any in-flight reader; ordinary transport handling owns retry + // and status, including the no-replay rule for actions. + _ = c.conn.Close() + return + } + } + } + }() +} + +// The owning Client serializes session disposal with rpcMu. +func (c *trueNASRPCClient) close() { + if c.keepaliveStop != nil { + close(c.keepaliveStop) + } + _ = c.conn.Close() + if c.keepaliveDone != nil { + <-c.keepaliveDone + } +} + func (c *Client) recordTransportError(err error) { message := c.sanitizeTransportError(err) c.updateTransportStatus(func(status *TransportStatus) { diff --git a/internal/truenas/transport_test.go b/internal/truenas/transport_test.go index e068861e8..34d16e2d2 100644 --- a/internal/truenas/transport_test.go +++ b/internal/truenas/transport_test.go @@ -993,3 +993,105 @@ func TestIssue1631HTTPSUpgradeTargetRules(t *testing.T) { } } } + +// Model an authenticated appliance with a transport idle limit. This is not +// evidence of a particular appliance firmware's authenticated timeout policy. +func TestAuthenticatedRPCSurvivesIdleTransportTimeout(t *testing.T) { + var sessions, pings atomic.Int32 + upgrader := websocket.Upgrader{} + server := httptest.NewTLSServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + conn, err := upgrader.Upgrade(w, r, nil) + if err != nil { + return + } + defer conn.Close() + sessions.Add(1) + resetIdle := func() { _ = conn.SetReadDeadline(time.Now().Add(35 * time.Second)) } + conn.SetPingHandler(func(data string) error { + pings.Add(1) + resetIdle() + return conn.WriteControl(websocket.PongMessage, []byte(data), time.Now().Add(time.Second)) + }) + resetIdle() + for { + var request trueNASRPCRequest + if err := conn.ReadJSON(&request); err != nil { + return + } + resetIdle() + var result any = map[string]any{"version": "TrueNAS-SCALE-25.04.2"} + if request.Method == "auth.login_ex" { + result = map[string]any{"response_type": "SUCCESS"} + } + if err := conn.WriteJSON(map[string]any{"jsonrpc": "2.0", "id": request.ID, "result": result}); err != nil { + return + } + } + })) + defer server.Close() + client := protocolFixtureClient(t, server.URL, ClientConfig{APIKey: "fixture-key", Username: "fixture-user"}) + defer client.Close() + ctx, cancel := context.WithTimeout(context.Background(), 50*time.Second) + defer cancel() + if _, err := client.GetSystemInfo(ctx); err != nil { + t.Fatal(err) + } + client.rpcMu.Lock() + session := client.rpc + client.rpcMu.Unlock() + // Cancelling the opening call must not cancel the persistent session. + cancel() + ctx, cancel = context.WithTimeout(context.Background(), 50*time.Second) + defer cancel() + // Longer than the fixture's idle timeout, shorter than two production pings. + time.Sleep(40 * time.Second) + if _, err := client.GetSystemInfo(ctx); err != nil { + t.Fatal(err) + } + if got := sessions.Load(); got != 1 { + t.Errorf("authenticated sessions = %d, want 1 (idle connection was lost)", got) + } + if pings.Load() == 0 { + t.Error("no websocket keepalive received") + } + client.Close() + select { + case <-session.keepaliveDone: + default: + t.Error("Close returned before keepalive exited") + } +} + +func TestRPCSessionKeepaliveConcurrentCallsAndShutdown(t *testing.T) { + fixture := newProtocolFixture(t, func(_ int, _ trueNASRPCRequest) protocolFixtureReply { + time.Sleep(3 * time.Millisecond) + return protocolFixtureReply{result: "ok"} + }, nil) + client := protocolFixtureClient(t, fixture.server.URL, ClientConfig{}) + ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second) + defer cancel() + conn, err := client.dialRPC(ctx) + if err != nil { + t.Fatal(err) + } + rpc := &trueNASRPCClient{conn: conn, nextID: 1} + rpc.startKeepalive(time.Millisecond) + defer rpc.close() + for i := 0; i < 20; i++ { + var result string + if err := rpc.call(ctx, "fixture.read", nil, &result); err != nil { + t.Fatal(err) + } + if result != "ok" { + t.Fatalf("result = %q", result) + } + } + // A failed control write must terminate its sender and close the socket, + // rather than retaining a goroutine until a future poll/client shutdown. + _ = conn.Close() + select { + case <-rpc.keepaliveDone: + case <-time.After(time.Second): + t.Fatal("keepalive did not exit after socket failure") + } +}