From b50313c59998d0e1506ab5ca42730bbb2e386dcd Mon Sep 17 00:00:00 2001 From: "pulse-triage[bot]" <249995291+pulse-triage[bot]@users.noreply.github.com> Date: Fri, 2 Oct 2026 00:27:34 +0100 Subject: [PATCH 1/5] Decode native CORE reporting rows across live metrics and History Use external RRD timing, native legends, scoped device graphs and measured CPU temperature for issue #2077. Keep absent, empty and zero buckets distinct, normalize CORE CPU states and preserve ARC alignment, safe transport failures and canonical projection. Add native-shape and end-to-end monitoring controls without changing release selection. Change-source: pulse-maintainer Contract-Neutral: Restores existing canonical metrics and temperature behaviour for CORE input formats; no public schema, route, resource identity or alert-policy delta. --- .../v6/internal/subsystems/monitoring.md | 61 ++- .../v6/internal/subsystems/registry.json | 2 + internal/monitoring/truenas_poller_test.go | 84 +++- internal/truenas/client.go | 178 +++++--- internal/truenas/client_test.go | 2 +- internal/truenas/core_reporting.go | 276 ++++++++++++ internal/truenas/core_reporting_test.go | 408 ++++++++++++++++++ internal/truenas/provider.go | 3 + .../truenas/reporting_aggregation_test.go | 3 + .../truenas/testdata/core13_reporting.json | 266 ++++++++++++ internal/truenas/types.go | 1 + 11 files changed, 1211 insertions(+), 73 deletions(-) create mode 100644 internal/truenas/core_reporting.go create mode 100644 internal/truenas/core_reporting_test.go create mode 100644 internal/truenas/testdata/core13_reporting.json diff --git a/docs/release-control/v6/internal/subsystems/monitoring.md b/docs/release-control/v6/internal/subsystems/monitoring.md index d0764f50b..b6d5ff16c 100644 --- a/docs/release-control/v6/internal/subsystems/monitoring.md +++ b/docs/release-control/v6/internal/subsystems/monitoring.md @@ -242,7 +242,8 @@ claim that #2077's missing graph journey is resolved. Live telemetry and history retain the legacy REST graph/query contract and the single successful-batch fast path. A batch rejected with HTTP 400, 422 or 500 -falls back to one request per graph (at most six calls including the batch). +falls back to one request per selected graph. The native catalogue extension +below bounds the expanded selection and its split loop. A failed optional graph must not discard successful CPU, memory or ARC graphs. Authentication, rate-limit, missing-endpoint, service-unavailable, decoding, transport and cancellation errors are not graph fallback triggers; these also @@ -275,6 +276,64 @@ timestamps/units, appliance acceptance or the reporter's exact failure. Native response evidence remains required before claiming that this repair resolves the reported telemetry journey. +### CORE native reporting rows and device graphs — issue #2077 + +`internal/truenas/core_reporting.go` owns the native reporting envelope alongside +`client.go`, not a second transport or metrics store. The reporter's five incoming +CORE 13.0-U6.1 results now establish value-only rows with one element per legend, +external `start`/`end`/`step`, FreeBSD memory-class labels, disk-octet labels, +interface RX/TX/overlap and per-core `cputemp` readings. They do not establish that +Pulse's REST bridge, installed browser or complete appliance journey works. + +Value-only rows derive time from `start + row-index * step`, with positive step, +valid ordered bounds, exact legend width and overflow-safe end checks. Explicit +row timestamps and supported object aliases retain compatibility. A metric is +never a timestamp, a null is never zero and a window mean is never a replacement +for a missing current row. History retains measured buckets; native current +readings accept only the last two steps before the requested end, preserving +sample time rather than stamping an old value as newly measured. The five-state +CORE CPU vector is normalized by its sum (RRD state rates need not sum to 100). +All-zero CPU or complete memory-class vectors are empty buckets, not 100% usage; +explicit zero usage, all-idle CPU, zero-free RAM with used pages and zero I/O +remain measured zeros. Free memory comes from `memory-free_value`, not active +pages or the sum of memory classes. ARC subtraction uses the matching free-RAM +bucket. Reporting I/O values remain rates in the existing bytes/s contract; +step supplies time, not a second rate division; `overlap` is not extra traffic. + +Legacy REST reads `/reporting/graphs` once per live/History collection and sends +its disk/interface identifiers, omitting identifiers on global CPU, memory, +ARC and CPU-temperature graphs. Only missing/method-unsupported catalogue +endpoints retain the older unscoped compatibility request. All other catalogue +failures remain failures; no authentication, TLS, rate-limit, transport, decode +or cancellation retry is added. Selection is deduplicated and capped at 256 +entries, refusing oversize rather than publishing a truncated total. Successful +reporting remains one batch; graph-only 400/422/500 failures may split once per +selected graph. Responses are bound to the selection: an omitted device stays +absent. Device rates sum once per identifier only at simultaneous observed +buckets of every selected member, never by carrying a stale member forward or +assuming a missing member was idle. Other graphs remain independently usable. + +Native CPU temperature maps into the existing canonical host temperature and +`temperature` History key, using the same package-versus-core outlier safeguard +as current readings. Disk-temperature History uses the same envelope decoder. +No frontend route, resource identity, alert suppression policy, credential scope +or recognised-legacy transport boundary changes. + +Verification: `internal/truenas/core_reporting_test.go` and its anonymised +`testdata/core13_reporting.json` cover five source-bound native excerpts, +row/timing/alias/zero/null/empty/finite controls, device totals, stale buckets, +ARC alignment, current and canonical History projections, request validation, +catalogue limits and failure boundaries. The excerpt windows are deliberately +shortened controls, not complete native windows. Extended +`TestTrueNASPartialReportingPipeline` in +`internal/monitoring/truenas_poller_test.go` crosses pinned-certificate HTTP, +poller, canonical registry, writer, in-memory/persistent chart readbacks and +native fallback, including temperature. Existing CORE isolation, recognised +transport/auth/TLS and modern SCALE aggregation tests remain required. Source +proofs do not substitute for containing-release qualification or installed CORE +live/collapsed/expanded/History acceptance. + + **Availability backfill preserves concurrent discovery changes (7 September 2026)** diff --git a/docs/release-control/v6/internal/subsystems/registry.json b/docs/release-control/v6/internal/subsystems/registry.json index df14c951a..28e7733d8 100644 --- a/docs/release-control/v6/internal/subsystems/registry.json +++ b/docs/release-control/v6/internal/subsystems/registry.json @@ -6157,6 +6157,7 @@ "match_prefixes": [], "match_files": [ "internal/truenas/client.go", + "internal/truenas/core_reporting.go", "internal/truenas/disk_health.go", "internal/truenas/fixtures.go", "internal/truenas/provider.go", @@ -6173,6 +6174,7 @@ "internal/truenas/client_api_shapes_test.go", "internal/truenas/client_test.go", "internal/truenas/contract_test.go", + "internal/truenas/core_reporting_test.go", "internal/truenas/provider_pool_health_contract_test.go", "internal/truenas/provider_test.go", "internal/truenas/transport_test.go" diff --git a/internal/monitoring/truenas_poller_test.go b/internal/monitoring/truenas_poller_test.go index e0af04581..bf8fe12a8 100644 --- a/internal/monitoring/truenas_poller_test.go +++ b/internal/monitoring/truenas_poller_test.go @@ -2221,13 +2221,23 @@ func TestTrueNASPartialReportingPipeline(t *testing.T) { status int want map[string]float64 errorCategory string + native bool }{ - {`[{"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}, ""}, + {`[{"name":"memory","legend":["free"],"data":[[1789000060,8]]}]`, 200, map[string]float64{"memory": 50}, "", false}, + {`[{"name":"cpu","legend":["usage"],"data":[[1789000060,0]]}]`, 200, map[string]float64{"cpu": 0}, "", false}, + {`provider-private-text`, 401, nil, "authentication", false}, + {`[{"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}, "", false}, + {`[ + {"name":"cpu","legend":["interrupt","system","user","nice","idle"],"data":[[1,2,3,4,90]]}, + {"name":"memory","legend":["memory-active_value","memory-inactive_value","memory-wired_value","memory-laundry_value","memory-free_value"],"data":[[1,1,10,0,4]]}, + {"name":"arcsize","legend":["arcsize_value"],"data":[[2]]}, + {"name":"cputemp","legend":["cputemp0","cputemp1"],"data":[[41,42]]}, + {"name":"interface","identifier":"nic-a","legend":["rx","tx","overlap"],"data":[[8,16,8]]}, + {"name":"disk","identifier":"disk-a","legend":["disk_octets_read","disk_octets_write"],"data":[[0,32]]} + ]`, 200, map[string]float64{"cpu": 10, "memory": 62.5, "netin": 8, "netout": 16, "diskread": 0, "diskwrite": 32}, "", true}, } var current atomic.Int64 + var nativeEnd atomic.Int64 server := httptest.NewTLSServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { w.Header().Set("Content-Type", "application/json") switch r.URL.Path { @@ -2237,10 +2247,48 @@ func TestTrueNASPartialReportingPipeline(t *testing.T) { _, _ = 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/graphs": + if cycles[current.Load()].native { + _, _ = w.Write([]byte(`[{"name":"disk","identifiers":["disk-a"]},{"name":"interface","identifiers":["nic-a"]}]`)) + } else { + http.NotFound(w, r) + } case "/api/v2.0/reporting/get_data": cycle := cycles[current.Load()] w.WriteHeader(cycle.status) - _, _ = w.Write([]byte(cycle.body)) + if !cycle.native { + _, _ = w.Write([]byte(cycle.body)) + break + } + var query struct { + Query map[string]any `json:"query"` + Graphs []map[string]any `json:"graphs"` + } + if err := json.NewDecoder(r.Body).Decode(&query); err != nil { + t.Error(err) + return + } + var rows []map[string]any + if err := json.Unmarshal([]byte(cycle.body), &rows); err != nil { + t.Error(err) + return + } + end := int64(query.Query["end"].(float64)) + nativeEnd.Store(end) + for _, row := range rows { + row["start"] = end + row["end"] = end + row["step"] = 10 + } + for _, graph := range query.Graphs { + if name := graph["name"]; (name == "interface" || name == "disk") && graph["identifier"] == nil { + t.Error("native device requested without identifier") + return + } + } + if err := json.NewEncoder(w).Encode(rows); err != nil { + t.Error(err) + } default: http.NotFound(w, r) } @@ -2318,6 +2366,9 @@ func TestTrueNASPartialReportingPipeline(t *testing.T) { if host.Agent.Memory.UsageUnavailable != !availability.Memory { t.Fatalf("cycle %d stale/fabricated memory metadata: %+v", index, host.Agent.Memory) } + if cycle.native && (host.Temperature == nil || *host.Temperature != 42) { + t.Fatalf("native collapsed/expanded temperature missing: %+v", host) + } var writes []metrics.WriteMetric monitor.syncUnifiedAgentMetrics(resourceStore, &writes) seen := make(map[string]bool) @@ -2350,14 +2401,33 @@ func TestTrueNASPartialReportingPipeline(t *testing.T) { } continue } - if err != nil || id != instance.ID || len(native) != len(cycle.want) { + wantNative := len(cycle.want) + if cycle.native { + wantNative++ + } + if err != nil || id != instance.ID || len(native) != wantNative { 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 { + if len(points) != 1 || points[0].Value != want || points[0].Timestamp.Unix() != func() int64 { + if cycle.native { + return nativeEnd.Load() + } + return 1789000060 + }() { t.Errorf("cycle %d %s native data/timestamp altered: %+v", index, key, points) } } + if cycle.native { + points := native["temperature"] + if len(points) != 1 || points[0].Value != 42 || points[0].Timestamp.Unix() != nativeEnd.Load() { + t.Fatalf("native temperature History missing: %+v", points) + } + shared := poller.GuestMetricHistory(nil, "default", "agent", time.Hour)[instance.ID] + if len(shared) != 7 || len(shared["temperature"]) != 1 { + t.Fatalf("shared History fallback dropped native panels: %+v", shared) + } + } } } diff --git a/internal/truenas/client.go b/internal/truenas/client.go index 81671d09b..769c78e74 100644 --- a/internal/truenas/client.go +++ b/internal/truenas/client.go @@ -288,7 +288,7 @@ func (c *Client) getSystemTelemetryREST(ctx context.Context) (*SystemInfo, error if err != nil { return nil, err } - history := parseSystemMetricHistory(response) + history := parseSystemMetricHistory(liveReportingResponses(response, end)) if history == nil { return nil, errNoSystemTelemetry } @@ -302,7 +302,27 @@ func systemInfoFromMetricHistory(history *SystemMetricHistory) *SystemInfo { if history == nil { return nil } - system := &SystemInfo{CollectedAt: time.Now().UTC(), Telemetry: &SystemTelemetryAvailability{}} + system := &SystemInfo{Telemetry: &SystemTelemetryAvailability{}} + for _, series := range [][]TimeSeriesPoint{history.CPUPercent, history.MemoryAvailableBytes, history.NetInRate, history.NetOutRate, history.DiskReadRate, history.DiskWriteRate} { + for _, point := range series { + if point.Timestamp.After(system.CollectedAt) { + system.CollectedAt = point.Timestamp + } + } + } + for key, series := range history.TemperatureCelsius { + if value, ok := latestTimeSeriesValue(series); ok { + if system.TemperatureCelsius == nil { + system.TemperatureCelsius = make(map[string]float64) + } + system.TemperatureCelsius[key] = value + } + for _, point := range series { + if point.Timestamp.After(system.CollectedAt) { + system.CollectedAt = point.Timestamp + } + } + } if value, ok := latestTimeSeriesValue(history.CPUPercent); ok { system.CPUPercent = value system.Telemetry.CPU = true @@ -314,9 +334,20 @@ func systemInfoFromMetricHistory(history *SystemMetricHistory) *SystemInfo { system.MemoryAvailableBytes = int64(value) system.Telemetry.Memory = value >= 0 } - if value, ok := latestTimeSeriesValue(history.ARCSizeBytes); ok { - system.ARCSizeBytes = int64(value) + // Cache and free RAM must describe the same bucket. An old ARC sample + // must not make a newer full-RAM reading look healthy. + var latestFree time.Time + for _, free := range history.MemoryAvailableBytes { + if free.Timestamp.After(latestFree) { + latestFree = free.Timestamp + } } + for _, arc := range history.ARCSizeBytes { + if arc.Timestamp.Equal(latestFree) { + system.ARCSizeBytes = int64(arc.Value) + } + } + if value, ok := latestTimeSeriesValue(history.NetInRate); ok { system.NetInRate = value system.Telemetry.NetIn = true @@ -388,6 +419,7 @@ func legacyRESTReportingGraphs() []map[string]any { reportingGraph("arcsize", ""), reportingGraph("interface", ""), reportingGraph("disk", ""), + reportingGraph("cputemp", ""), } } @@ -396,10 +428,16 @@ func legacyRESTReportingGraphs() []map[string]any { // graph-validation/server errors, never authentication, rate-limit, endpoint, // transport or cancellation failures. The query and graph set remain unchanged. func (c *Client) getLegacySystemReportingData(ctx context.Context, query map[string]any) ([]trueNASReportingGetDataResponse, error) { - graphs := legacyRESTReportingGraphs() + graphs, err := c.legacyReportingGraphs(ctx) + if err != nil { + return nil, err + } response, err := c.getReportingDataREST(ctx, graphs, query) - if err == nil || !isReportingGraphFailure(err) { - return response, err + if err == nil { + return bindLegacyReportingResponses(graphs, response), nil + } + if !isReportingGraphFailure(err) { + return nil, err } batchErr := err var collected []trueNASReportingGetDataResponse @@ -419,7 +457,7 @@ func (c *Client) getLegacySystemReportingData(ctx context.Context, query map[str if len(collected) == 0 { return nil, batchErr } - return collected, nil + return bindLegacyReportingResponses(graphs, collected), nil } func isReportingGraphFailure(err error) bool { @@ -2378,6 +2416,7 @@ type trueNASReportingGetDataResponse struct { Aggregations trueNASReportingAggregations `json:"aggregations"` Start int64 `json:"start"` End int64 `json:"end"` + Step int64 `json:"step"` Legend []string `json:"legend"` } @@ -3178,34 +3217,25 @@ func parseSystemTelemetry(fields map[string]any, intervalSeconds int, collectedA } func parseSystemTemperatures(responses []trueNASReportingGetDataResponse) map[string]float64 { - if len(responses) == 0 { - return nil - } - temperatures := make(map[string]float64) for _, response := range responses { - if strings.TrimSpace(strings.ToLower(response.Name)) != "cputemp" || len(response.Legend) == 0 { + if !strings.EqualFold(strings.TrimSpace(response.Name), "cputemp") { continue } - - values := extractReportingLegendFloatValues(response.Aggregations.Mean, response.Legend) - if len(values) == 0 && len(response.Data) > 0 { - values = extractReportingLegendFloatValues(response.Data[len(response.Data)-1], response.Legend) + values := latestReportingValues(response) + // Retain aggregate-only compatibility, never use a window mean to fill a + // missing/null current row or an invalid native timing envelope. + if len(response.Data) == 0 && response.Step == 0 { + values = finiteReportingValues(extractReportingLegendFloatValues(response.Aggregations.Mean, response.Legend)) } - if len(values) == 0 { - continue - } - for index, legend := range response.Legend { value, ok := values[legend] if !ok || value <= 0 { continue } - key := canonicalSystemTemperatureKey(legend, index, len(response.Legend)) - if key == "" { - continue + if key := canonicalSystemTemperatureKey(legend, index, len(response.Legend)); key != "" { + temperatures[key] = value } - temperatures[key] = value } } if len(temperatures) == 0 { @@ -3218,75 +3248,93 @@ func parseSystemMetricHistory(responses []trueNASReportingGetDataResponse) *Syst if len(responses) == 0 { return nil } - - history := &SystemMetricHistory{} + history := &SystemMetricHistory{TemperatureCelsius: make(map[string][]TimeSeriesPoint)} + io := map[string]map[string][]TimeSeriesPoint{ + "netin": {}, "netout": {}, "diskread": {}, "diskwrite": {}, + } + seen := make(map[string]bool) for _, response := range responses { - name := strings.TrimSpace(strings.ToLower(response.Name)) + name := strings.ToLower(strings.TrimSpace(response.Name)) switch name { - case "cpu", "memory", "arcsize", "interface", "disk": + case "cpu", "memory", "arcsize", "cputemp", "interface", "disk": default: continue } - - for _, raw := range response.Data { - timestamp, values, ok := parseReportingSeriesValues(raw, response.Legend) - if !ok || timestamp.IsZero() || len(values) == 0 { + identifier := readStringAny(map[string]any{"identifier": response.Identifier}, "identifier") + key := name + "\x00" + identifier + if seen[key] { + continue + } + seen[key] = true + if name == "interface" { + io["netin"][key] = nil + io["netout"][key] = nil + } + if name == "disk" { + io["diskread"][key] = nil + io["diskwrite"][key] = nil + } + for index := range response.Data { + timestamp, values, ok := reportingRow(response, index) + if !ok { continue } - switch name { case "cpu": - if value, ok := parseSystemCPUPercent(values); ok { + if value, ok := reportingCPUPercent(values, response.Legend); ok { history.CPUPercent = appendTimeSeriesPoint(history.CPUPercent, timestamp, value) } case "memory": - // FreeBSD active pages are only one used-memory component, not - // total usage. Without an explicit used series, let the provider - // derive usage from system capacity, free memory and ARC. + if !reportingMemoryRowValid(values, response.Legend) { + continue + } if value, ok := parseSystemMemoryPercent(values); ok { history.MemoryPercent = appendTimeSeriesPoint(history.MemoryPercent, timestamp, value) } - if value, ok := pickReportingValue(values, "used", "memory_used", "used_bytes", "memory"); ok { + if value, ok := pickReportingValue(values, "used", "memory_used", "used_bytes", "memory"); ok && value >= 0 { history.MemoryUsedBytes = appendTimeSeriesPoint(history.MemoryUsedBytes, timestamp, value) } - if value, ok := pickReportingValue(values, "available", "free", "available_bytes", "free_bytes"); ok { + if value, ok := pickReportingValue(values, "available", "free", "available_bytes", "free_bytes", "memory-free_value"); ok && value >= 0 { history.MemoryAvailableBytes = appendTimeSeriesPoint(history.MemoryAvailableBytes, timestamp, value) } - if value, ok := pickReportingValue(values, "total", "memory_total", "total_bytes", "physical_memory_total"); ok { + if value, ok := pickReportingValue(values, "total", "memory_total", "total_bytes", "physical_memory_total"); ok && value > 0 { history.MemoryTotalBytes = appendTimeSeriesPoint(history.MemoryTotalBytes, timestamp, value) } case "arcsize": - if value, ok := pickReportingValue(values, "arc_size", "size", "arcsize"); ok { + if value, ok := pickReportingValue(values, "arc_size", "size", "arcsize", "arcsize_value"); ok && value >= 0 { history.ARCSizeBytes = appendTimeSeriesPoint(history.ARCSizeBytes, timestamp, value) } - case "interface": - if value, ok := pickReportingValue(values, "received", "received_bytes", "rx", "rx_bytes", "netin", "in"); ok { - history.NetInRate = appendTimeSeriesPoint(history.NetInRate, timestamp, value) + case "cputemp": + for i, legend := range response.Legend { + if value, ok := values[legend]; ok && value > 0 { + key := canonicalSystemTemperatureKey(legend, i, len(response.Legend)) + if key != "" { + history.TemperatureCelsius[key] = appendTimeSeriesPoint(history.TemperatureCelsius[key], timestamp, value) + } + } } - if value, ok := pickReportingValue(values, "sent", "sent_bytes", "tx", "tx_bytes", "netout", "out"); ok { - history.NetOutRate = appendTimeSeriesPoint(history.NetOutRate, timestamp, value) + case "interface": + if value, ok := pickReportingValue(values, "received", "received_bytes", "rx", "rx_bytes", "netin", "in"); ok && value >= 0 { + io["netin"][key] = appendTimeSeriesPoint(io["netin"][key], timestamp, value) + } + if value, ok := pickReportingValue(values, "sent", "sent_bytes", "tx", "tx_bytes", "netout", "out"); ok && value >= 0 { + io["netout"][key] = appendTimeSeriesPoint(io["netout"][key], timestamp, value) } case "disk": - if value, ok := pickReportingValue(values, "read", "read_bytes", "diskread", "bytes_read"); ok { - history.DiskReadRate = appendTimeSeriesPoint(history.DiskReadRate, timestamp, value) + if value, ok := pickReportingValue(values, "read", "read_bytes", "diskread", "bytes_read", "disk_octets_read"); ok && value >= 0 { + io["diskread"][key] = appendTimeSeriesPoint(io["diskread"][key], timestamp, value) } - if value, ok := pickReportingValue(values, "write", "write_bytes", "diskwrite", "bytes_written"); ok { - history.DiskWriteRate = appendTimeSeriesPoint(history.DiskWriteRate, timestamp, value) + if value, ok := pickReportingValue(values, "write", "write_bytes", "diskwrite", "bytes_written", "disk_octets_write"); ok && value >= 0 { + io["diskwrite"][key] = appendTimeSeriesPoint(io["diskwrite"][key], timestamp, value) } } } } - - if len(history.CPUPercent) == 0 && - len(history.MemoryPercent) == 0 && - len(history.MemoryUsedBytes) == 0 && - len(history.MemoryAvailableBytes) == 0 && - len(history.MemoryTotalBytes) == 0 && - len(history.ARCSizeBytes) == 0 && - len(history.NetInRate) == 0 && - len(history.NetOutRate) == 0 && - len(history.DiskReadRate) == 0 && - len(history.DiskWriteRate) == 0 { + history.NetInRate = sumReportingDevices(io["netin"]) + history.NetOutRate = sumReportingDevices(io["netout"]) + history.DiskReadRate = sumReportingDevices(io["diskread"]) + history.DiskWriteRate = sumReportingDevices(io["diskwrite"]) + if len(history.CPUPercent)+len(history.MemoryPercent)+len(history.MemoryUsedBytes)+len(history.MemoryAvailableBytes)+len(history.MemoryTotalBytes)+len(history.ARCSizeBytes)+len(history.NetInRate)+len(history.NetOutRate)+len(history.DiskReadRate)+len(history.DiskWriteRate)+len(history.TemperatureCelsius) == 0 { return nil } return history @@ -3337,8 +3385,10 @@ func parseReportingDiskTemperatureHistory(responses []trueNASReportingGetDataRes } points := make([]TimeSeriesPoint, 0, len(response.Data)) - for _, raw := range response.Data { - timestamp, value, ok := parseReportingSeriesPoint(raw, response.Legend) + for index := range response.Data { + timestamp, values, ok := reportingRow(response, index) + value, hasValue := pickReportingValue(values, response.Legend...) + ok = ok && hasValue && value > 0 if !ok || timestamp.IsZero() { continue } diff --git a/internal/truenas/client_test.go b/internal/truenas/client_test.go index f5d338e39..ab68c5d3d 100644 --- a/internal/truenas/client_test.go +++ b/internal/truenas/client_test.go @@ -2222,7 +2222,7 @@ func TestRESTReportingGraphFailureBoundaries(t *testing.T) { } want := 1 if status == 400 || status == 422 || status == 500 { - want = 6 + want = 7 } if transport.calls != want { t.Fatalf("calls = %d, want %d", transport.calls, want) diff --git a/internal/truenas/core_reporting.go b/internal/truenas/core_reporting.go new file mode 100644 index 000000000..ac3eb46cc --- /dev/null +++ b/internal/truenas/core_reporting.go @@ -0,0 +1,276 @@ +package truenas + +import ( + "context" + "errors" + "fmt" + "math" + "net/http" + "sort" + "strings" + "time" +) + +// CORE's RRD graphs are device-scoped. Ask the reporting catalogue for its +// identifiers, not disk serials or configured (possibly renamed) interfaces. +// Older REST bridges without the catalogue retain the unscoped compatibility +// request. Authentication, transport and malformed replies are never retried. +func (c *Client) legacyReportingGraphs(ctx context.Context) ([]map[string]any, error) { + var catalogue []struct { + Name string `json:"name"` + Identifiers []string `json:"identifiers"` + } + if err := c.getJSON(ctx, http.MethodGet, "/reporting/graphs", &catalogue); err != nil { + if isAPIStatus(err, http.StatusNotFound) || isAPIStatus(err, http.StatusMethodNotAllowed) { + return legacyRESTReportingGraphs(), nil + } + return nil, err + } + graphs := []map[string]any{ + reportingGraph("cpu", ""), reportingGraph("memory", ""), + reportingGraph("arcsize", ""), reportingGraph("cputemp", ""), + } + seen := make(map[string]bool) + for _, graph := range catalogue { + name := strings.ToLower(strings.TrimSpace(graph.Name)) + if name != "disk" && name != "interface" { + continue + } + for _, identifier := range dedupeStrings(graph.Identifiers) { + identifier = strings.TrimSpace(identifier) + key := name + "\x00" + identifier + if identifier == "" || seen[key] { + continue + } + seen[key] = true + // Refuse an oversized catalogue, rather than silently publishing a + // truncated host total or allowing an unbounded split retry loop. + if len(graphs) >= 256 { + return nil, fmt.Errorf("truenas reporting catalogue exceeds 256 graphs") + } + graphs = append(graphs, reportingGraph(name, identifier)) + } + } + return graphs, nil +} + +func isAPIStatus(err error, status int) bool { + var apiErr *APIError + return errors.As(err, &apiErr) && apiErr.StatusCode == status +} + +func finiteReportingValues(values map[string]float64) map[string]float64 { + for key, value := range values { + if math.IsNaN(value) || math.IsInf(value, 0) { + delete(values, key) + } + } + return values +} + +// A native row has exactly one value per legend; its timestamp is external. +// An explicit-timestamp row has one extra element. Never guess an epoch from a +// value-only row with missing/invalid timing or accept a misaligned short row. +func reportingRow(response trueNASReportingGetDataResponse, index int) (time.Time, map[string]float64, bool) { + raw := response.Data[index] + if row, ok := raw.([]any); ok { + if len(response.Legend) == 0 { + return time.Time{}, nil, false + } + if len(row) == len(response.Legend) { + if response.Start <= 0 || response.End < response.Start || response.Step <= 0 || + int64(index) > (response.End-response.Start)/response.Step { + return time.Time{}, nil, false + } + timestamp := time.Unix(response.Start+int64(index)*response.Step, 0).UTC() + values := finiteReportingValues(extractReportingLegendFloatValues(row, response.Legend)) + return timestamp, values, len(values) > 0 + } + if len(row) != len(response.Legend)+1 { + return time.Time{}, nil, false + } + } + timestamp, values, ok := parseReportingSeriesValues(raw, response.Legend) + values = finiteReportingValues(values) + return timestamp, values, ok && timestamp.Unix() > 0 && len(values) > 0 +} + +// Keep native History intact, but do not label an old RRD point as a current +// reading. Two returned steps allow the usual unfinished/null RRD bucket; +// older points stay available only through History. Preserve row indices. +func liveReportingResponses(responses []trueNASReportingGetDataResponse, end int64) []trueNASReportingGetDataResponse { + live := append([]trueNASReportingGetDataResponse(nil), responses...) + for i, response := range live { + if response.Step <= 0 { + continue // timestamped legacy/modern compatibility payload + } + live[i].Data = append([]any(nil), response.Data...) + for index := range response.Data { + timestamp, _, ok := reportingRow(response, index) + gap := end - timestamp.Unix() + if !ok || timestamp.Unix() > end || (gap > response.Step && gap-response.Step > response.Step) { + live[i].Data[index] = nil + } + } + } + return live +} + +func latestReportingValues(response trueNASReportingGetDataResponse) map[string]float64 { + var latest time.Time + var values map[string]float64 + for index := range response.Data { + timestamp, row, ok := reportingRow(response, index) + if ok && (values == nil || timestamp.After(latest)) { + latest, values = timestamp, row + } + } + return values +} + +var coreCPUStates = []string{"interrupt", "system", "user", "nice", "idle"} +var coreMemoryClasses = []string{"memory-active_value", "memory-inactive_value", "memory-wired_value", "memory-laundry_value", "memory-free_value"} + +func hasReportingLegends(legends, required []string) bool { + present := make(map[string]bool, len(legends)) + for _, legend := range legends { + present[normalizeTemperatureLegendLabel(legend)] = true + } + for _, legend := range required { + if !present[normalizeTemperatureLegendLabel(legend)] { + return false + } + } + return true +} + +// RRD CPU state rates are not necessarily percentages (the supplied states sum +// above 100). Normalize the complete state vector; an all-zero bucket is not +// 100% busy. Explicit usage=0 and a measured all-idle vector remain real zeros. +func reportingCPUPercent(values map[string]float64, legends []string) (float64, bool) { + if !hasReportingLegends(legends, coreCPUStates) { + return parseSystemCPUPercent(values) + } + var total, idle float64 + for _, state := range coreCPUStates { + value, ok := pickReportingValue(values, state) + if !ok || value < 0 { + return 0, false + } + total += value + if state == "idle" { + idle = value + } + } + if total <= 0 || math.IsInf(total, 0) { + return 0, false + } + return (total - idle) / total * 100, true +} + +// Five zero memory classes cannot describe physical RAM. Keep this native RRD +// empty-bucket sentinel distinct from a real zero-free reading with used pages. +func reportingMemoryRowValid(values map[string]float64, legends []string) bool { + if !hasReportingLegends(legends, coreMemoryClasses) { + return true + } + var total float64 + for _, class := range coreMemoryClasses { + value, ok := pickReportingValue(values, class) + if !ok || value < 0 { + return false + } + total += value + } + return total > 0 && !math.IsInf(total, 0) +} + +// Sum only simultaneous observations of every returned device. Null/missing +// members are not zero, duplicate device replies are not additional traffic, +// and values are already rates: step is for time, not another division. +func sumReportingDevices(devices map[string][]TimeSeriesPoint) []TimeSeriesPoint { + if len(devices) == 0 { + return nil + } + var totals map[int64]float64 + for _, series := range devices { + points := make(map[int64]float64, len(series)) + for _, point := range series { + points[point.Timestamp.Unix()] = point.Value + } + if totals == nil { + totals = points + continue + } + for timestamp, total := range totals { + value, ok := points[timestamp] + if !ok { + delete(totals, timestamp) + } else { + totals[timestamp] = total + value + } + } + } + timestamps := make([]int64, 0, len(totals)) + for timestamp := range totals { + timestamps = append(timestamps, timestamp) + } + sort.Slice(timestamps, func(i, j int) bool { return timestamps[i] < timestamps[j] }) + points := make([]TimeSeriesPoint, 0, len(timestamps)) + for _, timestamp := range timestamps { + if !math.IsInf(totals[timestamp], 0) { + points = appendTimeSeriesPoint(points, time.Unix(timestamp, 0).UTC(), totals[timestamp]) + } + } + return points +} + +// Bind device responses to the requested catalogue. A failed/omitted device +// graph remains an absent member, not a smaller apparently complete host sum. +func bindLegacyReportingResponses(graphs []map[string]any, responses []trueNASReportingGetDataResponse) []trueNASReportingGetDataResponse { + byKey := make(map[string]trueNASReportingGetDataResponse) + for _, response := range responses { + name := strings.ToLower(strings.TrimSpace(response.Name)) + identifier := readStringAny(map[string]any{"identifier": response.Identifier}, "identifier") + key := name + "\x00" + identifier + if _, exists := byKey[key]; !exists { + byKey[key] = response + } + } + bound := make([]trueNASReportingGetDataResponse, 0, len(graphs)) + for _, graph := range graphs { + name := readStringAny(graph, "name") + identifier := readStringAny(graph, "identifier") + response, ok := byKey[name+"\x00"+identifier] + if !ok { + response = trueNASReportingGetDataResponse{Name: name, Identifier: identifier} + } + bound = append(bound, response) + } + return bound +} + +func systemTemperatureHistory(sensors map[string][]TimeSeriesPoint) []TimeSeriesPoint { + at := make(map[int64]map[string]float64) + for key, series := range sensors { + for _, point := range series { + timestamp := point.Timestamp.Unix() + if at[timestamp] == nil { + at[timestamp] = make(map[string]float64) + } + at[timestamp][key] = point.Value + } + } + timestamps := make([]int64, 0, len(at)) + for timestamp := range at { + timestamps = append(timestamps, timestamp) + } + sort.Slice(timestamps, func(i, j int) bool { return timestamps[i] < timestamps[j] }) + var points []TimeSeriesPoint + for _, timestamp := range timestamps { + if value := maxTrueNASSystemTemperature(SystemInfo{TemperatureCelsius: at[timestamp]}); value != nil { + points = appendTimeSeriesPoint(points, time.Unix(timestamp, 0).UTC(), *value) + } + } + return points +} diff --git a/internal/truenas/core_reporting_test.go b/internal/truenas/core_reporting_test.go new file mode 100644 index 000000000..07ee611c7 --- /dev/null +++ b/internal/truenas/core_reporting_test.go @@ -0,0 +1,408 @@ +package truenas + +import ( + "context" + "encoding/json" + "errors" + "math" + "net/http" + "os" + "reflect" + "strings" + "testing" + "time" + + "github.com/rcourtman/pulse-go-rewrite/internal/unifiedresources" +) + +// Five native result excerpts supplied in #2077 on 1 October 2026. Identifiers +// are anonymised and the SHORT EXCERPT windows made coherent. This is decoder +// evidence, not a complete time window, REST capture or appliance acceptance. +func core13ReportingFixture(t *testing.T) []trueNASReportingGetDataResponse { + t.Helper() + raw, err := os.ReadFile("testdata/core13_reporting.json") + if err != nil { + t.Fatal(err) + } + var replies []trueNASReportingGetDataResponse + if err := json.Unmarshal(raw, &replies); err != nil { + t.Fatal(err) + } + return replies +} + +func assertCOREPoint(t *testing.T, points []TimeSeriesPoint, index int, timestamp int64, value float64) { + t.Helper() + if len(points) <= index || points[index].Timestamp.Unix() != timestamp || math.Abs(points[index].Value-value) > 0.000001 { + t.Fatalf("point %d: %+v, want %d / %v", index, points, timestamp, value) + } +} + +// These five subtests also compile with the predecessor decoder and are the +// predetermined adverse controls; do not substitute compiler failure for them. +func TestCORE13NativeResultControls(t *testing.T) { + replies := core13ReportingFixture(t) + start := int64(1789000000) + t.Run("CPU", func(t *testing.T) { + history := parseSystemMetricHistory(replies[:1]) + if history == nil { + t.Fatal("native CPU rows produced no history") + } + want := (0.033358538468 + 4.4805217137 + 7.8754519386 + 7.8754519386) / (100 + 0.033358538468 + 4.4805217137 + 7.8754519386 + 7.8754519386) * 100 + assertCOREPoint(t, history.CPUPercent, 0, start, want) + if len(history.CPUPercent) != 1 { + t.Fatal("all-zero state buckets became CPU measurements") + } + }) + t.Run("memory", func(t *testing.T) { + history := parseSystemMetricHistory(replies[2:3]) + if history == nil { + t.Fatal("native memory aliases produced no history") + } + assertCOREPoint(t, history.MemoryAvailableBytes, 0, start, 2178945024) + if len(history.MemoryAvailableBytes) != 1 || len(history.MemoryUsedBytes) != 0 { + t.Fatal("empty bucket or active pages became total usage") + } + }) + t.Run("disk", func(t *testing.T) { + history := parseSystemMetricHistory(replies[3:4]) + if history == nil { + t.Fatal("native disk rows produced no history") + } + assertCOREPoint(t, history.DiskReadRate, 0, start, 484.54070578) + assertCOREPoint(t, history.DiskReadRate, 1, start+10, 0) + assertCOREPoint(t, history.DiskWriteRate, 1, start+10, 1276361.9904) + if len(history.DiskReadRate) != 2 { + t.Fatal("null disk bucket became zero") + } + }) + t.Run("network", func(t *testing.T) { + history := parseSystemMetricHistory(replies[4:5]) + if history == nil { + t.Fatal("native network rows produced no history") + } + assertCOREPoint(t, history.NetInRate, 0, start, 9728.2270494) + assertCOREPoint(t, history.NetOutRate, 0, start, 45551.708899) + if len(history.NetInRate) != 1 { + t.Fatal("null or overlap became RX") + } + }) + t.Run("temperature request", func(t *testing.T) { + found := false + for _, graph := range legacyRESTReportingGraphs() { + found = found || graph["name"] == "cputemp" + } + if !found { + t.Fatal("legacy graph set never requests CPU temperature") + } + }) +} + +func TestCORE13TemperatureUsesSamplesNotMeans(t *testing.T) { + replies := core13ReportingFixture(t)[1:2] + replies[0].Aggregations.Mean = []any{999., 999., 999., 999., 999., 999., 999., 999.} + temperatures := parseSystemTemperatures(replies) + if len(temperatures) != 8 || temperatures["cputemp2"] != 30.95 { + t.Fatalf("temperature samples lost: %+v", temperatures) + } + history := parseSystemMetricHistory(replies) + if history == nil { + t.Fatal("temperature-only history lost") + } + system := systemInfoFromMetricHistory(history) + if got := maxTrueNASSystemTemperature(*system); got == nil || *got != 30.95 { + t.Fatalf("current temperature = %v", got) + } + points := systemTemperatureHistory(history.TemperatureCelsius) + assertCOREPoint(t, points, 0, 1789000000, 30.95) + replies[0].Data[0] = []any{nil, nil, nil, nil, nil, nil, nil, nil} + if got := parseSystemTemperatures(replies); len(got) != 0 { + t.Fatalf("all-null rows substituted a mean: %v", got) + } +} + +func TestCORE13TimingAndMissingValueBoundaries(t *testing.T) { + for _, tc := range []struct { + name, body string + count int + last int64 + }{ + {"external zero", `{"start":1789000000,"end":1789000020,"step":10,"legend":["usage"],"data":[[0],[null],[12]]}`, 2, 1789000020}, + {"missing step", `{"start":1789000000,"end":1789000020,"legend":["usage"],"data":[[42]]}`, 0, 0}, + {"zero step", `{"start":1789000000,"end":1789000020,"step":0,"legend":["usage"],"data":[[42]]}`, 0, 0}, + {"negative step", `{"start":1789000000,"end":1789000020,"step":-10,"legend":["usage"],"data":[[42]]}`, 0, 0}, + {"missing start", `{"end":1789000020,"step":10,"legend":["usage"],"data":[[42]]}`, 0, 0}, + {"reversed range", `{"start":1789000020,"end":1789000000,"step":10,"legend":["usage"],"data":[[42]]}`, 0, 0}, + {"past end", `{"start":1789000000,"end":1789000000,"step":10,"legend":["usage"],"data":[[42],[43]]}`, 1, 1789000000}, + {"overflow", `{"start":9223372036854775797,"end":9223372036854775807,"step":9223372036854775807,"legend":["usage"],"data":[[42],[43]]}`, 1, 9223372036854775797}, + {"timestamped", `{"legend":["usage"],"data":[[1789000000,0],[1789000010000,12]]}`, 2, 1789000010}, + {"object", `{"legend":["usage"],"data":[{"timestamp":"2026-09-10T01:46:40Z","usage":0}]}`, 1, 1789004800}, + {"short row", `{"legend":["rx","tx","overlap"],"data":[[1789000000,12]]}`, 0, 0}, + {"nonfinite", `{"start":1789000000,"end":1789000000,"step":10,"legend":["usage"],"data":[["NaN"],["+Inf"]]}`, 0, 0}, + } { + t.Run(tc.name, func(t *testing.T) { + var response trueNASReportingGetDataResponse + if err := json.Unmarshal([]byte(tc.body), &response); err != nil { + t.Fatal(err) + } + response.Name = "cpu" + history := parseSystemMetricHistory([]trueNASReportingGetDataResponse{response}) + var points []TimeSeriesPoint + if history != nil { + points = history.CPUPercent + } + if len(points) != tc.count || (tc.count > 0 && points[len(points)-1].Timestamp.Unix() != tc.last) { + t.Fatalf("invalid timing/value handling: %+v", points) + } + }) + } +} + +func TestCORE13CPUAndMemoryEmptyBuckets(t *testing.T) { + for _, tc := range []struct { + row []any + valid bool + want float64 + }{ + {[]any{0., 0., 0., 0., 100.}, true, 0}, {[]any{0., 0., 0., 0., 0.}, false, 0}, + {[]any{nil, 0., 0., 0., 100.}, false, 0}, {[]any{-1., 0., 0., 0., 100.}, false, 0}, + {[]any{10., 10., 10., 10., 160.}, true, 20}, + } { + response := trueNASReportingGetDataResponse{Name: "cpu", Legend: coreCPUStates, Start: 1789000000, End: 1789000000, Step: 10, Data: []any{tc.row}} + history := parseSystemMetricHistory([]trueNASReportingGetDataResponse{response}) + if (history != nil) != tc.valid { + t.Fatalf("CPU row %v became %+v", tc.row, history) + } + if tc.valid { + assertCOREPoint(t, history.CPUPercent, 0, response.Start, tc.want) + } + } + for _, tc := range []struct { + row []any + valid bool + }{ + {[]any{0., 0., 0., 16., 0.}, true}, {[]any{0., 0., 0., 0., 0.}, false}, {[]any{1., 1., 1., nil, 0.}, false}, + } { + response := trueNASReportingGetDataResponse{Name: "memory", Legend: coreMemoryClasses, Start: 1789000000, End: 1789000000, Step: 10, Data: []any{tc.row}} + history := parseSystemMetricHistory([]trueNASReportingGetDataResponse{response}) + if (history != nil) != tc.valid { + t.Fatalf("memory row %v became %+v", tc.row, history) + } + if tc.valid { + assertCOREPoint(t, history.MemoryAvailableBytes, 0, response.Start, 0) + } + } +} + +func TestCORE13MultipleDevicesAndFreshness(t *testing.T) { + base := trueNASReportingGetDataResponse{Name: "interface", Identifier: "nic-a", Legend: []string{"rx", "tx", "overlap"}, Start: 1789000000, End: 1789000020, Step: 10, Data: []any{[]any{10., 20., 10.}, []any{0., 5., 0.}, []any{nil, nil, nil}}} + second := base + second.Identifier = "nic-b" + second.Data = []any{[]any{100., 200., 100.}, []any{0., nil, 0.}, []any{4., 8., 4.}} + for _, replies := range [][]trueNASReportingGetDataResponse{{base, second, base}, {second, base, base}} { + history := parseSystemMetricHistory(replies) + assertCOREPoint(t, history.NetInRate, 0, base.Start, 110) + assertCOREPoint(t, history.NetInRate, 1, base.Start+10, 0) + assertCOREPoint(t, history.NetOutRate, 0, base.Start, 220) + if len(history.NetInRate) != 2 || len(history.NetOutRate) != 1 { + t.Fatalf("missing device reading treated as zero: %+v", history) + } + } + if got := parseSystemMetricHistory(liveReportingResponses([]trueNASReportingGetDataResponse{base}, base.End+31)); got != nil { + t.Fatalf("stale sample became current: %+v", got) + } + fresh := parseSystemMetricHistory(liveReportingResponses([]trueNASReportingGetDataResponse{base}, base.End)) + assertCOREPoint(t, fresh.NetInRate, 1, base.Start+10, 0) + if at := systemInfoFromMetricHistory(fresh).CollectedAt.Unix(); at != base.Start+10 { + t.Fatalf("sample falsely stamped as current: %d", at) + } + if !reflect.DeepEqual(base.Data[0], []any{10., 20., 10.}) { + t.Fatal("live filtering mutated native History") + } +} + +// A request-validating REST bridge models catalogue -> per-device requests -> +// reporting rows -> complete snapshot and canonical host/history. Only the five +// excerpts are native evidence; catalogue, ARC, HTTP and timing are controls. +type core13ReportingTransport struct { + routes alertArgsTransport + replies []trueNASReportingGetDataResponse + requested []map[string]any + catalogueCalls, reportingCalls int +} + +func (s *core13ReportingTransport) RoundTrip(r *http.Request) (*http.Response, error) { + if r.URL.Path == "/api/v2.0/reporting/graphs" { + s.catalogueCalls++ + return (alertArgsTransport{r.URL.Path: {body: `[{"name":"disk","identifiers":["disk-a","disk-a",""]},{"name":"interface","identifiers":["nic-a"]}]`}}).RoundTrip(r) + } + if r.URL.Path != "/api/v2.0/reporting/get_data" { + return s.routes.RoundTrip(r) + } + s.reportingCalls++ + var request struct { + Graphs []map[string]any `json:"graphs"` + Query map[string]any `json:"query"` + } + if err := json.NewDecoder(r.Body).Decode(&request); err != nil { + return nil, err + } + s.requested = append(s.requested, request.Graphs...) + var responses []trueNASReportingGetDataResponse + for _, graph := range request.Graphs { + name := readStringAny(graph, "name") + identifier := readStringAny(graph, "identifier") + if (name == "disk" && identifier != "disk-a") || (name == "interface" && identifier != "nic-a") || (name != "disk" && name != "interface" && graph["identifier"] != nil) { + return (alertArgsTransport{r.URL.Path: {status: 422, body: `{"error":"native graph needs the correct device identifier"}`}}).RoundTrip(r) + } + for _, reply := range s.replies { + if reply.Name != name { + continue + } + reply.End = int64(request.Query["end"].(float64)) + reply.Start = reply.End - int64(len(reply.Data)-1)*reply.Step + responses = append(responses, reply) + } + if name == "arcsize" { + end := int64(request.Query["end"].(float64)) + responses = append(responses, trueNASReportingGetDataResponse{Name: name, Start: end - 10, End: end, Step: 10, Legend: []string{"arcsize_value"}, Data: []any{[]any{float64(8 << 30)}, []any{nil}}}) + } + } + raw, err := json.Marshal(responses) + if err != nil { + return nil, err + } + return (alertArgsTransport{r.URL.Path: {body: string(raw)}}).RoundTrip(r) +} + +func TestCORE13SnapshotAndCanonicalHistoryJourney(t *testing.T) { + previous := IsFeatureEnabled() + SetFeatureEnabled(true) + t.Cleanup(func() { SetFeatureEnabled(previous) }) + routes := alertArgsTransport(defaultAPIResponses()) + routes["/api/v2.0/system/info"] = apiResponse{body: `{"hostname":"core-fixture","version":"TrueNAS-13.0-U6.1","physmem":34359738368,"cores":8}`} + transport := &core13ReportingTransport{routes: routes, replies: core13ReportingFixture(t)} + client := newLegacyRESTReportingClient(t, routes) + client.httpClient.Transport = transport + t.Cleanup(client.Close) + provider := NewLiveProviderForConnection(&APIFetcher{Client: client}, "core-fixture-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 != 34359738368 || len(snapshot.Disks) != 2 || len(snapshot.Pools) != 1 { + t.Fatal("native telemetry lost inventory/capacity") + } + var host *unifiedresources.Resource + for _, record := range provider.Records() { + if record.Resource.Type == unifiedresources.ResourceTypeAgent { + copy := record.Resource + host = © + } + } + if host == nil || host.Metrics.CPU == nil || host.Metrics.Memory == nil || host.Temperature == nil || host.Agent.Memory == nil { + t.Fatalf("collapsed/expanded row lost native readings: %+v", host) + } + wantUsed := int64(34359738368-2178945024) - int64(8<<30) + if *host.Metrics.Memory.Used != wantUsed || host.Agent.Memory.Used != wantUsed || *host.Temperature != 30.95 || host.Agent.Memory.Cache != 8<<30 { + t.Fatalf("ARC-aware current values wrong: %+v %+v", host.Metrics, host.Agent) + } + for key, metric := range map[string]*unifiedresources.MetricValue{"netin": host.Metrics.NetIn, "netout": host.Metrics.NetOut, "diskread": host.Metrics.DiskRead, "diskwrite": host.Metrics.DiskWrite} { + want := map[string]float64{"netin": 9728.2270494, "netout": 45551.708899, "diskread": 0, "diskwrite": 1276361.9904}[key] + if metric == nil || metric.Unit != "bytes/s" || math.Abs(metric.Value-want) > 0.000001 { + t.Fatalf("%s current rate: %+v", key, metric) + } + } + id, history, err := provider.SystemMetricHistory(context.Background(), time.Hour) + if err != nil { + t.Fatal(err) + } + if id != "core-fixture-connection" || len(history) != 7 { + t.Fatalf("canonical History panels: %q %+v", id, history) + } + if len(history["cpu"]) != 1 || len(history["memory"]) != 1 || len(history["temperature"]) != 1 || len(history["diskread"]) != 2 { + t.Fatal("empty buckets filled with means or zeros") + } + if math.Abs(history["memory"][0].Value-host.Metrics.Memory.Percent) > 0.000001 { + t.Fatal("live/History disagree on free RAM/ARC") + } + } + if transport.catalogueCalls != 4 || transport.reportingCalls != 4 { + t.Fatalf("successful catalogue/batch unexpectedly retried: %+v", transport) + } + for _, graph := range transport.requested { + if name := readStringAny(graph, "name"); (name == "disk" || name == "interface") && readStringAny(graph, "identifier") == "" { + t.Fatal("device graph queried without identifier") + } + } +} + +func TestCORE13CatalogueAndResponseSafety(t *testing.T) { + for _, tc := range []struct { + status int + body string + fallback bool + }{ + {404, `{}`, true}, {405, `{}`, true}, {401, `{}`, false}, {403, `{}`, false}, {429, `{}`, false}, {500, `{}`, false}, {200, `{`, false}, + } { + routes := alertArgsTransport{"/api/v2.0/reporting/graphs": {status: tc.status, body: tc.body}} + client := newLegacyRESTReportingClient(t, routes) + graphs, err := client.legacyReportingGraphs(context.Background()) + if tc.fallback { + if err != nil || len(graphs) != 6 { + t.Fatalf("missing catalogue fallback: %v %v", graphs, err) + } + } else if err == nil { + t.Fatalf("catalogue failure concealed: %d", tc.status) + } + client.Close() + } + client := newLegacyRESTReportingClient(t, nil) + ctx, cancel := context.WithCancel(context.Background()) + cancel() + if _, err := client.GetSystemTelemetry(ctx); !errors.Is(err, context.Canceled) { + t.Fatalf("cancelled catalogue: %v", err) + } + replies := core13ReportingFixture(t) + missing := bindLegacyReportingResponses([]map[string]any{reportingGraph("interface", "nic-a"), reportingGraph("interface", "nic-b")}, replies) + if got := parseSystemMetricHistory(missing); got != nil { + t.Fatalf("omitted device became a complete host rate: %+v", got) + } + oversized := `[{"name":"disk","identifiers":[` + strings.Repeat(`"disk-a",`, 300) + `"disk-b"]}]` + // Duplicate catalogue entries do not consume the graph bound. + client = newLegacyRESTReportingClient(t, alertArgsTransport{"/api/v2.0/reporting/graphs": {body: oversized}}) + if graphs, err := client.legacyReportingGraphs(context.Background()); err != nil || len(graphs) != 6 { + t.Fatalf("catalogue deduplication: %v %v", graphs, err) + } + client.Close() +} + +func TestCORE13ARCAlignedWithFreeMemory(t *testing.T) { + replies := core13ReportingFixture(t) + history := parseSystemMetricHistory(replies) + history.ARCSizeBytes = []TimeSeriesPoint{{Timestamp: time.Unix(1789000000-10, 0), Value: 8 << 30}} + system := systemInfoFromMetricHistory(history) + if system.ARCSizeBytes != 0 { + t.Fatal("old ARC subtracted from newer memory") + } + history.ARCSizeBytes = append(history.ARCSizeBytes, TimeSeriesPoint{Timestamp: time.Unix(1789000000, 0), Value: 8 << 30}) + system = systemInfoFromMetricHistory(history) + system.MemoryTotalBytes = 32 << 30 + if system.ARCSizeBytes != 8<<30 { + t.Fatal("same-bucket ARC lost") + } + points := systemMemoryPercentHistory(history, system.MemoryTotalBytes) + want := metricsFromTrueNASSystem(*system, 0, 0).Memory.Percent + assertCOREPoint(t, points, 0, 1789000000, want) + history.ARCSizeBytes[1].Value = 64 << 30 + points = systemMemoryPercentHistory(history, system.MemoryTotalBytes) + assertCOREPoint(t, points, 0, 1789000000, 0) +} + +func TestCORE13DiskTemperatureHistoryExternalTiming(t *testing.T) { + response := trueNASReportingGetDataResponse{Name: "disktemp", Identifier: "disk-a", Legend: []string{"temperature"}, Start: 1789000000, End: 1789000020, Step: 10, Data: []any{[]any{40.}, []any{nil}, []any{42.}}} + points := parseReportingDiskTemperatureHistory([]trueNASReportingGetDataResponse{response})["disk-a"] + assertCOREPoint(t, points, 0, response.Start, 40) + assertCOREPoint(t, points, 1, response.Start+20, 42) +} diff --git a/internal/truenas/provider.go b/internal/truenas/provider.go index 2e909fcaa..f24f0efe4 100644 --- a/internal/truenas/provider.go +++ b/internal/truenas/provider.go @@ -369,6 +369,9 @@ func (p *Provider) SystemMetricHistory(ctx context.Context, duration time.Durati if len(nativeHistory.DiskWriteRate) > 0 { metricMap["diskwrite"] = cloneTimeSeriesPoints(nativeHistory.DiskWriteRate) } + if temperatures := systemTemperatureHistory(nativeHistory.TemperatureCelsius); len(temperatures) > 0 { + metricMap["temperature"] = temperatures + } if len(metricMap) == 0 { return resourceID, nil, nil } diff --git a/internal/truenas/reporting_aggregation_test.go b/internal/truenas/reporting_aggregation_test.go index b819a3319..1fcecb8b7 100644 --- a/internal/truenas/reporting_aggregation_test.go +++ b/internal/truenas/reporting_aggregation_test.go @@ -221,6 +221,9 @@ func TestRESTReportingDefaultAggregationsPreserveSamples(t *testing.T) { var calls int client := newLegacyRESTReportingClient(t, nil) client.httpClient.Transport = reportingAggregationRESTFunc(func(r *http.Request) (*http.Response, error) { + if r.URL.Path == "/api/v2.0/reporting/graphs" { + return (alertArgsTransport{}).RoundTrip(r) + } calls++ var request struct { Graphs []map[string]any `json:"graphs"` diff --git a/internal/truenas/testdata/core13_reporting.json b/internal/truenas/testdata/core13_reporting.json new file mode 100644 index 000000000..093fde369 --- /dev/null +++ b/internal/truenas/testdata/core13_reporting.json @@ -0,0 +1,266 @@ +[ + { + "name": "cpu", + "identifier": null, + "data": [ + [ + 0.033358538468, + 4.4805217137, + 7.8754519386, + 7.8754519386, + 100.0 + ], + [ + 0.0, + 0.0, + 0.0, + 0.0, + 0.0 + ], + [ + 0.0, + 0.0, + 0.0, + 0.0, + 0.0 + ] + ], + "start": 1789000000, + "end": 1789000020, + "step": 10, + "legend": [ + "interrupt", + "system", + "user", + "nice", + "idle" + ], + "aggregations": { + "min": [ + 0.0, + 0.0, + 0.0, + 0.0, + 0.0 + ], + "mean": [ + 0.019921323668173407, + 1.734965537597784, + 2.9715496785945983, + 2.9715496785945983, + 99.44598337950139 + ], + "max": [ + 0.088098918119, + 4.7359432381, + 8.3218063112, + 8.3218063112, + 100.0 + ] + } + }, + { + "name": "cputemp", + "identifier": null, + "data": [ + [ + 29.95, + 29.95, + 30.95, + 30.95, + 30.95, + 30.95, + 29.95, + 29.95 + ], + [ + null, + null, + null, + null, + null, + null, + null, + null + ] + ], + "start": 1789000000, + "end": 1789000010, + "step": 10, + "legend": [ + "cputemp0", + "cputemp1", + "cputemp2", + "cputemp3", + "cputemp4", + "cputemp5", + "cputemp6", + "cputemp7" + ], + "aggregations": { + "min": [ + 29.95, + 29.95, + 30.95, + 30.95, + 30.95, + 30.95, + 29.95, + 29.95 + ], + "mean": [ + 29.95, + 29.95, + 30.95, + 30.95, + 30.95, + 30.95, + 29.95, + 29.95 + ], + "max": [ + 29.95, + 29.95, + 30.95, + 30.95, + 30.95, + 30.95, + 29.95, + 29.95 + ] + } + }, + { + "name": "memory", + "identifier": null, + "data": [ + [ + 188961888.74, + 1918028677.6, + 29016502272.0, + 1335296.0, + 2178945024.0 + ], + [ + 0.0, + 0.0, + 0.0, + 0.0, + 0.0 + ] + ], + "start": 1789000000, + "end": 1789000010, + "step": 10, + "legend": [ + "memory-active_value", + "memory-inactive_value", + "memory-wired_value", + "memory-laundry_value", + "memory-free_value" + ], + "aggregations": { + "min": [ + 0.0, + 0.0, + 0.0, + 0.0, + 0.0 + ], + "mean": [ + 189003434.29373962, + 1911812503.3603878, + 28935538170.930748, + 1331597.1191135733, + 2179191845.1520777 + ], + "max": [ + 193937266.75, + 1921061077.8, + 29038160280.0, + 1335296.0, + 2199964529.7 + ] + } + }, + { + "name": "disk", + "identifier": "disk-a", + "data": [ + [ + 484.54070578, + 538546.26401 + ], + [ + 0.0, + 1276361.9904 + ], + [ + null, + null + ] + ], + "start": 1789000000, + "end": 1789000020, + "step": 10, + "legend": [ + "disk_octets_read", + "disk_octets_write" + ], + "aggregations": { + "min": [ + 0.0, + 383856.71231 + ], + "mean": [ + 336.1788029753333, + 1030387.1104943056 + ], + "max": [ + 2563.1187967, + 3593935.6719 + ] + } + }, + { + "name": "interface", + "identifier": "nic-a", + "data": [ + [ + 9728.2270494, + 45551.708899, + 9728.2270494 + ], + [ + null, + null, + null + ] + ], + "start": 1789000000, + "end": 1789000010, + "step": 10, + "legend": [ + "rx", + "tx", + "overlap" + ], + "aggregations": { + "min": [ + 3905.0068688, + 2298.9182121, + 2298.9182121 + ], + "mean": [ + 11905.4645962075, + 16747.828344191388, + 8876.474186066389 + ], + "max": [ + 61200.017929, + 86660.554326, + 61200.017929 + ] + } + } +] diff --git a/internal/truenas/types.go b/internal/truenas/types.go index 26c6fb6fc..dfffdaa91 100644 --- a/internal/truenas/types.go +++ b/internal/truenas/types.go @@ -204,6 +204,7 @@ type TimeSeriesPoint struct { // SystemMetricHistory stores provider-native TrueNAS system history before it // is normalized onto the canonical monitoring guest-chart surface. type SystemMetricHistory struct { + TemperatureCelsius map[string][]TimeSeriesPoint CPUPercent []TimeSeriesPoint MemoryPercent []TimeSeriesPoint MemoryUsedBytes []TimeSeriesPoint From c78840e3717b7dd3794ed401529e5f4bdcab2e50 Mon Sep 17 00:00:00 2001 From: "pulse-triage[bot]" <249995291+pulse-triage[bot]@users.noreply.github.com> Date: Fri, 2 Oct 2026 00:32:55 +0100 Subject: [PATCH 2/5] Classify the native CORE reporting fixture for canonical verification Bind the anonymised JSON regression input to the existing TrueNAS runtime proof policy. The repair restores existing canonical behavior and has no public contract shape change; its source, tests and substantive monitoring account remain in the preceding commit. Change-source: pulse-maintainer --- docs/release-control/v6/internal/subsystems/registry.json | 1 + 1 file changed, 1 insertion(+) diff --git a/docs/release-control/v6/internal/subsystems/registry.json b/docs/release-control/v6/internal/subsystems/registry.json index 28e7733d8..2b60224a0 100644 --- a/docs/release-control/v6/internal/subsystems/registry.json +++ b/docs/release-control/v6/internal/subsystems/registry.json @@ -6162,6 +6162,7 @@ "internal/truenas/fixtures.go", "internal/truenas/provider.go", "internal/truenas/telemetry.go", + "internal/truenas/testdata/core13_reporting.json", "internal/truenas/transport.go", "internal/truenas/types.go" ], From bab5cf77d60f5252311bad05f28aecb4d76f1de1 Mon Sep 17 00:00:00 2001 From: "pulse-triage[bot]" <249995291+pulse-triage[bot]@users.noreply.github.com> Date: Fri, 2 Oct 2026 00:48:46 +0100 Subject: [PATCH 3/5] Prove native CORE catalogue overflow fails without partial totals Exercise the unique-device selection ceiling separately from deduplication, and keep the graph-isolation account aligned with catalogue-selected requests. No runtime logic or public contract shape is changed. Change-source: pulse-maintainer --- .../v6/internal/subsystems/monitoring.md | 4 +++- internal/truenas/client.go | 3 ++- internal/truenas/core_reporting_test.go | 17 +++++++++++++++++ 3 files changed, 22 insertions(+), 2 deletions(-) diff --git a/docs/release-control/v6/internal/subsystems/monitoring.md b/docs/release-control/v6/internal/subsystems/monitoring.md index b6d5ff16c..28d198cc5 100644 --- a/docs/release-control/v6/internal/subsystems/monitoring.md +++ b/docs/release-control/v6/internal/subsystems/monitoring.md @@ -323,7 +323,9 @@ Verification: `internal/truenas/core_reporting_test.go` and its anonymised `testdata/core13_reporting.json` cover five source-bound native excerpts, row/timing/alias/zero/null/empty/finite controls, device totals, stale buckets, ARC alignment, current and canonical History projections, request validation, -catalogue limits and failure boundaries. The excerpt windows are deliberately +catalogue limits and failure boundaries. Unique overflow is rejected rather than +truncated, while duplicate identifiers consume only one selection entry. The +excerpt windows are deliberately shortened controls, not complete native windows. Extended `TestTrueNASPartialReportingPipeline` in `internal/monitoring/truenas_poller_test.go` crosses pinned-certificate HTTP, diff --git a/internal/truenas/client.go b/internal/truenas/client.go index 769c78e74..3c932617b 100644 --- a/internal/truenas/client.go +++ b/internal/truenas/client.go @@ -426,7 +426,8 @@ func legacyRESTReportingGraphs() []map[string]any { // getLegacySystemReportingData keeps a rejected optional graph from discarding // usable CPU/memory readings. Keep the successful batch fast path; only split // graph-validation/server errors, never authentication, rate-limit, endpoint, -// transport or cancellation failures. The query and graph set remain unchanged. +// transport or cancellation failures. Bind the catalogue-selected graph set once +// for both the batch and split; the requested query window remains unchanged. func (c *Client) getLegacySystemReportingData(ctx context.Context, query map[string]any) ([]trueNASReportingGetDataResponse, error) { graphs, err := c.legacyReportingGraphs(ctx) if err != nil { diff --git a/internal/truenas/core_reporting_test.go b/internal/truenas/core_reporting_test.go index 07ee611c7..f9cc2f586 100644 --- a/internal/truenas/core_reporting_test.go +++ b/internal/truenas/core_reporting_test.go @@ -4,6 +4,7 @@ import ( "context" "encoding/json" "errors" + "fmt" "math" "net/http" "os" @@ -378,6 +379,22 @@ func TestCORE13CatalogueAndResponseSafety(t *testing.T) { client.Close() } +func TestCORE13CatalogueRejectsUniqueOverflow(t *testing.T) { + identifiers := make([]string, 253) + for i := range identifiers { + identifiers[i] = fmt.Sprintf("disk-%d", i) + } + raw, err := json.Marshal([]map[string]any{{"name": "disk", "identifiers": identifiers}}) + if err != nil { + t.Fatal(err) + } + client := newLegacyRESTReportingClient(t, alertArgsTransport{"/api/v2.0/reporting/graphs": {body: string(raw)}}) + defer client.Close() + if graphs, err := client.legacyReportingGraphs(context.Background()); err == nil || graphs != nil { + t.Fatal("oversized unique catalogue became a partial host total") + } +} + func TestCORE13ARCAlignedWithFreeMemory(t *testing.T) { replies := core13ReportingFixture(t) history := parseSystemMetricHistory(replies) From 579eee290759f5518f6ef16d0b9189a2dec7d856 Mon Sep 17 00:00:00 2001 From: "pulse-triage[bot]" <249995291+pulse-triage[bot]@users.noreply.github.com> Date: Fri, 2 Oct 2026 01:49:25 +0100 Subject: [PATCH 4/5] Preserve independent ARC and stop cancelled CORE reporting before I/O Complete the native CORE #2077 repair after its whole-suite adverse result: retain ARC without inventing free RAM, enforce live-window freshness even with coarse RRD steps, and assert bounded graph splitting including CPU temperature. Keep empty/zero, aligned ARC and transport boundaries intact. Change-source: pulse-maintainer --- .../v6/internal/subsystems/monitoring.md | 7 ++++-- internal/truenas/client.go | 8 +++++++ internal/truenas/client_test.go | 6 ++--- internal/truenas/core_reporting.go | 9 ++++++-- internal/truenas/core_reporting_test.go | 22 +++++++++++++++++++ 5 files changed, 45 insertions(+), 7 deletions(-) diff --git a/docs/release-control/v6/internal/subsystems/monitoring.md b/docs/release-control/v6/internal/subsystems/monitoring.md index 28d198cc5..312bb58e6 100644 --- a/docs/release-control/v6/internal/subsystems/monitoring.md +++ b/docs/release-control/v6/internal/subsystems/monitoring.md @@ -291,13 +291,16 @@ row timestamps and supported object aliases retain compatibility. A metric is never a timestamp, a null is never zero and a window mean is never a replacement for a missing current row. History retains measured buckets; native current readings accept only the last two steps before the requested end, preserving -sample time rather than stamping an old value as newly measured. The five-state +sample time rather than stamping an old value as newly measured. This allowance +never exceeds the live query window, even with a coarse returned step. The five-state CORE CPU vector is normalized by its sum (RRD state rates need not sum to 100). All-zero CPU or complete memory-class vectors are empty buckets, not 100% usage; explicit zero usage, all-idle CPU, zero-free RAM with used pages and zero I/O remain measured zeros. Free memory comes from `memory-free_value`, not active pages or the sum of memory classes. ARC subtraction uses the matching free-RAM -bucket. Reporting I/O values remain rates in the existing bytes/s contract; +bucket when free RAM is present; independently reported ARC remains available +without claiming measured memory usage. Pre-cancelled collection stops before +reading the catalogue or posting a reporting query. Reporting I/O values remain rates in the existing bytes/s contract; step supplies time, not a second rate division; `overlap` is not extra traffic. Legacy REST reads `/reporting/graphs` once per live/History collection and sends diff --git a/internal/truenas/client.go b/internal/truenas/client.go index 3c932617b..a6ff48f5a 100644 --- a/internal/truenas/client.go +++ b/internal/truenas/client.go @@ -342,6 +342,14 @@ func systemInfoFromMetricHistory(history *SystemMetricHistory) *SystemInfo { latestFree = free.Timestamp } } + // ARC can be observed independently when no free-RAM sample exists; + // retain it without implying known memory usage. With free RAM, require + // cache and free to describe the same bucket before deriving usage. + if latestFree.IsZero() { + if value, ok := latestTimeSeriesValue(history.ARCSizeBytes); ok { + system.ARCSizeBytes = int64(value) + } + } for _, arc := range history.ARCSizeBytes { if arc.Timestamp.Equal(latestFree) { system.ARCSizeBytes = int64(arc.Value) diff --git a/internal/truenas/client_test.go b/internal/truenas/client_test.go index ab68c5d3d..b39886f44 100644 --- a/internal/truenas/client_test.go +++ b/internal/truenas/client_test.go @@ -2195,8 +2195,8 @@ func TestRESTReportingGraphFailurePreservesSnapshotTelemetry(t *testing.T) { t.Fatal("inventory lost") } } - if transport.calls != 12 { - t.Fatalf("reporting calls = %d, want bounded batch + five graphs per snapshot", transport.calls) + if transport.calls != 2*(1+len(legacyRESTReportingGraphs())) { + t.Fatalf("reporting calls = %d, want bounded batch + selected graphs per snapshot", transport.calls) } history, err := client.GetSystemMetricHistory(context.Background(), time.Hour) if err != nil { @@ -2464,7 +2464,7 @@ func TestRESTSnapshotRetainsSanitizedTelemetryFailure(t *testing.T) { } wantCalls := 1 if tc.status == 422 { - wantCalls = 6 + wantCalls = 1 + len(legacyRESTReportingGraphs()) } if transport.calls != wantCalls { t.Fatalf("reporting calls=%d, want %d", transport.calls, wantCalls) diff --git a/internal/truenas/core_reporting.go b/internal/truenas/core_reporting.go index ac3eb46cc..0e735e6b4 100644 --- a/internal/truenas/core_reporting.go +++ b/internal/truenas/core_reporting.go @@ -16,6 +16,9 @@ import ( // Older REST bridges without the catalogue retain the unscoped compatibility // request. Authentication, transport and malformed replies are never retried. func (c *Client) legacyReportingGraphs(ctx context.Context) ([]map[string]any, error) { + if err := ctx.Err(); err != nil { + return nil, err + } var catalogue []struct { Name string `json:"name"` Identifiers []string `json:"identifiers"` @@ -97,7 +100,9 @@ func reportingRow(response trueNASReportingGetDataResponse, index int) (time.Tim // Keep native History intact, but do not label an old RRD point as a current // reading. Two returned steps allow the usual unfinished/null RRD bucket; -// older points stay available only through History. Preserve row indices. +// older points stay available only through History. Never exceed the live +// query window even if a coarse/malformed step allows older buckets. Preserve +// row indices. func liveReportingResponses(responses []trueNASReportingGetDataResponse, end int64) []trueNASReportingGetDataResponse { live := append([]trueNASReportingGetDataResponse(nil), responses...) for i, response := range live { @@ -108,7 +113,7 @@ func liveReportingResponses(responses []trueNASReportingGetDataResponse, end int for index := range response.Data { timestamp, _, ok := reportingRow(response, index) gap := end - timestamp.Unix() - if !ok || timestamp.Unix() > end || (gap > response.Step && gap-response.Step > response.Step) { + if !ok || timestamp.Unix() > end || gap > legacyRESTTelemetryWindowSeconds || (gap > response.Step && gap-response.Step > response.Step) { live[i].Data[index] = nil } } diff --git a/internal/truenas/core_reporting_test.go b/internal/truenas/core_reporting_test.go index f9cc2f586..08a2594fd 100644 --- a/internal/truenas/core_reporting_test.go +++ b/internal/truenas/core_reporting_test.go @@ -417,6 +417,28 @@ func TestCORE13ARCAlignedWithFreeMemory(t *testing.T) { assertCOREPoint(t, points, 0, 1789000000, 0) } +func TestCORE13CoarseStepDoesNotMakeOldRowsCurrent(t *testing.T) { + response := trueNASReportingGetDataResponse{Name: "cpu", Legend: []string{"usage"}, Start: 1789000000, End: 1789000000, Step: 3600, Data: []any{[]any{12.}}} + if history := parseSystemMetricHistory(liveReportingResponses([]trueNASReportingGetDataResponse{response}, 1789000400)); history != nil { + t.Fatal("coarse step made an old CPU sample current") + } + if history := parseSystemMetricHistory([]trueNASReportingGetDataResponse{response}); history == nil { + t.Fatal("live age filtering destroyed native History") + } +} + +func TestCORE13ARCOnlyDoesNotInventMemoryUsage(t *testing.T) { + history := &SystemMetricHistory{ARCSizeBytes: []TimeSeriesPoint{{Timestamp: time.Unix(1789000000, 0), Value: 8 << 30}}} + system := systemInfoFromMetricHistory(history) + if system.ARCSizeBytes != 8<<30 || system.Telemetry.Memory { + t.Fatal("ARC-only telemetry lost or claimed free RAM") + } + system.MemoryTotalBytes = 32 << 30 + if metricsFromTrueNASSystem(*system, 0, 0).Memory != nil { + t.Fatal("ARC-only reading invented usage") + } +} + func TestCORE13DiskTemperatureHistoryExternalTiming(t *testing.T) { response := trueNASReportingGetDataResponse{Name: "disktemp", Identifier: "disk-a", Legend: []string{"temperature"}, Start: 1789000000, End: 1789000020, Step: 10, Data: []any{[]any{40.}, []any{nil}, []any{42.}}} points := parseReportingDiskTemperatureHistory([]trueNASReportingGetDataResponse{response})["disk-a"] From 45eec9a20eac55a356acfb1c17833cd6c8d067ad Mon Sep 17 00:00:00 2001 From: "pulse-triage[bot]" <249995291+pulse-triage[bot]@users.noreply.github.com> Date: Fri, 2 Oct 2026 02:09:48 +0100 Subject: [PATCH 5/5] Keep native TrueNAS Thermals in local and persisted host History Finish the connected #2077 outcome: persist an observed TrueNAS host temperature through the existing canonical writer, so local chart coverage cannot suppress the Thermals panel. Extend the pinned REST-to-chart control to cover the write, persistent readback and full local-window fast path; preserve source and absence gates. Change-source: pulse-maintainer Contract-Neutral: Restores existing TrueNAS Thermals history through its canonical writer; no agent lifecycle, schema, route, resource identity or alert-policy delta. --- .../v6/internal/subsystems/monitoring.md | 8 +++++++- internal/monitoring/monitor.go | 12 ++++++++++++ internal/monitoring/truenas_poller_test.go | 17 ++++++++++++----- 3 files changed, 31 insertions(+), 6 deletions(-) diff --git a/docs/release-control/v6/internal/subsystems/monitoring.md b/docs/release-control/v6/internal/subsystems/monitoring.md index 312bb58e6..af091c647 100644 --- a/docs/release-control/v6/internal/subsystems/monitoring.md +++ b/docs/release-control/v6/internal/subsystems/monitoring.md @@ -319,6 +319,11 @@ assuming a missing member was idle. Other graphs remain independently usable. Native CPU temperature maps into the existing canonical host temperature and `temperature` History key, using the same package-versus-core outlier safeguard as current readings. Disk-temperature History uses the same envelope decoder. +The canonical writer also records each positive finite TrueNAS host temperature +in local and persistent `temperature` History, under the existing source/dedup +gates. A full local CPU window must not make the existing Thermals panel depend +on native fallback or lose its observed temperature. Missing temperature is not +written as zero, and no additional sensor writes are added to other sources. No frontend route, resource identity, alert suppression policy, credential scope or recognised-legacy transport boundary changes. @@ -333,7 +338,8 @@ shortened controls, not complete native windows. Extended `TestTrueNASPartialReportingPipeline` in `internal/monitoring/truenas_poller_test.go` crosses pinned-certificate HTTP, poller, canonical registry, writer, in-memory/persistent chart readbacks and -native fallback, including temperature. Existing CORE isolation, recognised +native fallback, including temperature and the sufficiently covered local +chart fast path. Existing CORE isolation, recognised transport/auth/TLS and modern SCALE aggregation tests remain required. Source proofs do not substitute for containing-release qualification or installed CORE live/collapsed/expanded/History acceptance. diff --git a/internal/monitoring/monitor.go b/internal/monitoring/monitor.go index b0f3e2e5f..e2cf45886 100644 --- a/internal/monitoring/monitor.go +++ b/internal/monitoring/monitor.go @@ -5539,6 +5539,18 @@ func (m *Monitor) syncUnifiedAgentMetrics(store ResourceStoreInterface, sinks .. metricKey := fmt.Sprintf("agent:%s", targetID) observedAt := unifiedResourceObservedAt(resource, now) + // Native TrueNAS temperature must also survive the local History path. + // A sufficiently covered ring/store read need not query native fallback. + if monitorHasSource(resource.Sources, unifiedresources.SourceTrueNAS) && resource.Temperature != nil { + value := *resource.Temperature + if value > 0 && !math.IsNaN(value) && !math.IsInf(value, 0) { + if m.metricsHistory != nil { + m.metricsHistory.AddGuestMetric(metricKey, "temperature", value, observedAt) + } + appendStoreWrite("agent", targetID, "temperature", value, observedAt) + } + } + if metric := resource.Metrics.CPU; metric != nil { value := metric.Percent if value == 0 { diff --git a/internal/monitoring/truenas_poller_test.go b/internal/monitoring/truenas_poller_test.go index bf8fe12a8..bd99a72e1 100644 --- a/internal/monitoring/truenas_poller_test.go +++ b/internal/monitoring/truenas_poller_test.go @@ -2234,7 +2234,7 @@ func TestTrueNASPartialReportingPipeline(t *testing.T) { {"name":"cputemp","legend":["cputemp0","cputemp1"],"data":[[41,42]]}, {"name":"interface","identifier":"nic-a","legend":["rx","tx","overlap"],"data":[[8,16,8]]}, {"name":"disk","identifier":"disk-a","legend":["disk_octets_read","disk_octets_write"],"data":[[0,32]]} - ]`, 200, map[string]float64{"cpu": 10, "memory": 62.5, "netin": 8, "netout": 16, "diskread": 0, "diskwrite": 32}, "", true}, + ]`, 200, map[string]float64{"cpu": 10, "memory": 62.5, "netin": 8, "netout": 16, "diskread": 0, "diskwrite": 32, "temperature": 42}, "", true}, } var current atomic.Int64 var nativeEnd atomic.Int64 @@ -2389,7 +2389,7 @@ func TestTrueNASPartialReportingPipeline(t *testing.T) { 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"} { + for _, key := range []string{"cpu", "memory", "disk", "netin", "netout", "diskread", "diskwrite", "temperature"} { 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]) } @@ -2402,9 +2402,6 @@ func TestTrueNASPartialReportingPipeline(t *testing.T) { continue } wantNative := len(cycle.want) - if cycle.native { - wantNative++ - } if err != nil || id != instance.ID || len(native) != wantNative { t.Fatalf("cycle %d native History=%+v id=%q err=%v", index, native, id, err) } @@ -2430,4 +2427,14 @@ func TestTrueNASPartialReportingPipeline(t *testing.T) { } } } + // Once local CPU History covers the window, the chart's fast path must + // still retain Thermals without relying on another appliance request. + now := time.Now() + for i := 0; i <= 60; i++ { + monitor.metricsHistory.AddGuestMetric("agent:"+instance.ID, "cpu", 10, now.Add(time.Duration(i-60)*time.Minute)) + } + chart := monitor.GetGuestMetricsForChart("agent:"+instance.ID, "agent", instance.ID, time.Hour) + if points := chart["temperature"]; len(points) != 1 || points[0].Value != 42 { + t.Fatalf("sufficiently covered local History lost Thermals: %+v", points) + } }