diff --git a/docs/release-control/v6/internal/subsystems/monitoring.md b/docs/release-control/v6/internal/subsystems/monitoring.md index d0764f50b..af091c647 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,75 @@ 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. 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 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 +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. +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. + +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. 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, +poller, canonical registry, writer, in-memory/persistent chart readbacks and +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. + + **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..2b60224a0 100644 --- a/docs/release-control/v6/internal/subsystems/registry.json +++ b/docs/release-control/v6/internal/subsystems/registry.json @@ -6157,10 +6157,12 @@ "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", "internal/truenas/telemetry.go", + "internal/truenas/testdata/core13_reporting.json", "internal/truenas/transport.go", "internal/truenas/types.go" ], @@ -6173,6 +6175,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/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 e0af04581..bd99a72e1 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, "temperature": 42}, "", 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) @@ -2338,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]) } @@ -2350,14 +2401,40 @@ func TestTrueNASPartialReportingPipeline(t *testing.T) { } continue } - if err != nil || id != instance.ID || len(native) != len(cycle.want) { + wantNative := len(cycle.want) + 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) + } + } + } + // 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) } } diff --git a/internal/truenas/client.go b/internal/truenas/client.go index 81671d09b..a6ff48f5a 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,28 @@ 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 + } } + // 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) + } + } + if value, ok := latestTimeSeriesValue(history.NetInRate); ok { system.NetInRate = value system.Telemetry.NetIn = true @@ -388,18 +427,26 @@ func legacyRESTReportingGraphs() []map[string]any { reportingGraph("arcsize", ""), reportingGraph("interface", ""), reportingGraph("disk", ""), + reportingGraph("cputemp", ""), } } // 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 := 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 +466,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 +2425,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 +3226,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 +3257,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 +3394,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..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 { @@ -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) @@ -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 new file mode 100644 index 000000000..0e735e6b4 --- /dev/null +++ b/internal/truenas/core_reporting.go @@ -0,0 +1,281 @@ +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) { + if err := ctx.Err(); err != nil { + return nil, err + } + 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. 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 { + 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 > legacyRESTTelemetryWindowSeconds || (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..08a2594fd --- /dev/null +++ b/internal/truenas/core_reporting_test.go @@ -0,0 +1,447 @@ +package truenas + +import ( + "context" + "encoding/json" + "errors" + "fmt" + "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 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) + 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 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"] + 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