From 7e369d418f74da8023dea087ed434ec4ecd30764 Mon Sep 17 00:00:00 2001 From: "pulse-triage[bot]" <249995291+pulse-triage[bot]@users.noreply.github.com> Date: Fri, 2 Oct 2026 12:27:35 +0100 Subject: [PATCH] Bound TrueNAS RPC operations and keep polling after silent peers Apply the configured timeout to serialized JSON-RPC operations, subscriptions and permitted read retries. Let cancelled waiters leave without dispatching or poisoning the active session. Preserve modern transport selection and action no-replay; distinguish caller deadlines from successful bounded log tails. A required-method timeout must preserve cached host identity and previous success while allowing the next connection and recovery poll to proceed. Optional telemetry remains unavailable without blocking usable inventory. This repairs reproduced source paths, not a verified diagnosis or native resolution of issue #2382. Change-source: pulse-maintainer --- .../v6/internal/subsystems/monitoring.md | 36 ++ internal/monitoring/truenas_poller.go | 4 +- internal/monitoring/truenas_poller_test.go | 139 ++++++++ internal/truenas/client.go | 30 +- internal/truenas/transport.go | 68 +++- internal/truenas/transport_test.go | 315 ++++++++++++++++++ 6 files changed, 576 insertions(+), 16 deletions(-) diff --git a/docs/release-control/v6/internal/subsystems/monitoring.md b/docs/release-control/v6/internal/subsystems/monitoring.md index af091c647..d30789902 100644 --- a/docs/release-control/v6/internal/subsystems/monitoring.md +++ b/docs/release-control/v6/internal/subsystems/monitoring.md @@ -133,6 +133,42 @@ precede marker production. Legacy reports without the marker retain the existing default policy; matched agent and server support is required. This proves ingestion behaviour, not reporter installation or release delivery. +### TrueNAS RPC operation budgets and poll isolation + +The configured client timeout (30 seconds by default) bounds each logical +JSON-RPC operation, including its session-lock wait, handshake/authentication, +method exchange or subscription, and the permitted read retry/backoff. They +share one budget; a shorter caller deadline remains effective. A cancelled +waiter neither dispatches a request nor alters the current owner's socket or +transport status. RPC/stream readers receive the bounded context, and a +timed-out socket is discarded before a subsequent operation authenticates a +fresh session. Modern appliances never downgrade to REST on timeout, and an +action timeout after dispatch retains its unknown outcome without replay. +Keepalive still belongs to the session, not any completed operation's context. + +Optional live-telemetry timeouts do not prevent otherwise usable inventory; +unobserved metrics remain unavailable. A required-method timeout returns a +truthful failure, preserves the prior successful observation/cached identity, +and lets the shared poll cycle continue to other connections. Socket deadline +errors are classified as timeout, not credential or generic connection failures. +The next normal poll can recover; completion-based failure backoff is unchanged. +This is an operation budget, not a new total-snapshot deadline or parallel poller. +The app-log idle window still yields a bounded tail; an earlier caller/operation +deadline is an error rather than a successful empty log response. + +`TestJSONRPCConfiguredTimeout*`, +`TestJSONRPCWaitingDeadlineDoesNotInterruptSessionOwner`, +`TestJSONRPCPreCancelledCallDoesNotDispatch` and +`TestJSONRPCSnapshotContinuesAfterTelemetryTimeout` in +`internal/truenas/transport_test.go` cover silent handshake/authentication, +read/retry budgets, all three stream readers, queued cancellation, caller +deadlines, recovery and action no-replay. +`TestTrueNASPollerUnresponsiveRPCDoesNotFreezeOtherConnections` and +`TestClassifyTrueNASError` in `internal/monitoring/truenas_poller_test.go` +cover actual TLS/RPC clients through two successful/failed/recovered poll cycles, +cached identities, truthful health and unchanged backoff. These are controlled +runtime proofs, not native SCALE acceptance or a diagnosis/resolution of #2382. + ### TrueNAS persistent-session liveness and successful poll cadence Authenticated JSON-RPC WebSocket sessions send transport-only PING controls diff --git a/internal/monitoring/truenas_poller.go b/internal/monitoring/truenas_poller.go index 1f18061c1..51bf22ffc 100644 --- a/internal/monitoring/truenas_poller.go +++ b/internal/monitoring/truenas_poller.go @@ -1540,8 +1540,8 @@ func classifyTrueNASError(err error, connectionID string) *internalerrors.Monito } } // Transport-level errors: timeout takes precedence over generic connection failures. - var urlErr *url.Error - if (errors.As(err, &urlErr) && urlErr.Timeout()) || errors.Is(err, context.DeadlineExceeded) { + var netErr net.Error + if (errors.As(err, &netErr) && netErr.Timeout()) || errors.Is(err, context.DeadlineExceeded) { errType = internalerrors.ErrorTypeTimeout } else { var netOpErr *net.OpError diff --git a/internal/monitoring/truenas_poller_test.go b/internal/monitoring/truenas_poller_test.go index bd99a72e1..a263712ac 100644 --- a/internal/monitoring/truenas_poller_test.go +++ b/internal/monitoring/truenas_poller_test.go @@ -17,6 +17,7 @@ import ( "testing" "time" + "github.com/gorilla/websocket" "github.com/rcourtman/pulse-go-rewrite/internal/config" "github.com/rcourtman/pulse-go-rewrite/internal/models" "github.com/rcourtman/pulse-go-rewrite/internal/truenas" @@ -1815,6 +1816,14 @@ func TestClassifyTrueNASError(t *testing.T) { expectedType: "timeout", expectedRetry: true, }, + { + name: "WebSocket socket deadline classifies as timeout", + err: &truenas.RPCTransportError{Method: "system.info", Phase: "read", Err: &net.OpError{ + Op: "read", Net: "tcp", Err: os.ErrDeadlineExceeded, + }}, + expectedType: "timeout", + expectedRetry: true, + }, { name: "net.OpError classifies as connection", err: &net.OpError{Op: "dial", Net: "tcp", Addr: nil, Err: fmt.Errorf("connection refused")}, @@ -1866,6 +1875,136 @@ func TestClassifyTrueNASError(t *testing.T) { } } +func TestTrueNASPollerUnresponsiveRPCDoesNotFreezeOtherConnections(t *testing.T) { + var stalled atomic.Bool + release := make(chan struct{}) + newServer := func(hostname string, mayStall bool) *httptest.Server { + upgrader := websocket.Upgrader{} + return httptest.NewTLSServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if r.URL.Path != "/api/current" { + t.Error("modern polling fell back to REST") + http.NotFound(w, r) + return + } + conn, err := upgrader.Upgrade(w, r, nil) + if err != nil { + return + } + defer conn.Close() + for { + var request struct { + ID int64 `json:"id"` + Method string `json:"method"` + } + if err := conn.ReadJSON(&request); err != nil { + return + } + var result any = []any{} + switch request.Method { + case "auth.login_ex": + result = map[string]any{"response_type": "SUCCESS"} + case "system.info": + if mayStall && stalled.Load() { + <-release + } + result = map[string]any{"hostname": hostname, "version": "TrueNAS-SCALE-25.10.7", "system_serial": hostname} + case "core.subscribe": + result = "fixture-realtime" + case "core.unsubscribe": + result = nil + } + if err := conn.WriteJSON(map[string]any{"jsonrpc": "2.0", "id": request.ID, "result": result}); err != nil { + return + } + if request.Method == "core.subscribe" { + if err := conn.WriteJSON(map[string]any{ + "jsonrpc": "2.0", "method": "collection_update", + "params": map[string]any{"collection": "reporting.realtime", "fields": map[string]any{"cpu": map[string]any{"usage": 12}}}, + }); err != nil { + return + } + } + } + })) + } + brokenServer := newServer("timeout-nas", true) + healthyServer := newServer("healthy-nas", false) + t.Cleanup(brokenServer.Close) + t.Cleanup(healthyServer.Close) + poller := NewTrueNASPoller(nil, 0, nil) + instances := []config.TrueNASInstance{ + {ID: "timeout-connection", Host: brokenServer.URL, Enabled: true}, + {ID: "healthy-connection", Host: healthyServer.URL, Enabled: true}, + } + poller.providersByOrg["default"] = make(map[string]*truenas.Provider) + poller.configsByOrg["default"] = make(map[string]config.TrueNASInstance) + for i, instance := range instances { + server := []*httptest.Server{brokenServer, healthyServer}[i] + client, err := truenas.NewClient(truenas.ClientConfig{ + Host: server.URL, APIKey: "fixture-key", Username: "fixture-user", Timeout: 200 * time.Millisecond, + InsecureSkipVerify: true, Fingerprint: fmt.Sprintf("%x", sha256.Sum256(server.Certificate().Raw)), + }) + if err != nil { + t.Fatal(err) + } + t.Cleanup(client.Close) + poller.providersByOrg["default"][instance.ID] = truenas.NewLiveProviderForConnection(&truenas.APIFetcher{Client: client}, instance.ID) + poller.configsByOrg["default"][instance.ID] = instance + } + t.Cleanup(func() { close(release) }) + due := func() { + for _, instance := range instances { + poller.ensureConnectionRuntimeStatusLocked("default", instance.ID).nextPollAt = time.Now().Add(-time.Second) + } + } + poller.pollAll(context.Background()) + before := poller.ConnectionSummaries("default", instances) + for _, instance := range instances { + if before[instance.ID].Poll.LastSuccessAt == nil { + t.Fatal("initial successful poll was not established") + } + } + stalled.Store(true) + due() + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + done := make(chan struct{}) + go func() { poller.pollAll(ctx); close(done) }() + select { + case <-done: + case <-time.After(time.Second): + t.Error("one unresponsive RPC froze the shared poll cycle") + cancel() // Safely terminate the pre-repair adverse control. + select { + case <-done: + case <-time.After(5 * time.Second): + t.Fatal("poll cycle did not stop after cancellation") + } + } + after := poller.ConnectionSummaries("default", instances) + broken, healthy := after[instances[0].ID], after[instances[1].ID] + if broken.Poll.LastError == nil || broken.Poll.LastError.Category != "timeout" || broken.Poll.ConsecutiveFailures != 1 || !broken.Poll.LastSuccessAt.Equal(*before[instances[0].ID].Poll.LastSuccessAt) { + t.Fatalf("failed connection lost truthful timeout or previous success: %+v", broken) + } + if healthy.Poll.LastError != nil || !healthy.Poll.LastSuccessAt.After(*before[instances[1].ID].Poll.LastSuccessAt) { + t.Fatalf("healthy connection did not continue polling: %+v", healthy) + } + if !hasTrueNASHostForOrg(poller, "default", "timeout-nas") || !hasTrueNASHostForOrg(poller, "default", "healthy-nas") { + t.Fatal("timeout discarded a cached host identity") + } + status := poller.statusByOrg["default"][instances[0].ID] + if !status.nextPollAt.Equal(status.lastAttemptAt.Add(defaultTrueNASPollInterval)) { + t.Fatal("timeout changed completion-based failure backoff") + } + stalled.Store(false) + due() + poller.pollAll(context.Background()) + recovered := poller.ConnectionSummaries("default", instances)[instances[0].ID] + if recovered.Poll.LastError != nil || recovered.Poll.ConsecutiveFailures != 0 || !recovered.Poll.LastSuccessAt.After(*broken.Poll.LastSuccessAt) || recovered.Transport == nil || !recovered.Transport.Connected || recovered.Transport.Mode != truenas.TransportJSONRPC { + t.Fatalf("next poll did not recover the original connection over RPC: %+v", recovered) + } +} + type trueNASMockServer struct { server *httptest.Server requests atomic.Int64 diff --git a/internal/truenas/client.go b/internal/truenas/client.go index a6ff48f5a..6ad3e1104 100644 --- a/internal/truenas/client.go +++ b/internal/truenas/client.go @@ -66,7 +66,7 @@ type Client struct { baseURL string rpcURL string - rpcMu sync.Mutex + rpcMu rpcSessionMutex rpc *trueNASRPCClient mode TransportMode closed bool @@ -251,7 +251,7 @@ func (c *Client) GetSystemTelemetry(ctx context.Context) (*SystemInfo, error) { var telemetry *SystemInfo var temperatures map[string]float64 - err = c.withRPC(ctx, func(rpc *trueNASRPCClient) error { + err = c.withRPC(ctx, func(ctx context.Context, rpc *trueNASRPCClient) error { temperatures, _ = rpc.getSystemTemperatures(ctx) subscriptionName := fmt.Sprintf("reporting.realtime:{\"interval\":%d}", defaultRealtimeIntervalSeconds) subscriptionID, err := rpc.subscribe(ctx, subscriptionName) @@ -506,7 +506,7 @@ func (c *Client) GetSystemMetricHistory(ctx context.Context, duration time.Durat } var history *SystemMetricHistory - err = c.withRPC(ctx, func(rpc *trueNASRPCClient) error { + err = c.withRPC(ctx, func(ctx context.Context, rpc *trueNASRPCClient) error { var err error history, err = rpc.getSystemMetricHistory(ctx, duration) return err @@ -1332,7 +1332,7 @@ func (c *Client) GetDiskTemperatureHistory(ctx context.Context, identifiers []st } var history map[string][]TimeSeriesPoint - err := c.withRPC(ctx, func(rpc *trueNASRPCClient) error { + err := c.withRPC(ctx, func(ctx context.Context, rpc *trueNASRPCClient) error { var err error history, err = rpc.getDiskTemperatureHistory(ctx, identifiers, duration) return err @@ -1434,7 +1434,7 @@ func (c *Client) getDiskTemperaturesFromReporting(ctx context.Context, identifie } var temperatures map[string]int - err := c.withRPC(ctx, func(rpc *trueNASRPCClient) error { + err := c.withRPC(ctx, func(ctx context.Context, rpc *trueNASRPCClient) error { var err error temperatures, err = rpc.getDiskTemperatures(ctx, identifiers) return err @@ -1449,7 +1449,7 @@ func (c *Client) getDiskTemperatureAggregates(ctx context.Context, identifiers [ } var aggregates map[string]DiskTemperatureAggregate - err := c.withRPC(ctx, func(rpc *trueNASRPCClient) error { + err := c.withRPC(ctx, func(ctx context.Context, rpc *trueNASRPCClient) error { var err error aggregates, err = rpc.getDiskTemperatureAggregates(ctx, identifiers, windowDays) return err @@ -1805,7 +1805,7 @@ func (c *Client) parseAppsWithStats(ctx context.Context, response []map[string]a // failure. func (c *Client) GetAppStats(ctx context.Context) (map[string]AppStats, error) { var stats map[string]AppStats - err := c.withRPC(ctx, func(rpc *trueNASRPCClient) error { + err := c.withRPC(ctx, func(ctx context.Context, rpc *trueNASRPCClient) error { subscriptionName := fmt.Sprintf("app.stats:{\"interval\":%d}", defaultAppStatsIntervalSeconds) subscriptionID, err := rpc.subscribe(ctx, subscriptionName) if err != nil { @@ -1861,7 +1861,7 @@ func (c *Client) GetAppLogs(ctx context.Context, appName, containerID string, ta } subscriptionName := fmt.Sprintf("app.container_log_follow:%s", string(subscriptionJSON)) var lines []AppLogLine - err = c.withRPC(ctx, func(rpc *trueNASRPCClient) error { + err = c.withRPC(ctx, func(ctx context.Context, rpc *trueNASRPCClient) error { subscriptionID, err := rpc.subscribe(ctx, subscriptionName) if err != nil { return err @@ -2615,6 +2615,11 @@ func (c *trueNASRPCClient) call(ctx context.Context, method string, params any, if c == nil || c.conn == nil { return fmt.Errorf("truenas rpc connection is nil") } + // A cancelled waiter must not dispatch an action or poison a healthy + // session by setting an already-expired socket deadline. + if err := ctx.Err(); err != nil { + return err + } stopContext := c.armContext(ctx) defer stopContext() @@ -2759,6 +2764,15 @@ func (c *trueNASRPCClient) readAppLogEvents(ctx context.Context, tailLines int) var message trueNASRPCResponse if err := c.conn.ReadJSON(&message); err != nil { if isTimeoutError(err) { + // The normal idle window completes a bounded tail. A caller or + // operation budget expiring first is a failure, not evidence that + // the log stream successfully returned no data. + if ctxErr := ctx.Err(); ctxErr != nil { + return nil, false, &RPCTransportError{Method: "app.container_log_follow", Phase: "read", Err: ctxErr} + } + if ctxDeadline, ok := ctx.Deadline(); ok && !time.Now().Before(ctxDeadline) { + return nil, false, &RPCTransportError{Method: "app.container_log_follow", Phase: "read", Err: context.DeadlineExceeded} + } // Gorilla WebSocket documents a timed-out read as terminal for // the connection. Preserve the collected log data, then make // the caller discard this stream session instead of reusing it. diff --git a/internal/truenas/transport.go b/internal/truenas/transport.go index 68ef382ce..721bdd2ab 100644 --- a/internal/truenas/transport.go +++ b/internal/truenas/transport.go @@ -8,6 +8,7 @@ import ( "net/http" "strconv" "strings" + "sync" "time" "github.com/gorilla/websocket" @@ -15,6 +16,49 @@ import ( var errRPCStreamSessionConsumed = errors.New("truenas rpc stream session cannot be reused") +// rpcSessionMutex serializes the sole WebSocket reader/writer, while allowing +// a caller to abandon its wait without interrupting the current session owner. +// Its zero value is usable, including by Client.Close and protocol fixtures. +type rpcSessionMutex struct { + once sync.Once + token chan struct{} +} + +func (m *rpcSessionMutex) LockContext(ctx context.Context) error { + m.once.Do(func() { m.token = make(chan struct{}, 1) }) + if err := ctx.Err(); err != nil { + return err + } + select { + case <-ctx.Done(): + return ctx.Err() + case m.token <- struct{}{}: + if err := ctx.Err(); err != nil { + m.Unlock() + return err + } + return nil + } +} + +func (m *rpcSessionMutex) Lock() { _ = m.LockContext(context.Background()) } +func (m *rpcSessionMutex) Unlock() { <-m.token } + +// rpcOperationContext applies the same configured timeout as HTTP requests to +// a whole RPC operation: lock wait, negotiation/authentication, exchange or +// subscription, and its one permitted read retry share a single budget. A +// shorter caller deadline is never extended. Keepalive has its own lifetime. +func (c *Client) rpcOperationContext(ctx context.Context) (context.Context, context.CancelFunc) { + if ctx == nil { + ctx = context.Background() + } + timeout := c.config.Timeout + if timeout <= 0 { + timeout = defaultHTTPTimeout + } + return context.WithTimeout(ctx, timeout) +} + type discardRPCSessionError struct { err error } @@ -170,7 +214,11 @@ func (c *Client) ensureTransport(ctx context.Context) (TransportMode, error) { if c == nil { return TransportUnknown, fmt.Errorf("truenas client is nil") } - c.rpcMu.Lock() + ctx, cancel := c.rpcOperationContext(ctx) + defer cancel() + if err := c.rpcMu.LockContext(ctx); err != nil { + return TransportUnknown, err + } defer c.rpcMu.Unlock() return c.ensureTransportLocked(ctx) } @@ -286,11 +334,15 @@ func (c *Client) callRPC(ctx context.Context, method string, params any, result return c.callRPCWithRetry(ctx, method, params, result, true) } -func (c *Client) withRPC(ctx context.Context, operation func(*trueNASRPCClient) error) error { +func (c *Client) withRPC(ctx context.Context, operation func(context.Context, *trueNASRPCClient) error) error { if c == nil { return fmt.Errorf("truenas client is nil") } - c.rpcMu.Lock() + ctx, cancel := c.rpcOperationContext(ctx) + defer cancel() + if err := c.rpcMu.LockContext(ctx); err != nil { + return err + } defer c.rpcMu.Unlock() mode, err := c.ensureTransportLocked(ctx) @@ -300,7 +352,7 @@ func (c *Client) withRPC(ctx context.Context, operation func(*trueNASRPCClient) if mode != TransportJSONRPC || c.rpc == nil { return fmt.Errorf("truenas JSON-RPC operation is unavailable over negotiated transport %s", mode) } - if err := operation(c.rpc); err != nil { + if err := operation(ctx, c.rpc); err != nil { if errors.Is(err, errRPCStreamSessionConsumed) { c.closeRPCLocked() c.reconnect = 0 @@ -337,7 +389,7 @@ func (c *Client) withRPC(ctx context.Context, operation func(*trueNASRPCClient) status.LastConnectedAt = &now status.LastError = "" }) - if err := operation(c.rpc); err != nil { + if err := operation(ctx, c.rpc); err != nil { var secondTransportErr *RPCTransportError if errors.As(err, &secondTransportErr) { c.closeRPCLocked() @@ -360,7 +412,11 @@ func (c *Client) callRPCWithRetry(ctx context.Context, method string, params any if c == nil { return fmt.Errorf("truenas client is nil") } - c.rpcMu.Lock() + ctx, cancel := c.rpcOperationContext(ctx) + defer cancel() + if err := c.rpcMu.LockContext(ctx); err != nil { + return err + } defer c.rpcMu.Unlock() mode, err := c.ensureTransportLocked(ctx) diff --git a/internal/truenas/transport_test.go b/internal/truenas/transport_test.go index 34d16e2d2..a31578f71 100644 --- a/internal/truenas/transport_test.go +++ b/internal/truenas/transport_test.go @@ -1095,3 +1095,318 @@ func TestRPCSessionKeepaliveConcurrentCallsAndShutdown(t *testing.T) { t.Fatal("keepalive did not exit after socket failure") } } + +// The production poll loop has a cancellation context, not a per-request +// deadline. A silent peer must still be bounded by ClientConfig.Timeout. These +// fixtures never infer an appliance's reason for failing to reply. +func awaitConfiguredRPCBudget(t *testing.T, operation func(context.Context) error) error { + t.Helper() + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + done := make(chan error, 1) + go func() { done <- operation(ctx) }() + select { + case err := <-done: + return err + case <-time.After(time.Second): + t.Error("RPC operation ignored the configured timeout") + cancel() // Also lets the pre-repair adverse control terminate safely. + select { + case err := <-done: + return err + case <-time.After(5 * time.Second): + t.Fatal("cancelled RPC operation did not stop") + return nil + } + } +} + +func assertRPCTimeout(t *testing.T, err error) { + t.Helper() + if err == nil || (!errors.Is(err, context.DeadlineExceeded) && !isTimeoutError(err)) { + t.Errorf("expected a timeout, got %v", err) + } +} + +func TestJSONRPCConfiguredTimeoutBoundsHandshake(t *testing.T) { + release := make(chan struct{}) + var restRequests atomic.Int32 + server := httptest.NewTLSServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if r.URL.Path != "/api/current" { + restRequests.Add(1) + } + <-release // TCP/TLS succeeds, but no WebSocket handshake response follows. + })) + t.Cleanup(server.Close) + client := protocolFixtureClient(t, server.URL, ClientConfig{APIKey: "fixture-key", Username: "fixture-user", Timeout: 200 * time.Millisecond}) + t.Cleanup(func() { close(release) }) + err := awaitConfiguredRPCBudget(t, client.TestConnection) + assertRPCTimeout(t, err) + var handshake *RPCHandshakeError + if !errors.As(err, &handshake) || client.TransportStatus().Connected || restRequests.Load() != 0 { + t.Fatalf("handshake timeout lost its type/status or downgraded to REST: %v", err) + } +} + +func TestJSONRPCConfiguredTimeoutBoundsAuthentication(t *testing.T) { + release := make(chan struct{}) + var logins atomic.Int32 + fixture := newProtocolFixture(t, func(_ int, request trueNASRPCRequest) protocolFixtureReply { + if request.Method == "auth.login_ex" { + logins.Add(1) + <-release + } + return protocolFixtureReply{result: map[string]any{"response_type": "SUCCESS"}} + }, nil) + client := protocolFixtureClient(t, fixture.server.URL, ClientConfig{APIKey: "fixture-key", Username: "fixture-user", Timeout: 200 * time.Millisecond}) + t.Cleanup(func() { close(release) }) + err := awaitConfiguredRPCBudget(t, client.TestConnection) + assertRPCTimeout(t, err) + if fixture.sessions.Load() != 1 || logins.Load() != 1 || fixture.restRequests.Load() != 0 || client.TransportStatus().Connected { + t.Fatalf("silent authentication retried, downgraded or appeared connected: %+v", client.TransportStatus()) + } +} + +func TestJSONRPCConfiguredTimeoutBoundsReadAndRecovers(t *testing.T) { + release := make(chan struct{}) + var reads atomic.Int32 + fixture := newProtocolFixture(t, func(_ int, request trueNASRPCRequest) protocolFixtureReply { + if request.Method == "auth.login_ex" { + return protocolFixtureReply{result: map[string]any{"response_type": "SUCCESS"}} + } + if request.Method == "system.info" && reads.Add(1) == 1 { + <-release + } + return protocolFixtureReply{result: map[string]any{"hostname": "timeout-fixture", "version": "TrueNAS-SCALE-25.10.7"}} + }, nil) + client := protocolFixtureClient(t, fixture.server.URL, ClientConfig{APIKey: "fixture-key", Username: "fixture-user", Timeout: 200 * time.Millisecond}) + t.Cleanup(func() { close(release) }) + err := awaitConfiguredRPCBudget(t, client.TestConnection) + assertRPCTimeout(t, err) + if fixture.sessions.Load() != 1 || reads.Load() != 1 || client.TransportStatus().Connected { + t.Fatal("expired operation gained a retry budget or retained its timed-out session") + } + if err := client.TestConnection(context.Background()); err != nil { + t.Fatalf("subsequent operation did not recover: %v", err) + } + if fixture.sessions.Load() != 2 || reads.Load() != 2 || fixture.restRequests.Load() != 0 || !client.TransportStatus().Connected { + t.Fatalf("subsequent operation did not negotiate one fresh RPC session: %+v", client.TransportStatus()) + } +} + +func TestJSONRPCConfiguredTimeoutIncludesReadRetryBackoff(t *testing.T) { + release := make(chan struct{}) + var reads atomic.Int32 + fixture := newProtocolFixture(t, func(_ int, request trueNASRPCRequest) protocolFixtureReply { + if request.Method == "auth.login_ex" { + return protocolFixtureReply{result: map[string]any{"response_type": "SUCCESS"}} + } + if reads.Add(1) == 1 { + // The 100 ms reconnect backoff cannot fit in the remaining part of + // a 200 ms operation. It must not create a new deadline on retry. + time.Sleep(150 * time.Millisecond) + return protocolFixtureReply{close: true} + } + <-release + return protocolFixtureReply{result: nil} + }, nil) + client := protocolFixtureClient(t, fixture.server.URL, ClientConfig{APIKey: "fixture-key", Username: "fixture-user", Timeout: 200 * time.Millisecond}) + t.Cleanup(func() { close(release) }) + err := awaitConfiguredRPCBudget(t, client.TestConnection) + assertRPCTimeout(t, err) + if fixture.sessions.Load() != 1 || reads.Load() != 1 || fixture.restRequests.Load() != 0 { + t.Fatal("read retry restarted an exhausted operation budget") + } +} + +func TestJSONRPCConfiguredTimeoutDoesNotExtendCallerDeadline(t *testing.T) { + release := make(chan struct{}) + fixture := newProtocolFixture(t, func(_ int, request trueNASRPCRequest) protocolFixtureReply { + if request.Method == "auth.login_ex" { + return protocolFixtureReply{result: map[string]any{"response_type": "SUCCESS"}} + } + <-release + return protocolFixtureReply{result: nil} + }, nil) + client := protocolFixtureClient(t, fixture.server.URL, ClientConfig{APIKey: "fixture-key", Username: "fixture-user", Timeout: 5 * time.Second}) + t.Cleanup(func() { close(release) }) + err := awaitConfiguredRPCBudget(t, func(parent context.Context) error { + ctx, cancel := context.WithTimeout(parent, 100*time.Millisecond) + defer cancel() + return client.TestConnection(ctx) + }) + assertRPCTimeout(t, err) + if fixture.sessions.Load() != 1 || fixture.restRequests.Load() != 0 { + t.Fatal("caller deadline gained a retry or REST downgrade") + } +} + +func TestJSONRPCConfiguredTimeoutDoesNotReplayAction(t *testing.T) { + release := make(chan struct{}) + var actions atomic.Int32 + fixture := newProtocolFixture(t, func(_ int, request trueNASRPCRequest) protocolFixtureReply { + if request.Method == "auth.login_ex" { + return protocolFixtureReply{result: map[string]any{"response_type": "SUCCESS"}} + } + if request.Method == "app.start" { + actions.Add(1) + <-release + } + return protocolFixtureReply{result: nil} + }, nil) + client := protocolFixtureClient(t, fixture.server.URL, ClientConfig{APIKey: "fixture-key", Username: "fixture-user", Timeout: 200 * time.Millisecond}) + t.Cleanup(func() { close(release) }) + err := awaitConfiguredRPCBudget(t, func(ctx context.Context) error { return client.StartApp(ctx, "fixture-app") }) + assertRPCTimeout(t, err) + if !strings.Contains(err.Error(), "outcome is unknown") || actions.Load() != 1 || fixture.sessions.Load() != 1 || fixture.restRequests.Load() != 0 { + t.Fatalf("timed-out action was replayed or lost its unknown outcome: %v", err) + } +} + +func TestJSONRPCConfiguredTimeoutBoundsSubscriptions(t *testing.T) { + for _, stream := range []string{"telemetry", "app-stats", "app-logs"} { + t.Run(stream, func(t *testing.T) { + fixture := newProtocolFixture(t, func(_ int, request trueNASRPCRequest) protocolFixtureReply { + switch request.Method { + case "auth.login_ex": + return protocolFixtureReply{result: map[string]any{"response_type": "SUCCESS"}} + case "core.subscribe": + return protocolFixtureReply{result: "fixture-subscription"} // No stream event ever arrives. + default: + return protocolFixtureReply{result: []any{}} + } + }, nil) + client := protocolFixtureClient(t, fixture.server.URL, ClientConfig{APIKey: "fixture-key", Username: "fixture-user", Timeout: 200 * time.Millisecond}) + err := awaitConfiguredRPCBudget(t, func(ctx context.Context) error { + switch stream { + case "telemetry": + _, err := client.GetSystemTelemetry(ctx) + return err + case "app-stats": + _, err := client.GetAppStats(ctx) + return err + default: + _, err := client.GetAppLogs(ctx, "fixture-app", "fixture-container", 100) + return err + } + }) + assertRPCTimeout(t, err) + if fixture.sessions.Load() != 1 || fixture.restRequests.Load() != 0 || client.TransportStatus().Connected { + t.Fatal("expired stream was retried, downgraded or left reusable") + } + }) + } +} + +func TestJSONRPCWaitingDeadlineDoesNotInterruptSessionOwner(t *testing.T) { + release := make(chan struct{}) + started := make(chan struct{}) + var reads atomic.Int32 + fixture := newProtocolFixture(t, func(_ int, request trueNASRPCRequest) protocolFixtureReply { + if request.Method == "auth.login_ex" { + return protocolFixtureReply{result: map[string]any{"response_type": "SUCCESS"}} + } + if reads.Add(1) == 1 { + close(started) + <-release + } + return protocolFixtureReply{result: map[string]any{"hostname": "session-owner"}} + }, nil) + client := protocolFixtureClient(t, fixture.server.URL, ClientConfig{APIKey: "fixture-key", Username: "fixture-user", Timeout: 5 * time.Second}) + var released atomic.Bool + unblock := func() { + if released.CompareAndSwap(false, true) { + close(release) + } + } + t.Cleanup(unblock) + owner := make(chan error, 1) + go func() { owner <- client.TestConnection(context.Background()) }() + select { + case <-started: + case <-time.After(2 * time.Second): + t.Fatal("session owner did not dispatch") + } + ctx, cancel := context.WithTimeout(context.Background(), 100*time.Millisecond) + defer cancel() + waiter := make(chan error, 1) + go func() { waiter <- client.TestConnection(ctx) }() + select { + case err := <-waiter: + if !errors.Is(err, context.DeadlineExceeded) { + t.Errorf("waiter error = %v, want caller deadline", err) + } + case <-time.After(time.Second): + t.Error("caller deadline could not abandon the session lock wait") + unblock() // Do not strand the predecessor's uncancellable mutex waiter. + <-waiter + } + if reads.Load() != 1 { + t.Error("cancelled waiter dispatched a request") + } + unblock() + if err := <-owner; err != nil { + t.Fatalf("waiter interrupted the session owner: %v", err) + } + if err := client.TestConnection(context.Background()); err != nil || fixture.sessions.Load() != 1 { + t.Fatalf("cancelled waiter poisoned the shared session: %v, sessions=%d", err, fixture.sessions.Load()) + } +} + +func TestJSONRPCPreCancelledCallDoesNotDispatch(t *testing.T) { + var actions atomic.Int32 + fixture := newProtocolFixture(t, func(_ int, request trueNASRPCRequest) protocolFixtureReply { + if request.Method == "auth.login_ex" { + return protocolFixtureReply{result: map[string]any{"response_type": "SUCCESS"}} + } + if request.Method == "app.start" { + actions.Add(1) + } + return protocolFixtureReply{result: map[string]any{"hostname": "cancel-fixture"}} + }, nil) + client := protocolFixtureClient(t, fixture.server.URL, ClientConfig{APIKey: "fixture-key", Username: "fixture-user"}) + if err := client.TestConnection(context.Background()); err != nil { + t.Fatal(err) + } + ctx, cancel := context.WithCancel(context.Background()) + cancel() + if err := client.StartApp(ctx, "fixture-app"); !errors.Is(err, context.Canceled) { + t.Fatalf("pre-cancelled action error = %v", err) + } + if err := client.TestConnection(context.Background()); err != nil || fixture.sessions.Load() != 1 || actions.Load() != 0 { + t.Fatalf("pre-cancelled action dispatched or poisoned the session: %v", err) + } +} + +func TestJSONRPCSnapshotContinuesAfterTelemetryTimeout(t *testing.T) { + fixture := newProtocolFixture(t, func(_ int, request trueNASRPCRequest) protocolFixtureReply { + switch request.Method { + case "auth.login_ex": + return protocolFixtureReply{result: map[string]any{"response_type": "SUCCESS"}} + case "system.info": + return protocolFixtureReply{result: map[string]any{"hostname": "inventory-timeout", "version": "TrueNAS-SCALE-25.10.7", "system_serial": "FIXTURE-1", "physmem": 1024}} + case "core.subscribe": + return protocolFixtureReply{result: "silent-realtime"} + case "pool.query": + return protocolFixtureReply{result: []map[string]any{{"id": 1, "name": "tank", "status": "ONLINE"}}} + default: + return protocolFixtureReply{result: []any{}} + } + }, nil) + client := protocolFixtureClient(t, fixture.server.URL, ClientConfig{APIKey: "fixture-key", Username: "fixture-user", Timeout: 200 * time.Millisecond}) + var snapshot *FixtureSnapshot + err := awaitConfiguredRPCBudget(t, func(ctx context.Context) error { + var err error + snapshot, err = client.FetchSnapshot(ctx) + return err + }) + if err != nil || snapshot == nil || snapshot.System.Hostname != "inventory-timeout" || len(snapshot.Pools) != 1 { + t.Fatalf("optional telemetry timeout prevented usable inventory: %v, %+v", err, snapshot) + } + if snapshot.System.Telemetry == nil || snapshot.System.Telemetry.ErrorCategory != "timeout" || snapshot.System.Telemetry.CPU || snapshot.System.Telemetry.Memory { + t.Fatalf("silent telemetry appeared fresh or lost its timeout: %+v", snapshot.System.Telemetry) + } + if fixture.sessions.Load() != 2 || fixture.restRequests.Load() != 0 { + t.Fatalf("inventory did not recover on RPC after discarding the stream: sessions=%d rest=%d", fixture.sessions.Load(), fixture.restRequests.Load()) + } +}