Integrate reviewed native TrueNAS CORE telemetry and History repair

Change-source: pulse-maintainer
This commit is contained in:
pulse-triage[bot] 2026-10-02 02:56:46 +01:00
commit b32a4369e7
12 changed files with 1300 additions and 78 deletions

View file

@ -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)**

View file

@ -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"

View file

@ -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 {

View file

@ -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)
}
}

View file

@ -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
}

View file

@ -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)

View file

@ -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
}

View file

@ -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 = &copy
}
}
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)
}

View file

@ -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
}

View file

@ -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"`

View file

@ -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
]
}
}
]

View file

@ -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