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
This commit is contained in:
pulse-triage[bot] 2026-10-02 12:27:35 +01:00
parent 6cf72e6d88
commit 7e369d418f
6 changed files with 576 additions and 16 deletions

View file

@ -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

View file

@ -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

View file

@ -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

View file

@ -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.

View file

@ -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)

View file

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