Keep missing TrueNAS samples distinct from observed zero

Carry per-metric presence through REST and realtime snapshots, canonical host rows and shared History writes. Do not interpret arbitrary numeric object fields as reporting values. Preserve bounded telemetry failure diagnostics without blocking inventory or exposing provider text.

Change-source: pulse-maintainer
This commit is contained in:
pulse-triage[bot] 2026-10-01 03:25:28 +01:00
parent 1c8f51bf1f
commit 22ca07ba71
9 changed files with 543 additions and 33 deletions

View file

@ -150,9 +150,9 @@ host-chart history are reconstructed from the REST reporting API
`query` parameters). The reporting `memory` and `arcsize` graphs supply the
available and ARC readings that the JSON-RPC realtime subscription otherwise
provides. When no available-memory reading can be established, the provider
must not derive usage from a zero available value: the system memory metric is
must not derive usage from an absent available value: the system memory metric is
omitted and agent memory is projected as usage-unavailable with its known total,
so a CORE appliance is never reported as 100% used. This is a behavioral
so missing data never reports a CORE appliance as 100% used. This is a behavioral
correctness repair with no public API or schema delta.
`TestRESTSystemTelemetryReadsReportingMemory`,
`TestRESTSystemMetricHistoryUsesReporting` and
@ -163,6 +163,43 @@ and projection proofs, not native CORE appliance acceptance or reporter
confirmation.
### TrueNAS partial samples and nonblocking telemetry diagnostics — issue #2077
Live system identity and both telemetry parsers carry explicit per-metric
availability. A reporting timestamp or realtime interval is not evidence that
CPU, memory or any sibling I/O series was measured. Each canonical metric is
omitted unless its own reading exists; observed zeros remain valid, including
zero free memory. Capacity-only snapshots keep hardware capacity and mark
usage unavailable. Older static snapshots retain their compatibility projection;
no live client uses that fallback. Null/malformed single-series objects never
substitute a timestamp or another arbitrary numeric field for the missing value.
Supported generic `value`, `y` and `temperature` row aliases remain supported.
Optional telemetry failure must not discard inventory or degrade a successful
inventory poll. Existing connection diagnostics expose optional
`observed.telemetry` with six availability booleans (`cpu`, `memory`, `netIn`,
`netOut`, `diskRead`, `diskWrite`), a fixed `errorCategory` and, for HTTP failures,
`httpStatus`. No raw response body, endpoint, key, hostname or arbitrary provider
error text is copied into this projection. Error categories record the observed
collection result, not its cause. A later refresh replaces availability and
clears obsolete failures; snapshot and connection-summary copies cannot mutate
shared state. This is an additive diagnostic field on the existing read surface,
not a separate resource, route, telemetry store or alert policy.
`TestRESTReportingObservationPresence`,
`TestRESTSnapshotRetainsSanitizedTelemetryFailure`, `TestRealtimeObservationPresence`
and `TestSystemTelemetryFailureCategoriesAreBounded` in
`internal/truenas/client_test.go` pin repeated full snapshots, canonical rows,
native History/zero semantics and bounded failure/recovery controls.
`TestTrueNASPartialReportingPipeline` in
`internal/monitoring/truenas_poller_test.go` exercises the HTTP client, poller,
registry, memory metadata, shared metrics writer and in-memory/persisted chart
readbacks across partial, missing and zero-valued cycles. These synthetic proofs
are not the reporter's native response schema, installed CORE acceptance or a
claim that #2077's missing graph journey is resolved.
**Legacy reporting failure isolation and memory components — issue #2077**
Live telemetry and history retain the legacy REST graph/query contract and the

View file

@ -6105,6 +6105,7 @@
"internal/truenas/disk_health.go",
"internal/truenas/fixtures.go",
"internal/truenas/provider.go",
"internal/truenas/telemetry.go",
"internal/truenas/transport.go",
"internal/truenas/types.go"
],

View file

@ -57,6 +57,9 @@ type TrueNASConnectionObservedSummary struct {
Shares int `json:"shares"`
Disks int `json:"disks"`
RecoveryArtifacts int `json:"recoveryArtifacts"`
// Telemetry is separate from poll success: inventory can be current even
// when optional CPU/memory/IO collection failed or only partly succeeded.
Telemetry *truenas.SystemTelemetryAvailability `json:"telemetry,omitempty"`
}
// TrueNASConnectionSummary merges poll health with the most recent discovered
@ -930,6 +933,7 @@ func buildTrueNASObservedSummary(snapshot *truenas.FixtureSnapshot) *TrueNASConn
Shares: len(snapshot.Shares),
Disks: len(snapshot.Disks),
RecoveryArtifacts: len(snapshot.ZFSSnapshots) + len(snapshot.ReplicationTasks),
Telemetry: cloneTrueNASTelemetryAvailability(snapshot.System.Telemetry),
}
if host != "" || resourceID != "" {
summary.Systems = 1
@ -989,9 +993,18 @@ func cloneTrueNASObservedSummary(value *TrueNASConnectionObservedSummary) *TrueN
Shares: value.Shares,
Disks: value.Disks,
RecoveryArtifacts: value.RecoveryArtifacts,
Telemetry: cloneTrueNASTelemetryAvailability(value.Telemetry),
}
}
func cloneTrueNASTelemetryAvailability(value *truenas.SystemTelemetryAvailability) *truenas.SystemTelemetryAvailability {
if value == nil {
return nil
}
cloned := *value
return &cloned
}
func (p *TrueNASPoller) ingestRecoveryPoints(ctx context.Context, orgID string, connectionID string, provider *truenas.Provider) {
if p == nil || p.recoveryManager == nil || provider == nil {
return

View file

@ -2,6 +2,7 @@ package monitoring
import (
"context"
"encoding/json"
"errors"
"fmt"
"net"
@ -16,8 +17,10 @@ import (
"time"
"github.com/rcourtman/pulse-go-rewrite/internal/config"
"github.com/rcourtman/pulse-go-rewrite/internal/models"
"github.com/rcourtman/pulse-go-rewrite/internal/truenas"
"github.com/rcourtman/pulse-go-rewrite/internal/unifiedresources"
"github.com/rcourtman/pulse-go-rewrite/pkg/metrics"
)
func TestTrueNASPollerPollsConfiguredConnections(t *testing.T) {
@ -2205,3 +2208,148 @@ func TestTrueNASSuccessfulPollCadenceIncludesBoundedIdleGap(t *testing.T) {
})
}
}
// HTTP/client -> poller -> registry -> row/metric writer -> chart readback.
// Synthetic response shapes prove omission/zero handling, not CORE acceptance.
func TestTrueNASPartialReportingPipeline(t *testing.T) {
previous := truenas.IsFeatureEnabled()
truenas.SetFeatureEnabled(true)
t.Cleanup(func() { truenas.SetFeatureEnabled(previous) })
cycles := []struct {
body string
status int
want map[string]float64
errorCategory string
}{
{`[{"name":"memory","legend":["free"],"data":[[1789000060,8]]}]`, 200, map[string]float64{"memory": 50}, ""},
{`[{"name":"cpu","legend":["usage"],"data":[[1789000060,0]]}]`, 200, map[string]float64{"cpu": 0}, ""},
{`provider-private-text`, 401, nil, "authentication"},
{`[{"name":"cpu","legend":["usage"],"data":[[1789000060,0]]},{"name":"memory","legend":["free"],"data":[[1789000060,0]]},{"name":"interface","legend":["received"],"data":[[1789000060,0]]},{"name":"disk","legend":["write"],"data":[[1789000060,0]]}]`, 200, map[string]float64{"cpu": 0, "memory": 100, "netin": 0, "diskwrite": 0}, ""},
}
var current atomic.Int64
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
w.Header().Set("Content-Type", "application/json")
switch r.URL.Path {
case "/api/v2.0/system/info":
_, _ = w.Write([]byte(`{"hostname":"synthetic-core","version":"TrueNAS-13.0-U6.1","cores":8,"physmem":16}`))
case "/api/v2.0/pool":
_, _ = w.Write([]byte(`[{"id":1,"name":"tank","status":"ONLINE","size":100,"allocated":50,"free":50}]`))
case "/api/v2.0/pool/dataset", "/api/v2.0/disk", "/api/v2.0/alert/list":
_, _ = w.Write([]byte(`[]`))
case "/api/v2.0/reporting/get_data":
cycle := cycles[current.Load()]
w.WriteHeader(cycle.status)
_, _ = w.Write([]byte(cycle.body))
default:
http.NotFound(w, r)
}
}))
defer server.Close()
client, err := truenas.NewClient(truenas.ClientConfig{Host: server.URL, APIKey: "synthetic"})
if err != nil {
t.Fatal(err)
}
defer client.Close()
provider := truenas.NewLiveProviderForConnection(&truenas.APIFetcher{Client: client}, "synthetic-connection")
instance := config.TrueNASInstance{ID: "synthetic-connection", Host: server.URL, Enabled: true}
poller := NewTrueNASPoller(nil, 0, nil)
poller.providersByOrg["default"] = map[string]*truenas.Provider{instance.ID: provider}
poller.configsByOrg["default"] = map[string]config.TrueNASInstance{instance.ID: instance}
resourceStore := unifiedresources.NewMonitorAdapter(unifiedresources.NewRegistry(unifiedresources.NewMemoryStore()))
persistent, err := metrics.NewStore(metrics.DefaultConfig(t.TempDir()))
if err != nil {
t.Fatal(err)
}
defer func() { _ = persistent.Close() }()
monitor := &Monitor{metricsHistory: NewMetricsHistory(1024, 24*time.Hour), metricsStore: persistent}
counts := make(map[string]int)
observedAt := time.Now().UTC().Add(-10 * time.Minute).Truncate(time.Second)
for index, cycle := range cycles {
current.Store(int64(index))
poller.ensureConnectionRuntimeStatusLocked("default", instance.ID).nextPollAt = time.Time{}
poller.pollAll(context.Background())
summary := poller.ConnectionSummaries("default", []config.TrueNASInstance{instance})[instance.ID]
if summary.Poll == nil || summary.Poll.LastSuccessAt == nil || summary.Poll.LastError != nil || summary.Poll.ConsecutiveFailures != 0 || summary.Observed == nil || summary.Observed.StoragePools != 1 {
t.Fatalf("cycle %d lost successful inventory health: %+v", index, summary)
}
availability := summary.Observed.Telemetry
if availability == nil || availability.ErrorCategory != cycle.errorCategory {
t.Fatalf("cycle %d lost telemetry result: %+v", index, availability)
}
for key, present := range map[string]bool{"cpu": availability.CPU, "memory": availability.Memory, "netin": availability.NetIn, "netout": availability.NetOut, "diskread": availability.DiskRead, "diskwrite": availability.DiskWrite} {
_, want := cycle.want[key]
if present != want {
t.Errorf("cycle %d diagnostics %s presence=%v, want %v", index, key, present, want)
}
}
serialized, err := json.Marshal(summary)
if err != nil || strings.Contains(string(serialized), "provider-private-text") || strings.Contains(string(serialized), "synthetic\"") {
t.Fatalf("cycle %d secret-bearing diagnostic: %s %v", index, serialized, err)
}
availability.ErrorCategory = "mutated"
if poller.ConnectionSummaries("default", []config.TrueNASInstance{instance})[instance.ID].Observed.Telemetry.ErrorCategory != cycle.errorCategory {
t.Fatal("connection diagnostics share mutable provider state")
}
records := poller.GetCurrentRecordsForOrg("default")
// Space the modeled poll observations apart without a wall-clock sleep;
// the writer/store deliberately coalesce identical second timestamps.
for i := range records {
records[i].Resource.LastSeen = observedAt.Add(time.Duration(index) * time.Minute)
records[i].Resource.UpdatedAt = records[i].Resource.LastSeen
}
resourceStore.PopulateSnapshotAndSupplemental(models.StateSnapshot{}, map[unifiedresources.DataSource][]unifiedresources.IngestRecord{unifiedresources.SourceTrueNAS: records})
var host *unifiedresources.Resource
for _, resource := range resourceStore.GetAll() {
if resource.Type == unifiedresources.ResourceTypeAgent {
host = &resource
}
}
if host == nil || host.Agent == nil || host.Agent.CPUCount != 8 || host.Agent.Memory == nil || host.Agent.Memory.Total != 16 {
t.Fatalf("cycle %d lost host/capacity: %+v", index, host)
}
if host.Agent.Memory.UsageUnavailable != !availability.Memory {
t.Fatalf("cycle %d stale/fabricated memory metadata: %+v", index, host.Agent.Memory)
}
var writes []metrics.WriteMetric
monitor.syncUnifiedAgentMetrics(resourceStore, &writes)
seen := make(map[string]bool)
for _, write := range writes {
want, present := cycle.want[write.MetricType]
if write.MetricType == "disk" {
want, present = 50, true // independently collected pool capacity
}
if !present || write.Value != want || write.ResourceID != instance.ID || write.ResourceType != "agent" || seen[write.MetricType] {
t.Errorf("cycle %d wrong/duplicate/missing-source write: %+v", index, write)
}
seen[write.MetricType] = true
counts[write.MetricType]++
}
if len(writes) != len(cycle.want)+1 {
t.Fatalf("cycle %d writes=%+v, want only observed metrics plus pool usage", index, writes)
}
persistent.WriteBatchSync(writes)
inMemory := monitor.GetGuestMetrics("agent:"+instance.ID, time.Hour)
chart := monitor.GetGuestMetricsForChart("agent:"+instance.ID, "agent", instance.ID, time.Hour)
for _, key := range []string{"cpu", "memory", "disk", "netin", "netout", "diskread", "diskwrite"} {
if len(inMemory[key]) != counts[key] || len(chart[key]) != counts[key] {
t.Errorf("cycle %d %s chart contains fabricated/lost points: memory=%d chart=%d want=%d", index, key, len(inMemory[key]), len(chart[key]), counts[key])
}
}
id, native, err := provider.SystemMetricHistory(context.Background(), time.Hour)
if cycle.status != 200 {
if err == nil {
t.Fatal("failed native History read was disguised as success")
}
continue
}
if err != nil || id != instance.ID || len(native) != len(cycle.want) {
t.Fatalf("cycle %d native History=%+v id=%q err=%v", index, native, id, err)
}
for key, want := range cycle.want {
points := native[key]
if len(points) != 1 || points[0].Value != want || points[0].Timestamp.Unix() != 1789000060 {
t.Errorf("cycle %d %s native data/timestamp altered: %+v", index, key, points)
}
}
}
}

View file

@ -232,6 +232,7 @@ func systemInfoFromResponse(response systemInfoResponse) *SystemInfo {
MachineID: machineID,
CPUCount: cpuCount,
MemoryTotalBytes: response.Physmem,
Telemetry: &SystemTelemetryAvailability{},
}
}
@ -289,42 +290,48 @@ func (c *Client) getSystemTelemetryREST(ctx context.Context) (*SystemInfo, error
}
history := parseSystemMetricHistory(response)
if history == nil {
return nil, fmt.Errorf("truenas legacy REST reporting returned no system telemetry")
return nil, errNoSystemTelemetry
}
return systemInfoFromMetricHistory(history), nil
}
// systemInfoFromMetricHistory maps the latest reporting sample onto live system
// telemetry. Series that the appliance did not report stay zero so callers can
// distinguish "absent" from a genuine zero.
// telemetry. Numeric fields alone cannot distinguish absent from zero; retain
// presence independently so one successful graph cannot fabricate its siblings.
func systemInfoFromMetricHistory(history *SystemMetricHistory) *SystemInfo {
if history == nil {
return nil
}
system := &SystemInfo{CollectedAt: time.Now().UTC()}
system := &SystemInfo{CollectedAt: time.Now().UTC(), Telemetry: &SystemTelemetryAvailability{}}
if value, ok := latestTimeSeriesValue(history.CPUPercent); ok {
system.CPUPercent = value
system.Telemetry.CPU = true
}
if value, ok := latestTimeSeriesValue(history.MemoryTotalBytes); ok {
system.MemoryTotalBytes = int64(value)
}
if value, ok := latestTimeSeriesValue(history.MemoryAvailableBytes); ok {
system.MemoryAvailableBytes = int64(value)
system.Telemetry.Memory = value >= 0
}
if value, ok := latestTimeSeriesValue(history.ARCSizeBytes); ok {
system.ARCSizeBytes = int64(value)
}
if value, ok := latestTimeSeriesValue(history.NetInRate); ok {
system.NetInRate = value
system.Telemetry.NetIn = true
}
if value, ok := latestTimeSeriesValue(history.NetOutRate); ok {
system.NetOutRate = value
system.Telemetry.NetOut = true
}
if value, ok := latestTimeSeriesValue(history.DiskReadRate); ok {
system.DiskReadRate = value
system.Telemetry.DiskRead = true
}
if value, ok := latestTimeSeriesValue(history.DiskWriteRate); ok {
system.DiskWriteRate = value
system.Telemetry.DiskWrite = true
}
return system
}
@ -2036,7 +2043,11 @@ func (c *Client) FetchSnapshot(ctx context.Context) (*FixtureSnapshot, error) {
if err != nil {
return nil, fmt.Errorf("fetch truenas system info: %w", err)
}
if telemetry, err := c.GetSystemTelemetry(ctx); err == nil && telemetry != nil {
if telemetry, err := c.GetSystemTelemetry(ctx); err != nil {
// Inventory remains usable. Retain only a fixed error category/status,
// never the provider's body, endpoint, credentials or raw error text.
system.Telemetry = systemTelemetryFailure(err)
} else if telemetry != nil {
mergeSystemTelemetry(system, telemetry)
}
@ -2101,7 +2112,7 @@ func mergeSystemTelemetry(system *SystemInfo, telemetry *SystemInfo) {
if telemetry.MemoryTotalBytes > 0 {
system.MemoryTotalBytes = telemetry.MemoryTotalBytes
}
if telemetry.MemoryAvailableBytes > 0 {
if telemetry.MemoryAvailableBytes > 0 || (telemetry.Telemetry != nil && telemetry.Telemetry.Memory) {
system.MemoryAvailableBytes = telemetry.MemoryAvailableBytes
}
if telemetry.ARCSizeBytes > 0 {
@ -2121,6 +2132,10 @@ func mergeSystemTelemetry(system *SystemInfo, telemetry *SystemInfo) {
if !telemetry.CollectedAt.IsZero() {
system.CollectedAt = telemetry.CollectedAt
}
if telemetry.Telemetry != nil {
availability := *telemetry.Telemetry
system.Telemetry = &availability
}
}
func temperatureForTrueNASDisk(temperatures map[string]int, item diskResponse) int {
@ -3084,20 +3099,14 @@ func parseRealtimeFields(message trueNASRPCResponse, collectionPrefix string) (m
}
func parseSystemTelemetry(fields map[string]any, intervalSeconds int, collectedAt time.Time) *SystemInfo {
if len(fields) == 0 {
return &SystemInfo{
IntervalSeconds: intervalSeconds,
CollectedAt: collectedAt,
}
}
telemetry := &SystemInfo{
IntervalSeconds: intervalSeconds,
CollectedAt: collectedAt,
Telemetry: &SystemTelemetryAvailability{},
}
cpu := readMapAny(fields, "cpu")
cpuPercent := readFloatAny(cpu,
cpuPercent, hasCPU := readFloatValueAny(cpu,
"usage",
"percent",
"usage_percent",
@ -3106,18 +3115,21 @@ func parseSystemTelemetry(fields map[string]any, intervalSeconds int, collectedA
"total",
"overall",
)
if cpuPercent == 0 {
if !hasCPU {
if usage := readMapAny(cpu, "usage", "total"); len(usage) > 0 {
cpuPercent = readFloatAny(usage, "percent", "value", "usage")
cpuPercent, hasCPU = readFloatValueAny(usage, "percent", "value", "usage")
}
}
telemetry.CPUPercent = cpuPercent
telemetry.Telemetry.CPU = hasCPU
memory := readMapAny(fields, "memory")
total := readInt64Any(memory, "physical_memory_total", "total", "memory_total", "total_bytes")
available := readInt64Any(memory, "physical_memory_available", "available", "free", "available_bytes", "free_bytes")
_, hasMemory := readFloatValueAny(memory, "physical_memory_available", "available", "free", "available_bytes", "free_bytes")
telemetry.MemoryTotalBytes = total
telemetry.MemoryAvailableBytes = available
telemetry.Telemetry.Memory = hasMemory && available >= 0
telemetry.ARCSizeBytes = readInt64Any(memory, "arc_size")
interfaces := readMapAny(fields, "interfaces")
@ -3126,7 +3138,7 @@ func parseSystemTelemetry(fields map[string]any, intervalSeconds int, collectedA
if !ok {
continue
}
telemetry.NetInRate += readFloatAny(record,
in, hasIn := readFloatValueAny(record,
"rx_bytes",
"received_bytes",
"received_bytes_rate",
@ -3134,7 +3146,7 @@ func parseSystemTelemetry(fields map[string]any, intervalSeconds int, collectedA
"bytes_recv",
"bytes_received",
)
telemetry.NetOutRate += readFloatAny(record,
out, hasOut := readFloatValueAny(record,
"tx_bytes",
"sent_bytes",
"sent_bytes_rate",
@ -3142,6 +3154,10 @@ func parseSystemTelemetry(fields map[string]any, intervalSeconds int, collectedA
"bytes_sent",
"bytes_transmitted",
)
telemetry.NetInRate += in
telemetry.NetOutRate += out
telemetry.Telemetry.NetIn = telemetry.Telemetry.NetIn || hasIn
telemetry.Telemetry.NetOut = telemetry.Telemetry.NetOut || hasOut
}
disks := readMapAny(fields, "disks", "disls")
@ -3150,8 +3166,12 @@ func parseSystemTelemetry(fields map[string]any, intervalSeconds int, collectedA
if !ok {
continue
}
telemetry.DiskReadRate += readFloatAny(record, "read_bytes", "read_bytes_rate", "bytes_read")
telemetry.DiskWriteRate += readFloatAny(record, "write_bytes", "write_bytes_rate", "bytes_written")
read, hasRead := readFloatValueAny(record, "read_bytes", "read_bytes_rate", "bytes_read")
write, hasWrite := readFloatValueAny(record, "write_bytes", "write_bytes_rate", "bytes_written")
telemetry.DiskReadRate += read
telemetry.DiskWriteRate += write
telemetry.Telemetry.DiskRead = telemetry.Telemetry.DiskRead || hasRead
telemetry.Telemetry.DiskWrite = telemetry.Telemetry.DiskWrite || hasWrite
}
return telemetry
@ -3634,11 +3654,11 @@ func extractReportingLegendFloatValues(raw any, legends []string) map[string]flo
}
}
if len(values) == 0 && len(legends) == 1 {
for _, value := range typed {
if parsed, ok := parseFloat64Any(value); ok {
values[legends[0]] = parsed
break
}
// A missing/null legend value must not fall back to an unrelated
// numeric field (notably the row's timestamp). Keep only supported
// generic single-series value aliases.
if parsed, ok := readFloatValueAny(typed, "value", "y", "temperature"); ok {
values[legends[0]] = parsed
}
}
}

View file

@ -15,6 +15,7 @@ import (
"time"
"github.com/gorilla/websocket"
"github.com/rcourtman/pulse-go-rewrite/internal/unifiedresources"
)
type apiResponse struct {
@ -2344,3 +2345,201 @@ func TestRESTReportingRequestMatchesNativeGraphShape(t *testing.T) {
}
}
}
// These are synthetic response controls, not a capture from #2077's appliance.
// Pin the complete snapshot -> canonical row -> native History path so an
// optional graph cannot turn a missing measurement into a zero-valued series.
func TestRESTReportingObservationPresence(t *testing.T) {
previous := IsFeatureEnabled()
SetFeatureEnabled(true)
t.Cleanup(func() { SetFeatureEnabled(previous) })
for _, tc := range []struct {
name string
body string
want map[string]float64
}{
{"memory_only", `[{"name":"memory","legend":["free"],"data":[[1789000060,8]]}]`, map[string]float64{"memory": 50}},
{"arc_only", `[{"name":"arcsize","legend":["size"],"data":[[1789000060,8]]}]`, nil},
{"capacity_only", `[{"name":"memory","legend":["total"],"data":[[1789000060,16]]}]`, nil},
{"cpu_zero", `[{"name":"cpu","legend":["idle"],"data":[[1789000060,100]]}]`, map[string]float64{"cpu": 0}},
{"free_zero", `[{"name":"memory","legend":["free"],"data":[[1789000060,0]]}]`, map[string]float64{"memory": 100}},
{"network_in_zero", `[{"name":"interface","legend":["received"],"data":[[1789000060,0]]}]`, map[string]float64{"netin": 0}},
{"disk_write_zero", `[{"name":"disk","legend":["write"],"data":[[1789000060,0]]}]`, map[string]float64{"diskwrite": 0}},
{"null_array", `[{"name":"cpu","legend":["usage"],"data":[[1789000060,null]]}]`, nil},
{"null_object", `[{"name":"cpu","legend":["usage"],"data":[{"timestamp":1789000060,"usage":null}]}]`, nil},
{"malformed_object", `[{"name":"cpu","legend":["usage"],"data":[{"timestamp":1789000060,"usage":"missing","unrelated":10}]}]`, nil},
{"generic_zero", `[{"name":"cpu","legend":["usage"],"data":[{"timestamp":1789000060,"value":0}]}]`, map[string]float64{"cpu": 0}},
{"empty", `[]`, nil},
} {
t.Run(tc.name, func(t *testing.T) {
routes := alertArgsTransport(defaultAPIResponses())
routes["/api/v2.0/system/info"] = apiResponse{body: `{"hostname":"synthetic-core","version":"TrueNAS-13.0-U6.1","physmem":16,"cores":8}`}
routes["/api/v2.0/reporting/get_data"] = apiResponse{body: tc.body}
client := newLegacyRESTReportingClient(t, routes)
defer client.Close()
provider := NewLiveProviderForConnection(&APIFetcher{Client: client}, "synthetic-connection")
for cycle := 0; cycle < 2; cycle++ {
if err := provider.Refresh(context.Background()); err != nil {
t.Fatal(err)
}
snapshot := provider.Snapshot()
if snapshot.System.CPUCount != 8 || snapshot.System.MemoryTotalBytes != 16 || len(snapshot.Pools) == 0 || len(snapshot.Disks) == 0 {
t.Fatal("partial/unavailable telemetry must not discard hardware or inventory")
}
var metrics *unifiedresources.ResourceMetrics
for _, record := range provider.Records() {
if record.Resource.Type == unifiedresources.ResourceTypeAgent {
metrics = record.Resource.Metrics
}
}
if metrics == nil {
t.Fatal("missing canonical host row")
}
for key, metric := range map[string]*unifiedresources.MetricValue{
"cpu": metrics.CPU, "memory": metrics.Memory, "netin": metrics.NetIn,
"netout": metrics.NetOut, "diskread": metrics.DiskRead, "diskwrite": metrics.DiskWrite,
} {
want, present := tc.want[key]
if (metric != nil) != present || (metric != nil && metric.Value != want) {
t.Errorf("cycle %d %s: metric=%+v, want present=%v value=%v", cycle, key, metric, present, want)
}
}
id, history, err := provider.SystemMetricHistory(context.Background(), time.Hour)
if err != nil || id != "synthetic-connection" {
t.Fatalf("native History target = %q, err=%v", id, err)
}
for _, key := range []string{"cpu", "memory", "netin", "netout", "diskread", "diskwrite"} {
want, present := tc.want[key]
points := history[key]
if (len(points) > 0) != present || (len(points) > 0 && (len(points) != 1 || points[0].Value != want || points[0].Timestamp.Unix() != 1789000060)) {
t.Errorf("cycle %d %s: history=%+v, want present=%v value=%v at original timestamp", cycle, key, points, present, want)
}
}
}
})
}
}
func TestRESTSnapshotRetainsSanitizedTelemetryFailure(t *testing.T) {
for _, tc := range []struct {
name string
status int
body, category string
}{
{"no_samples", 200, `[]`, "no_samples"},
{"bad_response", 200, `{`, "invalid_response"},
{"unauthorized", 401, `provider-private-text`, "authentication"},
{"forbidden", 403, `provider-private-text`, "authentication"},
{"missing_endpoint", 404, `provider-private-text`, "unsupported_endpoint"},
{"rejected_graphs", 422, `provider-private-text`, "request_rejected"},
{"rate_limited", 429, `provider-private-text`, "rate_limited"},
{"unavailable", 503, `provider-private-text`, "http_error"},
} {
t.Run(tc.name, func(t *testing.T) {
routes := alertArgsTransport(defaultAPIResponses())
routes["/api/v2.0/system/info"] = apiResponse{body: `{"hostname":"synthetic-core","version":"TrueNAS-13.0-U6.1","physmem":16}`}
transport := &isolatedReportingTransport{routes: routes, status: tc.status, body: tc.body}
client := newLegacyRESTReportingClient(t, routes)
client.httpClient.Transport = transport
defer client.Close()
provider := NewLiveProvider(&APIFetcher{Client: client})
if err := provider.Refresh(context.Background()); err != nil {
t.Fatal(err)
}
snapshot := provider.Snapshot()
availability := snapshot.System.Telemetry
if availability == nil || availability.ErrorCategory != tc.category {
t.Fatalf("missing optional telemetry failure: %+v", availability)
}
wantStatus := tc.status
if tc.status == 200 {
wantStatus = 0
}
if availability.HTTPStatus != wantStatus || len(snapshot.Pools) == 0 || len(snapshot.Disks) == 0 {
t.Fatalf("unexpected status/inventory: %+v", snapshot)
}
serialized, err := json.Marshal(availability)
if err != nil || strings.Contains(string(serialized), "provider-private-text") || strings.Contains(string(serialized), "truenas.invalid") || strings.Contains(string(serialized), "synthetic") {
t.Fatalf("diagnostics exposed provider text: %s, err=%v", serialized, err)
}
wantCalls := 1
if tc.status == 422 {
wantCalls = 6
}
if transport.calls != wantCalls {
t.Fatalf("reporting calls=%d, want %d", transport.calls, wantCalls)
}
availability.ErrorCategory = "mutated"
if provider.Snapshot().System.Telemetry.ErrorCategory != tc.category {
t.Fatal("returned snapshot shares mutable telemetry diagnostics")
}
transport.status, transport.body = 200, `[{"name":"cpu","legend":["usage"],"data":[[1789000060,0]]}]`
if err := provider.Refresh(context.Background()); err != nil {
t.Fatal(err)
}
recovered := provider.Snapshot().System.Telemetry
if recovered == nil || !recovered.CPU || recovered.Memory || recovered.NetIn || recovered.NetOut || recovered.DiskRead || recovered.DiskWrite || recovered.ErrorCategory != "" || recovered.HTTPStatus != 0 {
t.Fatalf("recovery retained failure or invented readings: %+v", recovered)
}
})
}
}
func TestRealtimeObservationPresence(t *testing.T) {
for _, tc := range []struct {
name string
fields map[string]any
want SystemTelemetryAvailability
}{
{"empty", nil, SystemTelemetryAvailability{}},
{"cpu_zero", map[string]any{"cpu": map[string]any{"usage": 0}}, SystemTelemetryAvailability{CPU: true}},
{"cpu_nested_zero", map[string]any{"cpu": map[string]any{"usage": map[string]any{"percent": 0}}}, SystemTelemetryAvailability{CPU: true}},
{"cpu_null", map[string]any{"cpu": map[string]any{"usage": nil}}, SystemTelemetryAvailability{}},
{"capacity_only", map[string]any{"memory": map[string]any{"total": 16}}, SystemTelemetryAvailability{}},
{"free_zero", map[string]any{"memory": map[string]any{"total": 16, "available": 0}}, SystemTelemetryAvailability{Memory: true}},
{"free_null", map[string]any{"memory": map[string]any{"total": 16, "available": nil}}, SystemTelemetryAvailability{}},
{"network_in_zero", map[string]any{"interfaces": map[string]any{"eth0": map[string]any{"rx_bytes": 0}}}, SystemTelemetryAvailability{NetIn: true}},
{"disk_write_zero", map[string]any{"disks": map[string]any{"sda": map[string]any{"write_bytes": 0}}}, SystemTelemetryAvailability{DiskWrite: true}},
} {
t.Run(tc.name, func(t *testing.T) {
telemetry := parseSystemTelemetry(tc.fields, 2, time.Now().UTC())
if telemetry.Telemetry == nil || *telemetry.Telemetry != tc.want {
t.Fatalf("presence=%+v, want %+v", telemetry.Telemetry, tc.want)
}
system := systemInfoFromResponse(systemInfoResponse{Physmem: 16})
mergeSystemTelemetry(system, telemetry)
metrics := metricsFromTrueNASSystem(*system, 0, 0)
if (metrics.CPU != nil) != tc.want.CPU || (metrics.Memory != nil) != tc.want.Memory ||
(metrics.NetIn != nil) != tc.want.NetIn || (metrics.NetOut != nil) != tc.want.NetOut ||
(metrics.DiskRead != nil) != tc.want.DiskRead || (metrics.DiskWrite != nil) != tc.want.DiskWrite {
t.Fatalf("timestamp/interval fabricated missing siblings: %+v", metrics)
}
telemetry.Telemetry.CPU = !tc.want.CPU
if system.Telemetry.CPU != tc.want.CPU {
t.Fatal("merge retained caller-owned presence pointer")
}
})
}
}
func TestSystemTelemetryFailureCategoriesAreBounded(t *testing.T) {
for _, tc := range []struct {
err error
want string
}{
{fmt.Errorf("wrapped: %w", context.Canceled), "cancelled"},
{fmt.Errorf("wrapped: %w", context.DeadlineExceeded), "timeout"},
{&RPCAuthError{Mechanism: "provider-private-text"}, "authentication"},
{&RPCError{Message: "provider-private-text", Reason: "provider-private-text"}, "method_error"},
{fmt.Errorf("provider-private-text"), "collection_error"},
} {
failure := systemTelemetryFailure(tc.err)
if failure.ErrorCategory != tc.want || failure.HTTPStatus != 0 || failure.CPU || failure.Memory {
t.Fatalf("failure=%+v, want %s", failure, tc.want)
}
body, err := json.Marshal(failure)
if err != nil || strings.Contains(string(body), "provider-private-text") {
t.Fatalf("unbounded failure detail: %s %v", body, err)
}
}
}

View file

@ -918,10 +918,10 @@ func trueNASReclaimableARC(system SystemInfo) int64 {
// trueNASMemoryUsageKnown reports whether the appliance supplied an available
// memory reading. Without it the ARC-aware used figure cannot be derived, and
// treating the zero value as a measurement reports total capacity as 100% used
// (#2077, legacy REST transport on TrueNAS CORE).
// an unobserved numeric zero must not report total capacity as 100% used. A
// genuinely observed zero remains valid (#2077, legacy REST on TrueNAS CORE).
func trueNASMemoryUsageKnown(system SystemInfo) bool {
return system.MemoryAvailableBytes > 0
return systemTelemetryAvailability(system).Memory && system.MemoryAvailableBytes >= 0
}
// trueNASEffectiveMemoryUsed treats the ZFS ARC as reclaimable cache, the
@ -938,12 +938,12 @@ func trueNASEffectiveMemoryUsed(system SystemInfo) int64 {
}
func metricsFromTrueNASSystem(system SystemInfo, totalCapacity, totalUsed int64) *unifiedresources.ResourceMetrics {
hasRealtimeTelemetry := !system.CollectedAt.IsZero() || system.IntervalSeconds > 0
availability := systemTelemetryAvailability(system)
metrics := &unifiedresources.ResourceMetrics{
Disk: diskMetric(totalCapacity, totalUsed),
}
if hasRealtimeTelemetry {
if availability.CPU {
metrics.CPU = &unifiedresources.MetricValue{
Value: system.CPUPercent,
Percent: system.CPUPercent,
@ -965,22 +965,28 @@ func metricsFromTrueNASSystem(system SystemInfo, totalCapacity, totalUsed int64)
metrics.Memory = memory
}
if hasRealtimeTelemetry {
if availability.NetIn {
metrics.NetIn = &unifiedresources.MetricValue{
Value: system.NetInRate,
Unit: "bytes/s",
Source: unifiedresources.SourceTrueNAS,
}
}
if availability.NetOut {
metrics.NetOut = &unifiedresources.MetricValue{
Value: system.NetOutRate,
Unit: "bytes/s",
Source: unifiedresources.SourceTrueNAS,
}
}
if availability.DiskRead {
metrics.DiskRead = &unifiedresources.MetricValue{
Value: system.DiskReadRate,
Unit: "bytes/s",
Source: unifiedresources.SourceTrueNAS,
}
}
if availability.DiskWrite {
metrics.DiskWrite = &unifiedresources.MetricValue{
Value: system.DiskWriteRate,
Unit: "bytes/s",
@ -3046,6 +3052,10 @@ func clonePools(pools []Pool) []Pool {
func cloneSystemInfo(system SystemInfo) SystemInfo {
cloned := system
if system.Telemetry != nil {
availability := *system.Telemetry
cloned.Telemetry = &availability
}
if len(system.TemperatureCelsius) > 0 {
cloned.TemperatureCelsius = make(map[string]float64, len(system.TemperatureCelsius))
for key, value := range system.TemperatureCelsius {

View file

@ -0,0 +1,64 @@
package truenas
import (
"context"
"encoding/json"
"errors"
"io"
"net/http"
)
var errNoSystemTelemetry = errors.New("truenas legacy REST reporting returned no system telemetry")
func systemTelemetryAvailability(system SystemInfo) SystemTelemetryAvailability {
if system.Telemetry != nil {
return *system.Telemetry
}
// Compatibility for static snapshots written before per-metric presence.
// GetSystemInfo and both live telemetry parsers always set Telemetry, so
// their collection timestamps never imply that missing samples were zero.
hasTelemetry := !system.CollectedAt.IsZero() || system.IntervalSeconds > 0
return SystemTelemetryAvailability{
CPU: hasTelemetry, Memory: system.MemoryAvailableBytes > 0,
NetIn: hasTelemetry, NetOut: hasTelemetry,
DiskRead: hasTelemetry, DiskWrite: hasTelemetry,
}
}
func systemTelemetryFailure(err error) *SystemTelemetryAvailability {
availability := &SystemTelemetryAvailability{ErrorCategory: "collection_error"}
var apiErr *APIError
var authErr *RPCAuthError
var rpcErr *RPCError
var syntaxErr *json.SyntaxError
var typeErr *json.UnmarshalTypeError
switch {
case errors.Is(err, context.Canceled):
availability.ErrorCategory = "cancelled"
case errors.Is(err, context.DeadlineExceeded):
availability.ErrorCategory = "timeout"
case errors.Is(err, errNoSystemTelemetry):
availability.ErrorCategory = "no_samples"
case errors.As(err, &apiErr):
availability.HTTPStatus = apiErr.StatusCode
switch apiErr.StatusCode {
case http.StatusUnauthorized, http.StatusForbidden:
availability.ErrorCategory = "authentication"
case http.StatusNotFound, http.StatusMethodNotAllowed:
availability.ErrorCategory = "unsupported_endpoint"
case http.StatusBadRequest, http.StatusUnprocessableEntity:
availability.ErrorCategory = "request_rejected"
case http.StatusTooManyRequests:
availability.ErrorCategory = "rate_limited"
default:
availability.ErrorCategory = "http_error"
}
case errors.As(err, &authErr):
availability.ErrorCategory = "authentication"
case errors.As(err, &rpcErr):
availability.ErrorCategory = "method_error"
case errors.As(err, &syntaxErr), errors.As(err, &typeErr), errors.Is(err, io.ErrUnexpectedEOF):
availability.ErrorCategory = "invalid_response"
}
return availability
}

View file

@ -41,6 +41,24 @@ type SystemInfo struct {
TemperatureCelsius map[string]float64
IntervalSeconds int
CollectedAt time.Time
// Telemetry distinguishes an observed zero from an absent measurement.
// Live clients always populate it, including when collection fails. A nil
// value is retained for older static snapshots, not inferred by live reads.
Telemetry *SystemTelemetryAvailability
}
// SystemTelemetryAvailability is a bounded, non-secret observation summary.
// Memory means an available/free reading was reported; capacity alone is not
// usage. Errors describe the collection result, never an appliance diagnosis.
type SystemTelemetryAvailability struct {
CPU bool `json:"cpu"`
Memory bool `json:"memory"`
NetIn bool `json:"netIn"`
NetOut bool `json:"netOut"`
DiskRead bool `json:"diskRead"`
DiskWrite bool `json:"diskWrite"`
ErrorCategory string `json:"errorCategory,omitempty"`
HTTPStatus int `json:"httpStatus,omitempty"`
}
// Pool mirrors the subset of TrueNAS pool fields needed for unified mapping.