mirror of
https://github.com/rcourtman/Pulse.git
synced 2026-10-03 12:47:49 +00:00
Handle exact TrueNAS subscription termination without stalled reads
Recognise matching notify_unsubscribed events before and after acknowledgement, preserve unrelated streams and clean log completion, and retain only numeric errno with fixed diagnostics. Keep nonempty Apps inventory and poll-ledger progress while rejected stats recover on the next ordinary poll. The reporter supplied a source-derived envelope, not a retained capture; native acceptance and delivery remain separate. Change-source: pulse-maintainer
This commit is contained in:
parent
038579a28d
commit
f0ea29480e
6 changed files with 437 additions and 8 deletions
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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",
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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)
|
||||
|
|
|
|||
50
internal/truenas/subscription.go
Normal file
50
internal/truenas/subscription.go
Normal file
|
|
@ -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"}
|
||||
}
|
||||
|
|
@ -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)
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue