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()) + } +}