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.