mirror of
https://github.com/rcourtman/Pulse.git
synced 2026-10-03 04:38:48 +00:00
Merge reviewed polling write repairs on release/v6.4
Change-source: pulse-maintainer
This commit is contained in:
commit
76ca665470
18 changed files with 557 additions and 6 deletions
|
|
@ -31,7 +31,7 @@ Pulse automatically captures the following events:
|
|||
| `agent_profile_assigned` | Agent profile assignments | Profile `production` assigned to agent |
|
||||
| `agent_profile_unassigned` | Agent profile removals | Profile removed from agent |
|
||||
| `user_roles_updated` | RBAC role assignments changed | Updated roles for user jane: [operator] |
|
||||
| `agent_config_fetch` | Agent configuration retrieval attempts | Agent config fetched successfully |
|
||||
| `agent_config_fetch` | Every failed agent configuration fetch. A successful fetch is recorded on the agent's first delivery after Pulse starts, when its token or delivered configuration changes, and otherwise once a day; agents poll every minute, so repeat polls are not recorded | `agent_id=… token_id=… config=sha256:… reason=config_changed` |
|
||||
|
||||
Each event includes:
|
||||
- **Timestamp** (UTC)
|
||||
|
|
|
|||
|
|
@ -15,6 +15,24 @@
|
|||
|
||||
## Purpose
|
||||
|
||||
### Release-line routine polling write repairs — issues #2319/#2320
|
||||
|
||||
A linked Agent-only SMART disk's inherited PVE instance is presentation scope,
|
||||
not a PVE inventory observation. Skipped PVE polls require an actual Proxmox
|
||||
source before copying a disk into PVE-owned state. Agent source admission,
|
||||
identity, collected readings and command authority remain unchanged.
|
||||
`TestPhysicalDiskSkippedPollDoesNotPromoteAgentOnlySMARTToPVEInventory` pins
|
||||
repeated full/skip cycles; genuine PVE readback retains its existing continuity.
|
||||
|
||||
Routine successful agent config fetches are audited on the first delivery
|
||||
since startup, a token or desired-config hash change, or after 24 hours since
|
||||
the previous audit. Every failed fetch remains audited. Signing, payloads,
|
||||
scopes, response delivery and existing security rows are unchanged.
|
||||
`TestAgentConfigFetchAuditsNewDeliveriesAndEveryFailure` and the tracker
|
||||
regressions in `internal/api/unified_agent_handlers_test.go` cover the rule.
|
||||
These are synthetic source controls, not installed field acceptance.
|
||||
|
||||
|
||||
### MD RAID required members and spares — issue #2369
|
||||
|
||||
Host RAID reports carry optional `requiredDevices`: the configured member count
|
||||
|
|
|
|||
|
|
@ -20,6 +20,23 @@
|
|||
|
||||
## Purpose
|
||||
|
||||
### Successful agent-config audit information — issue #2320
|
||||
|
||||
`/api/agents/agent/{id}/config` still signs, scopes and delivers the same config.
|
||||
Every failed fetch is audited. A successful fetch is new audit information
|
||||
only on first delivery after startup, a different token or desired-config hash,
|
||||
or 24 hours after the last audited success. Recorded successes include
|
||||
`config=<desiredConfig hash>` and
|
||||
`reason=first_since_start|token_changed|config_changed|daily`. Tracking is
|
||||
isolated by organization and agent; a nil tracker audits every success.
|
||||
Existing security log rows age under retention, never deletion by this repair.
|
||||
`TestAgentConfigFetchAuditsNewDeliveriesAndEveryFailure` and
|
||||
`TestAgentConfigFetchAuditTrackerRecordsOnlyNewDeliveries` verify runtime
|
||||
scope failures and the per-agent/organization/token/config/daily decisions.
|
||||
The additional day/restart/concurrency control compares 1,440 unchanged polls
|
||||
with one audit event. This is source proof, not a natural installed day.
|
||||
|
||||
|
||||
### Ollama credential lifecycle
|
||||
Settings return only `ollama_username` and `ollama_password_set`, never the
|
||||
password. An omitted password preserves it, a supplied password replaces it
|
||||
|
|
|
|||
|
|
@ -17,6 +17,41 @@
|
|||
|
||||
## Purpose
|
||||
|
||||
### No-op identity lists do not create change history — issue #2319 companion
|
||||
|
||||
Resource change emission compares hostname/IP/MAC lists order-independently,
|
||||
including equivalent nil/empty lists, while scalar machine/DMI/cluster/guest
|
||||
identifiers and actual list-member changes remain exact change evidence.
|
||||
This carries reviewed main `a5aa881ac530` through this line's existing
|
||||
`recordRegistryChanges` List boundary; it does not import main's newer locked-
|
||||
generation/clone optimization or change canonical matching, wire fields,
|
||||
source authority, collection timing or existing persisted rows.
|
||||
`TestResourceChangeIdentityIgnoresSetOrderAndEmptySlices` and
|
||||
`TestRegistryListComparisonTreatsIdentityListsAsSets` pin no-op order/nil cases
|
||||
and one real new-address row. The serial-bearing SMART guard-removal fixture
|
||||
then reports one tags-only row; final guarded full/skip cycles add no rows.
|
||||
These are synthetic in-memory journal controls, not installed writes or CPU.
|
||||
|
||||
|
||||
### Agent-only SMART readback provenance — issue #2319
|
||||
|
||||
A linked Agent SMART disk may inherit its PVE node's instance for presentation.
|
||||
That scope must not invent a Proxmox source during a skipped physical-disk poll.
|
||||
`physicalDisksForInstanceFromReadState` admits only actual Proxmox observations;
|
||||
Agent-only disks remain visible through their Agent source. Explicit failed-
|
||||
query Agent fallback still supplies PVE inventory, while permission failures
|
||||
retain prior inventory and genuine PVE disks preserve identity and readings.
|
||||
The v6.4 prerequisite `PhysicalDiskView.SourceStatus` exposes the existing
|
||||
source map by value: explicit entry presence is provenance, not a healthy or
|
||||
freshness verdict. Its nil/Agent-only/Proxmox/mutation view regression pins
|
||||
the boundary before monitoring uses it; no source metadata is synthesized.
|
||||
No historical rows are deleted and no collection interval is changed.
|
||||
`TestPhysicalDiskSkippedPollDoesNotPromoteAgentOnlySMARTToPVEInventory` pins
|
||||
three empty-inventory/full-skip cycles with no new journal rows, and the
|
||||
existing source-identity, SMART-enrichment and failed-query controls retain
|
||||
provider continuity. Synthetic proof is not the reporter's installed row rate.
|
||||
|
||||
|
||||
### RAID count provenance across ingestion — issue #2369
|
||||
|
||||
Agent ingestion, host snapshots and canonical read-state projection retain
|
||||
|
|
|
|||
|
|
@ -21,6 +21,22 @@
|
|||
|
||||
## Purpose
|
||||
|
||||
### Polling-write repair preserves recovery evidence — issues #2319/#2320
|
||||
|
||||
Routine successful agent config polling no longer creates one security-audit
|
||||
row per minute: first delivery after startup, token/config changes and daily
|
||||
access remain recorded, as does every failure. Recovery actions keep their
|
||||
existing action audits, receipts and verification. No audit or change-history
|
||||
rows are removed. Skipped PVE disk readback requires an actual Proxmox source;
|
||||
Agent-only SMART presentation scope is not source-owned inventory or recovery
|
||||
state. Storage targets, samples, health, restore and command authority are
|
||||
unchanged. `internal/api/host_agent_removal_lifecycle_integration_test.go`
|
||||
verifies config delivery/failure auditing and
|
||||
`internal/monitoring/physical_disk_roundtrip_test.go` verifies unchanged disk
|
||||
identity and no fabricated journal changes. Native/installed acceptance is
|
||||
separate from these synthetic controls.
|
||||
|
||||
|
||||
### Independent physical disk temperature
|
||||
|
||||
Disk Overview displays a finite positive reported temperature even when the
|
||||
|
|
|
|||
|
|
@ -15,6 +15,35 @@
|
|||
|
||||
## Purpose
|
||||
|
||||
### No-op identity lists do not create change history — issue #2319 companion
|
||||
|
||||
Resource change emission compares hostname/IP/MAC lists order-independently,
|
||||
including equivalent nil/empty lists, while scalar machine/DMI/cluster/guest
|
||||
identifiers and actual list-member changes remain exact change evidence.
|
||||
This carries reviewed main `a5aa881ac530` through this line's existing
|
||||
`recordRegistryChanges` List boundary; it does not import main's newer locked-
|
||||
generation/clone optimization or change canonical matching, wire fields,
|
||||
source authority, collection timing or existing persisted rows.
|
||||
`TestResourceChangeIdentityIgnoresSetOrderAndEmptySlices` and
|
||||
`TestRegistryListComparisonTreatsIdentityListsAsSets` pin no-op order/nil cases
|
||||
and one real new-address row. The serial-bearing SMART guard-removal fixture
|
||||
then reports one tags-only row; final guarded full/skip cycles add no rows.
|
||||
These are synthetic in-memory journal controls, not installed writes or CPU.
|
||||
|
||||
|
||||
### Physical disk observation provenance on v6.4 — issue #2319
|
||||
|
||||
The read-only `PhysicalDiskView.SourceStatus` accessor returns the existing
|
||||
per-source observation and its actual presence, including an explicit zero
|
||||
status; nil or absent source metadata stays absent. A disk's inherited PVE
|
||||
instance is only presentation scope and cannot manufacture a Proxmox source.
|
||||
This small prerequisite of main's Agent-only SMART readback repair changes no
|
||||
resource wire shape, ID, collector, schedule or freshness policy. Status is
|
||||
returned by value, not as mutable source state. `TestView_PhysicalDiskSourceStatusRequiresAnObservation`
|
||||
in `internal/unifiedresources/views_test.go` covers Agent-only, explicit
|
||||
Proxmox, nil and caller-mutation boundaries; monitoring's repeated full/skip
|
||||
roundtrip test verifies the consuming no-churn path.
|
||||
|
||||
### Canonical RAID configured-member evidence — issue #2369
|
||||
|
||||
Host RAID metadata and read views retain optional `requiredDevices` (configured
|
||||
|
|
|
|||
|
|
@ -31,7 +31,7 @@ Pulse automatically captures the following events:
|
|||
| `agent_profile_assigned` | Agent profile assignments | Profile `production` assigned to agent |
|
||||
| `agent_profile_unassigned` | Agent profile removals | Profile removed from agent |
|
||||
| `user_roles_updated` | RBAC role assignments changed | Updated roles for user jane: [operator] |
|
||||
| `agent_config_fetch` | Agent configuration retrieval attempts | Agent config fetched successfully |
|
||||
| `agent_config_fetch` | Every failed agent configuration fetch. A successful fetch is recorded on the agent's first delivery after Pulse starts, when its token or delivered configuration changes, and otherwise once a day; agents poll every minute, so repeat polls are not recorded | `agent_id=… token_id=… config=sha256:… reason=config_changed` |
|
||||
|
||||
Each event includes:
|
||||
- **Timestamp** (UTC)
|
||||
|
|
|
|||
80
internal/api/agent_config_fetch_audit.go
Normal file
80
internal/api/agent_config_fetch_audit.go
Normal file
|
|
@ -0,0 +1,80 @@
|
|||
package api
|
||||
|
||||
import (
|
||||
"sync"
|
||||
"time"
|
||||
)
|
||||
|
||||
// agentConfigFetchAuditInterval re-records an unchanged successful config
|
||||
// delivery at most this often per agent, so the audit trail still shows
|
||||
// continued access.
|
||||
const agentConfigFetchAuditInterval = 24 * time.Hour
|
||||
|
||||
// maxAgentConfigFetchAudits bounds remembered deliveries. Beyond it, entries
|
||||
// older than agentConfigFetchAuditInterval are dropped; they would be
|
||||
// re-recorded on their next fetch anyway.
|
||||
const maxAgentConfigFetchAudits = 4096
|
||||
|
||||
// agentConfigFetchAuditTracker decides which successful agent config fetches
|
||||
// are audit events. Agents poll their config every minute, and recording each
|
||||
// poll made agent_config_fetch over 99.9% of audit rows (about 1,440 a day per
|
||||
// agent), burying the logins and failures the log exists to show. Failed
|
||||
// fetches are always audited and never pass through here.
|
||||
type agentConfigFetchAuditTracker struct {
|
||||
mu sync.Mutex
|
||||
last map[agentConfigFetchAuditKey]agentConfigFetchAuditEntry
|
||||
}
|
||||
|
||||
type agentConfigFetchAuditKey struct {
|
||||
orgID string
|
||||
agentID string
|
||||
}
|
||||
|
||||
type agentConfigFetchAuditEntry struct {
|
||||
tokenID string
|
||||
configHash string
|
||||
auditedAt time.Time
|
||||
}
|
||||
|
||||
func newAgentConfigFetchAuditTracker() *agentConfigFetchAuditTracker {
|
||||
return &agentConfigFetchAuditTracker{last: make(map[agentConfigFetchAuditKey]agentConfigFetchAuditEntry)}
|
||||
}
|
||||
|
||||
// observe records a successful delivery and reports whether it is new audit
|
||||
// information, with the reason: the agent's first delivery since startup, a
|
||||
// different token or delivered config, or an unchanged delivery last audited
|
||||
// at least agentConfigFetchAuditInterval ago. A nil tracker audits every
|
||||
// delivery.
|
||||
func (t *agentConfigFetchAuditTracker) observe(orgID, agentID, tokenID, configHash string, now time.Time) (string, bool) {
|
||||
if t == nil {
|
||||
return "", true
|
||||
}
|
||||
key := agentConfigFetchAuditKey{orgID: orgID, agentID: agentID}
|
||||
entry := agentConfigFetchAuditEntry{tokenID: tokenID, configHash: configHash, auditedAt: now}
|
||||
|
||||
t.mu.Lock()
|
||||
defer t.mu.Unlock()
|
||||
previous, seen := t.last[key]
|
||||
reason := ""
|
||||
switch {
|
||||
case !seen:
|
||||
reason = "first_since_start"
|
||||
case previous.tokenID != tokenID:
|
||||
reason = "token_changed"
|
||||
case previous.configHash != configHash:
|
||||
reason = "config_changed"
|
||||
case now.Sub(previous.auditedAt) >= agentConfigFetchAuditInterval:
|
||||
reason = "daily"
|
||||
default:
|
||||
return "", false
|
||||
}
|
||||
if !seen && len(t.last) >= maxAgentConfigFetchAudits {
|
||||
for staleKey, stale := range t.last {
|
||||
if now.Sub(stale.auditedAt) >= agentConfigFetchAuditInterval {
|
||||
delete(t.last, staleKey)
|
||||
}
|
||||
}
|
||||
}
|
||||
t.last[key] = entry
|
||||
return reason, true
|
||||
}
|
||||
|
|
@ -43,6 +43,8 @@ type UnifiedAgentHandlers struct {
|
|||
// server upgrade on their next report instead of their next hourly update
|
||||
// check. Empty (the default) omits it from acks.
|
||||
serverVersion string
|
||||
|
||||
configFetchAudits *agentConfigFetchAuditTracker
|
||||
}
|
||||
|
||||
// SetServerVersion supplies the running server version to include on report
|
||||
|
|
@ -63,7 +65,10 @@ func trimUnifiedAgentRoutePath(path string) string {
|
|||
|
||||
// NewUnifiedAgentHandlers constructs a new handler set for Pulse Unified Agent ingest.
|
||||
func NewUnifiedAgentHandlers(mtm *monitoring.MultiTenantMonitor, m *monitoring.Monitor, hub *websocket.Hub) *UnifiedAgentHandlers {
|
||||
return &UnifiedAgentHandlers{baseAgentHandlers: newBaseAgentHandlers(mtm, m, hub)}
|
||||
return &UnifiedAgentHandlers{
|
||||
baseAgentHandlers: newBaseAgentHandlers(mtm, m, hub),
|
||||
configFetchAudits: newAgentConfigFetchAuditTracker(),
|
||||
}
|
||||
}
|
||||
|
||||
// HandleReport ingests Pulse Unified Agent reports.
|
||||
|
|
@ -607,8 +612,18 @@ func (h *UnifiedAgentHandlers) handleGetConfig(w http.ResponseWriter, r *http.Re
|
|||
return
|
||||
}
|
||||
|
||||
LogAuditEventForTenant(GetOrgID(r.Context()), "agent_config_fetch", auth.GetUser(r.Context()), GetClientIP(r), r.URL.Path, true,
|
||||
fmt.Sprintf("agent_id=%s token_id=%s", agentID, tokenID(record)))
|
||||
configHash := ""
|
||||
if signedConfig.DesiredConfig != nil {
|
||||
configHash = signedConfig.DesiredConfig.Hash
|
||||
}
|
||||
orgID := GetOrgID(r.Context())
|
||||
if reason, audit := h.configFetchAudits.observe(orgID, agentID, tokenID(record), configHash, time.Now()); audit {
|
||||
details := fmt.Sprintf("agent_id=%s token_id=%s config=%s", agentID, tokenID(record), configHash)
|
||||
if reason != "" {
|
||||
details += " reason=" + reason
|
||||
}
|
||||
LogAuditEventForTenant(orgID, "agent_config_fetch", auth.GetUser(r.Context()), GetClientIP(r), r.URL.Path, true, details)
|
||||
}
|
||||
}
|
||||
|
||||
func tokenID(record *config.APITokenRecord) string {
|
||||
|
|
|
|||
|
|
@ -16,6 +16,7 @@ import (
|
|||
"github.com/rcourtman/pulse-go-rewrite/internal/models"
|
||||
"github.com/rcourtman/pulse-go-rewrite/internal/monitoring"
|
||||
agentshost "github.com/rcourtman/pulse-go-rewrite/pkg/agents/host"
|
||||
"github.com/rcourtman/pulse-go-rewrite/pkg/audit"
|
||||
)
|
||||
|
||||
type hostRemovalLifecycleHTTPRuntime struct {
|
||||
|
|
@ -647,3 +648,73 @@ func TestHostAgentRenameKeepsCommandGateAlignedWithChannelAdmission(t *testing.T
|
|||
t.Fatal("channel admission admitted a different agent ID on a matching hostname")
|
||||
}
|
||||
}
|
||||
|
||||
// Agents poll their config every minute. Only deliveries that tell an auditor
|
||||
// something new are recorded, while every failed fetch still is.
|
||||
func TestAgentConfigFetchAuditsNewDeliveriesAndEveryFailure(t *testing.T) {
|
||||
capture := &auditCaptureLogger{}
|
||||
prevLogger := audit.GetLogger()
|
||||
prevManager := GetTenantAuditManager()
|
||||
audit.SetLogger(capture)
|
||||
SetTenantAuditManager(nil)
|
||||
t.Cleanup(func() {
|
||||
audit.SetLogger(prevLogger)
|
||||
SetTenantAuditManager(prevManager)
|
||||
})
|
||||
|
||||
handler, monitor := newUnifiedAgentHandlers(t, nil)
|
||||
hostID := seedUnifiedAgentHost(t, monitor)
|
||||
monitorState(t, monitor).UpsertHost(models.Host{ID: hostID, Hostname: "node-1", TokenID: "runtime-token"})
|
||||
fetch := func(scopes ...string) int {
|
||||
t.Helper()
|
||||
req := httptest.NewRequest(http.MethodGet, "/api/agents/agent/"+hostID+"/config", nil)
|
||||
attachAPITokenRecord(req, &config.APITokenRecord{ID: "runtime-token", Scopes: scopes})
|
||||
rec := httptest.NewRecorder()
|
||||
handler.HandleConfig(rec, req)
|
||||
return rec.Code
|
||||
}
|
||||
|
||||
for i := 0; i < 3; i++ {
|
||||
if code := fetch(config.ScopeAgentConfigRead, config.ScopeAgentReport); code != http.StatusOK {
|
||||
t.Fatalf("config fetch %d status = %d, want 200", i, code)
|
||||
}
|
||||
}
|
||||
commandsEnabled := true
|
||||
if err := monitor.UpdateHostAgentConfig(hostID, &commandsEnabled); err != nil {
|
||||
t.Fatalf("UpdateHostAgentConfig: %v", err)
|
||||
}
|
||||
for i := 0; i < 2; i++ {
|
||||
if code := fetch(config.ScopeAgentConfigRead, config.ScopeAgentReport); code != http.StatusOK {
|
||||
t.Fatalf("config fetch after change %d status = %d, want 200", i, code)
|
||||
}
|
||||
}
|
||||
for i := 0; i < 2; i++ {
|
||||
if code := fetch(config.ScopeMonitoringRead); code == http.StatusOK {
|
||||
t.Fatalf("config fetch with only %s succeeded", config.ScopeMonitoringRead)
|
||||
}
|
||||
}
|
||||
|
||||
capture.mu.Lock()
|
||||
events := append([]audit.Event(nil), capture.events...)
|
||||
capture.mu.Unlock()
|
||||
var successes []string
|
||||
failures := 0
|
||||
for _, event := range events {
|
||||
if event.EventType != "agent_config_fetch" {
|
||||
continue
|
||||
}
|
||||
if event.Success {
|
||||
successes = append(successes, event.Details)
|
||||
} else {
|
||||
failures++
|
||||
}
|
||||
}
|
||||
if len(successes) != 2 ||
|
||||
!strings.Contains(successes[0], "reason=first_since_start") ||
|
||||
!strings.Contains(successes[1], "reason=config_changed") {
|
||||
t.Fatalf("successful fetch audit events = %q, want the first delivery and the config change", successes)
|
||||
}
|
||||
if failures != 2 {
|
||||
t.Fatalf("failed fetch audit events = %d, want every failure", failures)
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -7,6 +7,7 @@ import (
|
|||
"net/http"
|
||||
"net/http/httptest"
|
||||
"reflect"
|
||||
"sync"
|
||||
"testing"
|
||||
"time"
|
||||
"unsafe"
|
||||
|
|
@ -552,3 +553,85 @@ func TestUnifiedAgentHandlers_HandleLinkUnlink(t *testing.T) {
|
|||
t.Fatalf("unlink status = %d, want 200: %s", unlinkRec.Code, unlinkRec.Body.String())
|
||||
}
|
||||
}
|
||||
|
||||
func TestAgentConfigFetchAuditTrackerRecordsOnlyNewDeliveries(t *testing.T) {
|
||||
tracker := newAgentConfigFetchAuditTracker()
|
||||
start := time.Date(2026, 9, 24, 12, 0, 0, 0, time.UTC)
|
||||
for _, step := range []struct {
|
||||
name string
|
||||
token string
|
||||
hash string
|
||||
at time.Duration
|
||||
reason string
|
||||
audit bool
|
||||
}{
|
||||
{"first delivery", "tok-a", "sha256:1", 0, "first_since_start", true},
|
||||
{"unchanged poll", "tok-a", "sha256:1", time.Minute, "", false},
|
||||
{"config change", "tok-a", "sha256:2", 2 * time.Minute, "config_changed", true},
|
||||
{"unchanged after change", "tok-a", "sha256:2", 3 * time.Minute, "", false},
|
||||
{"token rotation", "tok-b", "sha256:2", 4 * time.Minute, "token_changed", true},
|
||||
{"just under a day", "tok-b", "sha256:2", 4*time.Minute + agentConfigFetchAuditInterval - time.Second, "", false},
|
||||
{"a day since last audit", "tok-b", "sha256:2", 4*time.Minute + agentConfigFetchAuditInterval, "daily", true},
|
||||
} {
|
||||
reason, audit := tracker.observe("default", "agent-1", step.token, step.hash, start.Add(step.at))
|
||||
if audit != step.audit || reason != step.reason {
|
||||
t.Fatalf("%s: audit=%v reason=%q, want audit=%v reason=%q", step.name, audit, reason, step.audit, step.reason)
|
||||
}
|
||||
}
|
||||
if reason, audit := tracker.observe("other-org", "agent-1", "tok-b", "sha256:2", start.Add(5*time.Minute)); !audit || reason != "first_since_start" {
|
||||
t.Fatalf("same agent ID in another org: audit=%v reason=%q, want its own first delivery", audit, reason)
|
||||
}
|
||||
var unset *agentConfigFetchAuditTracker
|
||||
if _, audit := unset.observe("default", "agent-1", "tok-a", "sha256:1", start); !audit {
|
||||
t.Fatal("a handler without a tracker must audit every delivery")
|
||||
}
|
||||
}
|
||||
|
||||
func TestAgentConfigFetchAuditTrackerDayRestartAndConcurrency(t *testing.T) {
|
||||
tracker := newAgentConfigFetchAuditTracker()
|
||||
start := time.Date(2026, 10, 2, 0, 0, 0, 0, time.UTC)
|
||||
recorded := 0
|
||||
for minute := 0; minute < 1440; minute++ {
|
||||
if _, audit := tracker.observe("org", "agent", "token", "sha256:config", start.Add(time.Duration(minute)*time.Minute)); audit {
|
||||
recorded++
|
||||
}
|
||||
}
|
||||
if recorded != 1 {
|
||||
t.Fatalf("1,440 unchanged minute polls recorded %d audits, want 1", recorded)
|
||||
}
|
||||
if reason, audit := tracker.observe("org", "agent", "token", "sha256:config", start.Add(24*time.Hour)); !audit || reason != "daily" {
|
||||
t.Fatalf("daily access audit = %v %q, want daily", audit, reason)
|
||||
}
|
||||
if reason, audit := tracker.observe("org", "agent", "token", "sha256:config", start.Add(23*time.Hour)); audit || reason != "" {
|
||||
t.Fatalf("a backwards clock recorded unchanged access: %v %q", audit, reason)
|
||||
}
|
||||
restarted := newAgentConfigFetchAuditTracker()
|
||||
if reason, audit := restarted.observe("org", "agent", "token", "sha256:config", start.Add(24*time.Hour)); !audit || reason != "first_since_start" {
|
||||
t.Fatalf("restart audit = %v %q, want first_since_start", audit, reason)
|
||||
}
|
||||
|
||||
// Concurrent deliveries must atomically share the same remembered success.
|
||||
concurrent := newAgentConfigFetchAuditTracker()
|
||||
results := make(chan bool, 32)
|
||||
var workers sync.WaitGroup
|
||||
for i := 0; i < cap(results); i++ {
|
||||
workers.Add(1)
|
||||
go func() {
|
||||
defer workers.Done()
|
||||
_, audit := concurrent.observe("org", "agent", "token", "sha256:config", start)
|
||||
results <- audit
|
||||
}()
|
||||
}
|
||||
workers.Wait()
|
||||
close(results)
|
||||
concurrentAudits := 0
|
||||
for audit := range results {
|
||||
if audit {
|
||||
concurrentAudits++
|
||||
}
|
||||
}
|
||||
if concurrentAudits != 1 {
|
||||
t.Fatalf("32 concurrent identical deliveries recorded %d audits, want 1", concurrentAudits)
|
||||
}
|
||||
t.Logf("unchanged minute polls=1440 recorded=%d; concurrent polls=32 recorded=%d; daily and restart access preserved", recorded, concurrentAudits)
|
||||
}
|
||||
|
|
|
|||
|
|
@ -851,6 +851,13 @@ func physicalDisksForInstanceFromReadState(readState unifiedresources.ReadState,
|
|||
if disk == nil || disk.Instance() != instance {
|
||||
continue
|
||||
}
|
||||
// A linked host Agent's SMART disk inherits the PVE instance for
|
||||
// presentation, but that does not make it a PVE inventory record.
|
||||
// Feeding it back into State.PhysicalDisks invents a PVE source and
|
||||
// alternates tags/history whenever the real disks/list is empty.
|
||||
if _, observedByPVE := disk.SourceStatus(unifiedresources.SourceProxmox); !observedByPVE {
|
||||
continue
|
||||
}
|
||||
out = append(out, physicalDiskFromReadStateView(disk))
|
||||
}
|
||||
return out
|
||||
|
|
|
|||
|
|
@ -178,6 +178,64 @@ func TestPhysicalDiskSkippedPollPreservesSourceIdentity(t *testing.T) {
|
|||
}
|
||||
}
|
||||
|
||||
// An agent disk may inherit its linked PVE node's instance for presentation,
|
||||
// even when the PVE disks/list endpoint did not report that disk. The skipped
|
||||
// poll must not turn that presentation scope into a PVE inventory observation:
|
||||
// an empty PVE inventory on the next full poll would then remove the invented
|
||||
// observation and record spurious configuration changes on every cycle (#2319).
|
||||
func TestPhysicalDiskSkippedPollDoesNotPromoteAgentOnlySMARTToPVEInventory(t *testing.T) {
|
||||
state := models.NewState()
|
||||
now := time.Now().UTC()
|
||||
state.UpdateNodesForInstance("pve", []models.Node{{
|
||||
ID: "pve-node", Name: "node", Instance: "pve", LinkedAgentID: "agent",
|
||||
Status: "online", LastSeen: now,
|
||||
}})
|
||||
state.UpsertHost(models.Host{
|
||||
ID: "agent", Hostname: "node", LinkedNodeID: "pve-node",
|
||||
Status: "online", LastSeen: now,
|
||||
Sensors: models.HostSensorSummary{SMART: []models.HostDiskSMART{{
|
||||
// A stable serial keeps the canonical identity unchanged; without
|
||||
// the source guard, the false PVE observation produces exactly the
|
||||
// tags-only change reported in #2319 rather than a tags+identity row.
|
||||
Device: "sda", Serial: "disk-serial", Type: "sata", Health: "PASSED",
|
||||
}}},
|
||||
})
|
||||
store := unifiedresources.NewMemoryStore()
|
||||
adapter := unifiedresources.NewMonitorAdapter(unifiedresources.NewRegistry(store))
|
||||
monitor := &Monitor{
|
||||
state: state, resourceStore: adapter,
|
||||
lastPhysicalDiskPoll: map[string]time.Time{"pve": now},
|
||||
}
|
||||
for cycle := 0; cycle < 3; cycle++ {
|
||||
// A successful full PVE inventory read found no disks on this node.
|
||||
state.UpdatePhysicalDisks("pve", nil)
|
||||
adapter.PopulateFromSnapshot(state.GetSnapshot())
|
||||
disks := adapter.PhysicalDisks()
|
||||
if len(disks) != 1 || disks[0].Instance() != "pve" {
|
||||
t.Fatalf("cycle %d: linked Agent SMART disk missing from presentation: %+v", cycle, disks)
|
||||
}
|
||||
if _, hasPVE := disks[0].SourceStatus(unifiedresources.SourceProxmox); hasPVE {
|
||||
t.Fatalf("cycle %d: Agent-only disk unexpectedly has PVE source", cycle)
|
||||
}
|
||||
before, err := store.GetRecentChanges(disks[0].ID(), time.Time{}, 100)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
monitor.maybePollPhysicalDisksAsync(context.Background(), "pve", &config.PVEInstance{}, nil, nil, nil, nil)
|
||||
adapter.PopulateFromSnapshot(state.GetSnapshot())
|
||||
after, err := store.GetRecentChanges(disks[0].ID(), time.Time{}, 100)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if len(after) != len(before) {
|
||||
t.Fatalf("cycle %d: unchanged SMART disk emitted %d new history rows: %+v", cycle, len(after)-len(before), after)
|
||||
}
|
||||
if got := state.GetSnapshot().PhysicalDisks; len(got) != 0 {
|
||||
t.Fatalf("cycle %d: Agent-only disk was written into PVE inventory: %+v", cycle, got)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestPhysicalDiskReadbackSourceIDFallback(t *testing.T) {
|
||||
for _, resource := range []unifiedresources.Resource{
|
||||
{ID: "canonical"},
|
||||
|
|
|
|||
|
|
@ -180,7 +180,7 @@ func resourceChangedFields(before, after Resource) []string {
|
|||
if before.CustomURL != after.CustomURL {
|
||||
changed = append(changed, "customUrl")
|
||||
}
|
||||
if !reflect.DeepEqual(before.Identity, after.Identity) {
|
||||
if !resourceIdentityEquivalent(before.Identity, after.Identity) {
|
||||
changed = append(changed, "identity")
|
||||
}
|
||||
if (isAvailabilityOwnedResource(before) || isAvailabilityOwnedResource(after)) &&
|
||||
|
|
@ -489,6 +489,20 @@ func changeSourceAdapterForDataSource(source DataSource) ChangeSourceAdapter {
|
|||
}
|
||||
}
|
||||
|
||||
// Identity addresses, hostnames and MACs are sets, not ordered observations.
|
||||
// Registry merges can reorder them or turn a nil slice into an empty one while
|
||||
// the actual identity remains unchanged; those differences must not create
|
||||
// configuration-history rows on every poll.
|
||||
func resourceIdentityEquivalent(a, b ResourceIdentity) bool {
|
||||
return a.MachineID == b.MachineID &&
|
||||
a.DMIUUID == b.DMIUUID &&
|
||||
a.ClusterName == b.ClusterName &&
|
||||
a.ProxmoxGuestKey == b.ProxmoxGuestKey &&
|
||||
sameStringSet(a.Hostnames, b.Hostnames) &&
|
||||
sameStringSet(a.IPAddresses, b.IPAddresses) &&
|
||||
sameStringSet(a.MACAddresses, b.MACAddresses)
|
||||
}
|
||||
|
||||
func sameStringSet(a, b []string) bool {
|
||||
if len(a) != len(b) {
|
||||
return false
|
||||
|
|
|
|||
|
|
@ -483,6 +483,36 @@ func TestBuildResourceChange_ClassifiesConfigUpdate(t *testing.T) {
|
|||
}
|
||||
}
|
||||
|
||||
func TestResourceChangeIdentityIgnoresSetOrderAndEmptySlices(t *testing.T) {
|
||||
before := Resource{
|
||||
ID: "disk-1", Type: ResourceTypePhysicalDisk, Status: StatusOnline,
|
||||
Identity: ResourceIdentity{
|
||||
MachineID: "serial-1", Hostnames: []string{"node", "node.example"},
|
||||
},
|
||||
}
|
||||
after := before
|
||||
after.Identity = ResourceIdentity{
|
||||
MachineID: "serial-1", Hostnames: []string{"node.example", "node"},
|
||||
IPAddresses: []string{}, MACAddresses: []string{},
|
||||
}
|
||||
if change := buildResourceChange(before, true, after, true, time.Now().UTC(), nil, SourcePulseDiff, ""); change != nil {
|
||||
t.Fatalf("unchanged identity set emitted history: %+v", change)
|
||||
}
|
||||
|
||||
after.Identity.MachineID = "serial-2"
|
||||
change := buildResourceChange(before, true, after, true, time.Now().UTC(), nil, SourcePulseDiff, "")
|
||||
if change == nil || !sameStringSet(mustChangedFields(t, change), []string{"identity"}) {
|
||||
t.Fatalf("real identity change not recorded: %+v", change)
|
||||
}
|
||||
|
||||
after.Identity.MachineID = before.Identity.MachineID
|
||||
after.Identity.IPAddresses = []string{"192.0.2.10"}
|
||||
change = buildResourceChange(before, true, after, true, time.Now().UTC(), nil, SourcePulseDiff, "")
|
||||
if change == nil || !sameStringSet(mustChangedFields(t, change), []string{"identity"}) {
|
||||
t.Fatalf("new identity set member not recorded: %+v", change)
|
||||
}
|
||||
}
|
||||
|
||||
func mustChangedFields(t *testing.T, change *ResourceChange) []string {
|
||||
t.Helper()
|
||||
raw, ok := change.Metadata["changedFields"]
|
||||
|
|
|
|||
|
|
@ -6536,3 +6536,36 @@ func TestIssue2076USBLinkedDiskAmbiguity(t *testing.T) {
|
|||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestRegistryListComparisonTreatsIdentityListsAsSets(t *testing.T) {
|
||||
store := NewMemoryStore()
|
||||
now := time.Date(2026, 9, 29, 17, 0, 0, 0, time.UTC)
|
||||
entry := IngestRecord{
|
||||
SourceID: "disk-1",
|
||||
Resource: Resource{Type: ResourceTypePhysicalDisk, Name: "disk", Status: StatusOnline, LastSeen: now},
|
||||
Identity: ResourceIdentity{MachineID: "serial-1", Hostnames: []string{"node", "node.example"}},
|
||||
}
|
||||
before := NewRegistry(store)
|
||||
before.IngestRecords(SourceAgent, []IngestRecord{entry})
|
||||
|
||||
reordered := entry
|
||||
reordered.Identity = ResourceIdentity{
|
||||
MachineID: "serial-1", Hostnames: []string{"node.example", "node"},
|
||||
IPAddresses: []string{}, MACAddresses: []string{},
|
||||
}
|
||||
after := NewRegistry(store)
|
||||
after.IngestRecords(SourceAgent, []IngestRecord{reordered})
|
||||
recordRegistryChanges(store, before.List(), after.List(), now.Add(time.Minute), nil, SourcePulseDiff, "")
|
||||
if len(store.changes) != 0 {
|
||||
t.Fatalf("equivalent identity lists emitted %d history rows: %+v", len(store.changes), store.changes)
|
||||
}
|
||||
|
||||
changed := reordered
|
||||
changed.Identity.IPAddresses = []string{"192.0.2.10"}
|
||||
later := NewRegistry(store)
|
||||
later.IngestRecords(SourceAgent, []IngestRecord{changed})
|
||||
recordRegistryChanges(store, after.List(), later.List(), now.Add(2*time.Minute), nil, SourcePulseDiff, "")
|
||||
if len(store.changes) != 1 || !sameStringSet(mustChangedFields(t, &store.changes[0]), []string{"identity"}) {
|
||||
t.Fatalf("real address change not recorded once: %+v", store.changes)
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -2158,6 +2158,14 @@ func (v PhysicalDiskView) Status() ResourceStatus {
|
|||
return v.r.Status
|
||||
}
|
||||
|
||||
func (v PhysicalDiskView) SourceStatus(source DataSource) (SourceStatus, bool) {
|
||||
if v.r == nil {
|
||||
return SourceStatus{}, false
|
||||
}
|
||||
status, ok := v.r.SourceStatus[source]
|
||||
return status, ok
|
||||
}
|
||||
|
||||
func (v PhysicalDiskView) DevPath() string {
|
||||
if v.r == nil || v.r.PhysicalDisk == nil {
|
||||
return ""
|
||||
|
|
|
|||
|
|
@ -1000,6 +1000,43 @@ func TestView_PhysicalDiskViewNodeFallsBackToIdentityHostnames(t *testing.T) {
|
|||
}
|
||||
}
|
||||
|
||||
func TestView_PhysicalDiskSourceStatusRequiresAnObservation(t *testing.T) {
|
||||
var nilView PhysicalDiskView
|
||||
if _, ok := nilView.SourceStatus(SourceProxmox); ok {
|
||||
t.Fatal("nil view invented a Proxmox observation")
|
||||
}
|
||||
r := &Resource{
|
||||
Type: ResourceTypePhysicalDisk,
|
||||
Sources: []DataSource{SourceAgent},
|
||||
Proxmox: &ProxmoxData{Instance: "pve"},
|
||||
PhysicalDisk: &PhysicalDiskMeta{},
|
||||
SourceStatus: map[DataSource]SourceStatus{
|
||||
SourceAgent: {Status: "online"},
|
||||
},
|
||||
}
|
||||
v := NewPhysicalDiskView(r)
|
||||
if v.Instance() != "pve" {
|
||||
t.Fatal("fixture did not inherit its presentation instance")
|
||||
}
|
||||
if _, ok := v.SourceStatus(SourceProxmox); ok {
|
||||
t.Fatal("Agent presentation scope invented a Proxmox observation")
|
||||
}
|
||||
if status, ok := v.SourceStatus(SourceAgent); !ok || status.Status != "online" {
|
||||
t.Fatalf("actual Agent observation lost: %+v, %v", status, ok)
|
||||
}
|
||||
// Observation provenance is independent of freshness or appliance health.
|
||||
// Even a zero-valued entry records that source; an absent map entry does not.
|
||||
r.SourceStatus[SourceProxmox] = SourceStatus{}
|
||||
status, ok := v.SourceStatus(SourceProxmox)
|
||||
if !ok {
|
||||
t.Fatal("explicit Proxmox observation was treated as absent")
|
||||
}
|
||||
status.Status = "offline"
|
||||
if r.SourceStatus[SourceProxmox].Status != "" {
|
||||
t.Fatal("caller modified the stored source status through the view")
|
||||
}
|
||||
}
|
||||
|
||||
func TestView_DockerHostViewAccessors(t *testing.T) {
|
||||
now := time.Date(2026, 2, 10, 12, 4, 0, 0, time.UTC)
|
||||
temp := 44.4
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue