diff --git a/docs/release-control/v6/internal/subsystems/monitoring.md b/docs/release-control/v6/internal/subsystems/monitoring.md index c149fa5c1..6265060be 100644 --- a/docs/release-control/v6/internal/subsystems/monitoring.md +++ b/docs/release-control/v6/internal/subsystems/monitoring.md @@ -17,6 +17,32 @@ ## Purpose +### Exact TrueNAS subscription termination — issue #2396 + +JSON-RPC stream readers and the subscription-acknowledgement wait recognise +`notify_unsubscribed` only for the exact requested collection, including its +arguments. An unrelated interval or event source cannot reject the current +subscription. A matching rejection ends the operation promptly, discards the +stream session and is not retried or downgraded to REST. A subsequent ordinary +read may authenticate a fresh session within the existing operation budget. +The notification uses middleware errno fields, not a JSON-RPC response error; +only its numeric errno and a fixed message survive, never provider reason, +trace, extra fields or private collection arguments. Malformed termination +cannot leave a stream reusable. A clean end without telemetry remains +unavailable, not a zero sample; a clean log end returns the bounded collected +tail and retains the socket without cancelling an already-ended subscription. + +`TestJSONRPCSubscriptionTerminationRejectsWithoutWaitingOrRetry` covers all +three readers before/after acknowledgement, cleanup, confidentiality and fresh +read recovery. The adjacent termination controls pin exact-collection +isolation, malformed events, absent samples and clean log completion. +`TestTrueNASPollerSubscriptionRejectionKeepsPollingAndRecovers` exercises +three connected protocol/poller snapshots: usable inventory and poll-ledger +progress during rejected nonempty Apps enrichment, native alert disappearance, +then ordinary stats recovery with stable identity. These are source controls +using the reporter's source-derived envelope, not a retained wire capture, +native appliance/incident acceptance or release availability. + ### RAID count provenance across ingestion — issue #2369 Agent ingestion, host snapshots and canonical read-state projection retain diff --git a/docs/release-control/v6/internal/subsystems/registry.json b/docs/release-control/v6/internal/subsystems/registry.json index 222a9fd47..5b5955934 100644 --- a/docs/release-control/v6/internal/subsystems/registry.json +++ b/docs/release-control/v6/internal/subsystems/registry.json @@ -6162,6 +6162,7 @@ "internal/truenas/disk_health.go", "internal/truenas/fixtures.go", "internal/truenas/provider.go", + "internal/truenas/subscription.go", "internal/truenas/telemetry.go", "internal/truenas/testdata/core13_reporting.json", "internal/truenas/transport.go", diff --git a/internal/monitoring/truenas_poller_test.go b/internal/monitoring/truenas_poller_test.go index 9074c9cdd..fcf3a34c6 100644 --- a/internal/monitoring/truenas_poller_test.go +++ b/internal/monitoring/truenas_poller_test.go @@ -2015,6 +2015,133 @@ func TestTrueNASPollerUnresponsiveRPCDoesNotFreezeOtherConnections(t *testing.T) } } +// #2396's nonempty Apps path: an accepted subscription is then rejected by +// middlewared. Keep the inventory, advance the poll ledger, observe native +// alert disappearance and recover stats on the next ordinary poll. This is a +// protocol/poller control, not a native appliance or incident-delivery proof. +func TestTrueNASPollerSubscriptionRejectionKeepsPollingAndRecovers(t *testing.T) { + var reject atomic.Bool + var sessions, statsSubscriptions atomic.Int32 + upgrader := websocket.Upgrader{} + server := 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() + sessions.Add(1) + for { + var request struct { + ID int64 `json:"id"` + Method string `json:"method"` + Params []json.RawMessage `json:"params"` + } + if conn.ReadJSON(&request) != nil { + return + } + var result any = []any{} + var collection string + switch request.Method { + case "auth.login_ex": + result = map[string]any{"response_type": "SUCCESS"} + case "system.info": + result = map[string]any{"hostname": "reject-nas", "version": "TrueNAS-SCALE-25.04.2.6", "system_serial": "REJECT-FIXTURE"} + case "pool.query": + result = []map[string]any{{"id": 1, "name": "tank", "status": "ONLINE", "size": 1000, "allocated": 400}} + case "app.query": + result = []map[string]any{{"id": "fixture-app", "name": "fixture-app", "state": "RUNNING"}} + case "alert.list": + if !reject.Load() { + result = []map[string]any{{"id": "fixture-alert", "level": "WARNING", "formatted": "fixture warning", "source": "fixture"}} + } + case "core.subscribe": + if len(request.Params) != 1 || json.Unmarshal(request.Params[0], &collection) != nil { + t.Error("invalid subscription") + return + } + result = "fixture-sub" + case "core.unsubscribe": + result = nil + } + if conn.WriteJSON(map[string]any{"jsonrpc": "2.0", "id": request.ID, "result": result}) != nil { + return + } + if collection == "" { + continue + } + var params any + method := "collection_update" + if strings.HasPrefix(collection, "app.stats:") { + statsSubscriptions.Add(1) + if reject.Load() { + method = "notify_unsubscribed" + params = map[string]any{"collection": collection, "error": map[string]any{"error": 14, "errname": "EFAULT", "reason": "[EFAULT] Apps are not available", "trace": nil, "extra": nil}} + } else { + params = map[string]any{"collection": collection, "fields": []any{map[string]any{"app_name": "fixture-app", "cpu_usage": 12}}} + } + } else { + params = map[string]any{"collection": collection, "fields": map[string]any{"cpu": map[string]any{"usage": 10}}} + } + if conn.WriteJSON(map[string]any{"jsonrpc": "2.0", "method": method, "params": params}) != nil { + return + } + } + })) + t.Cleanup(server.Close) + client, err := truenas.NewClient(truenas.ClientConfig{ + Host: server.URL, APIKey: "fixture-key", Username: "fixture-user", Timeout: 2 * time.Second, + InsecureSkipVerify: true, Fingerprint: fmt.Sprintf("%x", sha256.Sum256(server.Certificate().Raw)), + }) + if err != nil { + t.Fatal(err) + } + t.Cleanup(client.Close) + instance := config.TrueNASInstance{ID: "reject-connection", Host: server.URL, Enabled: true} + provider := truenas.NewLiveProviderForConnection(&truenas.APIFetcher{Client: client}, instance.ID) + poller := NewTrueNASPoller(nil, 0, nil) + poller.providersByOrg["default"] = map[string]*truenas.Provider{instance.ID: provider} + poller.configsByOrg["default"] = map[string]config.TrueNASInstance{instance.ID: instance} + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + var lastSuccess time.Time + for poll := 0; poll < 3; poll++ { + reject.Store(poll == 1) + poller.ensureConnectionRuntimeStatusLocked("default", instance.ID).nextPollAt = time.Now().Add(-time.Second) + started := time.Now() + poller.pollAll(ctx) + if elapsed := time.Since(started); elapsed >= time.Second { + t.Fatalf("rejection stalled poll %d for %s", poll, elapsed) + } + summary := poller.ConnectionSummaries("default", []config.TrueNASInstance{instance})[instance.ID] + if summary.Poll.LastError != nil || summary.Poll.LastSuccessAt == nil || !summary.Poll.LastSuccessAt.After(lastSuccess) || summary.Poll.LastAttemptAt == nil { + t.Fatalf("optional stats rejection prevented poll progress: %+v", summary) + } + lastSuccess = *summary.Poll.LastSuccessAt + snapshot := provider.Snapshot() + if snapshot == nil || snapshot.System.Hostname != "reject-nas" || len(snapshot.Pools) != 1 || len(snapshot.Apps) != 1 || snapshot.System.CPUPercent != 10 { + t.Fatalf("poll %d lost usable inventory/telemetry: %+v", poll, snapshot) + } + if poll == 1 { + if snapshot.Apps[0].Stats != nil || len(snapshot.Alerts) != 0 { + t.Fatal("rejected stats fabricated data or left the former native alert") + } + } else if snapshot.Apps[0].Stats == nil || snapshot.Apps[0].Stats.CPUPercent != 12 || len(snapshot.Alerts) != 1 { + t.Fatal("healthy stats/native alerts did not return") + } + if !hasTrueNASHostForOrg(poller, "default", "reject-nas") { + t.Fatal("poll lost the appliance identity") + } + } + if statsSubscriptions.Load() != 3 || sessions.Load() != 2 { + t.Fatalf("rejection replayed or churned later polls: subscriptions=%d sessions=%d", statsSubscriptions.Load(), sessions.Load()) + } +} + type trueNASMockServer struct { server *httptest.Server requests atomic.Int64 diff --git a/internal/truenas/client.go b/internal/truenas/client.go index ada44b00f..b057e40d1 100644 --- a/internal/truenas/client.go +++ b/internal/truenas/client.go @@ -258,7 +258,7 @@ func (c *Client) GetSystemTelemetry(ctx context.Context) (*SystemInfo, error) { if err != nil { return err } - telemetry, err = rpc.readSystemTelemetryEvent(ctx, defaultRealtimeIntervalSeconds) + telemetry, err = rpc.readSystemTelemetryEvent(ctx, defaultRealtimeIntervalSeconds, subscriptionName) if err != nil { return discardRPCSessionForStreamError("reporting.realtime", err) } @@ -1817,7 +1817,7 @@ func (c *Client) GetAppStats(ctx context.Context) (map[string]AppStats, error) { if err != nil { return err } - stats, err = rpc.readAppStatsEvent(ctx, defaultAppStatsIntervalSeconds) + stats, err = rpc.readAppStatsEvent(ctx, defaultAppStatsIntervalSeconds, subscriptionName) if err != nil { return discardRPCSessionForStreamError("app.stats", err) } @@ -1873,7 +1873,10 @@ func (c *Client) GetAppLogs(ctx context.Context, appName, containerID string, ta return err } var reusable bool - lines, reusable, err = rpc.readAppLogEvents(ctx, tailLines) + lines, reusable, err = rpc.readAppLogEvents(ctx, tailLines, subscriptionName) + if errors.Is(err, errRPCSubscriptionComplete) { + return nil + } if err != nil { return discardRPCSessionForStreamError("app.container_log_follow", err) } @@ -2656,6 +2659,21 @@ func (c *trueNASRPCClient) call(ctx context.Context, method string, params any, return &RPCTransportError{Method: method, Phase: "read", Err: err} } if message.Method != "" { + // A rejection can race with the subscription acknowledgement. + // Do not drop it while waiting for the core.subscribe response. + if method == "core.subscribe" { + if args, ok := params.([]any); ok && len(args) == 1 { + if collection, ok := args[0].(string); ok { + ended, err := rpcSubscriptionTermination(message, collection, method) + if ended || err != nil { + if err == nil { + err = rpcSubscriptionEndedBeforeData(method) + } + return discardRPCSessionForStreamError(method, err) + } + } + } + } continue } if message.ID != request.ID { @@ -2674,7 +2692,7 @@ func (c *trueNASRPCClient) call(ctx context.Context, method string, params any, } } -func (c *trueNASRPCClient) readAppStatsEvent(ctx context.Context, intervalSeconds int) (map[string]AppStats, error) { +func (c *trueNASRPCClient) readAppStatsEvent(ctx context.Context, intervalSeconds int, collection string) (map[string]AppStats, error) { if c == nil || c.conn == nil { return nil, fmt.Errorf("truenas rpc connection is nil") } @@ -2690,6 +2708,12 @@ func (c *trueNASRPCClient) readAppStatsEvent(ctx context.Context, intervalSecond if err := c.conn.ReadJSON(&message); err != nil { return nil, &RPCTransportError{Method: "app.stats", Phase: "read", Err: err} } + if ended, err := rpcSubscriptionTermination(message, collection, "app.stats"); ended || err != nil { + if err == nil { + err = rpcSubscriptionEndedBeforeData("app.stats") + } + return nil, err + } if message.Method == "" { if message.Error != nil { return nil, rpcErrorFromWire("app.stats", message.Error) @@ -2742,7 +2766,7 @@ func (c *trueNASRPCClient) readAppStatsEvent(ctx context.Context, intervalSecond } } -func (c *trueNASRPCClient) readAppLogEvents(ctx context.Context, tailLines int) ([]AppLogLine, bool, error) { +func (c *trueNASRPCClient) readAppLogEvents(ctx context.Context, tailLines int, collection string) ([]AppLogLine, bool, error) { if c == nil || c.conn == nil { return nil, false, fmt.Errorf("truenas rpc connection is nil") } @@ -2788,6 +2812,12 @@ func (c *trueNASRPCClient) readAppLogEvents(ctx context.Context, tailLines int) } return nil, false, &RPCTransportError{Method: "app.container_log_follow", Phase: "read", Err: err} } + if ended, err := rpcSubscriptionTermination(message, collection, "app.container_log_follow"); ended || err != nil { + if err != nil { + return nil, false, err + } + return trimAppLogLines(lines, tailLines), true, errRPCSubscriptionComplete + } if message.Method == "" { if message.Error != nil { return nil, false, rpcErrorFromWire("app.container_log_follow", message.Error) @@ -3096,7 +3126,7 @@ func (c *trueNASRPCClient) getReportingDataWithQuery(ctx context.Context, graphs return response, nil } -func (c *trueNASRPCClient) readSystemTelemetryEvent(ctx context.Context, intervalSeconds int) (*SystemInfo, error) { +func (c *trueNASRPCClient) readSystemTelemetryEvent(ctx context.Context, intervalSeconds int, collection string) (*SystemInfo, error) { if c == nil || c.conn == nil { return nil, fmt.Errorf("truenas rpc connection is nil") } @@ -3112,6 +3142,12 @@ func (c *trueNASRPCClient) readSystemTelemetryEvent(ctx context.Context, interva if err := c.conn.ReadJSON(&message); err != nil { return nil, &RPCTransportError{Method: "reporting.realtime", Phase: "read", Err: err} } + if ended, err := rpcSubscriptionTermination(message, collection, "reporting.realtime"); ended || err != nil { + if err == nil { + err = rpcSubscriptionEndedBeforeData("reporting.realtime") + } + return nil, err + } if message.Method == "" { if message.Error != nil { return nil, rpcErrorFromWire("reporting.realtime", message.Error) diff --git a/internal/truenas/subscription.go b/internal/truenas/subscription.go new file mode 100644 index 000000000..cd60b1864 --- /dev/null +++ b/internal/truenas/subscription.go @@ -0,0 +1,50 @@ +package truenas + +import ( + "encoding/json" + "errors" + "fmt" +) + +// A clean terminal event completes a log tail without leaving a subscription +// to cancel. It is not a sample for a telemetry stream. +var errRPCSubscriptionComplete = errors.New("truenas rpc subscription completed") + +// rpcSubscriptionTermination recognises only the exact collection we subscribed +// to, including its arguments. Other subscriptions (even another interval for +// the same event source) must not terminate the current reader. middlewared's +// notification error is not a JSON-RPC response error: it uses errno/reason +// fields, and may contain private names and traces. Retain only its numeric +// errno and a fixed message, never that wire text or the collection arguments. +func rpcSubscriptionTermination(message trueNASRPCResponse, collection, method string) (bool, error) { + if message.Method != "notify_unsubscribed" { + return false, nil + } + var notification struct { + Collection string `json:"collection"` + Error json.RawMessage `json:"error"` + } + if err := json.Unmarshal(message.Params, ¬ification); err != nil { + return false, fmt.Errorf("invalid truenas %s termination notification", method) + } + if notification.Collection != collection { + return false, nil + } + if len(notification.Error) == 0 { + return true, fmt.Errorf("invalid truenas %s termination notification", method) + } + if string(notification.Error) == "null" { + return true, nil + } + var failure struct { + Errno int `json:"error"` + } + if err := json.Unmarshal(notification.Error, &failure); err != nil { + return true, fmt.Errorf("invalid truenas %s termination error", method) + } + return true, &RPCError{Method: method, Code: failure.Errno, Message: "subscription rejected"} +} + +func rpcSubscriptionEndedBeforeData(method string) error { + return &RPCError{Method: method, Message: "subscription ended before a sample"} +} diff --git a/internal/truenas/transport_test.go b/internal/truenas/transport_test.go index cdbf58351..4dc03527f 100644 --- a/internal/truenas/transport_test.go +++ b/internal/truenas/transport_test.go @@ -21,10 +21,192 @@ import ( "github.com/rs/zerolog/log" ) +// Subscription terminal-message regressions for #2396. The envelope is from +// the reporter's source-derived 25.04.2.6 example, not a retained wire capture. +func TestJSONRPCSubscriptionTerminationRejectsWithoutWaitingOrRetry(t *testing.T) { + for _, stream := range []string{"app.stats", "reporting.realtime", "app.container_log_follow"} { + for _, beforeReply := range []bool{false, true} { + t.Run(fmt.Sprintf("%s/before_reply=%t", stream, beforeReply), func(t *testing.T) { + var subscribes, unsubscribes atomic.Int32 + 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 "reporting.get_data": + return protocolFixtureReply{result: []any{}} + case "core.subscribe": + subscribes.Add(1) + collection := request.Params.([]any)[0].(string) + return protocolFixtureReply{result: "private-sub-id", beforeReply: beforeReply, notifications: []protocolFixtureNotification{{ + method: "notify_unsubscribed", params: map[string]any{"collection": collection, "error": map[string]any{ + "error": 14, "errname": "EFAULT", "reason": "private-app-name fixture-key fixture-user https://private.invalid", + "trace": "private-trace", "extra": "private-extra", + }}, + }}} + case "core.unsubscribe": + unsubscribes.Add(1) + return protocolFixtureReply{result: nil} + case "pool.query": + return protocolFixtureReply{result: []any{}} + default: + t.Errorf("unexpected method %s", request.Method) + return protocolFixtureReply{close: true} + } + }, nil) + client := protocolFixtureClient(t, fixture.server.URL, ClientConfig{APIKey: "fixture-key", Username: "fixture-user", Timeout: 2 * time.Second}) + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + started := time.Now() + var err error + switch stream { + case "app.stats": + var stats map[string]AppStats + stats, err = client.GetAppStats(ctx) + if stats != nil { + t.Error("rejected stats fabricated a sample") + } + case "reporting.realtime": + var telemetry *SystemInfo + telemetry, err = client.GetSystemTelemetry(ctx) + if telemetry != nil { + t.Error("rejected telemetry fabricated a sample") + } + default: + var lines []AppLogLine + lines, err = client.GetAppLogs(ctx, "private-app-name", "private-container", 100) + if lines != nil { + t.Error("rejected logs appeared successful") + } + } + var rpcErr *RPCError + if !errors.As(err, &rpcErr) || rpcErr.Code != 14 || rpcErr.Message != "subscription rejected" { + t.Fatalf("termination = %v, want structured errno=14 rejection", err) + } + if elapsed := time.Since(started); elapsed >= time.Second { + t.Fatalf("terminal event waited for the operation budget: %s", elapsed) + } + for _, private := range []string{"private-app-name", "fixture-key", "fixture-user", "private.invalid", "private-trace", "private-extra", "private-sub-id", "private-container"} { + if strings.Contains(err.Error(), private) || strings.Contains(client.TransportStatus().LastError, private) { + t.Fatalf("termination exposed %s", private) + } + } + if client.rpc != nil || client.TransportStatus().Connected || subscribes.Load() != 1 || unsubscribes.Load() != 0 || fixture.sessions.Load() != 1 { + t.Fatal("terminal rejection retried, reused the stream or tried to cancel its ended subscription") + } + if _, err := client.GetPools(ctx); err != nil || fixture.sessions.Load() != 2 || fixture.restRequests.Load() != 0 { + t.Fatalf("next read did not recover on a fresh RPC session: %v", err) + } + }) + } + } +} + +func TestJSONRPCSubscriptionTerminationIgnoresOtherCollections(t *testing.T) { + var unsubscribes atomic.Int32 + 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": + collection := request.Params.([]any)[0].(string) + return protocolFixtureReply{result: "stats-sub", notifications: []protocolFixtureNotification{ + {method: "notify_unsubscribed", params: map[string]any{"collection": "app.stats:{\"interval\":99}", "error": map[string]any{"error": 14}}}, + {method: "notify_unsubscribed", params: map[string]any{"collection": "app.stats.extra:{\"interval\":2}", "error": "malformed-unrelated-error"}}, + {method: "notify_unsubscribed", params: map[string]any{"collection": "reporting.realtime:{\"interval\":2}", "error": nil}}, + {method: "collection_update", params: map[string]any{"collection": collection, "fields": []any{map[string]any{"app_name": "fixture-app", "cpu_usage": 7}}}}, + }} + case "core.unsubscribe": + unsubscribes.Add(1) + return protocolFixtureReply{result: nil} + default: + return protocolFixtureReply{result: []any{}} + } + }, nil) + client := protocolFixtureClient(t, fixture.server.URL, ClientConfig{APIKey: "fixture-key", Username: "fixture-user"}) + for poll := 0; poll < 2; poll++ { + stats, err := client.GetAppStats(context.Background()) + if err != nil || stats["fixture-app"].CPUPercent != 7 { + t.Fatalf("unrelated terminal event interrupted stats: %v, %+v", err, stats) + } + } + if fixture.sessions.Load() != 1 || unsubscribes.Load() != 2 || !client.TransportStatus().Connected { + t.Fatal("unrelated event prevented ordinary cleanup and session reuse") + } +} + +func TestJSONRPCSubscriptionCleanTerminationDoesNotInventStats(t *testing.T) { + 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 == "core.subscribe" { + return protocolFixtureReply{result: "stats-sub", notifications: []protocolFixtureNotification{{method: "notify_unsubscribed", params: map[string]any{"collection": request.Params.([]any)[0], "error": nil}}}} + } + return protocolFixtureReply{result: []any{}} + }, nil) + client := protocolFixtureClient(t, fixture.server.URL, ClientConfig{APIKey: "fixture-key", Username: "fixture-user", Timeout: 2 * time.Second}) + stats, err := client.GetAppStats(context.Background()) + var rpcErr *RPCError + if stats != nil || !errors.As(err, &rpcErr) || rpcErr.Message != "subscription ended before a sample" || client.TransportStatus().Connected { + t.Fatalf("clean termination appeared sampled or left an unresolved stream: %v, %+v", err, stats) + } +} + +func TestJSONRPCSubscriptionCleanLogEndKeepsSession(t *testing.T) { + var unsubscribes 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 == "core.subscribe" { + collection := request.Params.([]any)[0] + return protocolFixtureReply{result: "log-sub", notifications: []protocolFixtureNotification{ + {method: "collection_update", params: map[string]any{"collection": collection, "fields": []any{map[string]any{"data": "fixture-line"}}}}, + {method: "notify_unsubscribed", params: map[string]any{"collection": "app.container_log_follow:{}", "error": map[string]any{"error": 14}}}, + {method: "notify_unsubscribed", params: map[string]any{"collection": collection, "error": nil}}, + }} + } + if request.Method == "core.unsubscribe" { + unsubscribes.Add(1) + } + return protocolFixtureReply{result: []any{}} + }, nil) + client := protocolFixtureClient(t, fixture.server.URL, ClientConfig{APIKey: "fixture-key", Username: "fixture-user"}) + lines, err := client.GetAppLogs(context.Background(), "fixture-app", "fixture-container", 100) + if err != nil || len(lines) != 1 || lines[0].Data != "fixture-line" { + t.Fatalf("clean log end lost the bounded tail: %v, %+v", err, lines) + } + if _, err := client.GetPools(context.Background()); err != nil || fixture.sessions.Load() != 1 || unsubscribes.Load() != 0 || !client.TransportStatus().Connected { + t.Fatalf("clean log end discarded or unsubscribed an already-ended stream: %v", err) + } +} + +func TestJSONRPCSubscriptionMalformedTerminalDiscardsWithoutWireText(t *testing.T) { + for _, terminal := range []any{map[string]any{"error": "private-invalid-errno", "reason": "private-error"}, "private-error", true} { + t.Run(fmt.Sprintf("%T", terminal), func(t *testing.T) { + 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 == "core.subscribe" { + return protocolFixtureReply{result: "stats-sub", notifications: []protocolFixtureNotification{{method: "notify_unsubscribed", params: map[string]any{"collection": request.Params.([]any)[0], "error": terminal}}}} + } + return protocolFixtureReply{result: []any{}} + }, nil) + client := protocolFixtureClient(t, fixture.server.URL, ClientConfig{APIKey: "fixture-key", Username: "fixture-user", Timeout: 2 * time.Second}) + _, err := client.GetAppStats(context.Background()) + if err == nil || strings.Contains(err.Error(), "private-") || !strings.Contains(err.Error(), "invalid truenas app.stats termination error") || client.rpc != nil || fixture.sessions.Load() != 1 { + t.Fatalf("malformed termination was lost, leaked or reused: %v", err) + } + }) + } +} + type protocolFixtureReply struct { result any err *trueNASRPCError notifications []protocolFixtureNotification + beforeReply bool close bool } @@ -89,8 +271,10 @@ func newProtocolFixture( } response.Result = raw } - if err := conn.WriteJSON(response); err != nil { - return + if !reply.beforeReply { + if err := conn.WriteJSON(response); err != nil { + return + } } for _, notification := range reply.notifications { params, err := json.Marshal(notification.params) @@ -106,6 +290,11 @@ func newProtocolFixture( return } } + if reply.beforeReply { + if err := conn.WriteJSON(response); err != nil { + return + } + } } })) t.Cleanup(fixture.server.Close)