From 22ca07ba7171686f1e0904d1d14eb1f08101e134 Mon Sep 17 00:00:00 2001 From: "pulse-triage[bot]" <249995291+pulse-triage[bot]@users.noreply.github.com> Date: Thu, 1 Oct 2026 03:25:28 +0100 Subject: [PATCH] Keep missing TrueNAS samples distinct from observed zero Carry per-metric presence through REST and realtime snapshots, canonical host rows and shared History writes. Do not interpret arbitrary numeric object fields as reporting values. Preserve bounded telemetry failure diagnostics without blocking inventory or exposing provider text. Change-source: pulse-maintainer --- .../v6/internal/subsystems/monitoring.md | 41 +++- .../v6/internal/subsystems/registry.json | 1 + internal/monitoring/truenas_poller.go | 13 ++ internal/monitoring/truenas_poller_test.go | 148 +++++++++++++ internal/truenas/client.go | 70 +++--- internal/truenas/client_test.go | 199 ++++++++++++++++++ internal/truenas/provider.go | 22 +- internal/truenas/telemetry.go | 64 ++++++ internal/truenas/types.go | 18 ++ 9 files changed, 543 insertions(+), 33 deletions(-) create mode 100644 internal/truenas/telemetry.go diff --git a/docs/release-control/v6/internal/subsystems/monitoring.md b/docs/release-control/v6/internal/subsystems/monitoring.md index 453d3a538..8f0ba47e6 100644 --- a/docs/release-control/v6/internal/subsystems/monitoring.md +++ b/docs/release-control/v6/internal/subsystems/monitoring.md @@ -150,9 +150,9 @@ host-chart history are reconstructed from the REST reporting API `query` parameters). The reporting `memory` and `arcsize` graphs supply the available and ARC readings that the JSON-RPC realtime subscription otherwise provides. When no available-memory reading can be established, the provider -must not derive usage from a zero available value: the system memory metric is +must not derive usage from an absent available value: the system memory metric is omitted and agent memory is projected as usage-unavailable with its known total, -so a CORE appliance is never reported as 100% used. This is a behavioral +so missing data never reports a CORE appliance as 100% used. This is a behavioral correctness repair with no public API or schema delta. `TestRESTSystemTelemetryReadsReportingMemory`, `TestRESTSystemMetricHistoryUsesReporting` and @@ -163,6 +163,43 @@ and projection proofs, not native CORE appliance acceptance or reporter confirmation. + +### TrueNAS partial samples and nonblocking telemetry diagnostics — issue #2077 + +Live system identity and both telemetry parsers carry explicit per-metric +availability. A reporting timestamp or realtime interval is not evidence that +CPU, memory or any sibling I/O series was measured. Each canonical metric is +omitted unless its own reading exists; observed zeros remain valid, including +zero free memory. Capacity-only snapshots keep hardware capacity and mark +usage unavailable. Older static snapshots retain their compatibility projection; +no live client uses that fallback. Null/malformed single-series objects never +substitute a timestamp or another arbitrary numeric field for the missing value. +Supported generic `value`, `y` and `temperature` row aliases remain supported. + +Optional telemetry failure must not discard inventory or degrade a successful +inventory poll. Existing connection diagnostics expose optional +`observed.telemetry` with six availability booleans (`cpu`, `memory`, `netIn`, +`netOut`, `diskRead`, `diskWrite`), a fixed `errorCategory` and, for HTTP failures, +`httpStatus`. No raw response body, endpoint, key, hostname or arbitrary provider +error text is copied into this projection. Error categories record the observed +collection result, not its cause. A later refresh replaces availability and +clears obsolete failures; snapshot and connection-summary copies cannot mutate +shared state. This is an additive diagnostic field on the existing read surface, +not a separate resource, route, telemetry store or alert policy. + +`TestRESTReportingObservationPresence`, +`TestRESTSnapshotRetainsSanitizedTelemetryFailure`, `TestRealtimeObservationPresence` +and `TestSystemTelemetryFailureCategoriesAreBounded` in +`internal/truenas/client_test.go` pin repeated full snapshots, canonical rows, +native History/zero semantics and bounded failure/recovery controls. +`TestTrueNASPartialReportingPipeline` in +`internal/monitoring/truenas_poller_test.go` exercises the HTTP client, poller, +registry, memory metadata, shared metrics writer and in-memory/persisted chart +readbacks across partial, missing and zero-valued cycles. These synthetic proofs +are not the reporter's native response schema, installed CORE acceptance or a +claim that #2077's missing graph journey is resolved. + + **Legacy reporting failure isolation and memory components — issue #2077** Live telemetry and history retain the legacy REST graph/query contract and the diff --git a/docs/release-control/v6/internal/subsystems/registry.json b/docs/release-control/v6/internal/subsystems/registry.json index 7f8b616f6..feeb999fb 100644 --- a/docs/release-control/v6/internal/subsystems/registry.json +++ b/docs/release-control/v6/internal/subsystems/registry.json @@ -6105,6 +6105,7 @@ "internal/truenas/disk_health.go", "internal/truenas/fixtures.go", "internal/truenas/provider.go", + "internal/truenas/telemetry.go", "internal/truenas/transport.go", "internal/truenas/types.go" ], diff --git a/internal/monitoring/truenas_poller.go b/internal/monitoring/truenas_poller.go index b44e26049..1f18061c1 100644 --- a/internal/monitoring/truenas_poller.go +++ b/internal/monitoring/truenas_poller.go @@ -57,6 +57,9 @@ type TrueNASConnectionObservedSummary struct { Shares int `json:"shares"` Disks int `json:"disks"` RecoveryArtifacts int `json:"recoveryArtifacts"` + // Telemetry is separate from poll success: inventory can be current even + // when optional CPU/memory/IO collection failed or only partly succeeded. + Telemetry *truenas.SystemTelemetryAvailability `json:"telemetry,omitempty"` } // TrueNASConnectionSummary merges poll health with the most recent discovered @@ -930,6 +933,7 @@ func buildTrueNASObservedSummary(snapshot *truenas.FixtureSnapshot) *TrueNASConn Shares: len(snapshot.Shares), Disks: len(snapshot.Disks), RecoveryArtifacts: len(snapshot.ZFSSnapshots) + len(snapshot.ReplicationTasks), + Telemetry: cloneTrueNASTelemetryAvailability(snapshot.System.Telemetry), } if host != "" || resourceID != "" { summary.Systems = 1 @@ -989,9 +993,18 @@ func cloneTrueNASObservedSummary(value *TrueNASConnectionObservedSummary) *TrueN Shares: value.Shares, Disks: value.Disks, RecoveryArtifacts: value.RecoveryArtifacts, + Telemetry: cloneTrueNASTelemetryAvailability(value.Telemetry), } } +func cloneTrueNASTelemetryAvailability(value *truenas.SystemTelemetryAvailability) *truenas.SystemTelemetryAvailability { + if value == nil { + return nil + } + cloned := *value + return &cloned +} + func (p *TrueNASPoller) ingestRecoveryPoints(ctx context.Context, orgID string, connectionID string, provider *truenas.Provider) { if p == nil || p.recoveryManager == nil || provider == nil { return diff --git a/internal/monitoring/truenas_poller_test.go b/internal/monitoring/truenas_poller_test.go index caa605a7e..af4425c40 100644 --- a/internal/monitoring/truenas_poller_test.go +++ b/internal/monitoring/truenas_poller_test.go @@ -2,6 +2,7 @@ package monitoring import ( "context" + "encoding/json" "errors" "fmt" "net" @@ -16,8 +17,10 @@ import ( "time" "github.com/rcourtman/pulse-go-rewrite/internal/config" + "github.com/rcourtman/pulse-go-rewrite/internal/models" "github.com/rcourtman/pulse-go-rewrite/internal/truenas" "github.com/rcourtman/pulse-go-rewrite/internal/unifiedresources" + "github.com/rcourtman/pulse-go-rewrite/pkg/metrics" ) func TestTrueNASPollerPollsConfiguredConnections(t *testing.T) { @@ -2205,3 +2208,148 @@ func TestTrueNASSuccessfulPollCadenceIncludesBoundedIdleGap(t *testing.T) { }) } } + +// HTTP/client -> poller -> registry -> row/metric writer -> chart readback. +// Synthetic response shapes prove omission/zero handling, not CORE acceptance. +func TestTrueNASPartialReportingPipeline(t *testing.T) { + previous := truenas.IsFeatureEnabled() + truenas.SetFeatureEnabled(true) + t.Cleanup(func() { truenas.SetFeatureEnabled(previous) }) + cycles := []struct { + body string + status int + want map[string]float64 + errorCategory string + }{ + {`[{"name":"memory","legend":["free"],"data":[[1789000060,8]]}]`, 200, map[string]float64{"memory": 50}, ""}, + {`[{"name":"cpu","legend":["usage"],"data":[[1789000060,0]]}]`, 200, map[string]float64{"cpu": 0}, ""}, + {`provider-private-text`, 401, nil, "authentication"}, + {`[{"name":"cpu","legend":["usage"],"data":[[1789000060,0]]},{"name":"memory","legend":["free"],"data":[[1789000060,0]]},{"name":"interface","legend":["received"],"data":[[1789000060,0]]},{"name":"disk","legend":["write"],"data":[[1789000060,0]]}]`, 200, map[string]float64{"cpu": 0, "memory": 100, "netin": 0, "diskwrite": 0}, ""}, + } + var current atomic.Int64 + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + w.Header().Set("Content-Type", "application/json") + switch r.URL.Path { + case "/api/v2.0/system/info": + _, _ = w.Write([]byte(`{"hostname":"synthetic-core","version":"TrueNAS-13.0-U6.1","cores":8,"physmem":16}`)) + case "/api/v2.0/pool": + _, _ = w.Write([]byte(`[{"id":1,"name":"tank","status":"ONLINE","size":100,"allocated":50,"free":50}]`)) + case "/api/v2.0/pool/dataset", "/api/v2.0/disk", "/api/v2.0/alert/list": + _, _ = w.Write([]byte(`[]`)) + case "/api/v2.0/reporting/get_data": + cycle := cycles[current.Load()] + w.WriteHeader(cycle.status) + _, _ = w.Write([]byte(cycle.body)) + default: + http.NotFound(w, r) + } + })) + defer server.Close() + client, err := truenas.NewClient(truenas.ClientConfig{Host: server.URL, APIKey: "synthetic"}) + if err != nil { + t.Fatal(err) + } + defer client.Close() + provider := truenas.NewLiveProviderForConnection(&truenas.APIFetcher{Client: client}, "synthetic-connection") + instance := config.TrueNASInstance{ID: "synthetic-connection", Host: server.URL, Enabled: true} + 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} + resourceStore := unifiedresources.NewMonitorAdapter(unifiedresources.NewRegistry(unifiedresources.NewMemoryStore())) + persistent, err := metrics.NewStore(metrics.DefaultConfig(t.TempDir())) + if err != nil { + t.Fatal(err) + } + defer func() { _ = persistent.Close() }() + monitor := &Monitor{metricsHistory: NewMetricsHistory(1024, 24*time.Hour), metricsStore: persistent} + counts := make(map[string]int) + observedAt := time.Now().UTC().Add(-10 * time.Minute).Truncate(time.Second) + for index, cycle := range cycles { + current.Store(int64(index)) + poller.ensureConnectionRuntimeStatusLocked("default", instance.ID).nextPollAt = time.Time{} + poller.pollAll(context.Background()) + summary := poller.ConnectionSummaries("default", []config.TrueNASInstance{instance})[instance.ID] + if summary.Poll == nil || summary.Poll.LastSuccessAt == nil || summary.Poll.LastError != nil || summary.Poll.ConsecutiveFailures != 0 || summary.Observed == nil || summary.Observed.StoragePools != 1 { + t.Fatalf("cycle %d lost successful inventory health: %+v", index, summary) + } + availability := summary.Observed.Telemetry + if availability == nil || availability.ErrorCategory != cycle.errorCategory { + t.Fatalf("cycle %d lost telemetry result: %+v", index, availability) + } + for key, present := range map[string]bool{"cpu": availability.CPU, "memory": availability.Memory, "netin": availability.NetIn, "netout": availability.NetOut, "diskread": availability.DiskRead, "diskwrite": availability.DiskWrite} { + _, want := cycle.want[key] + if present != want { + t.Errorf("cycle %d diagnostics %s presence=%v, want %v", index, key, present, want) + } + } + serialized, err := json.Marshal(summary) + if err != nil || strings.Contains(string(serialized), "provider-private-text") || strings.Contains(string(serialized), "synthetic\"") { + t.Fatalf("cycle %d secret-bearing diagnostic: %s %v", index, serialized, err) + } + availability.ErrorCategory = "mutated" + if poller.ConnectionSummaries("default", []config.TrueNASInstance{instance})[instance.ID].Observed.Telemetry.ErrorCategory != cycle.errorCategory { + t.Fatal("connection diagnostics share mutable provider state") + } + records := poller.GetCurrentRecordsForOrg("default") + // Space the modeled poll observations apart without a wall-clock sleep; + // the writer/store deliberately coalesce identical second timestamps. + for i := range records { + records[i].Resource.LastSeen = observedAt.Add(time.Duration(index) * time.Minute) + records[i].Resource.UpdatedAt = records[i].Resource.LastSeen + } + resourceStore.PopulateSnapshotAndSupplemental(models.StateSnapshot{}, map[unifiedresources.DataSource][]unifiedresources.IngestRecord{unifiedresources.SourceTrueNAS: records}) + var host *unifiedresources.Resource + for _, resource := range resourceStore.GetAll() { + if resource.Type == unifiedresources.ResourceTypeAgent { + host = &resource + } + } + if host == nil || host.Agent == nil || host.Agent.CPUCount != 8 || host.Agent.Memory == nil || host.Agent.Memory.Total != 16 { + t.Fatalf("cycle %d lost host/capacity: %+v", index, host) + } + if host.Agent.Memory.UsageUnavailable != !availability.Memory { + t.Fatalf("cycle %d stale/fabricated memory metadata: %+v", index, host.Agent.Memory) + } + var writes []metrics.WriteMetric + monitor.syncUnifiedAgentMetrics(resourceStore, &writes) + seen := make(map[string]bool) + for _, write := range writes { + want, present := cycle.want[write.MetricType] + if write.MetricType == "disk" { + want, present = 50, true // independently collected pool capacity + } + if !present || write.Value != want || write.ResourceID != instance.ID || write.ResourceType != "agent" || seen[write.MetricType] { + t.Errorf("cycle %d wrong/duplicate/missing-source write: %+v", index, write) + } + seen[write.MetricType] = true + counts[write.MetricType]++ + } + if len(writes) != len(cycle.want)+1 { + t.Fatalf("cycle %d writes=%+v, want only observed metrics plus pool usage", index, writes) + } + persistent.WriteBatchSync(writes) + inMemory := monitor.GetGuestMetrics("agent:"+instance.ID, time.Hour) + chart := monitor.GetGuestMetricsForChart("agent:"+instance.ID, "agent", instance.ID, time.Hour) + for _, key := range []string{"cpu", "memory", "disk", "netin", "netout", "diskread", "diskwrite"} { + if len(inMemory[key]) != counts[key] || len(chart[key]) != counts[key] { + t.Errorf("cycle %d %s chart contains fabricated/lost points: memory=%d chart=%d want=%d", index, key, len(inMemory[key]), len(chart[key]), counts[key]) + } + } + id, native, err := provider.SystemMetricHistory(context.Background(), time.Hour) + if cycle.status != 200 { + if err == nil { + t.Fatal("failed native History read was disguised as success") + } + continue + } + if err != nil || id != instance.ID || len(native) != len(cycle.want) { + t.Fatalf("cycle %d native History=%+v id=%q err=%v", index, native, id, err) + } + for key, want := range cycle.want { + points := native[key] + if len(points) != 1 || points[0].Value != want || points[0].Timestamp.Unix() != 1789000060 { + t.Errorf("cycle %d %s native data/timestamp altered: %+v", index, key, points) + } + } + } +} diff --git a/internal/truenas/client.go b/internal/truenas/client.go index 6d51ff1c0..81671d09b 100644 --- a/internal/truenas/client.go +++ b/internal/truenas/client.go @@ -232,6 +232,7 @@ func systemInfoFromResponse(response systemInfoResponse) *SystemInfo { MachineID: machineID, CPUCount: cpuCount, MemoryTotalBytes: response.Physmem, + Telemetry: &SystemTelemetryAvailability{}, } } @@ -289,42 +290,48 @@ func (c *Client) getSystemTelemetryREST(ctx context.Context) (*SystemInfo, error } history := parseSystemMetricHistory(response) if history == nil { - return nil, fmt.Errorf("truenas legacy REST reporting returned no system telemetry") + return nil, errNoSystemTelemetry } return systemInfoFromMetricHistory(history), nil } // systemInfoFromMetricHistory maps the latest reporting sample onto live system -// telemetry. Series that the appliance did not report stay zero so callers can -// distinguish "absent" from a genuine zero. +// telemetry. Numeric fields alone cannot distinguish absent from zero; retain +// presence independently so one successful graph cannot fabricate its siblings. func systemInfoFromMetricHistory(history *SystemMetricHistory) *SystemInfo { if history == nil { return nil } - system := &SystemInfo{CollectedAt: time.Now().UTC()} + system := &SystemInfo{CollectedAt: time.Now().UTC(), Telemetry: &SystemTelemetryAvailability{}} if value, ok := latestTimeSeriesValue(history.CPUPercent); ok { system.CPUPercent = value + system.Telemetry.CPU = true } if value, ok := latestTimeSeriesValue(history.MemoryTotalBytes); ok { system.MemoryTotalBytes = int64(value) } if value, ok := latestTimeSeriesValue(history.MemoryAvailableBytes); ok { system.MemoryAvailableBytes = int64(value) + system.Telemetry.Memory = value >= 0 } if value, ok := latestTimeSeriesValue(history.ARCSizeBytes); ok { system.ARCSizeBytes = int64(value) } if value, ok := latestTimeSeriesValue(history.NetInRate); ok { system.NetInRate = value + system.Telemetry.NetIn = true } if value, ok := latestTimeSeriesValue(history.NetOutRate); ok { system.NetOutRate = value + system.Telemetry.NetOut = true } if value, ok := latestTimeSeriesValue(history.DiskReadRate); ok { system.DiskReadRate = value + system.Telemetry.DiskRead = true } if value, ok := latestTimeSeriesValue(history.DiskWriteRate); ok { system.DiskWriteRate = value + system.Telemetry.DiskWrite = true } return system } @@ -2036,7 +2043,11 @@ func (c *Client) FetchSnapshot(ctx context.Context) (*FixtureSnapshot, error) { if err != nil { return nil, fmt.Errorf("fetch truenas system info: %w", err) } - if telemetry, err := c.GetSystemTelemetry(ctx); err == nil && telemetry != nil { + if telemetry, err := c.GetSystemTelemetry(ctx); err != nil { + // Inventory remains usable. Retain only a fixed error category/status, + // never the provider's body, endpoint, credentials or raw error text. + system.Telemetry = systemTelemetryFailure(err) + } else if telemetry != nil { mergeSystemTelemetry(system, telemetry) } @@ -2101,7 +2112,7 @@ func mergeSystemTelemetry(system *SystemInfo, telemetry *SystemInfo) { if telemetry.MemoryTotalBytes > 0 { system.MemoryTotalBytes = telemetry.MemoryTotalBytes } - if telemetry.MemoryAvailableBytes > 0 { + if telemetry.MemoryAvailableBytes > 0 || (telemetry.Telemetry != nil && telemetry.Telemetry.Memory) { system.MemoryAvailableBytes = telemetry.MemoryAvailableBytes } if telemetry.ARCSizeBytes > 0 { @@ -2121,6 +2132,10 @@ func mergeSystemTelemetry(system *SystemInfo, telemetry *SystemInfo) { if !telemetry.CollectedAt.IsZero() { system.CollectedAt = telemetry.CollectedAt } + if telemetry.Telemetry != nil { + availability := *telemetry.Telemetry + system.Telemetry = &availability + } } func temperatureForTrueNASDisk(temperatures map[string]int, item diskResponse) int { @@ -3084,20 +3099,14 @@ func parseRealtimeFields(message trueNASRPCResponse, collectionPrefix string) (m } func parseSystemTelemetry(fields map[string]any, intervalSeconds int, collectedAt time.Time) *SystemInfo { - if len(fields) == 0 { - return &SystemInfo{ - IntervalSeconds: intervalSeconds, - CollectedAt: collectedAt, - } - } - telemetry := &SystemInfo{ IntervalSeconds: intervalSeconds, CollectedAt: collectedAt, + Telemetry: &SystemTelemetryAvailability{}, } cpu := readMapAny(fields, "cpu") - cpuPercent := readFloatAny(cpu, + cpuPercent, hasCPU := readFloatValueAny(cpu, "usage", "percent", "usage_percent", @@ -3106,18 +3115,21 @@ func parseSystemTelemetry(fields map[string]any, intervalSeconds int, collectedA "total", "overall", ) - if cpuPercent == 0 { + if !hasCPU { if usage := readMapAny(cpu, "usage", "total"); len(usage) > 0 { - cpuPercent = readFloatAny(usage, "percent", "value", "usage") + cpuPercent, hasCPU = readFloatValueAny(usage, "percent", "value", "usage") } } telemetry.CPUPercent = cpuPercent + telemetry.Telemetry.CPU = hasCPU memory := readMapAny(fields, "memory") total := readInt64Any(memory, "physical_memory_total", "total", "memory_total", "total_bytes") available := readInt64Any(memory, "physical_memory_available", "available", "free", "available_bytes", "free_bytes") + _, hasMemory := readFloatValueAny(memory, "physical_memory_available", "available", "free", "available_bytes", "free_bytes") telemetry.MemoryTotalBytes = total telemetry.MemoryAvailableBytes = available + telemetry.Telemetry.Memory = hasMemory && available >= 0 telemetry.ARCSizeBytes = readInt64Any(memory, "arc_size") interfaces := readMapAny(fields, "interfaces") @@ -3126,7 +3138,7 @@ func parseSystemTelemetry(fields map[string]any, intervalSeconds int, collectedA if !ok { continue } - telemetry.NetInRate += readFloatAny(record, + in, hasIn := readFloatValueAny(record, "rx_bytes", "received_bytes", "received_bytes_rate", @@ -3134,7 +3146,7 @@ func parseSystemTelemetry(fields map[string]any, intervalSeconds int, collectedA "bytes_recv", "bytes_received", ) - telemetry.NetOutRate += readFloatAny(record, + out, hasOut := readFloatValueAny(record, "tx_bytes", "sent_bytes", "sent_bytes_rate", @@ -3142,6 +3154,10 @@ func parseSystemTelemetry(fields map[string]any, intervalSeconds int, collectedA "bytes_sent", "bytes_transmitted", ) + telemetry.NetInRate += in + telemetry.NetOutRate += out + telemetry.Telemetry.NetIn = telemetry.Telemetry.NetIn || hasIn + telemetry.Telemetry.NetOut = telemetry.Telemetry.NetOut || hasOut } disks := readMapAny(fields, "disks", "disls") @@ -3150,8 +3166,12 @@ func parseSystemTelemetry(fields map[string]any, intervalSeconds int, collectedA if !ok { continue } - telemetry.DiskReadRate += readFloatAny(record, "read_bytes", "read_bytes_rate", "bytes_read") - telemetry.DiskWriteRate += readFloatAny(record, "write_bytes", "write_bytes_rate", "bytes_written") + read, hasRead := readFloatValueAny(record, "read_bytes", "read_bytes_rate", "bytes_read") + write, hasWrite := readFloatValueAny(record, "write_bytes", "write_bytes_rate", "bytes_written") + telemetry.DiskReadRate += read + telemetry.DiskWriteRate += write + telemetry.Telemetry.DiskRead = telemetry.Telemetry.DiskRead || hasRead + telemetry.Telemetry.DiskWrite = telemetry.Telemetry.DiskWrite || hasWrite } return telemetry @@ -3634,11 +3654,11 @@ func extractReportingLegendFloatValues(raw any, legends []string) map[string]flo } } if len(values) == 0 && len(legends) == 1 { - for _, value := range typed { - if parsed, ok := parseFloat64Any(value); ok { - values[legends[0]] = parsed - break - } + // A missing/null legend value must not fall back to an unrelated + // numeric field (notably the row's timestamp). Keep only supported + // generic single-series value aliases. + if parsed, ok := readFloatValueAny(typed, "value", "y", "temperature"); ok { + values[legends[0]] = parsed } } } diff --git a/internal/truenas/client_test.go b/internal/truenas/client_test.go index f5818d337..f5d338e39 100644 --- a/internal/truenas/client_test.go +++ b/internal/truenas/client_test.go @@ -15,6 +15,7 @@ import ( "time" "github.com/gorilla/websocket" + "github.com/rcourtman/pulse-go-rewrite/internal/unifiedresources" ) type apiResponse struct { @@ -2344,3 +2345,201 @@ func TestRESTReportingRequestMatchesNativeGraphShape(t *testing.T) { } } } + +// These are synthetic response controls, not a capture from #2077's appliance. +// Pin the complete snapshot -> canonical row -> native History path so an +// optional graph cannot turn a missing measurement into a zero-valued series. +func TestRESTReportingObservationPresence(t *testing.T) { + previous := IsFeatureEnabled() + SetFeatureEnabled(true) + t.Cleanup(func() { SetFeatureEnabled(previous) }) + for _, tc := range []struct { + name string + body string + want map[string]float64 + }{ + {"memory_only", `[{"name":"memory","legend":["free"],"data":[[1789000060,8]]}]`, map[string]float64{"memory": 50}}, + {"arc_only", `[{"name":"arcsize","legend":["size"],"data":[[1789000060,8]]}]`, nil}, + {"capacity_only", `[{"name":"memory","legend":["total"],"data":[[1789000060,16]]}]`, nil}, + {"cpu_zero", `[{"name":"cpu","legend":["idle"],"data":[[1789000060,100]]}]`, map[string]float64{"cpu": 0}}, + {"free_zero", `[{"name":"memory","legend":["free"],"data":[[1789000060,0]]}]`, map[string]float64{"memory": 100}}, + {"network_in_zero", `[{"name":"interface","legend":["received"],"data":[[1789000060,0]]}]`, map[string]float64{"netin": 0}}, + {"disk_write_zero", `[{"name":"disk","legend":["write"],"data":[[1789000060,0]]}]`, map[string]float64{"diskwrite": 0}}, + {"null_array", `[{"name":"cpu","legend":["usage"],"data":[[1789000060,null]]}]`, nil}, + {"null_object", `[{"name":"cpu","legend":["usage"],"data":[{"timestamp":1789000060,"usage":null}]}]`, nil}, + {"malformed_object", `[{"name":"cpu","legend":["usage"],"data":[{"timestamp":1789000060,"usage":"missing","unrelated":10}]}]`, nil}, + {"generic_zero", `[{"name":"cpu","legend":["usage"],"data":[{"timestamp":1789000060,"value":0}]}]`, map[string]float64{"cpu": 0}}, + {"empty", `[]`, nil}, + } { + t.Run(tc.name, func(t *testing.T) { + routes := alertArgsTransport(defaultAPIResponses()) + routes["/api/v2.0/system/info"] = apiResponse{body: `{"hostname":"synthetic-core","version":"TrueNAS-13.0-U6.1","physmem":16,"cores":8}`} + routes["/api/v2.0/reporting/get_data"] = apiResponse{body: tc.body} + client := newLegacyRESTReportingClient(t, routes) + defer client.Close() + provider := NewLiveProviderForConnection(&APIFetcher{Client: client}, "synthetic-connection") + for cycle := 0; cycle < 2; cycle++ { + if err := provider.Refresh(context.Background()); err != nil { + t.Fatal(err) + } + snapshot := provider.Snapshot() + if snapshot.System.CPUCount != 8 || snapshot.System.MemoryTotalBytes != 16 || len(snapshot.Pools) == 0 || len(snapshot.Disks) == 0 { + t.Fatal("partial/unavailable telemetry must not discard hardware or inventory") + } + var metrics *unifiedresources.ResourceMetrics + for _, record := range provider.Records() { + if record.Resource.Type == unifiedresources.ResourceTypeAgent { + metrics = record.Resource.Metrics + } + } + if metrics == nil { + t.Fatal("missing canonical host row") + } + for key, metric := range map[string]*unifiedresources.MetricValue{ + "cpu": metrics.CPU, "memory": metrics.Memory, "netin": metrics.NetIn, + "netout": metrics.NetOut, "diskread": metrics.DiskRead, "diskwrite": metrics.DiskWrite, + } { + want, present := tc.want[key] + if (metric != nil) != present || (metric != nil && metric.Value != want) { + t.Errorf("cycle %d %s: metric=%+v, want present=%v value=%v", cycle, key, metric, present, want) + } + } + id, history, err := provider.SystemMetricHistory(context.Background(), time.Hour) + if err != nil || id != "synthetic-connection" { + t.Fatalf("native History target = %q, err=%v", id, err) + } + for _, key := range []string{"cpu", "memory", "netin", "netout", "diskread", "diskwrite"} { + want, present := tc.want[key] + points := history[key] + if (len(points) > 0) != present || (len(points) > 0 && (len(points) != 1 || points[0].Value != want || points[0].Timestamp.Unix() != 1789000060)) { + t.Errorf("cycle %d %s: history=%+v, want present=%v value=%v at original timestamp", cycle, key, points, present, want) + } + } + } + }) + } +} + +func TestRESTSnapshotRetainsSanitizedTelemetryFailure(t *testing.T) { + for _, tc := range []struct { + name string + status int + body, category string + }{ + {"no_samples", 200, `[]`, "no_samples"}, + {"bad_response", 200, `{`, "invalid_response"}, + {"unauthorized", 401, `provider-private-text`, "authentication"}, + {"forbidden", 403, `provider-private-text`, "authentication"}, + {"missing_endpoint", 404, `provider-private-text`, "unsupported_endpoint"}, + {"rejected_graphs", 422, `provider-private-text`, "request_rejected"}, + {"rate_limited", 429, `provider-private-text`, "rate_limited"}, + {"unavailable", 503, `provider-private-text`, "http_error"}, + } { + t.Run(tc.name, func(t *testing.T) { + routes := alertArgsTransport(defaultAPIResponses()) + routes["/api/v2.0/system/info"] = apiResponse{body: `{"hostname":"synthetic-core","version":"TrueNAS-13.0-U6.1","physmem":16}`} + transport := &isolatedReportingTransport{routes: routes, status: tc.status, body: tc.body} + client := newLegacyRESTReportingClient(t, routes) + client.httpClient.Transport = transport + defer client.Close() + provider := NewLiveProvider(&APIFetcher{Client: client}) + if err := provider.Refresh(context.Background()); err != nil { + t.Fatal(err) + } + snapshot := provider.Snapshot() + availability := snapshot.System.Telemetry + if availability == nil || availability.ErrorCategory != tc.category { + t.Fatalf("missing optional telemetry failure: %+v", availability) + } + wantStatus := tc.status + if tc.status == 200 { + wantStatus = 0 + } + if availability.HTTPStatus != wantStatus || len(snapshot.Pools) == 0 || len(snapshot.Disks) == 0 { + t.Fatalf("unexpected status/inventory: %+v", snapshot) + } + serialized, err := json.Marshal(availability) + if err != nil || strings.Contains(string(serialized), "provider-private-text") || strings.Contains(string(serialized), "truenas.invalid") || strings.Contains(string(serialized), "synthetic") { + t.Fatalf("diagnostics exposed provider text: %s, err=%v", serialized, err) + } + wantCalls := 1 + if tc.status == 422 { + wantCalls = 6 + } + if transport.calls != wantCalls { + t.Fatalf("reporting calls=%d, want %d", transport.calls, wantCalls) + } + availability.ErrorCategory = "mutated" + if provider.Snapshot().System.Telemetry.ErrorCategory != tc.category { + t.Fatal("returned snapshot shares mutable telemetry diagnostics") + } + transport.status, transport.body = 200, `[{"name":"cpu","legend":["usage"],"data":[[1789000060,0]]}]` + if err := provider.Refresh(context.Background()); err != nil { + t.Fatal(err) + } + recovered := provider.Snapshot().System.Telemetry + if recovered == nil || !recovered.CPU || recovered.Memory || recovered.NetIn || recovered.NetOut || recovered.DiskRead || recovered.DiskWrite || recovered.ErrorCategory != "" || recovered.HTTPStatus != 0 { + t.Fatalf("recovery retained failure or invented readings: %+v", recovered) + } + }) + } +} + +func TestRealtimeObservationPresence(t *testing.T) { + for _, tc := range []struct { + name string + fields map[string]any + want SystemTelemetryAvailability + }{ + {"empty", nil, SystemTelemetryAvailability{}}, + {"cpu_zero", map[string]any{"cpu": map[string]any{"usage": 0}}, SystemTelemetryAvailability{CPU: true}}, + {"cpu_nested_zero", map[string]any{"cpu": map[string]any{"usage": map[string]any{"percent": 0}}}, SystemTelemetryAvailability{CPU: true}}, + {"cpu_null", map[string]any{"cpu": map[string]any{"usage": nil}}, SystemTelemetryAvailability{}}, + {"capacity_only", map[string]any{"memory": map[string]any{"total": 16}}, SystemTelemetryAvailability{}}, + {"free_zero", map[string]any{"memory": map[string]any{"total": 16, "available": 0}}, SystemTelemetryAvailability{Memory: true}}, + {"free_null", map[string]any{"memory": map[string]any{"total": 16, "available": nil}}, SystemTelemetryAvailability{}}, + {"network_in_zero", map[string]any{"interfaces": map[string]any{"eth0": map[string]any{"rx_bytes": 0}}}, SystemTelemetryAvailability{NetIn: true}}, + {"disk_write_zero", map[string]any{"disks": map[string]any{"sda": map[string]any{"write_bytes": 0}}}, SystemTelemetryAvailability{DiskWrite: true}}, + } { + t.Run(tc.name, func(t *testing.T) { + telemetry := parseSystemTelemetry(tc.fields, 2, time.Now().UTC()) + if telemetry.Telemetry == nil || *telemetry.Telemetry != tc.want { + t.Fatalf("presence=%+v, want %+v", telemetry.Telemetry, tc.want) + } + system := systemInfoFromResponse(systemInfoResponse{Physmem: 16}) + mergeSystemTelemetry(system, telemetry) + metrics := metricsFromTrueNASSystem(*system, 0, 0) + if (metrics.CPU != nil) != tc.want.CPU || (metrics.Memory != nil) != tc.want.Memory || + (metrics.NetIn != nil) != tc.want.NetIn || (metrics.NetOut != nil) != tc.want.NetOut || + (metrics.DiskRead != nil) != tc.want.DiskRead || (metrics.DiskWrite != nil) != tc.want.DiskWrite { + t.Fatalf("timestamp/interval fabricated missing siblings: %+v", metrics) + } + telemetry.Telemetry.CPU = !tc.want.CPU + if system.Telemetry.CPU != tc.want.CPU { + t.Fatal("merge retained caller-owned presence pointer") + } + }) + } +} + +func TestSystemTelemetryFailureCategoriesAreBounded(t *testing.T) { + for _, tc := range []struct { + err error + want string + }{ + {fmt.Errorf("wrapped: %w", context.Canceled), "cancelled"}, + {fmt.Errorf("wrapped: %w", context.DeadlineExceeded), "timeout"}, + {&RPCAuthError{Mechanism: "provider-private-text"}, "authentication"}, + {&RPCError{Message: "provider-private-text", Reason: "provider-private-text"}, "method_error"}, + {fmt.Errorf("provider-private-text"), "collection_error"}, + } { + failure := systemTelemetryFailure(tc.err) + if failure.ErrorCategory != tc.want || failure.HTTPStatus != 0 || failure.CPU || failure.Memory { + t.Fatalf("failure=%+v, want %s", failure, tc.want) + } + body, err := json.Marshal(failure) + if err != nil || strings.Contains(string(body), "provider-private-text") { + t.Fatalf("unbounded failure detail: %s %v", body, err) + } + } +} diff --git a/internal/truenas/provider.go b/internal/truenas/provider.go index e81c5ef02..2e909fcaa 100644 --- a/internal/truenas/provider.go +++ b/internal/truenas/provider.go @@ -918,10 +918,10 @@ func trueNASReclaimableARC(system SystemInfo) int64 { // trueNASMemoryUsageKnown reports whether the appliance supplied an available // memory reading. Without it the ARC-aware used figure cannot be derived, and -// treating the zero value as a measurement reports total capacity as 100% used -// (#2077, legacy REST transport on TrueNAS CORE). +// an unobserved numeric zero must not report total capacity as 100% used. A +// genuinely observed zero remains valid (#2077, legacy REST on TrueNAS CORE). func trueNASMemoryUsageKnown(system SystemInfo) bool { - return system.MemoryAvailableBytes > 0 + return systemTelemetryAvailability(system).Memory && system.MemoryAvailableBytes >= 0 } // trueNASEffectiveMemoryUsed treats the ZFS ARC as reclaimable cache, the @@ -938,12 +938,12 @@ func trueNASEffectiveMemoryUsed(system SystemInfo) int64 { } func metricsFromTrueNASSystem(system SystemInfo, totalCapacity, totalUsed int64) *unifiedresources.ResourceMetrics { - hasRealtimeTelemetry := !system.CollectedAt.IsZero() || system.IntervalSeconds > 0 + availability := systemTelemetryAvailability(system) metrics := &unifiedresources.ResourceMetrics{ Disk: diskMetric(totalCapacity, totalUsed), } - if hasRealtimeTelemetry { + if availability.CPU { metrics.CPU = &unifiedresources.MetricValue{ Value: system.CPUPercent, Percent: system.CPUPercent, @@ -965,22 +965,28 @@ func metricsFromTrueNASSystem(system SystemInfo, totalCapacity, totalUsed int64) metrics.Memory = memory } - if hasRealtimeTelemetry { + if availability.NetIn { metrics.NetIn = &unifiedresources.MetricValue{ Value: system.NetInRate, Unit: "bytes/s", Source: unifiedresources.SourceTrueNAS, } + } + if availability.NetOut { metrics.NetOut = &unifiedresources.MetricValue{ Value: system.NetOutRate, Unit: "bytes/s", Source: unifiedresources.SourceTrueNAS, } + } + if availability.DiskRead { metrics.DiskRead = &unifiedresources.MetricValue{ Value: system.DiskReadRate, Unit: "bytes/s", Source: unifiedresources.SourceTrueNAS, } + } + if availability.DiskWrite { metrics.DiskWrite = &unifiedresources.MetricValue{ Value: system.DiskWriteRate, Unit: "bytes/s", @@ -3046,6 +3052,10 @@ func clonePools(pools []Pool) []Pool { func cloneSystemInfo(system SystemInfo) SystemInfo { cloned := system + if system.Telemetry != nil { + availability := *system.Telemetry + cloned.Telemetry = &availability + } if len(system.TemperatureCelsius) > 0 { cloned.TemperatureCelsius = make(map[string]float64, len(system.TemperatureCelsius)) for key, value := range system.TemperatureCelsius { diff --git a/internal/truenas/telemetry.go b/internal/truenas/telemetry.go new file mode 100644 index 000000000..7a594cd90 --- /dev/null +++ b/internal/truenas/telemetry.go @@ -0,0 +1,64 @@ +package truenas + +import ( + "context" + "encoding/json" + "errors" + "io" + "net/http" +) + +var errNoSystemTelemetry = errors.New("truenas legacy REST reporting returned no system telemetry") + +func systemTelemetryAvailability(system SystemInfo) SystemTelemetryAvailability { + if system.Telemetry != nil { + return *system.Telemetry + } + // Compatibility for static snapshots written before per-metric presence. + // GetSystemInfo and both live telemetry parsers always set Telemetry, so + // their collection timestamps never imply that missing samples were zero. + hasTelemetry := !system.CollectedAt.IsZero() || system.IntervalSeconds > 0 + return SystemTelemetryAvailability{ + CPU: hasTelemetry, Memory: system.MemoryAvailableBytes > 0, + NetIn: hasTelemetry, NetOut: hasTelemetry, + DiskRead: hasTelemetry, DiskWrite: hasTelemetry, + } +} + +func systemTelemetryFailure(err error) *SystemTelemetryAvailability { + availability := &SystemTelemetryAvailability{ErrorCategory: "collection_error"} + var apiErr *APIError + var authErr *RPCAuthError + var rpcErr *RPCError + var syntaxErr *json.SyntaxError + var typeErr *json.UnmarshalTypeError + switch { + case errors.Is(err, context.Canceled): + availability.ErrorCategory = "cancelled" + case errors.Is(err, context.DeadlineExceeded): + availability.ErrorCategory = "timeout" + case errors.Is(err, errNoSystemTelemetry): + availability.ErrorCategory = "no_samples" + case errors.As(err, &apiErr): + availability.HTTPStatus = apiErr.StatusCode + switch apiErr.StatusCode { + case http.StatusUnauthorized, http.StatusForbidden: + availability.ErrorCategory = "authentication" + case http.StatusNotFound, http.StatusMethodNotAllowed: + availability.ErrorCategory = "unsupported_endpoint" + case http.StatusBadRequest, http.StatusUnprocessableEntity: + availability.ErrorCategory = "request_rejected" + case http.StatusTooManyRequests: + availability.ErrorCategory = "rate_limited" + default: + availability.ErrorCategory = "http_error" + } + case errors.As(err, &authErr): + availability.ErrorCategory = "authentication" + case errors.As(err, &rpcErr): + availability.ErrorCategory = "method_error" + case errors.As(err, &syntaxErr), errors.As(err, &typeErr), errors.Is(err, io.ErrUnexpectedEOF): + availability.ErrorCategory = "invalid_response" + } + return availability +} diff --git a/internal/truenas/types.go b/internal/truenas/types.go index dfbf41158..26c6fb6fc 100644 --- a/internal/truenas/types.go +++ b/internal/truenas/types.go @@ -41,6 +41,24 @@ type SystemInfo struct { TemperatureCelsius map[string]float64 IntervalSeconds int CollectedAt time.Time + // Telemetry distinguishes an observed zero from an absent measurement. + // Live clients always populate it, including when collection fails. A nil + // value is retained for older static snapshots, not inferred by live reads. + Telemetry *SystemTelemetryAvailability +} + +// SystemTelemetryAvailability is a bounded, non-secret observation summary. +// Memory means an available/free reading was reported; capacity alone is not +// usage. Errors describe the collection result, never an appliance diagnosis. +type SystemTelemetryAvailability struct { + CPU bool `json:"cpu"` + Memory bool `json:"memory"` + NetIn bool `json:"netIn"` + NetOut bool `json:"netOut"` + DiskRead bool `json:"diskRead"` + DiskWrite bool `json:"diskWrite"` + ErrorCategory string `json:"errorCategory,omitempty"` + HTTPStatus int `json:"httpStatus,omitempty"` } // Pool mirrors the subset of TrueNAS pool fields needed for unified mapping.