diff --git a/docs/qualification/PATROL_ASSISTANT_CUSTOMER_JOURNEY.md b/docs/qualification/PATROL_ASSISTANT_CUSTOMER_JOURNEY.md
index ce85608a2..29c6f571f 100644
--- a/docs/qualification/PATROL_ASSISTANT_CUSTOMER_JOURNEY.md
+++ b/docs/qualification/PATROL_ASSISTANT_CUSTOMER_JOURNEY.md
@@ -4,6 +4,105 @@ The customer job is: "Tell me what needs my attention, explain why, and help me
deal with it without creating more work." Patrol owns the issue and investigation.
Assistant explains that same issue and uses existing governed action contracts.
+## Active redesign plan, 2026-09-05
+
+The maintainer requested a whole-design assessment and an explicit goal to
+complete the resulting plan. The existing customer-outcomes candidate remains
+the active scope under `v6-product-lane-expansion`. This is an architecture and
+qualification effort, not further cosmetic refinement of the current surface.
+
+### Product and architecture decision
+
+Patrol owns proactive investigation of the user's infrastructure and retained
+operating intent. Assistant explains and continues that same issue and its
+existing investigation/action records. An observed signal, model hypothesis,
+proposed action, accepted plan, executed operation and independently verified
+outcome are different facts. No transition may promote one into another merely
+because a tool accepted a structurally valid record.
+
+Keep canonical resource identity, observation provenance, retained evidence,
+operator intent and existing governed action records. Keep permissions, approval,
+mutual exclusion, idempotency, explicit resource/cost limits, execution and
+independent postcondition verification deterministic. Interpretation, relevance,
+causal diagnosis, investigation choices and action judgment belong to the model.
+Every orchestration pass must identify its objective invariant. A pass whose
+purpose is to manufacture or force a diagnosis from proxy counts must be removed
+or replaced with better model context and an explicit model decision.
+
+The design review reproduced a concrete trust violation in
+`internal/ai/chat/agentic_investigation_budget.go`: accepted proposal rationale is
+called an evidence checkpoint, inserted into the final Root Cause section, and
+the completion prompt forbids downgrading it in the reviewed baseline. Acceptance establishes that a
+proposal was recorded, not that its causal claim is true. This mechanism must
+be corrected before the diagnosis/action journey can qualify. The current
+working change removes both insertion paths and the duplicated proposal prompt
+state. A full-service regression now preserves the exact uncertain conclusion in the
+stream, returned result and persisted session after evidence and proposal turns.
+The full live investigation journey remains to be qualified.
+
+### Baseline and measurement limits
+
+The recorded 2026-09-05 telemetry review in the owning coverage gap used latest
+reports from monitoring-active, multi-ping installations, excluding development
+and deployment proof. It contained 127 paid installations, 71 with Patrol enabled
+and 23 with Assistant calls. Fourteen reported verified resolutions came from one
+installation. These are the previously recorded aggregate review, not a fresh
+query made during this redesign. Paid includes all non-free tiers. Cooccurring
+usage does not establish a linked successful task. Schema 17 outcome/provider/cost
+fields had no adoption in that review, so those fields cannot establish present
+customer effectiveness or model cost.
+
+Local live proof has independently exposed lost evidence, incorrect history
+coordinates, premature compaction, unsupported causal claims and roughly
+three-minute interactive investigations. One maintainer installation is useful
+reproduction evidence, not a representative customer success rate.
+
+### Execution order and acceptance
+
+| Step | Work | Acceptance | Current state |
+|---|---|---|---|
+| 1. Product contract and baseline | Map the current loop and sources of judgment. Record telemetry populations and gaps. | Every identified decision has an owner. Activity is not labelled usefulness. | Complete for this redesign scope. Contract, ownership decisions and baseline limits are recorded. |
+| 2. Shared evidence | Preserve canonical risk reasons and SMART counters, source/time semantics and history across tools/turns. | Regression tests preserve unknown versus zero and all canonical evidence. Real responses can inspect the same facts as the product. | Implemented and qualified for the named shared-evidence defects. Canonical disk detail, risk and cadence pass real data-path proof. Full affected package, concurrency and retained-query performance checks pass. Real-model interpretation failures remain tracked in step 5. |
+| 3. Diagnostic orchestration | Correct proposal-as-proof. Audit triage budgets, unmatched-signal evaluation, assessment completion and investigation cutoffs. | No code-written causal conclusion. No quality inferred from tool, flag or finding counts. Each retained pass has an objective reason. Safety boundaries and incomplete outcomes remain explicit. | Proposal promotion and capture inference removed in the working change. Targeted tool, loop and full-service regressions pass. Remaining orchestration audit and live qualification are open. |
+| 4. Issue through verified outcome | Follow existing issue/investigation/action records into Assistant, approval, execution and independent readback. | Accepted proposal is visibly distinct from execution and verification. Rejected or unsupported actions do not become success. Uncertainty can survive an action proposal. | Existing foundation, full journey qualification pending. |
+| 5. Ground-truth qualification and landing | Extend existing qualification tooling only where necessary. Exercise healthy/unhealthy, dependency, missing-access, storage/backup and approved/rejected action cases. Inspect the final browser journey at desktop and narrow widths. | Record exact source/model/permissions, evidence, decisions, faults/misses, latency and verification. Fix in-scope failures, pass appropriate proofs and land scoped commits. | Pending. |
+
+Use one shared runtime and the existing qualification runner, not a second
+product intelligence engine or a new parallel lifecycle. Preserve independent
+negative controls and fault oracles. Recorded fixtures prove contracts and
+reproducibility, not live model competence. Do not tune success wording or
+scoring to make the model pass.
+
+### Diagnostic orchestration audit
+
+| Mechanism | Current implementation | Decision for implementation review |
+|---|---|---|
+| Proposal rationale inserted as Root Cause | The reviewed baseline amended the conclusion after acceptance in the agent loop and service. Both mutations are removed in the working change. | Prove uncertain prose survives accepted proposals through stream, persistence and linked Assistant display. Keep the proposal record as an attributed model decision and allow uncertainty in the diagnosis. |
+| Causal-resource validator | Removed the duplicate resource graph and name/status inference from the working capture boundary. Causal attribution is optional when unknown. | Prove capability/schema validation, parameter isolation and invocation integrity remain enforced. Dependency evidence remains available to the model through canonical queries. |
+| Flag-count turn ladder | `computeTriageMaxTurns` grants 5 + 3 turns per flag, bounded to 8–40, with a separate quick limit. | Replace quality/urgency proxies with explicit execution resource limits. More flags must not imply a better investigation budget. |
+| Unmatched-signal evaluation | `runAIAnalysisState` detects signals from tool output and triage, then starts a second model pass when they lack matching findings. | Audit for removal in favour of complete initial evidence and model-owned decisions. Preserve negative-control and missed-fault qualification rather than force reports. |
+| Missing-finding assessment sweep | A bounded continuation requests missing explicit present/resolved/uncertain verdicts for known findings. | Retain only the mechanical completeness obligation, with sufficient original evidence and no fabricated resolution or requirement to discover new issues. |
+| Investigation evidence-call floor | Completion is rejected unless at least one successful structured evidence call occurred. | Audit seed evidence and action freshness separately. A tool count does not prove grounding or quality. Preserve explicit failed/unavailable evidence. |
+| Authority and execution boundaries | Tenant identity, capability schemas, approvals, invocation IDs, parameter redaction and independent readback. | Keep and prove unchanged when diagnostic policy is simplified. These enforce objective invariants. |
+
+This is an audited change list, not a claim that the changes are already made.
+Each removal must run its focused regression and affected complete journey.
+
+### Completion and external dependencies
+
+The local implementation goal remains open until required qualification is
+performed. Ordinary Assistant requests work with the current subscription route,
+but autonomous Patrol has an explicit provider-policy refusal and remains
+blocked. Do not rephrase the refused probe, bypass the readiness boundary or
+count an interactive request as an autonomous Patrol pass. A supported provider
+path is required for that qualification. Prepare other work while resolving the
+provider dependency through supported configuration.
+
+Release publication and wider product readiness are separate. Independent
+volunteered Pro environments are still required before claiming repeatable
+customer value. Capture that rollout requirement in the owning qualification gap
+rather than presenting one homelab result as completion of it.
+
## Implemented interaction
`Explain with Assistant` on an attention item or active alert starts an explicit
@@ -394,3 +493,223 @@ gap. These residuals remain in the outcome qualification gap. The final
retained-canonical-coverage receipts bind the current code to Assistant
expansion, output scrolling, answer pixels and reload at both viewports, plus
the actual saved provider refusal and HTTP 409 gate.
+
+
+## Redesign baseline retest, 2026-09-05
+
+The unchanged performance question completed in 224 seconds with eleven visible
+tool controls after the local disk and timestamp-context changes. Playwright
+exercised `/patrol` at 1440x1000 and 390x1000, including keyboard/pointer tool
+expansion, output scrolling, Escape, session selection and reload. Initial answer
+pixels were inspected. Final conclusion pixels and the focused disk question
+remain pending, so this is not full browser qualification for the slice.
+
+The model explicitly repeated the timestamp caveat while still calling the
+returned series unbroken and inferring that the alert start almost certainly
+marked collection start. It also treated a present-time pressure read as a
+definitive resolution of the historical pressure question. These are failed
+diagnoses, not evidence that more caveat text is sufficient.
+
+The changed disk tool returned SMART `PASSED`, canonical `warning`, and remaining
+life 63, but neither SMART counters nor risk reasons. A direct canonical resource
+API read confirmed those fields are absent upstream too. The subsequent source-status read identified `proxmox.status=stale`, which
+explains a freshness warning without a hardware-risk reason. The disk tool had
+omitted that field too. Preserve the canonical source status alongside health
+and risk before qualifying this explanation. Do not fabricate a wear warning
+from 63 percent remaining.
+
+Artifacts remain private under the local
+`tmp/patrol-outcome-telemetry/homelab/diagnostic-evidence-performance` directory.
+An initial browser attempt never reached a provider because the configured SSD
+temporary directory was absent. The directory was restored before the recorded
+run. Installed Chrome was used because the expected cached Playwright browser
+was unavailable. No refused readiness probe or infrastructure mutation occurred.
+
+The qualification scorer also needs an explicit trust review: current root-cause
+grounding checks match resource names/IDs in a named Markdown section, and other
+diagnosis checks use required words. Combined with code-inserted proposal prose,
+this can reward a formatted assertion without independent diagnostic support.
+Keep the physical fault oracles and action readback checks, but do not treat
+those text checks as semantic proof of cause.
+
+### Regression work discovered during redesign
+
+PR #1920 CI exposed a retained-history performance regression, including the
+bounded chart query and shared store reads. Reconciliation must remain correct,
+but display aggregation belongs inside SQLite rather than scanning all retained
+points into Go. The working correction also reuses scan destinations. Initial paired
+worker benchmarks confirm reduced allocations but still show material runtime
+regressions. Performance correction and source-bound browser proof remain open.
+
+The resource API test selected the first canonical `agent` resource and assumed
+it was Agent-backed. A Proxmox-only node has the same canonical resource type,
+but different history storage coordinates. The test now selects the fixture's
+actual Agent identity and separately asserts the Proxmox node target.
+
+Disk warning investigation exposed another canonical gap. Proxmox physical disks
+poll every five minutes by default, but source freshness currently uses the
+shorter general Proxmox threshold. Registry merge status selection also requires
+a regression check for recovery from stale Proxmox data. Preserve source
+freshness in the tool now, and correct cadence/recovery at the owning registry
+and monitoring boundary before qualifying those warnings as useful diagnosis.
+
+### Disk live retest, 2026-09-06 local time
+
+The current local Assistant answered the disk question through Claude Opus 5
+in about three minutes. The source freshness field was visible and correctly
+identified as distinct from SMART failure. The answer still failed qualification:
+it invented a midnight polling restart from coincident timestamps, treated the
+remaining-life measurement as ambiguous, and hit three rejected disk-detail calls.
+The tool advertised `physical-disk` for get but its handler did not support it.
+
+The next working correction routes exact canonical disk get/health requests
+through the same shared projection and labels remaining life with the explicit
+`life_remaining_percent` field. The live answer is a recorded failure, not proof
+of completed diagnosis. The screenshots exercised `/patrol`, Assistant, tool
+expansions, keyboard activation, narrow wrapping, Escape and session reload at
+1440x1000 and 390x1000. Fresh browser proof is required after these later changes.
+
+### Canonical freshness correction under verification
+
+Physical disk collection now carries its independent polling interval into the
+canonical per-source freshness record. Staleness uses at least two expected
+intervals, while retaining a longer configured source threshold. A new observation
+from the highest-priority available source may refresh its resource status,
+including a Proxmox-only resource that previously kept its first warning.
+The targeted regression covers five- and fifteen-minute disk schedules, stale
+transition, identity-preserving recovery and measured hardware risk after recovery.
+The live collector exposed a further loss in the quick temperature refresh:
+`physicalDiskFromReadStateView` discarded cadence before writing disk records
+back to monitoring state. That conversion now preserves the typed source status
+schedule. Broader registry regressions and final browser proof remain required. This does not make timestamp coincidences evidence of collector restarts.
+
+The retained-query candidate also encountered a SQLite native fault during the
+full metrics package run. The same baseline package passed. This is unresolved
+until the failing context and current candidate are verified. Narrow benchmark
+and query-plan success do not close this failure.
+
+### Continued orchestration review
+
+The assessment sweep currently receives only the retained finding title, severity,
+resource and up to 500 characters of old evidence. It receives neither the main
+run's observations nor their collection times. That is insufficient context for a
+new present/resolved judgment. The completion obligation may remain mechanical,
+but any continuation must share the actual investigation context. It must not
+turn old evidence into a fresh assessment.
+
+The investigation evidence-call floor also misclassifies a failed or
+approval-blocked read as no evidence at all. Such a result establishes a limit
+on available access, while a successful tool call alone establishes no diagnostic
+quality. Remove the quality inference and forced start-repair turn. Preserve
+explicit call/turn/time limits and let an incomplete or uncertain conclusion
+record the unavailable evidence. Action authority remains independently enforced.
+
+The unmatched-signal follow-up rebuilds a reduced evidence list after the main
+model run and prohibits further investigation. Its candidate selection ranks
+health/reliability/backup/connectivity/anomaly categories and caps them at twenty.
+This duplicates interpretation outside the model and discards the main reasoning
+context. The planned replacement is the original model assessment with complete
+seed/tool evidence, measured against independent missed-fault controls.
+
+A fresh live canonical disk read after the cadence round-trip correction returned
+`online` and `expectedUpdateIntervalSeconds=300`. This verifies the actual
+collector-to-canonical-resource path. Full affected package tests and the current
+Assistant answer remain separate acceptance checks.
+
+### Current disk answer and browser evidence
+
+The 2026-09-06 00:28 BST Assistant retest completed through the configured Claude
+subscription route in 2m38s, with 12,436 reported tokens. All nine evidence calls
+completed, including the exact canonical disk detail request that failed before.
+The answer correctly interpreted 63 percent remaining life, source recency and
+the distinction between disk and node temperature evidence.
+
+The overall answer still fails qualification. It asserted that recovery exists
+because an online backup datastore has space, without reading actual backup
+coverage or restore evidence. It also offered host commands without checking
+that execution access exists. The conclusion should remain bounded to the
+observed disk state. Latency and unnecessary fleet-wide explanation remain
+product defects to address in the shared diagnostic contract.
+
+Playwright exercised `/patrol` and the Assistant drawer at 1440x1000 and
+390x1000, expanded all nine evidence calls with keyboard activation, scrolled
+long outputs, collapsed them, used Escape, reloaded and reopened the retained
+session. Root inspection covered actual disk-list/detail output pixels, retained
+summary output and the final answer through its last paragraph. The answer and
+evidence remain readable at both widths. The source-bound receipt and screenshots
+are private under `tmp/patrol-outcome-telemetry/homelab/diagnostic-evidence-disk-canonical-final`.
+This is interface and data-path proof, not a passing diagnostic outcome.
+
+### Retained-performance retest and transport correction
+
+The 00:32 BST read-only performance request completed in 3m36s with 17,677
+reported tokens. It correctly reported unavailable agent access and separated
+high memory occupancy from established pressure. It still called the retained
+slice unbroken and inferred collector startup from coincident timestamps despite
+explicitly repeating the evidence caveat. That diagnosis remains failed.
+
+This run also exposed raw serialized routing JSON in the visible conversation.
+The subscription fallback had accumulated native CLI text before the first
+declared tool call and forwarded it as Assistant content. It also retained only
+one call from a native batch. The working provider correction returns the entire
+first declared batch in order, emits no routing-turn prose, and keeps the
+structured final answer unchanged. Targeted provider tests pass. All explicit
+refusal and local-tool isolation boundaries remain. Fresh browser evidence is
+required after this runtime change.
+
+The transport retest exceeded the browser's five-minute limit and was cancelled
+when that browser closed. The saved session contains four ordered routing turns
+with 2, 5, 4 and 3 calls respectively. Their answer content is empty, so the
+protocol text leak is absent from the persisted transport result. Fourteen calls
+did not produce a final answer before cancellation. This is incomplete browser
+qualification and failed interactive latency, not a completed diagnostic run.
+The response instruction currently demands thorough investigation and suggested
+next steps even after the user's specific question can be answered. Review that
+shared product instruction before collecting another equally broad run.
+
+### Scoped response and rendering retest
+
+The 00:56 BST performance request completed in 2m42s with 13,021 reported
+tokens after the shared response prompt was scoped to the user's actual job.
+Routing JSON was absent, missing agent access remained explicit and memory
+pressure was correctly left unknown. The answer still overstated continuous
+coverage from retained timestamps, inferred history start causes, treated an
+alert on the same utilization measurement as corroboration, and mixed UTC/BST.
+This remains a failed diagnostic outcome. It does not qualify autonomous Patrol.
+
+That real answer exposed a six-column table expanding the narrow conversation.
+The shared sanitized renderer now gives tables their own keyboard-scrollable
+region. Playwright reopened the real session at 1440x1000, 900x1000 and 390x1000,
+reached the last column, verified fixed conversation width, and exercised Escape,
+reload and session reopening. Browser-only short and streaming table fixtures
+confirmed preserved focus, region identity and scroll position. Actual pixels
+were inspected. The source-bound frontend receipt is
+`frontend-modern/browser-verification.json`; private artifacts are under
+`tmp/patrol-outcome-telemetry/homelab/markdown-table-final`. These fixture results
+establish rendering behavior only. No additional model calls were made.
+
+### Final evidence-foundation proof
+
+The affected tools, chat, providers, models, unified resources, monitoring and
+resource API package suites pass on the worker. The final metrics implementation
+(`store.go` SHA256 `5ef4b22506ecc73131bcd881a5bbc72d2a45f43f393f2dccbfba52dc287b18e8`)
+passes the full metrics suite in 90.446s, the database suite in 0.131s and the
+focused concurrent read/write race proof in 19.625s. Tests cover new/deleted
+history visibility, resource isolation, bounded compiled-statement retention,
+one-connection operation and transaction-bound query instrumentation.
+
+The unchanged bounded chart benchmark measures 889.3 microseconds baseline
+versus 888.8 microseconds candidate (p=.796, n=10), with no significant slowdown.
+Final paired metrics benchmarks pass the existing regression checker. The
+single-metric downsample case reports +10.00% (p=.005, n=10), close to the
+checker's greater-than-10% threshold. Raw single-series and both multi-metric
+queries have no significant latency difference. These are bounded worker
+measurements, not a production latency guarantee. The earlier +20.34% chart
+regression is corrected by reusing compiled presence SQL while evaluating it
+inside every current read snapshot. No history or tier-presence result is cached.
+
+One earlier metrics package run terminated inside SQLite native binding code.
+Subsequent full and focused race runs passed, including the final implementation.
+The original fault has no established root cause and is not labelled harmless
+or explained by the successful reruns. Its private worker log is
+`paired-metrics-local-20260905T235637/full-metrics-final.log`.
diff --git a/docs/release-control/v6/internal/status.json b/docs/release-control/v6/internal/status.json
index 96f614a7a..a44d8bd1c 100644
--- a/docs/release-control/v6/internal/status.json
+++ b/docs/release-control/v6/internal/status.json
@@ -10201,7 +10201,7 @@
},
{
"id": "patrol-assistant-customer-outcome-qualification",
- "summary": "The 2026-09-05 product review found that Explain with Assistant opened a blank generic conversation and the attention workbench discarded richer canonical finding context. The explicit explanation dispatcher and canonical handoff repair address that frontend break, with scripted browser proof and regression tests. They do not qualify model reasoning, real infrastructure actions, or repeated customer value. Latest-report, monitoring-active, multi-ping telemetry excluding dev and deployment proof contained 127 paid installs, 71 with Patrol enabled and 23 with Assistant calls. Fourteen reported verified resolutions were concentrated in one install. Paid includes all non-free tiers, and activity cooccurrence is not a linked successful journey. Schema 17 outcome/provider/cost fields had no adoption in that review. Three real jobs remain to qualify across named provider/model/version configurations: unhealthy service diagnosis, backup or capacity risk, and supported VM/LXC plan, approval and verified result. Repeatability must be established across independent paid customer environments without pooling free-tier or ineligible installs. The same-day real maintainer homelab evaluation with Claude Opus 5 exposed missing Proxmox temperature-tool evidence, contradictory summary health, incomplete visible prose, alert-count confusion and an inaccurately classified provider refusal during Patrol readiness. The shared temperature projection and premature within-window evidence compaction are repaired, while the documented reasoning, completion, readiness and end-to-end outcome failures remain unqualified. The continuation repaired canonical resource lookup, operation-versus-disk filtering, Proxmox CPU topology and retained resource-scoped history across backend restarts. Unscoped summary/baseline modernization, wear semantics, disk risk explanation and complete retained-window coverage remain open. A named-resource repeat exposed overly restrictive current_resource recovery guidance, now corrected without weakening the attachment boundary. Live testing also exposed warning-based API execution after an incomplete provider check. Runtime and API now enforce the selected-mode verdict, including saved incomplete results, and the repaired live endpoint returns HTTP 409. One read-only test run started before the gate correction and was stopped, so it is not qualification evidence. The evidence interpretation continuation replaces model-facing heuristic summaries with retained statistics, explicit source scope, real observation coverage and preserved bucket extrema/timestamps. A typed provider_refusal diagnosis now survives CLI exit and legacy cache decoding without granting readiness or repeating the refused request. The final real summary distinguished utilisation from demonstrated pressure and explained extrema versus bucket averages, but still speculated about sensor artefacts and monitoring restarts, inferred continuous coverage from minute spacing, and conflated returned history with retained history. Reasoning remains unqualified. A local fixture confirmed that 24-hour Query and QueryAll can hide newer raw points when a preferred aggregate tier is non-empty. Shared temporal tier reconciliation now covers single and batch queries, with extrema preserved through rollups and indexed overlap probes. Metrics race and query SLO checks pass on the worker. The live retained summary reaches seven seconds before its query instead of returning a roughly 45-minute-stale tail. The adjacent node History tab exposed a separate canonical-target defect and an inferred Agent link from discovery routing. The unified PVE metrics target now matches the node store writer, the drawer consumes that target and explicit Agent linkage, and summary tools no longer carry a PVE-specific exception. Full diagnosis and autonomous Patrol outcomes remain unqualified.",
+ "summary": "The maintainer requested a whole-design assessment and an explicit completion goal for Patrol and Assistant. The ordered plan and source-bound evidence are in docs/qualification/PATROL_ASSISTANT_CUSTOMER_JOURNEY.md. The contract separates observations, model conclusions, proposals, execution and independently verified outcomes. Assistant continues the same issue and governed action records. The explicit explanation handoff is an implemented frontend foundation, not outcome qualification. The recorded 2026-09-05 latest-report, monitoring-active, multi-ping review excluded development and deployment proof and contained 127 paid installs, 71 with Patrol enabled and 23 with Assistant calls. Fourteen reported verified resolutions came from one install. Paid includes all non-free tiers. Activity cooccurrence does not establish a successful linked task, and schema 17 outcome/provider/cost fields had no adoption in that review. Real read-only homelab testing with Claude Opus 5 exposed missing canonical evidence, history target/tier gaps, premature context loss, unsupported causal claims and multi-minute responses. Prior history and handoff changes are recorded in the plan. The current working slice removes proposal rationale promotion into Root Cause and duplicate causal inference, preserves disk SMART/risk/source freshness and explicit remaining life, supports canonical disk get/health, and corrects collector cadence and status recovery. Targeted regressions and the live disk data path pass. Affected full worker package suites pass. The single-metric downsample regression and an adjacent bounded chart regression are corrected. Compiled presence statements are reused within a bounded set, but every query reads current data in one snapshot. The chart comparison shows no significant slowdown. Final metrics/database suites, focused concurrency race checks and the existing benchmark guard pass. Single-metric downsampling remains close to the performance threshold. Final landing checks remain pending. One earlier SQLite native binding fault has no established root cause, although subsequent final full and focused race proofs pass. The latest live answers still infer collector startup from timestamp coincidence, claim continuous coverage from retained spacing and infer backup recovery from an online datastore. They fail diagnostic qualification. A serialized routing-envelope leak prompted a subscription transport correction with targeted regression proof, verified without routing prose in a completed real Assistant answer. A narrow answer-table overflow defect is corrected in the shared sanitized renderer with desktop, intermediate and mobile browser proof. The selected subscription provider explicitly refused autonomous Patrol readiness. The runtime and API preserve that refusal and prevent execution. Ordinary Assistant success is not autonomous proof. Required healthy/unhealthy, dependency, missing-access, backup/capacity and approved/rejected action qualification remains incomplete. Independent volunteered Pro environments are required before claiming repeatable customer value or wider readiness.",
"owner": "project-owner",
"status": "planned",
"recorded_at": "2026-09-05",
@@ -10369,7 +10369,7 @@
{
"id": "patrol-assistant-customer-outcomes",
"name": "Patrol and Assistant Customer Outcomes",
- "summary": "Qualify unhealthy service diagnosis, backup or capacity risk, and supported VM/LXC planning through approval and independent verification on disposable infrastructure. Compare named model configurations for useful outcomes, missed faults, false positives, latency and cost, then establish repeated value across independent volunteered Pro environments. The explicit explanation handoff is the completed frontend foundation, not completion of this outcome qualification.",
+ "summary": "Simplify Patrol around model-owned investigation and one issue-to-verified-outcome journey shared with Assistant. Preserve canonical evidence and remove proposal-as-proof and proxy-driven diagnostic policy while retaining deterministic authority and independent verification. Execute the ordered redesign plan in docs/qualification/PATROL_ASSISTANT_CUSTOMER_JOURNEY.md, qualify unhealthy service diagnosis, backup or capacity risk and supported VM/LXC actions with negative controls, and record real model quality, latency and cost limits. Independent volunteered Pro qualification remains required for wider readiness.",
"status": "proposed",
"recorded_at": "2026-09-05",
"target_id": "v6-product-lane-expansion",
@@ -10392,7 +10392,8 @@
],
"demand_evidence": [
"user-direction: 2026-09-05 prioritise dependable Patrol and Assistant outcomes as core Pulse Pro value",
- "qualification-contract: docs/qualification/PATROL_ASSISTANT_CUSTOMER_JOURNEY.md"
+ "qualification-contract: docs/qualification/PATROL_ASSISTANT_CUSTOMER_JOURNEY.md",
+ "user-direction: 2026-09-05 assess the whole design, establish an executable plan and pursue its completion as an explicit goal"
]
}
],
diff --git a/docs/release-control/v6/internal/subsystems/agent-lifecycle.md b/docs/release-control/v6/internal/subsystems/agent-lifecycle.md
index ecc09bfbe..0560dc321 100644
--- a/docs/release-control/v6/internal/subsystems/agent-lifecycle.md
+++ b/docs/release-control/v6/internal/subsystems/agent-lifecycle.md
@@ -15,6 +15,13 @@
## Purpose
+The shared physical-disk model preserves an internal collector cadence through
+state snapshot cloning and canonical typed read-state reconstruction. This is
+monitoring metadata and does not change Agent registration, execution authority,
+or the Agent report wire contract. `internal/models/metrics_types_test.go` covers
+snapshot retention and isolation, and the monitoring collection-trust roundtrip
+covers merged Agent/Proxmox disks without losing the Proxmox source schedule.
+
Assistant historical metric wiring uses the current monitor's retained store
and registry metrics coordinates. Historical reads do not alter enrollment,
agent ownership, or command capabilities. Proxmox-only canonical agent details
diff --git a/docs/release-control/v6/internal/subsystems/ai-runtime.md b/docs/release-control/v6/internal/subsystems/ai-runtime.md
index b8a567f52..91344d70b 100644
--- a/docs/release-control/v6/internal/subsystems/ai-runtime.md
+++ b/docs/release-control/v6/internal/subsystems/ai-runtime.md
@@ -19,6 +19,23 @@
## Purpose
+Physical-disk reads in `pulse_metrics` and `pulse_query` list/get/health
+share `physicalDiskSummaryFromResource`. The projection preserves canonical
+resource status, per-source freshness, physical-disk risk reasons and optional SMART counters alongside
+the device-reported SMART health result. It never recomputes risk or collapses a
+SMART pass into an all-clear. A stale source may explain a warning without a
+hardware risk. Missing counters remain absent and observed zero
+remains present. The tool projection labels `life_remaining_percent` explicitly, rather than
+exposing the source-specific ambiguous `wearout` name, and distinguishes it from
+`smart.percentageUsed` as consumed endurance. The regression proof is
+`internal/ai/tools/physical_disk_evidence_test.go`.
+
+Retained summaries disclose that point and bucket timestamps describe returned
+history inside the requested window. Their spacing does not measure collection
+uptime or explain missing history. Collector lifecycle and retention settings
+are not part of the summary query. These are source semantics supplied to the
+model, not a harness-generated diagnosis or a claim of model qualification.
+
The manual Patrol API consumes the service's selected-mode runtime verdict as
an execution check. Dimension-level presentation warnings cannot authorize a
run after a provider failure leaves Watch-only readiness unassessed. Operator
@@ -148,25 +165,19 @@ results before asserting that no peer is implicated. The Docker `services`
operation is explicitly Swarm-only: an empty service list never proves the
absence of ordinary Docker containers or container dependencies, and its tool
result must direct the model to topology or search for that evidence.
-Every accepted investigation proposal carries two separately typed resource
-identities: the action target and the exact canonical causal resource
-established by collected evidence. They may be equal, but a cross-resource
-diagnosis must preserve the peer or dependency as the causal identity and the
-proposal reason must preserve its observed state and causal chain. The
-request-local proposal boundary retains a minimal canonical resource graph from
-successful structured query results (resource ID, name, state, and health-check
-targets only). When that graph uniquely resolves a claimed causal resource's
-health-check target to a different unavailable resource, a proposal that still
-names the affected resource as causal is rejected before capture and returned
-to the investigation for correction. Ambiguous names and non-structured or
-failed evidence never create a causal assertion; core must not guess among
-possible resources or retain arbitrary log text at this boundary. The
-tool-free completion turn receives that accepted record as an evidence
-checkpoint. Before persistence, core reconciles the Root Cause section against
-the checkpoint and restores the exact causal identity and recorded basis if the
-provider's final prose omits or downgrades them. This reconciliation adds no new
-model claim and grants no action authority; it prevents the narrative from
-contradicting the already accepted, evidence-bound proposal record.
+Every accepted investigation proposal preserves the canonical action target.
+An optional `causal_resource_id` records the model's attribution when supported
+by evidence. It is omitted when the cause is unknown. A recovery proposal may
+address an observed symptom without establishing its underlying cause. The
+reason carries observed evidence, expected benefit and remaining uncertainty.
+The capture boundary validates capability schemas, permissions, correlation,
+parameter isolation and invocation integrity. It does not maintain a duplicate
+resource graph or infer causality from dependency names and status vocabulary.
+The tool-free completion turn distinguishes proposal acceptance from diagnosis.
+The proposal remains an attributed model decision. Core preserves final model
+prose unchanged in the response, stream and transcript, including uncertainty
+or a correction of the earlier rationale. Neither structural validation nor
+proposal acceptance establishes a root cause or authorizes execution.
An investigation cannot complete or submit a typed action proposal before the
model has received at least one successful structured result from an advertised
evidence tool. Until then, core withholds proposal authority while leaving
@@ -8011,3 +8022,37 @@ reasoning and real remediation in
`docs/qualification/PATROL_ASSISTANT_CUSTOMER_JOURNEY.md`. The repeatable browser
proof is `scripts/check-patrol-assistant-journey.mjs`. A passing scripted
response does not establish a useful customer outcome or model qualification.
+
+### Subscription routing content boundary
+
+Subscription routing turns emit validated tool calls only. Native CLI narration
+and synthetic local errors are transport material, not Assistant answer text or
+infrastructure observations. For Claude's native-tool fallback, retain the entire
+first declared tool batch in order and discard later CLI continuations produced
+without Pulse tool results. Exact tool names, invocation identity, refusal,
+permission-denial and disabled local-tool boundaries remain enforced. Structured
+final answers retain their content unchanged.
+
+The regression uses a serialized routing envelope before a two-call native batch
+and a later synthetic continuation. It verifies empty routing prose, both
+original calls, original order and preserved final uncertainty. A completed
+real Assistant request after the transport change showed no routing envelope
+in the persisted conversation or visible answer. Its unsupported historical
+inferences remain failed diagnostic qualification.
+
+The shared Assistant response prompt is scoped to the user's resource, period
+and decision. It no longer demands general thoroughness or unsolicited next
+steps on every answer. Provenance follows the product's distinct source
+observation, model hypothesis, proposal, execution and verification records.
+This is model context, not a deterministic quality classifier. Real responses
+remain subject to factual review and latency qualification.
+
+### Readable answer tables
+
+The shared sanitized markdown renderer owns horizontal overflow for each table
+through a labelled, keyboard-focusable region. Application styling is added only
+after sanitization. Model-provided classes, styles and event handlers remain
+disallowed. Streaming DOM reconciliation preserves table region identity, focus
+and scroll position. The conversation itself must not acquire horizontal scroll
+from wide answer tables. Browser proof covers a persisted real answer and
+explicit renderer fixtures at 1440, 900 and 390 pixel widths.
diff --git a/docs/release-control/v6/internal/subsystems/monitoring.md b/docs/release-control/v6/internal/subsystems/monitoring.md
index e4ca25082..9e1793589 100644
--- a/docs/release-control/v6/internal/subsystems/monitoring.md
+++ b/docs/release-control/v6/internal/subsystems/monitoring.md
@@ -17,6 +17,16 @@
## Purpose
+Physical disk inventory has an independent collector schedule. The PVE poller
+carries its default five-minute or configured interval with each disk record,
+while keeping the last successful observation timestamp on retained records.
+Quick temperature refreshes and typed read-state conversions must preserve that
+cadence rather than falling back to the general node polling interval. It feeds
+the canonical source-freshness contract, not a model-generated diagnosis.
+Proof: the async fallback poll in `monitor_pve_disk_fallback_test.go`, the canonical
+roundtrip in `issue1595_collection_trust_test.go`, and registry cadence/recovery
+coverage. Live `/api/resources` and final Assistant browser proof are also required.
+
PBS node-status collection must reject HTTP-success responses whose `data`
is omitted or null (including a null response envelope). Absent status is
unavailable telemetry, not measured zero usage: the poller retains independently
diff --git a/docs/release-control/v6/internal/subsystems/performance-and-scalability.md b/docs/release-control/v6/internal/subsystems/performance-and-scalability.md
index cb3c81b40..e1dffead6 100644
--- a/docs/release-control/v6/internal/subsystems/performance-and-scalability.md
+++ b/docs/release-control/v6/internal/subsystems/performance-and-scalability.md
@@ -29,10 +29,23 @@ buckets whose stored timestamp lies in the requested window participate. This
avoids double-counting and prevents observations outside the requested range
from suppressing evidence inside it. Reconciliation precedes display
aggregation, which retains the existing unweighted mean and centred display
-bucket timestamp. Single and fleet queries use the same streaming aggregator.
+bucket timestamp. Single-resource charts aggregate in SQLite after reconciliation so only display
+buckets cross into Go. Fleet queries stream ordered observations into bounded
+display buckets without retaining every input point. Both preserve identical
+bucket, mean and extrema semantics, and reuse scan destinations across rows.
-The overlap probes use indexed series/time searches in one SQLite statement per
-bounded resource chunk. Query-plan tests exercise the runtime SQL builder,
+Each bounded resource chunk uses one read transaction. Indexed existence checks
+identify which retention tiers have any observations in that snapshot. Every
+present tier participates in the shared indexed overlap query. Presence never
+stands in for per-series coverage. The store retains at most 32 compiled
+presence-statement shapes and evaluates them again inside each read snapshot.
+It never caches tier presence, query results or timestamp windows. Less common
+shapes run uncached after the bound is reached. Preparation occurs before
+acquiring the transaction, including with a single-connection pool. The shared
+database instrumentation preserves timing for transaction-bound prepared
+statements, and database closure owns their lifetime. Empty tiers need no per-observation probe,
+and an all-raw series uses a direct range query. Fixed query dimensions need
+not be decoded again for every returned point. Query-plan tests exercise the runtime SQL builder,
including all tier priorities and metric filters. Multi-stage rollups preserve
recorded minima and maxima instead of taking extrema from bucket averages.
Already discarded historical extrema cannot be reconstructed by this change.
diff --git a/docs/release-control/v6/internal/subsystems/storage-recovery.md b/docs/release-control/v6/internal/subsystems/storage-recovery.md
index 99e1148c0..105942858 100644
--- a/docs/release-control/v6/internal/subsystems/storage-recovery.md
+++ b/docs/release-control/v6/internal/subsystems/storage-recovery.md
@@ -21,6 +21,12 @@
## Purpose
+Canonical disk source status may carry the collector's expected update interval.
+This freshness metadata remains separate from physical-disk risk, SMART values
+and recovery/action authority. Source status cloning preserves the optional field,
+and absent cadence retains the source default. Fresh observation recovery must
+not erase a current hardware-risk assessment.
+
Assistant physical-disk metrics use the canonical disk host projection and
separate the disk format filter from the metrics operation selector. Missing
SMART observations remain missing coverage, not disk failure or recovery
diff --git a/docs/release-control/v6/internal/subsystems/unified-resources.md b/docs/release-control/v6/internal/subsystems/unified-resources.md
index 5bf1decfb..6739cc88e 100644
--- a/docs/release-control/v6/internal/subsystems/unified-resources.md
+++ b/docs/release-control/v6/internal/subsystems/unified-resources.md
@@ -15,6 +15,16 @@
## Purpose
+Physical disk source freshness preserves the collector-authored
+`expectedUpdateIntervalSeconds` alongside the actual last observation. Registry
+ingest, merge, cloning and typed disk views preserve it. Staleness uses the
+longer of the configured source threshold and two collector intervals. A new
+observation from the highest-priority represented source may refresh canonical
+status, including Proxmox-only resources. Lower-priority platform data cannot
+overwrite a linked Agent status. Freshness changes never change resource identity.
+Proof: disk cadence/recovery tests in `internal/unifiedresources/registry_test.go`
+and source-status clone coverage in `internal/unifiedresources/clone_test.go`.
+
Canonical resource identity and metrics storage coordinates are distinct.
API-only Proxmox hosts remain `agent` resources but their metrics target is
`node` with the Proxmox source ID, matching the collector's writes. A real
diff --git a/frontend-modern/browser-verification.json b/frontend-modern/browser-verification.json
index 032c0022e..a5fa18a8a 100644
--- a/frontend-modern/browser-verification.json
+++ b/frontend-modern/browser-verification.json
@@ -1,41 +1,39 @@
{
"version": 1,
- "base_sha": "5288b64d40642b29cd27fa0ca67080cb261d37f1",
- "verified_at": "2026-09-05T22:04:00.286Z",
+ "base_sha": "3a189f31d44789a642b68a915d3d4a683ab81104",
+ "verified_at": "2026-09-06T00:41:42.025224Z",
"result": "passed",
"changed_paths": [
- "frontend-modern/src/components/Workloads/nodeDrawerModel.ts",
- "frontend-modern/src/types/api.ts",
- "frontend-modern/src/utils/resourceStateAdapters.ts"
+ "frontend-modern/src/components/AI/aiChatUtils.ts"
],
"content_sha256": {
- "frontend-modern/src/components/Workloads/nodeDrawerModel.ts": "2a46640db19fe358e9f0ab32c4819a1c30a7c91b71015104b53253296176bdd6",
- "frontend-modern/src/types/api.ts": "16e3c1c60450eb5bb8e4f0db57f595d29cc9dbfd2e01bffee1f0529262f9f60a",
- "frontend-modern/src/utils/resourceStateAdapters.ts": "98ee10447e5fa0fe122122727ce1b59a29c92afbcdcf7f046f4cea9f2edc9e80"
+ "frontend-modern/src/components/AI/aiChatUtils.ts": "c65a0f95d3e32da3e8f590de1d40067b5568e5d7d5573652024541e8f8709fb8"
},
"routes": [
- "/proxmox",
- "/patrol",
- "/settings/pulse-intelligence/patrol"
+ "/patrol"
],
"viewports": [
{
"width": 1440,
"height": 1000
},
+ {
+ "width": 900,
+ "height": 1000
+ },
{
"width": 390,
"height": 1000
}
],
"states": [
- "Current managed frontend and backend, mock monitoring disabled, real homelab and read-only model investigation. No infrastructure mutation or repeated refused readiness probe.",
- "Proxmox node overview and History, 1h/24h/7d ranges, populated utilization/network/thermal charts. API-only node does not claim an installed Agent or expose unavailable disk throughput.",
- "Assistant pending/completed turn, thirteen tool controls and full output/answer scrolling. Saved provider refusal remains visible and actual manual execution returns HTTP 409."
+ "Assistant persisted real retained-performance answer with a six-column table. Browser-only short-table and streaming-table fixtures use the current shared renderer and DOM morph implementation.",
+ "Wide table overflow belongs to its focusable region. Conversation and surrounding prose do not scroll horizontally. All columns remain reachable.",
+ "Saved subscription provider refusal remains visible. No provider probe, model request or infrastructure mutation was made for this rendering verification."
],
"interactions": [
- "Node expansion, History tab, native range selection, hover on each chart, Overview/History switching, reload and reopen. Normal and narrow chart pixels inspected.",
- "Assistant send and completion, keyboard Enter/pointer expansion and collapse, inner output scroll to end, answer pixels, Escape, history selection and reload.",
- "Local Playwright receipts: retained-canonical-coverage/history-receipt.json and receipt.json. Source hashes match the final implementation. The model outcome is a bounded evidence-path pass, not diagnosis or autonomous Patrol qualification."
+ "Open Assistant, select retained session, keyboard-focus table, ArrowRight scrolling, scroll to last column, vertical message scrolling, inspect start/end pixels at desktop and narrow widths.",
+ "Append a table row through the current renderer and DOM morph. Verify region identity, keyboard focus and horizontal scroll position are preserved. Inspect short and streaming table pixels at 1440, 900 and 390 widths.",
+ "Escape, reload, reopen Assistant and retained session. Private proof: tmp/patrol-outcome-telemetry/homelab/markdown-table-final. This is renderer qualification, not a passing diagnostic outcome."
]
}
diff --git a/frontend-modern/src/components/AI/__tests__/aiChatUtils.test.ts b/frontend-modern/src/components/AI/__tests__/aiChatUtils.test.ts
index c35073a3c..d11a3c0c5 100644
--- a/frontend-modern/src/components/AI/__tests__/aiChatUtils.test.ts
+++ b/frontend-modern/src/components/AI/__tests__/aiChatUtils.test.ts
@@ -115,6 +115,24 @@ describe('aiChatUtils', () => {
expect(output).not.toContain('alert("xss")');
});
+ it('contains table overflow without trusting model layout attributes', () => {
+ const output = utils.renderMarkdown(
+ '
',
+ );
+ const root = document.createElement('div');
+ root.innerHTML = output;
+ const table = root.querySelector('table');
+ const region = table?.parentElement;
+ expect(region?.getAttribute('role')).toBe('region');
+ expect(region?.getAttribute('aria-label')).toBe('Scrollable table');
+ expect(region?.tabIndex).toBe(0);
+ expect(region?.classList.contains('overflow-x-auto')).toBe(true);
+ expect(table?.textContent).toBe('MetricValueMemory88%');
+ expect(table?.hasAttribute('class')).toBe(false);
+ expect(table?.hasAttribute('style')).toBe(false);
+ expect(table?.hasAttribute('onclick')).toBe(false);
+ });
+
it('escapes HTML entities if markdown parsing fails', () => {
const spy = vi.spyOn(marked, 'parse').mockImplementation(() => {
throw new Error('boom');
diff --git a/frontend-modern/src/components/AI/aiChatUtils.ts b/frontend-modern/src/components/AI/aiChatUtils.ts
index 7bf9de19c..782caa17f 100644
--- a/frontend-modern/src/components/AI/aiChatUtils.ts
+++ b/frontend-modern/src/components/AI/aiChatUtils.ts
@@ -97,7 +97,7 @@ export const renderMarkdown = (content: unknown): string => {
// - 'target' and 'rel' are not in ALLOWED_ATTR because the
// afterSanitizeAttributes hook above sets them on every ; they
// don't need to be parsed in from the LLM-supplied HTML.
- return DOMPurify.sanitize(rawHtml, {
+ const sanitized = DOMPurify.sanitize(rawHtml, {
// Allow common formatting tags but block scripts, iframes, etc.
ALLOWED_TAGS: [
'p',
@@ -135,6 +135,21 @@ export const renderMarkdown = (content: unknown): string => {
// Force all links to open in new tab and prevent opener attacks
ADD_ATTR: ['target', 'rel'],
});
+ // Add application-owned scrolling only after sanitization. Model HTML
+ // cannot supply classes or attributes that escape the message layout.
+ const template = document.createElement('template');
+ template.innerHTML = sanitized;
+ for (const table of template.content.querySelectorAll('table')) {
+ const region = document.createElement('div');
+ region.className =
+ 'max-w-full overflow-x-auto focus-visible:outline focus-visible:outline-2 focus-visible:outline-blue-500';
+ region.tabIndex = 0;
+ region.setAttribute('role', 'region');
+ region.setAttribute('aria-label', 'Scrollable table');
+ table.replaceWith(region);
+ region.appendChild(table);
+ }
+ return template.innerHTML;
} catch {
// If parsing fails, escape HTML entities as fallback
return normalized.replace(/[&<>"']/g, (char) => {
diff --git a/internal/ai/chat/agentic.go b/internal/ai/chat/agentic.go
index 8fb69ec62..f8f03f29a 100644
--- a/internal/ai/chat/agentic.go
+++ b/internal/ai/chat/agentic.go
@@ -954,7 +954,6 @@ func (a *AgenticLoop) executeWithTools(ctx context.Context, sessionID string, me
investigationEvidenceStartRepairAttempted := false
toolBlockedLastTurn := false // When true, request final text after budget/loop block
investigationProposalCompleted := false
- var acceptedInvestigationProposalBasis *investigationProposalBasis
acceptedFindingReports := 0
// Patrol core normally establishes the exact-scope active-finding snapshot
// before the provider is invoked. Legacy/narrow adapters can still expose a
@@ -1210,7 +1209,7 @@ agenticLoop:
case investigationProposalCompleted:
req.Tools = nil
textOnlySafetyBrake = true
- req.System += investigationProposalCompletionSystemPrompt(acceptedInvestigationProposalBasis)
+ req.System += investigationProposalCompletionSystemPrompt
case a.maxEvidenceCalls > 0 && a.totalEvidenceCalls >= a.maxEvidenceCalls:
if a.successfulEvidenceCalls > 0 {
req.Tools = investigationTerminalTools(tools)
@@ -1761,23 +1760,6 @@ agenticLoop:
continue
}
- if isPatrolInvestigationExecution(a.currentExecutionProfile()) && investigationProposalCompleted && acceptedInvestigationProposalBasis != nil {
- grounded, addition := groundInvestigationConclusionInProposal(assistantMsg.Content, acceptedInvestigationProposalBasis)
- if grounded != assistantMsg.Content {
- assistantMsg.Content = grounded
- resultMessages[len(resultMessages)-1].Content = grounded
- providerMessages[len(providerMessages)-1].Content = grounded
- if addition != "" {
- jsonData, _ := json.Marshal(ContentData{Text: addition})
- callback(StreamEvent{Type: "content", Data: jsonData})
- }
- log.Warn().
- Str("session_id", sessionID).
- Str("causal_resource_id", acceptedInvestigationProposalBasis.CausalResourceID).
- Msg("[AgenticLoop] Restored accepted proposal evidence checkpoint in investigation conclusion")
- }
- }
-
// === ADVERTISED-ACTION GATE: an action request ends in pulse_control, not prose ===
// The field failure this pins: the operator asks to reboot N guests,
// the model resolves them, then writes a report with "next steps"
@@ -2467,9 +2449,6 @@ agenticLoop:
}
successfulInvestigationEvidence := isPatrolInvestigationExecution(a.currentExecutionProfile()) &&
isSuccessfulInvestigationEvidenceResult(tc.Name, resultText, isError)
- if successfulInvestigationEvidence && a.executor != nil {
- a.executor.RecordProposalEvidence(tc.Name, resultText)
- }
// Track pending recovery for strict resolution blocks
// (FSM blocks are tracked above; strict resolution blocks come from the executor)
@@ -2598,7 +2577,6 @@ agenticLoop:
}
if isPatrolInvestigationExecution(a.currentExecutionProfile()) && tc.Name == agentcapabilities.PatrolProposeActionToolName {
investigationProposalCompleted = true
- acceptedInvestigationProposalBasis = investigationProposalBasisFromToolCall(tc)
}
}
diff --git a/internal/ai/chat/agentic_investigation_budget.go b/internal/ai/chat/agentic_investigation_budget.go
index f3071e1f3..8ef7f7fb8 100644
--- a/internal/ai/chat/agentic_investigation_budget.go
+++ b/internal/ai/chat/agentic_investigation_budget.go
@@ -14,121 +14,11 @@ func isPatrolInvestigationExecution(profile aitools.ExecutionProfile) bool {
return profile == aitools.ProfilePatrolInvestigation
}
-type investigationProposalBasis struct {
- TargetResourceID string
- CausalResourceID string
- CapabilityName string
- Reason string
-}
+// Structural proposal acceptance records an intended action, not a verified
+// diagnosis. The model retains responsibility for interpreting tool evidence.
+const investigationProposalCompletionSystemPrompt = `
-func investigationProposalBasisFromToolCall(call providers.ToolCall) *investigationProposalBasis {
- if call.Name != agentcapabilities.PatrolProposeActionToolName {
- return nil
- }
- stringInput := func(key string) string {
- value, _ := call.Input[key].(string)
- return strings.TrimSpace(value)
- }
- basis := &investigationProposalBasis{
- TargetResourceID: stringInput("resource_id"),
- CausalResourceID: stringInput("causal_resource_id"),
- CapabilityName: stringInput("capability_name"),
- Reason: stringInput("reason"),
- }
- if basis.TargetResourceID == "" || basis.CausalResourceID == "" || basis.CapabilityName == "" || basis.Reason == "" {
- return nil
- }
- return basis
-}
-
-func investigationProposalCompletionSystemPrompt(basis *investigationProposalBasis) string {
- prompt := `
-
-INVESTIGATION COMPLETION: A typed action proposal has already been accepted for this run. Do not call more tools. Produce the required investigation summary from the evidence already collected and state that the proposal is pending governed policy or operator handling; never claim that it executed.`
- if basis == nil {
- return prompt
- }
- return prompt + fmt.Sprintf(`
-
-ACCEPTED PROPOSAL EVIDENCE CHECKPOINT: The governed proposal record identifies causal resource %q, action target %q, capability %q, and this exact recorded basis: %q. Preserve the causal resource identity and recorded basis in the Root Cause section, distinguish the action target in Affected Resources, and do not contradict or downgrade facts already used to justify the accepted proposal.`, basis.CausalResourceID, basis.TargetResourceID, basis.CapabilityName, basis.Reason)
-}
-
-// groundInvestigationConclusionInProposal makes the typed proposal record the
-// terminal source of truth for facts the model already committed as its action
-// rationale. Provider prose remains free-form, but it cannot omit the exact
-// causal identity or rationale after Patrol has accepted them structurally.
-func groundInvestigationConclusionInProposal(content string, basis *investigationProposalBasis) (string, string) {
- if basis == nil || strings.TrimSpace(content) == "" {
- return content, ""
- }
- reason := strings.Join(strings.Fields(basis.Reason), " ")
- checkpoint := fmt.Sprintf("Proposal evidence checkpoint: causal resource `%s`; recorded basis: %s", basis.CausalResourceID, reason)
- rootCause := markdownInvestigationSection(content, "Root Cause")
- if strings.Contains(strings.ToLower(rootCause), strings.ToLower(basis.CausalResourceID)) &&
- strings.Contains(strings.ToLower(rootCause), strings.ToLower(reason)) {
- return content, ""
- }
- grounded, ok := appendToMarkdownInvestigationSection(content, "Root Cause", checkpoint)
- if !ok {
- separator := "\n\n"
- if strings.TrimSpace(content) == "" {
- separator = ""
- }
- addition := separator + "### Root Cause\n" + checkpoint
- return content + addition, addition
- }
- return grounded, "\n\n" + checkpoint
-}
-
-func markdownInvestigationSection(content, heading string) string {
- lines := strings.Split(content, "\n")
- start := -1
- for index, line := range lines {
- trimmed := strings.TrimSpace(line)
- if start < 0 {
- if strings.EqualFold(trimmed, "### "+heading) {
- start = index + 1
- }
- continue
- }
- if strings.HasPrefix(trimmed, "### ") {
- return strings.Join(lines[start:index], "\n")
- }
- }
- if start >= 0 {
- return strings.Join(lines[start:], "\n")
- }
- return ""
-}
-
-func appendToMarkdownInvestigationSection(content, heading, addition string) (string, bool) {
- lines := strings.Split(content, "\n")
- start := -1
- insertAt := len(lines)
- for index, line := range lines {
- trimmed := strings.TrimSpace(line)
- if start < 0 {
- if strings.EqualFold(trimmed, "### "+heading) {
- start = index + 1
- }
- continue
- }
- if strings.HasPrefix(trimmed, "### ") {
- insertAt = index
- break
- }
- }
- if start < 0 {
- return content, false
- }
- prefix := append([]string(nil), lines[:insertAt]...)
- for len(prefix) > 0 && strings.TrimSpace(prefix[len(prefix)-1]) == "" {
- prefix = prefix[:len(prefix)-1]
- }
- prefix = append(prefix, "", addition, "")
- prefix = append(prefix, lines[insertAt:]...)
- return strings.Join(prefix, "\n"), true
-}
+INVESTIGATION COMPLETION: An action proposal has been recorded for this run. Do not call more tools. Summarize the evidence collected and any uncertainty. The proposal is pending governed policy or operator handling and has not executed. Proposal acceptance validates the action contract, not the rationale or root cause.`
const investigationEvidenceStartRepairSystemPrompt = `
diff --git a/internal/ai/chat/agentic_investigation_budget_w0716_coverage_test.go b/internal/ai/chat/agentic_investigation_budget_w0716_coverage_test.go
index 868c457a1..0a2cb244b 100644
--- a/internal/ai/chat/agentic_investigation_budget_w0716_coverage_test.go
+++ b/internal/ai/chat/agentic_investigation_budget_w0716_coverage_test.go
@@ -49,53 +49,6 @@ func TestInvestigationCompletionRequiresSupportedProposalBeforeProse(t *testing.
}
}
-func TestInvestigationProposalCompletionCarriesTypedCausalCheckpoint(t *testing.T) {
- call := providers.ToolCall{
- Name: agentcapabilities.PatrolProposeActionToolName,
- Input: map[string]interface{}{
- "resource_id": "app-container-client",
- "causal_resource_id": "app-container-dependency",
- "capability_name": "restart",
- "reason": "app-container-dependency is exited and the client health check is failing",
- },
- }
- basis := investigationProposalBasisFromToolCall(call)
- if basis == nil || basis.CausalResourceID != "app-container-dependency" {
- t.Fatalf("proposal basis = %#v, want exact causal resource", basis)
- }
- prompt := investigationProposalCompletionSystemPrompt(basis)
- for _, want := range []string{"app-container-dependency", "app-container-client", "is exited", "Root Cause"} {
- if !strings.Contains(prompt, want) {
- t.Fatalf("completion prompt missing %q: %s", want, prompt)
- }
- }
-}
-
-func TestGroundInvestigationConclusionInProposalRestoresOmittedCausalEvidence(t *testing.T) {
- basis := &investigationProposalBasis{
- TargetResourceID: "app-container-client",
- CausalResourceID: "app-container-dependency",
- CapabilityName: "restart",
- Reason: "app-container-dependency is exited and its client is unhealthy",
- }
- content := "### Investigation Summary\nClient health is failing.\n\n### Root Cause\nThe precise dependency state is unknown.\n\n### Affected Resources\n`app-container-client`.\n\n### Recommendation\nRestart proposal pending.\n\n### Conclusion\nNEEDS_ATTENTION: pending."
- grounded, addition := groundInvestigationConclusionInProposal(content, basis)
- rootCause := markdownInvestigationSection(grounded, "Root Cause")
- for _, want := range []string{"app-container-dependency", "is exited", "recorded basis"} {
- if !strings.Contains(rootCause, want) {
- t.Fatalf("grounded Root Cause missing %q:\n%s", want, grounded)
- }
- }
- if addition == "" {
- t.Fatal("expected a stream-visible grounding addition")
- }
-
- again, duplicate := groundInvestigationConclusionInProposal(grounded, basis)
- if again != grounded || duplicate != "" {
- t.Fatalf("grounding must be idempotent; duplicate=%q\n%s", duplicate, again)
- }
-}
-
func Test_w0716_budget_IsInvestigationEvidenceTool(t *testing.T) {
tests := []struct {
name string
diff --git a/internal/ai/chat/service.go b/internal/ai/chat/service.go
index aad96c31b..3f166c6d4 100644
--- a/internal/ai/chat/service.go
+++ b/internal/ai/chat/service.go
@@ -3871,26 +3871,11 @@ func (s *Service) buildSystemPromptWithToolGovernance(toolGovernance string) str
- Asking may not be your first action: a clarification request issued before you have attempted any tool call this turn will be rejected. For broad questions ("how is my infrastructure doing?", "any alerts?"), call the obvious read-only tool with defaults — pulse_summarize {"action":"fleet"} needs no parameters, and the alert tools list everything without a target.
- Never guess a target you did not resolve. Do not attempt a tool call with current_resource or another placeholder as a stand-in for a missing target — in autonomous mode the same rules apply: resolve with read-only tools first, then ask in normal assistant text if genuine ambiguity remains.
-## HOW TO RESPOND
-You are like a colleague doing pair programming on infrastructure tasks. Tool calls are your internal investigation — the user sees your final synthesized response.
-
-1. INVESTIGATE THOROUGHLY: Decide whether tool evidence is needed, then gather enough information to answer well. Don't stop after the first tool call if more context would help.
-
-2. SYNTHESIZE YOUR FINDINGS: After using tools, explain what you learned and did. Don't just confirm "done" — provide context that helps the user understand the outcome.
-
-3. SURFACE ISSUES PROACTIVELY: If you discover something during investigation that affects the user's goal (prerequisites missing, config issues, limitations), mention it. Don't hide problems.
-
-4. SUGGEST NEXT STEPS: If there's something the user might need to do next, or if you noticed a potential improvement, mention it.
-
-5. BE DIRECT: Acknowledge mistakes or complications honestly. If something won't work as the user expects, say so clearly.
-
-6. KEEP FORMATTING CLEAN: Use concise headings, bullets, and tables where they help. Do not use emoji, warning icons, or decorative symbols in normal operational answers unless the user explicitly asks for that tone.
+## USER JOB
+Help the user decide what needs attention, understand why, and take the next supported step through Pulse. Keep the investigation and answer focused on the resource, period and decision they asked about. Lead with the answer and the evidence needed to assess it. Use concise prose or a small table, without decorative symbols.
## GROUNDING & PROVENANCE
-Pulse context carries provenance: discovered facts include the source that produced them and a confidence, discovery context states when it was last gathered ("Last discovered: 2 days ago"), and metrics and events carry timestamps.
-- When you state a fact drawn from this context, attribute it briefly so the user can trust and verify it — name the source for a discovered fact ("Debian 12, per /etc/os-release") and note recency for time-sensitive claims ("as of the last poll", "the discovery is 2 days old").
-- Do not present stale or cached context as current. If the discovery is old or the underlying state may have changed since, say so and offer to re-check live.
-- Keep attribution concise and inline; do not append a citation to every sentence or clutter the answer.
+Pulse records source observations, model hypotheses, proposed actions, executed operations and verified outcomes separately. Preserve those distinctions when explaining an issue. Attribute measured facts to their source and observation time. A retained finding or cached discovery is historical evidence until a current observation supports it. A failed read establishes an access or collection limit, not a healthy result.
## TASK COMPLETION
- After successful control actions, respond once you have enough evidence to explain the result.
diff --git a/internal/ai/chat/service_investigation.go b/internal/ai/chat/service_investigation.go
index 73885d442..61383e1d5 100644
--- a/internal/ai/chat/service_investigation.go
+++ b/internal/ai/chat/service_investigation.go
@@ -2,7 +2,6 @@ package chat
import (
"context"
- "encoding/json"
"errors"
"fmt"
"strings"
@@ -232,19 +231,7 @@ func (s *Service) ExecuteInvestigationStream(ctx context.Context, req Investigat
}
}
content := contentBuilder.String()
- if runErr == nil && proposalErr == nil && proposal != nil {
- grounded, addition := groundInvestigationConclusionInProposal(content, &investigationProposalBasis{
- TargetResourceID: proposal.ResourceID,
- CausalResourceID: proposal.CausalResourceID,
- CapabilityName: proposal.CapabilityName,
- Reason: proposal.Reason,
- })
- content = grounded
- if addition != "" {
- data, _ := json.Marshal(ContentData{Text: addition})
- callback(StreamEvent{Type: "content", Data: data})
- }
- }
+
result := &InvestigationRunResult{
Content: content,
Proposal: proposal,
diff --git a/internal/ai/chat/service_tooling_test.go b/internal/ai/chat/service_tooling_test.go
index 87aac0717..d834d25eb 100644
--- a/internal/ai/chat/service_tooling_test.go
+++ b/internal/ai/chat/service_tooling_test.go
@@ -353,7 +353,7 @@ func TestBuildSystemPrompt_DoesNotClaimGenericVMControl(t *testing.T) {
if !strings.Contains(prompt, "pulse_kubernetes") {
t.Fatalf("expected system prompt to include the governed Kubernetes tool contract, got %q", prompt)
}
- if !strings.Contains(prompt, "Do not use emoji, warning icons, or decorative symbols") {
+ if !strings.Contains(prompt, "without decorative symbols") {
t.Fatalf("expected system prompt to keep Assistant formatting operational, got %q", prompt)
}
}
@@ -365,9 +365,9 @@ func TestBuildSystemPrompt_IncludesProvenanceGuidance(t *testing.T) {
for _, expected := range []string{
"## GROUNDING & PROVENANCE",
- "attribute it briefly so the user can trust and verify it",
- "Do not present stale or cached context as current",
- "Keep attribution concise and inline",
+ "source observations, model hypotheses, proposed actions, executed operations and verified outcomes separately",
+ "historical evidence until a current observation supports it",
+ "Attribute measured facts to their source and observation time",
} {
if !strings.Contains(prompt, expected) {
t.Fatalf("expected provenance guidance %q in system prompt, got %q", expected, prompt)
@@ -1182,6 +1182,7 @@ func TestNonInteractiveProfileBlocksQuestionPersistsPairAndContinues(t *testing.
}
func TestInvestigationLoopRedactsProposalParamsEverywhereDurable(t *testing.T) {
+ const conclusion = "### Investigation Summary\nRestart proposed.\n\n### Root Cause\nThe cause is unknown. The earlier dependency rationale is unconfirmed.\n\n### Affected Resources\n`vm:42`.\n\n### Recommendation\nRestart pending.\n\n### Conclusion\nNEEDS_ATTENTION: pending."
turn := 0
provider := &stubStreamingProvider{}
provider.chatStream = func(ctx context.Context, req providers.ChatRequest, callback providers.StreamCallback) error {
@@ -1214,7 +1215,7 @@ func TestInvestigationLoopRedactsProposalParamsEverywhereDurable(t *testing.T) {
}
}
}
- callback(providers.StreamEvent{Type: "content", Data: providers.ContentEvent{Text: "### Investigation Summary\nRestart proposed.\n\n### Root Cause\nDependency failure.\n\n### Affected Resources\n`vm:42`.\n\n### Recommendation\nRestart pending.\n\n### Conclusion\nNEEDS_ATTENTION: pending."}})
+ callback(providers.StreamEvent{Type: "content", Data: providers.ContentEvent{Text: conclusion}})
callback(providers.StreamEvent{Type: "done", Data: providers.DoneEvent{}})
return nil
}
@@ -1239,12 +1240,20 @@ func TestInvestigationLoopRedactsProposalParamsEverywhereDurable(t *testing.T) {
loop.successfulEvidenceCalls = 1
var streamedRawParam bool
+ var streamedConclusion strings.Builder
messages, err := loop.ExecuteWithTools(
context.Background(),
"session-investigation",
[]Message{{Role: "user", Content: "investigate finding f-9"}},
nil,
func(event StreamEvent) {
+ if event.Type == "content" {
+ var data ContentData
+ if err := json.Unmarshal(event.Data, &data); err != nil {
+ t.Errorf("decode content: %v", err)
+ }
+ streamedConclusion.WriteString(data.Text)
+ }
if strings.Contains(string(event.Data), "graceful") {
streamedRawParam = true
}
@@ -1257,6 +1266,19 @@ func TestInvestigationLoopRedactsProposalParamsEverywhereDurable(t *testing.T) {
t.Fatal("stream events must never expose proposal parameter values")
}
+ if streamedConclusion.String() != conclusion {
+ t.Fatalf("stream changed the model's uncertain diagnosis: %q", streamedConclusion.String())
+ }
+ var durableConclusion strings.Builder
+ for _, msg := range messages {
+ if msg.Role == "assistant" {
+ durableConclusion.WriteString(msg.Content)
+ }
+ }
+ if durableConclusion.String() != conclusion {
+ t.Fatalf("transcript changed the model's uncertain diagnosis: %q", durableConclusion.String())
+ }
+
// The durable transcript keeps the call but with redacted params.
var sawCall bool
for _, msg := range messages {
@@ -1287,6 +1309,80 @@ func TestInvestigationLoopRedactsProposalParamsEverywhereDurable(t *testing.T) {
}
}
+func TestInvestigationServicePreservesUncertainDiagnosisAfterProposal(t *testing.T) {
+ const conclusion = "### Root Cause\nUnknown. A restart may restore service, but the earlier dependency explanation is unconfirmed.\n\n### Recommendation\nProposal pending approval. No action has executed."
+ store, err := NewSessionStore(t.TempDir())
+ if err != nil {
+ t.Fatal(err)
+ }
+ turn := 0
+ provider := &stubStreamingProvider{}
+ provider.chatStream = func(ctx context.Context, req providers.ChatRequest, callback providers.StreamCallback) error {
+ turn++
+ switch turn {
+ case 1:
+ callback(providers.StreamEvent{Type: "done", Data: providers.DoneEvent{ToolCalls: []providers.ToolCall{{
+ ID: "evidence-1", Name: agentcapabilities.PulseQueryToolName,
+ Input: map[string]interface{}{"action": "health"},
+ }}}})
+ case 2:
+ callback(providers.StreamEvent{Type: "done", Data: providers.DoneEvent{ToolCalls: []providers.ToolCall{{
+ ID: "proposal-1", Name: agentcapabilities.PatrolProposeActionToolName,
+ Input: map[string]interface{}{"resource_id": "vm:42", "causal_resource_id": "vm:dependency", "capability_name": "restart", "reason": "The dependency may have stopped."},
+ }}}})
+ default:
+ callback(providers.StreamEvent{Type: "content", Data: providers.ContentEvent{Text: conclusion}})
+ callback(providers.StreamEvent{Type: "done", Data: providers.DoneEvent{}})
+ }
+ return nil
+ }
+ service := &Service{
+ started: true, sessions: store,
+ executor: tools.NewPulseToolExecutor(tools.ExecutorConfig{StateProvider: &mockStateProvider{}}),
+ cfg: &config.AIConfig{PatrolModel: "mock:model", ControlLevel: config.ControlLevelReadOnly},
+ patrolProviderFactory: func(string) (providers.StreamingProvider, error) { return provider, nil },
+ }
+ var streamed strings.Builder
+ result, err := service.ExecuteInvestigationStream(context.Background(), InvestigationRunRequest{
+ SessionID: "uncertain-investigation", Prompt: "Investigate the issue", SystemPrompt: "Investigate using current evidence.",
+ MaxTurns: 5, MaxEvidenceCalls: 3, ResourceType: "vm",
+ Identity: tools.ProposalIdentity{FindingID: "finding-1", InvestigationID: "investigation-1"},
+ Catalog: func(context.Context, string) ([]ur.ResourceCapability, error) {
+ return []ur.ResourceCapability{{Name: "restart"}}, nil
+ },
+ }, func(event StreamEvent) {
+ if event.Type == "content" {
+ var data ContentData
+ if err := json.Unmarshal(event.Data, &data); err != nil {
+ t.Errorf("content decode: %v", err)
+ }
+ streamed.WriteString(data.Text)
+ }
+ })
+ if err != nil {
+ t.Fatal(err)
+ }
+ if result.Proposal == nil || result.Proposal.CausalResourceID != "vm:dependency" {
+ t.Fatalf("proposal attribution lost: %+v", result.Proposal)
+ }
+ if result.Content != conclusion || streamed.String() != conclusion {
+ t.Fatalf("service promoted proposal rationale into diagnosis: result=%q stream=%q", result.Content, streamed.String())
+ }
+ messages, err := store.GetMessages("uncertain-investigation")
+ if err != nil {
+ t.Fatal(err)
+ }
+ var persisted strings.Builder
+ for _, msg := range messages {
+ if msg.Role == "assistant" {
+ persisted.WriteString(msg.Content)
+ }
+ }
+ if persisted.String() != conclusion {
+ t.Fatalf("stored diagnosis mutated: %q", persisted.String())
+ }
+}
+
func TestProposalRawInputOverrideNeverLeaksThroughProgressEvents(t *testing.T) {
var events []string
callback := func(event StreamEvent) { events = append(events, string(event.Data)) }
diff --git a/internal/ai/providers/subscription_agent.go b/internal/ai/providers/subscription_agent.go
index 1a41320a5..ee952d541 100644
--- a/internal/ai/providers/subscription_agent.go
+++ b/internal/ai/providers/subscription_agent.go
@@ -855,8 +855,8 @@ func decodeSubscriptionAgentTurn(agent SubscriptionAgent, raw []byte) (subscript
// all of its local tools are intentionally disabled. The CLI then fabricates a
// "No such tool available" result and lets the model continue, which destroys
// Pulse's one-tool-turn-at-a-time provider protocol. The stream preserves the
-// original intended call before that local error. Route the first declared
-// Pulse call back through Pulse's executor; never execute it inside Claude.
+// original intended calls before that local error. Route the first declared
+// Pulse tool batch back through Pulse's executor; never execute it inside Claude.
func decodeClaudeSubscriptionAgentResponse(req ChatRequest, raw []byte) (subscriptionAgentTurn, error) {
if !bytes.Contains(raw, []byte{'\n'}) {
return decodeSubscriptionAgentTurn(SubscriptionAgentClaude, raw)
@@ -869,8 +869,7 @@ func decodeClaudeSubscriptionAgentResponse(req ChatRequest, raw []byte) (subscri
var terminal claudePrintResponse
var terminalFound bool
- var content strings.Builder
- var routed *subscriptionAgentToolCall
+ var routed []subscriptionAgentToolCall
for _, line := range bytes.Split(raw, []byte{'\n'}) {
line = bytes.TrimSpace(line)
if len(line) == 0 {
@@ -882,12 +881,9 @@ func decodeClaudeSubscriptionAgentResponse(req ChatRequest, raw []byte) (subscri
}
switch event.Type {
case "assistant":
+ var batch []subscriptionAgentToolCall
for _, block := range event.Message.Content {
switch block.Type {
- case "text":
- if routed == nil && strings.TrimSpace(block.Text) != "" {
- content.WriteString(block.Text)
- }
case "tool_use":
if block.Name == "StructuredOutput" {
continue
@@ -895,12 +891,12 @@ func decodeClaudeSubscriptionAgentResponse(req ChatRequest, raw []byte) (subscri
if _, ok := allowed[block.Name]; !ok {
return subscriptionAgentTurn{}, fmt.Errorf("Claude subscription agent attempted undeclared native tool %q", block.Name)
}
- if routed == nil {
- call := subscriptionAgentToolCall{ID: block.ID, Name: block.Name, Input: block.Input}
- routed = &call
- }
+ batch = append(batch, subscriptionAgentToolCall{ID: block.ID, Name: block.Name, Input: block.Input})
}
}
+ if len(routed) == 0 && len(batch) > 0 {
+ routed = batch
+ }
case "result":
if err := json.Unmarshal(line, &terminal); err != nil {
return subscriptionAgentTurn{}, fmt.Errorf("decode Claude subscription result: %w", err)
@@ -919,8 +915,7 @@ func decodeClaudeSubscriptionAgentResponse(req ChatRequest, raw []byte) (subscri
}
if routed != nil {
return subscriptionAgentTurn{
- Content: content.String(),
- RawToolCalls: []subscriptionAgentToolCall{*routed},
+ RawToolCalls: routed,
InputTokens: terminal.Usage.InputTokens,
OutputTokens: terminal.Usage.OutputTokens,
}, nil
@@ -1102,6 +1097,10 @@ func validateSubscriptionAgentTurn(req ChatRequest, turn *subscriptionAgentTurn)
}
}
if len(turn.RawToolCalls) > 0 {
+ // Routing turns contain only validated calls in this transport contract.
+ // CLI narration can include serialized protocol or synthetic local tool
+ // errors. Neither is a user-facing answer or infrastructure evidence.
+ turn.Content = ""
turn.StopReason = "tool_use"
} else {
turn.StopReason = "end_turn"
diff --git a/internal/ai/providers/subscription_agent_test.go b/internal/ai/providers/subscription_agent_test.go
index 4dd8d2ed2..a613e71e4 100644
--- a/internal/ai/providers/subscription_agent_test.go
+++ b/internal/ai/providers/subscription_agent_test.go
@@ -636,7 +636,7 @@ func TestDecodeClaudeSubscriptionAgentRoutesNativePulseToolAttempt(t *testing.T)
if err != nil {
t.Fatal(err)
}
- if turn.Content != "Checking logs next." || len(turn.RawToolCalls) != 1 {
+ if turn.Content != "" || len(turn.RawToolCalls) != 1 {
t.Fatalf("routed turn = %#v", turn)
}
call := turn.RawToolCalls[0]
@@ -664,6 +664,51 @@ func TestDecodeClaudeSubscriptionAgentRoutesOnlyFirstNativePulseToolAttempt(t *t
}
}
+func TestDecodeClaudeSubscriptionAgentPreservesFirstBatchWithoutProtocolText(t *testing.T) {
+ raw := []byte(strings.Join([]string{
+ `{"type":"assistant","message":{"content":[{"type":"text","text":"{\"content\":\"\",\"stop_reason\":\"tool_use\",\"tool_calls\":[]}"}]}}`,
+ `{"type":"assistant","message":{"content":[{"type":"tool_use","id":"read-1","name":"pulse_read","input":{"action":"logs"}},{"type":"tool_use","id":"query-1","name":"pulse_query","input":{"action":"health"}}]}}`,
+ `{"type":"user","message":{"content":[{"type":"tool_result","tool_use_id":"read-1","is_error":true,"content":"No such tool available"}]}}`,
+ `{"type":"assistant","message":{"content":[{"type":"tool_use","id":"read-2","name":"pulse_read","input":{"action":"exec"}}]}}`,
+ `{"type":"result","subtype":"success","structured_output":{"content":"Synthetic local failure","stop_reason":"end_turn","tool_calls":[]},"permission_denials":[],"usage":{"input_tokens":11,"output_tokens":17}}`,
+ }, "\n"))
+ req := ChatRequest{Tools: []Tool{{Name: "pulse_read"}, {Name: "pulse_query"}}}
+ turn, err := decodeClaudeSubscriptionAgentResponse(req, raw)
+ if err != nil {
+ t.Fatal(err)
+ }
+ if err := validateSubscriptionAgentTurn(req, &turn); err != nil {
+ t.Fatal(err)
+ }
+ if turn.Content != "" || turn.StopReason != "tool_use" || len(turn.ProviderToolCalls) != 2 {
+ t.Fatalf("routing turn leaked text or lost batch: %#v", turn)
+ }
+ if turn.ProviderToolCalls[0].ID != "read-1" || turn.ProviderToolCalls[1].ID != "query-1" {
+ t.Fatalf("original ordered batch replaced by local continuation: %#v", turn.ProviderToolCalls)
+ }
+}
+
+func TestSubscriptionAgentStructuredRoutingTurnHasNoAnswerContent(t *testing.T) {
+ req := ChatRequest{Tools: []Tool{{Name: "pulse_query"}}}
+ turn := subscriptionAgentTurn{
+ Content: `{"content":"","stop_reason":"tool_use","tool_calls":[{"id":"q1","name":"pulse_query"}]}`,
+ RawToolCalls: []subscriptionAgentToolCall{{ID: "q1", Name: "pulse_query", Input: map[string]interface{}{"action": "health"}}},
+ }
+ if err := validateSubscriptionAgentTurn(req, &turn); err != nil {
+ t.Fatal(err)
+ }
+ if turn.Content != "" || len(turn.ProviderToolCalls) != 1 {
+ t.Fatalf("structured routing turn leaked transport content: %#v", turn)
+ }
+ final := subscriptionAgentTurn{Content: "The cause is unknown. Read access is unavailable."}
+ if err := validateSubscriptionAgentTurn(req, &final); err != nil {
+ t.Fatal(err)
+ }
+ if final.Content != "The cause is unknown. Read access is unavailable." || final.StopReason != "end_turn" {
+ t.Fatalf("final answer changed: %#v", final)
+ }
+}
+
func TestDecodeClaudeSubscriptionAgentRejectsUndeclaredNativeToolAttempt(t *testing.T) {
raw := []byte(strings.Join([]string{
`{"type":"assistant","message":{"content":[{"type":"tool_use","id":"toolu-1","name":"Bash","input":{"command":"true"}}]}}`,
diff --git a/internal/ai/tools/data_types.go b/internal/ai/tools/data_types.go
index 594b2eb28..24b3f76c7 100644
--- a/internal/ai/tools/data_types.go
+++ b/internal/ai/tools/data_types.go
@@ -2026,21 +2026,25 @@ func (r PhysicalDisksResponse) NormalizeCollections() PhysicalDisksResponse {
// PhysicalDiskSummary summarizes a physical disk with SMART health info
type PhysicalDiskSummary struct {
- ID string `json:"id"`
- Node string `json:"node"`
- Instance string `json:"instance"`
- DevPath string `json:"dev_path"`
- Model string `json:"model,omitempty"`
- Serial string `json:"serial,omitempty"`
- WWN string `json:"wwn,omitempty"`
- Type string `json:"type"` // nvme, sata, sas
- SizeBytes int64 `json:"size_bytes"`
- Health string `json:"health"` // PASSED, FAILED, UNKNOWN
- Wearout *int `json:"wearout,omitempty"` // SSD wear percentage (0-100), nil when unavailable
- Temperature *int `json:"temperature,omitempty"` // Celsius, nil when unavailable
- RPM *int `json:"rpm,omitempty"` // 0 for SSDs, nil when unavailable
- Used string `json:"used,omitempty"`
- LastChecked time.Time `json:"last_checked,omitempty"`
+ ID string `json:"id"`
+ Node string `json:"node"`
+ Instance string `json:"instance"`
+ DevPath string `json:"dev_path"`
+ Model string `json:"model,omitempty"`
+ Serial string `json:"serial,omitempty"`
+ WWN string `json:"wwn,omitempty"`
+ Type string `json:"type"` // nvme, sata, sas
+ SizeBytes int64 `json:"size_bytes"`
+ Health string `json:"health"` // Device-reported SMART result: PASSED, FAILED, UNKNOWN
+ Status unifiedresources.ResourceStatus `json:"status,omitempty"`
+ SourceStatus map[unifiedresources.DataSource]unifiedresources.SourceStatus `json:"source_status,omitempty"`
+ Risk *unifiedresources.PhysicalDiskRisk `json:"risk,omitempty"`
+ SMART *unifiedresources.SMARTMeta `json:"smart,omitempty"`
+ LifeRemainingPercent *int `json:"life_remaining_percent,omitempty"` // SSD life remaining percent (0-100), nil when unavailable
+ Temperature *int `json:"temperature,omitempty"` // Celsius, nil when unavailable
+ RPM *int `json:"rpm,omitempty"` // 0 for SSDs, nil when unavailable
+ Used string `json:"used,omitempty"`
+ LastChecked time.Time `json:"last_checked,omitempty"`
}
// ========== Host RAID Types ==========
diff --git a/internal/ai/tools/physical_disk_evidence_test.go b/internal/ai/tools/physical_disk_evidence_test.go
new file mode 100644
index 000000000..320ebff9c
--- /dev/null
+++ b/internal/ai/tools/physical_disk_evidence_test.go
@@ -0,0 +1,113 @@
+package tools
+
+import (
+ "context"
+ "encoding/json"
+ "testing"
+ "time"
+
+ "github.com/rcourtman/pulse-go-rewrite/internal/models"
+ "github.com/rcourtman/pulse-go-rewrite/internal/storagehealth"
+ "github.com/rcourtman/pulse-go-rewrite/internal/unifiedresources"
+ "github.com/stretchr/testify/require"
+)
+
+// A SMART pass does not erase a canonical warning. Both query entry points must
+// preserve its evidence, including the distinction between absent and zero.
+func TestPhysicalDiskEvidencePreservesCanonicalRisk(t *testing.T) {
+ zero, mediaErrors := int64(0), int64(12)
+ used := 37
+ risk := &unifiedresources.PhysicalDiskRisk{
+ Level: storagehealth.RiskWarning,
+ Reasons: []unifiedresources.PhysicalDiskRiskReason{{
+ Code: "media_errors", Severity: storagehealth.RiskWarning,
+ Summary: "12 media errors reported",
+ }},
+ }
+ observed := time.Date(2026, 9, 5, 21, 0, 0, 0, time.UTC)
+ resources := []unifiedresources.Resource{
+ {ID: "disk-warning", Name: "Warning disk", Type: unifiedresources.ResourceTypePhysicalDisk,
+ Status: unifiedresources.StatusWarning, ParentName: "node-one", LastSeen: observed,
+ PhysicalDisk: &unifiedresources.PhysicalDiskMeta{
+ Health: "PASSED", Wearout: 63, Risk: risk,
+ SMART: &unifiedresources.SMARTMeta{MediaErrors: &mediaErrors, UDMACRCErrors: &zero, PercentageUsed: &used},
+ }},
+ {ID: "disk-unknown", Name: "Unknown disk", Type: unifiedresources.ResourceTypePhysicalDisk,
+ Status: unifiedresources.StatusUnknown, ParentName: "node-one", LastSeen: observed,
+ PhysicalDisk: &unifiedresources.PhysicalDiskMeta{Health: "UNKNOWN", Wearout: -1}},
+ {ID: "disk-stale", Name: "Stale disk", Type: unifiedresources.ResourceTypePhysicalDisk,
+ Status: unifiedresources.StatusWarning, ParentName: "node-one", LastSeen: observed,
+ SourceStatus: map[unifiedresources.DataSource]unifiedresources.SourceStatus{
+ unifiedresources.SourceProxmox: {Status: "stale", LastSeen: observed},
+ },
+ PhysicalDisk: &unifiedresources.PhysicalDiskMeta{Health: "PASSED", Wearout: 63}},
+ }
+ executor := NewPulseToolExecutor(ExecutorConfig{
+ StateProvider: &mockStateProvider{state: models.StateSnapshot{}},
+ UnifiedResourceProvider: &stubUnifiedResourceProvider{resources: resources},
+ ControlLevel: ControlLevelReadOnly,
+ })
+ for _, tool := range []struct {
+ name string
+ args map[string]interface{}
+ key string
+ }{
+ {"pulse_metrics", map[string]interface{}{"type": "disks"}, "disks"},
+ {"pulse_query", map[string]interface{}{"action": "list", "type": "physical-disks"}, "physical_disks"},
+ } {
+ t.Run(tool.name, func(t *testing.T) {
+ result, err := executor.ExecuteTool(context.Background(), tool.name, tool.args)
+ require.NoError(t, err)
+ require.False(t, result.IsError, "%+v", result.Content)
+ var payload map[string]json.RawMessage
+ require.NoError(t, json.Unmarshal([]byte(result.Content[0].Text), &payload))
+ var disks []PhysicalDiskSummary
+ require.NoError(t, json.Unmarshal(payload[tool.key], &disks))
+ require.Len(t, disks, 3)
+ byID := map[string]PhysicalDiskSummary{}
+ for _, disk := range disks {
+ byID[disk.ID] = disk
+ }
+ warning := byID["disk-warning"]
+ require.Equal(t, "PASSED", warning.Health)
+ require.Equal(t, unifiedresources.StatusWarning, warning.Status)
+ require.Equal(t, risk, warning.Risk)
+ require.Equal(t, resources[0].PhysicalDisk.SMART, warning.SMART)
+ require.NotNil(t, warning.LifeRemainingPercent)
+ require.Equal(t, 63, *warning.LifeRemainingPercent)
+ require.Equal(t, observed, warning.LastChecked)
+ require.Equal(t, "node-one", warning.Node)
+ unknown := byID["disk-unknown"]
+ require.Equal(t, unifiedresources.StatusUnknown, unknown.Status)
+ require.Nil(t, unknown.SMART)
+ require.Nil(t, unknown.Risk)
+ require.Nil(t, unknown.LifeRemainingPercent)
+ stale := byID["disk-stale"]
+ require.Equal(t, "PASSED", stale.Health)
+ require.Equal(t, unifiedresources.StatusWarning, stale.Status)
+ require.Nil(t, stale.Risk)
+ require.Equal(t, resources[2].SourceStatus, stale.SourceStatus)
+ encoded, err := json.Marshal(warning.SMART)
+ require.NoError(t, err)
+ require.Contains(t, string(encoded), `"udmaCrcErrors":0`)
+ require.NotContains(t, string(encoded), `"pendingSectors"`)
+ })
+ }
+ for _, action := range []string{"get", "health"} {
+ t.Run(action, func(t *testing.T) {
+ result, err := executor.ExecuteTool(context.Background(), "pulse_query", map[string]interface{}{
+ "action": action, "resource_type": "physical-disk", "resource_id": "disk-warning",
+ })
+ require.NoError(t, err)
+ require.False(t, result.IsError, "%+v", result.Content)
+ var disk PhysicalDiskSummary
+ require.NoError(t, json.Unmarshal([]byte(result.Content[0].Text), &disk))
+ require.Equal(t, risk, disk.Risk)
+ require.Equal(t, resources[0].PhysicalDisk.SMART, disk.SMART)
+ require.Equal(t, 63, *disk.LifeRemainingPercent)
+ require.Contains(t, result.Content[0].Text, `"life_remaining_percent":63`)
+ require.NotContains(t, result.Content[0].Text, `"wearout"`)
+ })
+ }
+
+}
diff --git a/internal/ai/tools/proposal_capture.go b/internal/ai/tools/proposal_capture.go
index fd46cd46a..dd453d096 100644
--- a/internal/ai/tools/proposal_capture.go
+++ b/internal/ai/tools/proposal_capture.go
@@ -11,7 +11,6 @@ import (
"sync"
"github.com/rcourtman/pulse-go-rewrite/internal/actionplanner"
- "github.com/rcourtman/pulse-go-rewrite/internal/agentcapabilities"
unified "github.com/rcourtman/pulse-go-rewrite/internal/unifiedresources"
)
@@ -61,11 +60,9 @@ type CapturedProposal struct {
InvocationID string
Identity ProposalIdentity
ResourceID string
- // CausalResourceID is the exact canonical resource whose observed state
- // established the remediation rationale. It may equal ResourceID when the
- // affected/action target is also causal. Keeping this identity in the typed
- // proposal prevents the terminal narrative from collapsing a cross-resource
- // diagnosis back onto the symptom resource.
+ // CausalResourceID is the model's optional causal attribution. Empty means
+ // the cause is not established. Capture validates the action contract, not
+ // this diagnosis or its rationale.
CausalResourceID string
CapabilityName string
Params map[string]interface{}
@@ -114,14 +111,6 @@ type ProposalCapture struct {
proposal *CapturedProposal
fingerprint string
failedAttempts int
- evidence map[string]proposalEvidenceResource
-}
-
-type proposalEvidenceResource struct {
- ID string
- Name string
- Status string
- HealthcheckTargets []string
}
func (i ProposalIdentity) clone() ProposalIdentity {
@@ -136,123 +125,9 @@ func NewProposalCapture(identity ProposalIdentity, catalog ProposalCatalog) *Pro
return &ProposalCapture{
identity: identity.clone(),
catalog: catalog,
- evidence: make(map[string]proposalEvidenceResource),
}
}
-// RecordEvidence retains only the small canonical resource graph needed to
-// validate a later causal-resource claim. Raw evidence and arbitrary log text
-// are never retained here. The graph is populated from successful structured
-// query results after the model has observed them, so proposal validation can
-// reject a conclusion that contradicts the investigation's own evidence.
-func (c *ProposalCapture) RecordEvidence(toolName, content string) {
- if c == nil || strings.TrimSpace(toolName) != agentcapabilities.PulseQueryToolName {
- return
- }
- var decoded interface{}
- if err := json.Unmarshal([]byte(content), &decoded); err != nil {
- return
- }
- observed := make([]proposalEvidenceResource, 0)
- collectProposalEvidenceResources(decoded, &observed)
- if len(observed) == 0 {
- return
- }
-
- c.mu.Lock()
- defer c.mu.Unlock()
- if c.evidence == nil {
- c.evidence = make(map[string]proposalEvidenceResource)
- }
- for _, resource := range observed {
- resource.ID = unified.CanonicalResourceID(resource.ID)
- resource.Name = strings.TrimSpace(resource.Name)
- if resource.ID == "" || resource.Name == "" {
- continue
- }
- previous := c.evidence[resource.ID]
- if resource.Status == "" {
- resource.Status = previous.Status
- }
- if len(resource.HealthcheckTargets) == 0 {
- resource.HealthcheckTargets = previous.HealthcheckTargets
- }
- c.evidence[resource.ID] = resource
- }
-}
-
-func collectProposalEvidenceResources(value interface{}, resources *[]proposalEvidenceResource) {
- switch typed := value.(type) {
- case []interface{}:
- for _, item := range typed {
- collectProposalEvidenceResources(item, resources)
- }
- case map[string]interface{}:
- id, _ := typed["id"].(string)
- name, _ := typed["name"].(string)
- if strings.TrimSpace(id) != "" && strings.TrimSpace(name) != "" {
- status, _ := typed["status"].(string)
- if strings.TrimSpace(status) == "" {
- status, _ = typed["state"].(string)
- }
- resource := proposalEvidenceResource{ID: id, Name: name, Status: strings.TrimSpace(status)}
- if targets, ok := typed["healthcheck_targets"].([]interface{}); ok {
- for _, target := range targets {
- if text, ok := target.(string); ok && strings.TrimSpace(text) != "" {
- resource.HealthcheckTargets = append(resource.HealthcheckTargets, strings.TrimSpace(text))
- }
- }
- }
- *resources = append(*resources, resource)
- }
- for _, child := range typed {
- collectProposalEvidenceResources(child, resources)
- }
- }
-}
-
-func proposalEvidenceUnavailable(status string) bool {
- switch strings.ToLower(strings.TrimSpace(status)) {
- case "dead", "exited", "failed", "inactive", "not running", "offline", "stopped":
- return true
- default:
- return false
- }
-}
-
-func (c *ProposalCapture) validateCausalResource(causalResourceID string) error {
- causalResourceID = unified.CanonicalResourceID(causalResourceID)
- c.mu.Lock()
- defer c.mu.Unlock()
-
- claimed, ok := c.evidence[causalResourceID]
- if !ok || len(claimed.HealthcheckTargets) == 0 {
- return nil
- }
- for _, target := range claimed.HealthcheckTargets {
- matches := make([]proposalEvidenceResource, 0, 1)
- for _, candidate := range c.evidence {
- if strings.EqualFold(candidate.Name, target) {
- matches = append(matches, candidate)
- }
- }
- if len(matches) != 1 || !proposalEvidenceUnavailable(matches[0].Status) {
- continue
- }
- dependency := matches[0]
- return fmt.Errorf(
- "causal_resource_id %q conflicts with collected canonical evidence: health-check target %q is resource %q with status %q; investigate the implicated dependency and use its advertised capabilities before proposing bounded recovery on the affected resource",
- causalResourceID, target, dependency.ID, dependency.Status,
- )
- }
- return nil
-}
-
-// Capabilities resolves the current advertised action contract for a
-// canonical resource without changing proposal state. Investigations use
-// this to inspect a causal resource discovered after the run started; the
-// same catalog is used again when a proposal is validated, so lookup and
-// acceptance cannot drift.
func (c *ProposalCapture) Capabilities(ctx context.Context, resourceID string) ([]unified.ResourceCapability, error) {
if c == nil || c.catalog == nil {
return nil, errors.New("no capability catalog is wired for this investigation")
diff --git a/internal/ai/tools/proposal_capture_test.go b/internal/ai/tools/proposal_capture_test.go
index d75362bec..d68f042ae 100644
--- a/internal/ai/tools/proposal_capture_test.go
+++ b/internal/ai/tools/proposal_capture_test.go
@@ -146,58 +146,29 @@ func TestFailedAttemptsWithoutSuccessAreATypedError(t *testing.T) {
assert.Equal(t, 0, failed)
}
-func TestProposalRejectsCausalResourceContradictedByObservedHealthDependency(t *testing.T) {
- catalog := func(_ context.Context, resourceID string) ([]unified.ResourceCapability, error) {
- switch resourceID {
- case "app-container-client", "app-container-dependency":
- return []unified.ResourceCapability{{Name: "restart", MinimumApprovalLevel: unified.ApprovalAdmin}}, nil
- default:
- return nil, nil
- }
- }
- capture := NewProposalCapture(ProposalIdentity{InvestigationID: "inv-1"}, catalog)
+func TestProposalAllowsUncertainCauseWithoutInventingAttribution(t *testing.T) {
+ capture := NewProposalCapture(ProposalIdentity{InvestigationID: "inv-1"}, testProposalCatalog())
exec := newInvestigationExecutor(t, capture)
-
- exec.RecordProposalEvidence(agentcapabilities.PulseQueryToolName, `{
- "id":"app-container-client",
- "name":"client",
- "status":"running",
- "health":"unhealthy",
- "healthcheck_targets":["dependency"]
- }`)
- exec.RecordProposalEvidence(agentcapabilities.PulseQueryToolName, `{
- "docker":{"hosts":[{"containers":[
- {"id":"app-container-client","name":"client","state":"running","health":"unhealthy","healthcheck_targets":["dependency"]},
- {"id":"app-container-dependency","name":"dependency","state":"exited","health":"unhealthy"}
- ]}]}
- }`)
-
- wrong := map[string]interface{}{
- "resource_id": "app-container-client",
- "causal_resource_id": "app-container-client",
- "capability_name": "restart",
- "reason": "restart the unhealthy client",
- }
- rejected := executePropose(t, exec, "call-wrong", wrong)
- assert.True(t, rejected.IsError)
- assert.Contains(t, rejected.Content[0].Text, "app-container-dependency")
- assert.Contains(t, rejected.Content[0].Text, `status "exited"`)
-
- corrected := map[string]interface{}{
- "resource_id": "app-container-dependency",
- "causal_resource_id": "app-container-dependency",
- "capability_name": "restart",
- "reason": "restart the exited dependency required by the unhealthy client",
- }
- accepted := executePropose(t, exec, "call-corrected", corrected)
- assert.False(t, accepted.IsError)
-
+ args := proposeArgs()
+ delete(args, "causal_resource_id")
+ args["reason"] = "The service is stopped. A restart may restore service, but the cause is unknown."
+ result := executePropose(t, exec, "recovery-1", args)
+ require.False(t, result.IsError, "%+v", result)
proposal, failed, err := capture.Outcome()
require.NoError(t, err)
require.NotNil(t, proposal)
- assert.Equal(t, "app-container-dependency", proposal.ResourceID)
- assert.Equal(t, "app-container-dependency", proposal.CausalResourceID)
- assert.Equal(t, 1, failed)
+ assert.Empty(t, proposal.CausalResourceID)
+ assert.Equal(t, args["reason"], proposal.Reason)
+ assert.Equal(t, "vm:42", proposal.ResourceID)
+ assert.Zero(t, failed)
+ // Unknown cause does not weaken invocation integrity or grant execution.
+ changed := proposeArgs()
+ changed["reason"] = args["reason"]
+ result = executePropose(t, exec, "recovery-1", changed)
+ assert.True(t, result.IsError)
+ proposal, _, err = capture.Outcome()
+ assert.Nil(t, proposal)
+ assert.ErrorIs(t, err, ErrProposalIntegrity)
}
func TestSensitiveProposalParamsRejectedWithoutEcho(t *testing.T) {
@@ -515,7 +486,7 @@ func TestProposeActionSchemaExposesOnlyModelAuthoredFields(t *testing.T) {
[]string{"resource_id", "causal_resource_id", "capability_name", "params", "reason"},
mapKeys(tool.InputSchema.Properties),
)
- assert.Contains(t, tool.InputSchema.Required, "causal_resource_id")
+ assert.NotContains(t, tool.InputSchema.Required, "causal_resource_id")
return
}
t.Fatal("patrol_propose_action missing from investigation projection")
diff --git a/internal/ai/tools/tools_metrics.go b/internal/ai/tools/tools_metrics.go
index 810426106..8f987b69c 100644
--- a/internal/ai/tools/tools_metrics.go
+++ b/internal/ai/tools/tools_metrics.go
@@ -24,7 +24,7 @@ Types:
- temperatures: CPU, disk, and sensor temperatures from hosts
- network: Network interface statistics (rx/tx bytes, speed)
- diskio: Disk I/O statistics (read/write bytes, ops)
-- disks: Physical disk health (SMART, wearout, temperatures)
+- disks: Physical disk health, canonical status/risk reasons, source freshness, SMART counters and temperatures. health is the device-reported SMART result. risk carries hardware-risk reasons, while source_status describes collection freshness and may explain a status warning even when SMART is PASSED and no hardware risk is recorded. life_remaining_percent is percent life remaining (100 is unworn), while smart.percentageUsed is percent endurance consumed. Missing counters are unobserved, not zero.
- baselines: Learned normal behavior baselines for resources
- patterns: Detected operational patterns and predictions
@@ -762,34 +762,7 @@ func (e *PulseToolExecutor) executeListPhysicalDisks(_ context.Context, args map
continue
}
- summary := PhysicalDiskSummary{
- ID: r.ID,
- Node: node,
- DevPath: pd.DevPath,
- Model: pd.Model,
- Serial: pd.Serial,
- WWN: pd.WWN,
- Type: pd.DiskType,
- SizeBytes: pd.SizeBytes,
- Health: pd.Health,
- Used: pd.Used,
- LastChecked: r.LastSeen,
- }
-
- if pd.Wearout >= 0 {
- wearout := pd.Wearout
- summary.Wearout = &wearout
- }
- if pd.Temperature > 0 {
- temp := pd.Temperature
- summary.Temperature = &temp
- }
- if pd.RPM > 0 {
- rpm := pd.RPM
- summary.RPM = &rpm
- }
-
- disks = append(disks, summary)
+ disks = append(disks, physicalDiskSummaryFromResource(r))
}
if disks == nil {
diff --git a/internal/ai/tools/tools_propose.go b/internal/ai/tools/tools_propose.go
index c04f259d1..07dfcc626 100644
--- a/internal/ai/tools/tools_propose.go
+++ b/internal/ai/tools/tools_propose.go
@@ -48,7 +48,7 @@ func (e *PulseToolExecutor) registerProposeTools() {
Reference an advertised resource capability (see the resource's capability catalog) and fill only its declared parameters. Never place secrets in params - sensitive parameters are supplied by an operator at approval time.
-Identify the exact canonical causal resource established by the collected evidence. This may equal the action target, but when a dependency or peer caused the finding it must name that different resource. The reason must preserve the causal resource's observed state and explain the causal chain.
+If the evidence establishes a causal resource, identify it separately from the action target. Otherwise omit causal_resource_id. Explain the observed problem, why the proposed action should help, and any uncertainty. An action can address a symptom without establishing its underlying cause.
Submit at most one proposal per investigation. If no safe remediation exists, conclude without proposing.`,
InputSchema: InputSchema{
@@ -64,7 +64,7 @@ Submit at most one proposal per investigation. If no safe remediation exists, co
},
"causal_resource_id": {
Type: "string",
- Description: "Exact canonical resource ID whose observed state establishes the remediation rationale; may equal resource_id",
+ Description: "Optional canonical resource ID attributed as causal by the investigation. Omit when the cause is unknown.",
},
"params": {
Type: "object",
@@ -72,10 +72,10 @@ Submit at most one proposal per investigation. If no safe remediation exists, co
},
"reason": {
Type: "string",
- Description: "Why this action remediates the finding",
+ Description: "Observed evidence, why this action should help, and remaining uncertainty",
},
},
- Required: []string{"resource_id", "causal_resource_id", "capability_name", "reason"},
+ Required: []string{"resource_id", "capability_name", "reason"},
},
},
Handler: func(ctx context.Context, exec *PulseToolExecutor, args map[string]interface{}) (CallToolResult, error) {
@@ -164,13 +164,9 @@ func (e *PulseToolExecutor) executeProposeAction(ctx context.Context, args map[s
if params == nil {
params = map[string]interface{}{}
}
- if resourceID == "" || causalResourceID == "" || capabilityName == "" || reason == "" {
+ if resourceID == "" || capabilityName == "" || reason == "" {
capture.RecordFailedAttempt()
- return NewErrorResult(fmt.Errorf("resource_id, causal_resource_id, capability_name, and reason are required")), nil
- }
- if err := capture.validateCausalResource(causalResourceID); err != nil {
- capture.RecordFailedAttempt()
- return NewErrorResult(err), nil
+ return NewErrorResult(fmt.Errorf("resource_id, capability_name, and reason are required")), nil
}
if err := validateProposalAgainstCatalog(ctx, capture.catalog, resourceID, capabilityName, params); err != nil {
@@ -253,12 +249,3 @@ func stringArg(args map[string]interface{}, key string) string {
func (e *PulseToolExecutor) SetProposalCapture(capture *ProposalCapture) {
e.proposalCapture = capture
}
-
-// RecordProposalEvidence updates the request-local proposal validator with a
-// successful structured evidence result. It is a no-op outside investigations.
-func (e *PulseToolExecutor) RecordProposalEvidence(toolName, content string) {
- if e == nil || e.proposalCapture == nil {
- return
- }
- e.proposalCapture.RecordEvidence(toolName, content)
-}
diff --git a/internal/ai/tools/tools_query.go b/internal/ai/tools/tools_query.go
index 195052427..7654a9f8e 100644
--- a/internal/ai/tools/tools_query.go
+++ b/internal/ai/tools/tools_query.go
@@ -2265,6 +2265,8 @@ func canonicalQueryResourceType(resourceType string) string {
return "agent"
case "storage-pool":
return "storage"
+ case "physical_disk":
+ return "physical-disk"
default:
return strings.ToLower(strings.TrimSpace(resourceType))
}
@@ -3527,17 +3529,19 @@ func canonicalPhysicalDiskHost(resource unifiedresources.Resource) string {
func physicalDiskSummaryFromResource(resource unifiedresources.Resource) PhysicalDiskSummary {
pd := resource.PhysicalDisk
summary := PhysicalDiskSummary{
- ID: resource.ID,
- Node: canonicalPhysicalDiskHost(resource),
- DevPath: "",
- Model: "",
- Serial: "",
- WWN: "",
- Type: "",
- SizeBytes: 0,
- Health: "",
- Used: "",
- LastChecked: resource.LastSeen,
+ ID: resource.ID,
+ Node: canonicalPhysicalDiskHost(resource),
+ DevPath: "",
+ Model: "",
+ Serial: "",
+ WWN: "",
+ Type: "",
+ SizeBytes: 0,
+ Health: "",
+ Status: resource.Status,
+ SourceStatus: resource.SourceStatus,
+ Used: "",
+ LastChecked: resource.LastSeen,
}
if pd == nil {
return summary
@@ -3549,10 +3553,12 @@ func physicalDiskSummaryFromResource(resource unifiedresources.Resource) Physica
summary.Type = pd.DiskType
summary.SizeBytes = pd.SizeBytes
summary.Health = pd.Health
+ summary.Risk = pd.Risk
+ summary.SMART = pd.SMART
summary.Used = pd.Used
if pd.Wearout >= 0 {
wearout := pd.Wearout
- summary.Wearout = &wearout
+ summary.LifeRemainingPercent = &wearout
}
if pd.Temperature > 0 {
temp := pd.Temperature
@@ -4853,6 +4859,18 @@ func (e *PulseToolExecutor) executeGetResource(_ context.Context, args map[strin
return NewErrorResult(fmt.Errorf("resource_id is required")), nil
}
+ if resourceType == "physical-disk" {
+ if e.unifiedResourceProvider == nil {
+ return NewErrorResult(fmt.Errorf("physical disk inventory is unavailable")), nil
+ }
+ for _, resource := range e.unifiedResourceProvider.GetByType(unifiedresources.ResourceTypePhysicalDisk) {
+ if resource.ID == strings.TrimSpace(resourceID) {
+ return NewJSONResult(physicalDiskSummaryFromResource(resource)), nil
+ }
+ }
+ return NewErrorResult(fmt.Errorf("physical disk %q not found in canonical inventory", resourceID)), nil
+ }
+
rs, err := e.readStateForControl()
if err != nil {
return NewTextResult("State information not available."), nil
diff --git a/internal/ai/tools/tools_summarize.go b/internal/ai/tools/tools_summarize.go
index e322c2da4..b9a7b4ef7 100644
--- a/internal/ai/tools/tools_summarize.go
+++ b/internal/ai/tools/tools_summarize.go
@@ -18,7 +18,7 @@ func (e *PulseToolExecutor) registerSummarizeTools() {
e.registry.registerBuiltin(RegisteredTool{
Definition: Tool{
Name: agentcapabilities.PulseSummarizeToolName,
- Description: `Read retained metric evidence for one resource or a fleet over 24h, 7d, or 30d. Returns measured statistics, units, actual first/latest observation times, retained point counts and largest gaps. The requested window does not imply complete or fresh coverage. Means are unweighted means of returned retained points, which may already be retention aggregates.
+ Description: `Read retained metric evidence for one resource or a fleet over 24h, 7d, or 30d. Returns measured statistics, units, first/latest retained timestamps, retained point counts and largest timestamp gaps. The requested window does not imply complete or fresh coverage. Means are unweighted means of returned retained points, which may already be retention aggregates.
This tool reads metrics only. Alerts, findings, disk health, backup coverage and topology are not queried. Use the relevant tools for those sources before drawing conclusions about health or causes.
@@ -118,16 +118,18 @@ func (e *PulseToolExecutor) executeSummarize(ctx context.Context, args map[strin
// EvidenceScope states what was collected, independently from an empty result.
// In particular, an empty metrics map says nothing about alert or disk health.
type summarizeEvidenceScope struct {
- Source string `json:"source"`
- NotQueried []string `json:"not_queried"`
- Aggregation string `json:"aggregation"`
+ Source string `json:"source"`
+ NotQueried []string `json:"not_queried"`
+ Aggregation string `json:"aggregation"`
+ TimeSemantics string `json:"time_semantics"`
}
func retainedSummaryScope() summarizeEvidenceScope {
return summarizeEvidenceScope{
- Source: "retained_metrics",
- NotQueried: []string{"alerts", "findings", "disk_health", "backups", "topology"},
- Aggregation: "Mean and latest describe retained point values, which may be bucket averages. Min and max preserve recorded bucket extrema, with their bucket timestamps, so peaks can exceed the plotted averages. The mean is unweighted. First and last timestamps do not prove continuous coverage.",
+ Source: "retained_metrics",
+ NotQueried: []string{"alerts", "findings", "disk_health", "backups", "topology"},
+ Aggregation: "Mean and latest describe retained point values, which may be bucket averages. Min and max preserve recorded bucket extrema, with their bucket timestamps, so peaks can exceed the plotted averages. The mean is unweighted. First and last timestamps do not prove continuous coverage.",
+ TimeSemantics: "first_at, last_at and max_gap_seconds describe only returned point or retention-bucket timestamps inside the requested window. Buckets do not record sample completeness. These fields establish neither collection uptime nor why older observations are absent. Collector start/restart events and retention configuration are not queried.",
}
}
diff --git a/internal/ai/tools/tools_summarize_test.go b/internal/ai/tools/tools_summarize_test.go
index 821025b37..ae0bdebb2 100644
--- a/internal/ai/tools/tools_summarize_test.go
+++ b/internal/ai/tools/tools_summarize_test.go
@@ -152,6 +152,9 @@ func TestSummarizeTool_ResourceReturnsHeuristicNarrative(t *testing.T) {
if parsed.Scope.Source != "retained_metrics" || parsed.Evidence == nil || len(parsed.Evidence.Metrics) != 0 {
t.Fatalf("expected empty retained evidence, got %+v", parsed)
}
+ if parsed.Scope.TimeSemantics == "" {
+ t.Fatal("even empty results must disclose the retained timestamp semantics")
+ }
for _, field := range []string{"health_status", "health_message", "observations", "recommendations"} {
if strings.Contains(res.Content[0].Text, `"`+field+`"`) {
t.Fatalf("empty evidence invented %s: %s", field, res.Content[0].Text)
diff --git a/internal/api/resourceapi/resources_frontend_types_test.go b/internal/api/resourceapi/resources_frontend_types_test.go
index 8bc365dd1..a919ead40 100644
--- a/internal/api/resourceapi/resources_frontend_types_test.go
+++ b/internal/api/resourceapi/resources_frontend_types_test.go
@@ -405,7 +405,7 @@ func TestResourceListUsesCanonicalContractTypes(t *testing.T) {
// Agent-backed host resources should publish agent metrics targets.
var foundAgentHost *unified.Resource
for i := range resp.Data {
- if resp.Data[i].Type == "agent" {
+ if resp.Data[i].Type == "agent" && resp.Data[i].Agent != nil && resp.Data[i].Agent.AgentID == "agent-host-1" {
foundAgentHost = &resp.Data[i]
break
}
@@ -437,6 +437,9 @@ func TestResourceListUsesCanonicalContractTypes(t *testing.T) {
if foundNode.Proxmox.NodeName != "pve1" {
t.Fatalf("proxmox.nodeName = %q, want pve1", foundNode.Proxmox.NodeName)
}
+ if foundNode.MetricsTarget == nil || foundNode.MetricsTarget.ResourceType != "node" || foundNode.MetricsTarget.ResourceID != "instance-pve1" {
+ t.Fatalf("Proxmox-only resource must use its node storage coordinates, got %+v", foundNode.MetricsTarget)
+ }
// Test 2: canonical docker-host filter only returns docker-backed runtime resources.
dockerRec := httptest.NewRecorder()
diff --git a/internal/models/metrics_types_test.go b/internal/models/metrics_types_test.go
index 2b53b9354..7dc45cf5b 100644
--- a/internal/models/metrics_types_test.go
+++ b/internal/models/metrics_types_test.go
@@ -380,3 +380,16 @@ func TestPBSNodeMetricAvailabilityModelBoundary(t *testing.T) {
t.Fatal("wire input asserted monitoring-owned availability evidence")
}
}
+
+func TestPhysicalDiskSnapshotPreservesIndependentPollingCadence(t *testing.T) {
+ state := NewState()
+ state.UpdatePhysicalDisks("pve", []PhysicalDisk{{ID: "disk", Instance: "pve", ExpectedUpdateInterval: 15 * time.Minute}})
+ snapshot := state.GetSnapshot()
+ if len(snapshot.PhysicalDisks) != 1 || snapshot.PhysicalDisks[0].ExpectedUpdateInterval != 15*time.Minute {
+ t.Fatalf("snapshot lost cadence: %+v", snapshot.PhysicalDisks)
+ }
+ snapshot.PhysicalDisks[0].ExpectedUpdateInterval = time.Minute
+ if state.GetSnapshot().PhysicalDisks[0].ExpectedUpdateInterval != 15*time.Minute {
+ t.Fatal("snapshot mutation changed collector state")
+ }
+}
diff --git a/internal/models/models.go b/internal/models/models.go
index f0e979774..ee61bc8b9 100644
--- a/internal/models/models.go
+++ b/internal/models/models.go
@@ -2673,28 +2673,29 @@ type CephServiceStatus struct {
// PhysicalDisk represents a physical disk on a node
type PhysicalDisk struct {
- ID string `json:"id"` // "{instance}-{node}-{devpath}"
- Node string `json:"node"`
- Instance string `json:"instance"`
- DevPath string `json:"devPath"` // /dev/nvme0n1, /dev/sda
- Model string `json:"model"`
- Vendor string `json:"vendor,omitempty"`
- Serial string `json:"serial"`
- WWN string `json:"wwn"` // World Wide Name
- Type string `json:"type"` // nvme, sata, sas
- Controller string `json:"controller,omitempty"` // Controller association when reported
- Target string `json:"target,omitempty"` // Controller target/HCTL when reported
- Size int64 `json:"size"` // bytes
- Health string `json:"health"` // PASSED, FAILED, UNKNOWN
- Wearout int `json:"wearout"` // SSD wear metric from Proxmox (0-100, -1 when unavailable)
- Temperature int `json:"temperature"` // Celsius (if available)
- RPM int `json:"rpm"` // 0 for SSDs
- Used string `json:"used"` // Filesystem or partition usage
- StorageGroup string `json:"storageGroup"` // Pool/VG/array this disk belongs to (e.g. ZFS pool name); empty if not matched
- SmartAttributes *SMARTAttributes `json:"smartAttributes,omitempty"`
- IO *DiskIO `json:"io,omitempty"`
- Collection *diskinventory.CollectionStatus `json:"collection,omitempty"`
- LastChecked time.Time `json:"lastChecked"`
+ ID string `json:"id"` // "{instance}-{node}-{devpath}"
+ Node string `json:"node"`
+ Instance string `json:"instance"`
+ DevPath string `json:"devPath"` // /dev/nvme0n1, /dev/sda
+ Model string `json:"model"`
+ Vendor string `json:"vendor,omitempty"`
+ Serial string `json:"serial"`
+ WWN string `json:"wwn"` // World Wide Name
+ Type string `json:"type"` // nvme, sata, sas
+ Controller string `json:"controller,omitempty"` // Controller association when reported
+ Target string `json:"target,omitempty"` // Controller target/HCTL when reported
+ Size int64 `json:"size"` // bytes
+ Health string `json:"health"` // PASSED, FAILED, UNKNOWN
+ Wearout int `json:"wearout"` // SSD wear metric from Proxmox (0-100, -1 when unavailable)
+ Temperature int `json:"temperature"` // Celsius (if available)
+ RPM int `json:"rpm"` // 0 for SSDs
+ Used string `json:"used"` // Filesystem or partition usage
+ StorageGroup string `json:"storageGroup"` // Pool/VG/array this disk belongs to (e.g. ZFS pool name); empty if not matched
+ SmartAttributes *SMARTAttributes `json:"smartAttributes,omitempty"`
+ IO *DiskIO `json:"io,omitempty"`
+ Collection *diskinventory.CollectionStatus `json:"collection,omitempty"`
+ ExpectedUpdateInterval time.Duration `json:"-"` // Collector schedule, independent of general node polling
+ LastChecked time.Time `json:"lastChecked"`
}
// PBSInstance represents a Proxmox Backup Server instance
diff --git a/internal/monitoring/issue1595_collection_trust_test.go b/internal/monitoring/issue1595_collection_trust_test.go
index 8c6a79efc..83587605f 100644
--- a/internal/monitoring/issue1595_collection_trust_test.go
+++ b/internal/monitoring/issue1595_collection_trust_test.go
@@ -127,6 +127,7 @@ func TestIssue1595SASTopologySurvivesMergeRegistryAndReadState(t *testing.T) {
},
LastChecked: now,
}
+ providerDisk.ExpectedUpdateInterval = 15 * time.Minute
providerDisks = append(providerDisks, providerDisk)
want := providerDisk
@@ -191,6 +192,9 @@ func TestIssue1595SASTopologySurvivesMergeRegistryAndReadState(t *testing.T) {
t.Fatalf("disk %q metrics target = %q, want serial-stable target", want.Serial, view.MetricResourceID())
}
readBack := physicalDiskFromReadStateView(view)
+ if readBack.ExpectedUpdateInterval != want.ExpectedUpdateInterval {
+ t.Fatalf("disk cadence lost during canonical roundtrip: %s != %s", readBack.ExpectedUpdateInterval, want.ExpectedUpdateInterval)
+ }
assertIssue1595PhysicalDisk(t, readBack, want)
if readBack.Collection == nil ||
readBack.Collection.Serial.State != diskinventory.FieldAvailable ||
diff --git a/internal/monitoring/monitor.go b/internal/monitoring/monitor.go
index aa795b652..11a370b76 100644
--- a/internal/monitoring/monitor.go
+++ b/internal/monitoring/monitor.go
@@ -735,7 +735,7 @@ func physicalDiskFromReadStateView(view *unifiedresources.PhysicalDiskView) mode
return models.PhysicalDisk{Wearout: unifiedresources.WearoutUnreported}
}
- return models.PhysicalDisk{
+ disk := models.PhysicalDisk{
ID: view.ID(),
Node: view.Node(),
Instance: view.Instance(),
@@ -758,6 +758,11 @@ func physicalDiskFromReadStateView(view *unifiedresources.PhysicalDiskView) mode
Collection: diskinventory.CloneStatus(view.Collection()),
LastChecked: view.LastSeen(),
}
+ if source, ok := view.SourceStatus(unifiedresources.SourceProxmox); ok {
+ disk.ExpectedUpdateInterval = time.Duration(source.ExpectedUpdateIntervalSeconds) * time.Second
+ }
+ return disk
+
}
func physicalDiskIOFromUnifiedMeta(in *unifiedresources.PhysicalDiskIOMeta) *models.DiskIO {
diff --git a/internal/monitoring/monitor_pve.go b/internal/monitoring/monitor_pve.go
index f6903e7a3..1394b9a20 100644
--- a/internal/monitoring/monitor_pve.go
+++ b/internal/monitoring/monitor_pve.go
@@ -1183,6 +1183,11 @@ func (m *Monitor) maybePollPhysicalDisksAsync(
Int("diskCount", len(allDisks)).
Int("preservedCount", len(existingDisksMap)-len(polledNodes)).
Msg("Updating physical disks in state")
+ // Keep the observation timestamp from the last successful read, including
+ // preserved records, and attach the independent collector schedule.
+ for i := range allDisks {
+ allDisks[i].ExpectedUpdateInterval = pollingInterval
+ }
m.state.UpdatePhysicalDisks(inst, allDisks)
}(instanceName, client, nodes, nodeEffectiveStatus, modelNodes)
}
diff --git a/internal/monitoring/monitor_pve_disk_fallback_test.go b/internal/monitoring/monitor_pve_disk_fallback_test.go
index ec43665b3..86a3c3303 100644
--- a/internal/monitoring/monitor_pve_disk_fallback_test.go
+++ b/internal/monitoring/monitor_pve_disk_fallback_test.go
@@ -162,6 +162,9 @@ func TestMaybePollPhysicalDisksAsync_AgentFallbackWhenDiskQueryFails(t *testing.
if disk.DevPath != "/dev/sda" || disk.Node != "node1" || disk.Instance != "pve1" {
t.Fatalf("unexpected fallback disk identity: %+v", disk)
}
+ if disk.ExpectedUpdateInterval != 5*time.Minute {
+ t.Fatalf("default collector cadence lost: %s", disk.ExpectedUpdateInterval)
+ }
if disk.Health != "PASSED" || disk.Temperature != 30 {
t.Fatalf("unexpected fallback disk data: %+v", disk)
}
diff --git a/internal/unifiedresources/adapters.go b/internal/unifiedresources/adapters.go
index 935baf82f..5a8287413 100644
--- a/internal/unifiedresources/adapters.go
+++ b/internal/unifiedresources/adapters.go
@@ -2224,10 +2224,13 @@ func resourceFromPhysicalDisk(disk models.PhysicalDisk) (Resource, ResourceIdent
}
resource := Resource{
- Type: ResourceTypePhysicalDisk,
- Name: name,
- Status: physicalDiskStatus(disk.Model, disk.Health, assessment),
- LastSeen: disk.LastChecked,
+ Type: ResourceTypePhysicalDisk,
+ Name: name,
+ Status: physicalDiskStatus(disk.Model, disk.Health, assessment),
+ LastSeen: disk.LastChecked,
+ SourceStatus: map[DataSource]SourceStatus{
+ SourceProxmox: {ExpectedUpdateIntervalSeconds: int64(disk.ExpectedUpdateInterval / time.Second)},
+ },
UpdatedAt: time.Now().UTC(),
Metrics: metricsFromPhysicalDisk(disk),
PhysicalDisk: pdMeta,
diff --git a/internal/unifiedresources/clone_test.go b/internal/unifiedresources/clone_test.go
index c6d7ce567..5a0dd1ddd 100644
--- a/internal/unifiedresources/clone_test.go
+++ b/internal/unifiedresources/clone_test.go
@@ -155,10 +155,13 @@ func TestCloneResource_MutateSourceStatusMap(t *testing.T) {
original := &Resource{
ID: "r-1",
SourceStatus: map[DataSource]SourceStatus{
- SourceProxmox: {Status: "online"},
+ SourceProxmox: {Status: "online", ExpectedUpdateIntervalSeconds: 300},
},
}
cloned := cloneResource(original)
+ if cloned.SourceStatus[SourceProxmox].ExpectedUpdateIntervalSeconds != 300 {
+ t.Fatal("clone lost collector cadence")
+ }
cloned.SourceStatus[SourceDocker] = SourceStatus{Status: "online"}
if _, exists := original.SourceStatus[SourceDocker]; exists {
diff --git a/internal/unifiedresources/registry.go b/internal/unifiedresources/registry.go
index dea0c852b..bea193ae3 100644
--- a/internal/unifiedresources/registry.go
+++ b/internal/unifiedresources/registry.go
@@ -1613,6 +1613,11 @@ func (rr *ResourceRegistry) markStaleLocked(now time.Time, thresholds map[DataSo
if !ok {
threshold = 120 * time.Second
}
+ if status.ExpectedUpdateIntervalSeconds > 0 {
+ // Slow inventory polls have their own cadence. A source must miss
+ // two expected intervals before its retained observation is stale.
+ threshold = max(threshold, 2*time.Duration(status.ExpectedUpdateIntervalSeconds)*time.Second)
+ }
if status.LastSeen.IsZero() {
continue
}
@@ -2667,9 +2672,10 @@ func (rr *ResourceRegistry) ingest(source DataSource, sourceID string, resource
}
resource.Identity = identity
resource.Sources = []DataSource{source}
- resource.SourceStatus = map[DataSource]SourceStatus{
- source: {Status: sourceSightingStatus(resource.LastSeen), LastSeen: resource.LastSeen},
- }
+ sighting := resource.SourceStatus[source]
+ sighting.Status = sourceSightingStatus(resource.LastSeen)
+ sighting.LastSeen = resource.LastSeen
+ resource.SourceStatus = map[DataSource]SourceStatus{source: sighting}
resource.parentBySource = make(map[DataSource]string)
rr.setSourceParent(&resource, source, resource.ParentID)
@@ -3557,7 +3563,10 @@ func (rr *ResourceRegistry) mergeInto(existing *Resource, incoming Resource, sou
if existing.SourceStatus == nil {
existing.SourceStatus = make(map[DataSource]SourceStatus)
}
- existing.SourceStatus[source] = SourceStatus{Status: sourceSightingStatus(incoming.LastSeen), LastSeen: incoming.LastSeen}
+ sighting := incoming.SourceStatus[source]
+ sighting.Status = sourceSightingStatus(incoming.LastSeen)
+ sighting.LastSeen = incoming.LastSeen
+ existing.SourceStatus[source] = sighting
if incoming.LastSeen.After(existing.LastSeen) {
existing.LastSeen = incoming.LastSeen
@@ -3566,7 +3575,7 @@ func (rr *ResourceRegistry) mergeInto(existing *Resource, incoming Resource, sou
existing.UpdatedAt = now
existing.ParentID = rr.resolveCanonicalParentID(existing)
- existing.Status = chooseStatus(existing.Status, incoming.Status, source)
+ existing.Status = chooseStatus(existing.Status, incoming.Status, source, existing.Sources)
existing.Metrics = mergeMetrics(existing, existing.Metrics, incoming.Metrics, source, now, existing.SourceStatus, nil)
existing.Metrics = clearUnavailableSourceMemoryMetric(existing.Metrics, &incoming, source)
@@ -5483,14 +5492,18 @@ func sourcePriority(source DataSource) int {
}
}
-func chooseStatus(existing ResourceStatus, incoming ResourceStatus, source DataSource) ResourceStatus {
+func chooseStatus(existing ResourceStatus, incoming ResourceStatus, source DataSource, sources []DataSource) ResourceStatus {
if existing == "" || existing == StatusUnknown {
return incoming
}
- if sourcePriority(source) >= sourcePriority(SourceAgent) {
- return incoming
+ for _, observed := range sources {
+ if sourcePriority(observed) > sourcePriority(source) {
+ return existing
+ }
}
- return existing
+ // Refresh the highest-priority observed source, including a platform-only
+ // resource. Its first status must not remain sticky after recovery.
+ return incoming
}
func aggregateStatus(resource *Resource) ResourceStatus {
diff --git a/internal/unifiedresources/registry_test.go b/internal/unifiedresources/registry_test.go
index 9b0083a80..5f87f952f 100644
--- a/internal/unifiedresources/registry_test.go
+++ b/internal/unifiedresources/registry_test.go
@@ -6414,3 +6414,71 @@ func TestIssue1720ArrayVolumeMergesBareSerialWithPrefixedWWN(t *testing.T) {
t.Fatalf("sibling volumes collapsed into one identity: %+v", byPath)
}
}
+
+func TestPhysicalDiskFreshnessUsesCollectorCadenceAndRecovers(t *testing.T) {
+ now := time.Now().UTC()
+ for _, cadence := range []time.Duration{5 * time.Minute, 15 * time.Minute} {
+ t.Run(cadence.String(), func(t *testing.T) {
+ rr := NewRegistry(nil)
+ disk := models.PhysicalDisk{ID: "disk-source", Instance: "pve", Node: "node", DevPath: "/dev/nvme0n1", Serial: "serial-1", Health: "PASSED", Wearout: 63, LastChecked: now, ExpectedUpdateInterval: cadence}
+ rr.ingestPhysicalDisk(disk)
+ check := func(wantStatus ResourceStatus, wantFreshness string) Resource {
+ t.Helper()
+ resources := rr.ListByType(ResourceTypePhysicalDisk)
+ if len(resources) != 1 {
+ t.Fatalf("disks: %+v", resources)
+ }
+ got := resources[0]
+ if got.Status != wantStatus || got.SourceStatus[SourceProxmox].Status != wantFreshness {
+ t.Fatalf("status=%s freshness=%+v, want %s/%s", got.Status, got.SourceStatus, wantStatus, wantFreshness)
+ }
+ if got.SourceStatus[SourceProxmox].ExpectedUpdateIntervalSeconds != int64(cadence/time.Second) {
+ t.Fatalf("cadence lost: %+v", got.SourceStatus)
+ }
+ return got
+ }
+ rr.MarkStale(now.Add(cadence), map[DataSource]time.Duration{SourceProxmox: time.Minute})
+ original := check(StatusOnline, "online")
+ rr.MarkStale(now.Add(2*cadence+time.Second), map[DataSource]time.Duration{SourceProxmox: time.Minute})
+ check(StatusWarning, "stale")
+ disk.LastChecked = now.Add(2*cadence + 2*time.Second)
+ rr.ingestPhysicalDisk(disk)
+ recovered := check(StatusOnline, "online")
+ if recovered.ID != original.ID {
+ t.Fatal("freshness recovery changed identity")
+ }
+ disk.Wearout = 8
+ disk.LastChecked = disk.LastChecked.Add(cadence)
+ rr.ingestPhysicalDisk(disk)
+ warning := check(StatusWarning, "online")
+ if warning.PhysicalDisk.Risk == nil {
+ t.Fatal("fresh source hid measured disk risk")
+ }
+ })
+ }
+}
+
+func TestPlatformOnlyResourceStatusRecoversWithoutOverridingAgent(t *testing.T) {
+ now := time.Now().UTC()
+ rr := NewRegistry(nil)
+ identity := ResourceIdentity{MachineID: "0123456789abcdef0123456789abcdef", Hostnames: []string{"node-one"}}
+ source := Resource{Type: ResourceTypeAgent, Name: "node-one", Status: StatusOnline, LastSeen: now}
+ id := rr.ingest(SourceProxmox, "pve-node", source, identity)
+ rr.MarkStale(now.Add(3*time.Minute), nil)
+ source.LastSeen = now.Add(4 * time.Minute)
+ rr.ingest(SourceProxmox, "pve-node", source, identity)
+ got, _ := rr.Get(id)
+ if got.Status != StatusOnline {
+ t.Fatalf("fresh platform status remains %s", got.Status)
+ }
+ source.Status = StatusWarning
+ if mergedID := rr.ingest(SourceAgent, "agent-node", source, identity); mergedID != id {
+ t.Fatalf("fixture did not merge: %s != %s", mergedID, id)
+ }
+ source.Status = StatusOnline
+ rr.ingest(SourceProxmox, "pve-node", source, identity)
+ got, _ = rr.Get(id)
+ if got.Status != StatusWarning {
+ t.Fatalf("platform overwrote higher-priority agent status: %s", got.Status)
+ }
+}
diff --git a/internal/unifiedresources/types.go b/internal/unifiedresources/types.go
index 2c4055d04..f310a7bbe 100644
--- a/internal/unifiedresources/types.go
+++ b/internal/unifiedresources/types.go
@@ -290,9 +290,10 @@ const (
// SourceStatus describes the freshness of data from a source.
type SourceStatus struct {
- Status string `json:"status"` // online, stale, offline
- LastSeen time.Time `json:"lastSeen"`
- Error string `json:"error,omitempty"`
+ Status string `json:"status"` // online, stale, offline
+ LastSeen time.Time `json:"lastSeen"`
+ Error string `json:"error,omitempty"`
+ ExpectedUpdateIntervalSeconds int64 `json:"expectedUpdateIntervalSeconds,omitempty"` // Collector-authored cadence, zero uses source default
}
// ResourceIdentity holds identifiers used for matching.
diff --git a/internal/unifiedresources/views.go b/internal/unifiedresources/views.go
index 7ae22239b..952fe4318 100644
--- a/internal/unifiedresources/views.go
+++ b/internal/unifiedresources/views.go
@@ -2115,6 +2115,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 ""
diff --git a/pkg/db/slowlog.go b/pkg/db/slowlog.go
index 335f26a86..d4bd4f8ef 100644
--- a/pkg/db/slowlog.go
+++ b/pkg/db/slowlog.go
@@ -245,6 +245,12 @@ func (t *InstrumentedTx) Prepare(query string) (*InstrumentedStmt, error) {
return &InstrumentedStmt{Stmt: stmt, name: t.name, query: query}, nil
}
+// Stmt binds a database-level prepared statement to this transaction while
+// retaining the normal statement timing and slow-query instrumentation.
+func (t *InstrumentedTx) Stmt(stmt *InstrumentedStmt) *InstrumentedStmt {
+ return &InstrumentedStmt{Stmt: t.Tx.Stmt(stmt.Stmt), name: t.name, query: stmt.query}
+}
+
// Commit commits the transaction.
func (t *InstrumentedTx) Commit() error {
start := time.Now()
diff --git a/pkg/db/slowlog_test.go b/pkg/db/slowlog_test.go
index 8730d7703..06f5551b1 100644
--- a/pkg/db/slowlog_test.go
+++ b/pkg/db/slowlog_test.go
@@ -328,3 +328,43 @@ func slowQueryCount(t *testing.T, database, operation string) float64 {
}
return m.Counter.GetValue()
}
+
+func TestTransactionBoundPreparedStatementInstrumentation(t *testing.T) {
+ raw := openTestDB(t)
+ raw.SetMaxOpenConns(1)
+ idb := Wrap(raw, "bound_stmt_db")
+ if _, err := idb.Exec("CREATE TABLE t (val INTEGER)"); err != nil {
+ t.Fatal(err)
+ }
+ stmt, err := idb.Prepare("SELECT COUNT(*) FROM t")
+ if err != nil {
+ t.Fatal(err)
+ }
+ defer stmt.Close()
+ tx, err := idb.Begin()
+ if err != nil {
+ t.Fatal(err)
+ }
+ defer tx.Rollback()
+ if _, err := tx.Exec("INSERT INTO t VALUES (1)"); err != nil {
+ t.Fatal(err)
+ }
+ before := histogramSampleCount(t, "bound_stmt_db", "stmt_query_row")
+ bound := tx.Stmt(stmt)
+ var count int
+ if err := bound.QueryRow().Scan(&count); err != nil || count != 1 {
+ t.Fatalf("transaction snapshot: %d, %v", count, err)
+ }
+ if err := bound.Close(); err != nil {
+ t.Fatal(err)
+ }
+ if histogramSampleCount(t, "bound_stmt_db", "stmt_query_row")-before != 1 {
+ t.Fatal("bound statement lost query instrumentation")
+ }
+ if err := tx.Rollback(); err != nil {
+ t.Fatal(err)
+ }
+ if err := stmt.QueryRow().Scan(&count); err != nil || count != 0 {
+ t.Fatalf("database statement must survive bound close and rollback: %d, %v", count, err)
+ }
+}
diff --git a/pkg/metrics/store.go b/pkg/metrics/store.go
index 2f2812985..16f873f32 100644
--- a/pkg/metrics/store.go
+++ b/pkg/metrics/store.go
@@ -3,6 +3,7 @@
package metrics
import (
+ "context"
"database/sql"
"errors"
"fmt"
@@ -195,6 +196,11 @@ type Store struct {
db *pdb.InstrumentedDB
config StoreConfig
+ // Cache compiled presence SQL only. Results always come from the current
+ // read transaction. Bound the number of parameter-count shapes retained.
+ presenceMu sync.Mutex
+ presenceStatements map[string]*pdb.InstrumentedStmt
+
// Write buffer
bufferMu sync.Mutex
buffer []bufferedMetric
@@ -1337,20 +1343,31 @@ func normalizeMetricTypes(metricTypes []string) []string {
// buckets, and a coarser fallback is omitted if a preferred point overlaps it.
// A bucket's presence does not establish continuous underlying collection.
// All probes are scoped to the same series and requested timestamp window.
-func retainedQuerySQL(resourceType string, resourceIDs, metricTypes []string, start, end time.Time, tiers []Tier) (string, []interface{}) {
+func retainedQuerySQL(resourceType string, resourceIDs, metricTypes []string, start, end time.Time, stepSecs int64, tiers []Tier) (string, []interface{}) {
+ // Query dimensions fixed by the caller need not be decoded per point.
+ identityColumns := ""
+ if len(resourceIDs) != 1 {
+ identityColumns += "resource_id, "
+ }
+ if len(metricTypes) != 1 {
+ identityColumns += "metric_type, "
+ }
idSlots := strings.TrimSuffix(strings.Repeat("?,", len(resourceIDs)), ",")
metricClause := ""
if len(metricTypes) > 0 {
metricClause = " AND m.metric_type IN (" + strings.TrimSuffix(strings.Repeat("?,", len(metricTypes)), ",") + ")"
}
+ index := "idx_metrics_query_all"
+ if len(metricTypes) > 0 {
+ index = "idx_metrics_lookup"
+ }
+ // Keep tier and timestamp constraints in the index even when an ordering
+ // index could avoid a sort by walking unrelated retention tiers.
+ scope := `m.resource_type = ? AND m.resource_id IN (` + idSlots + `)` + metricClause + `
+ AND m.tier = ? AND m.timestamp >= ? AND m.timestamp <= ?`
branches := make([]string, 0, len(tiers))
var params []interface{}
- for i, tier := range tiers {
- branch := `SELECT m.resource_id, m.metric_type, m.timestamp, m.value,
- COALESCE(m.min_value, m.value) AS min_value, COALESCE(m.max_value, m.value) AS max_value
- FROM metrics AS m
- WHERE m.resource_type = ? AND m.resource_id IN (` + idSlots + `)` + metricClause + `
- AND m.tier = ? AND m.timestamp >= ? AND m.timestamp <= ?`
+ appendScopeParams := func(tier Tier) {
params = append(params, resourceType)
for _, id := range resourceIDs {
params = append(params, id)
@@ -1359,11 +1376,28 @@ func retainedQuerySQL(resourceType string, resourceIDs, metricTypes []string, st
params = append(params, metric)
}
params = append(params, string(tier), start.Unix(), end.Unix())
+ }
+ projection := identityColumns + `m.timestamp, m.value,
+ COALESCE(m.min_value, m.value) AS min_value, COALESCE(m.max_value, m.value) AS max_value`
+ directAggregate := stepSecs > 1 && len(tiers) == 1
+ if directAggregate {
+ // A snapshot containing one tier needs no intermediate projection.
+ // Aggregate its values directly, avoiding expression materialization
+ // for every input row while keeping the same bounded output.
+ projection = identityColumns + `(m.timestamp / ?) * ? + (? / 2) AS bucket_ts,
+ AVG(m.value), MIN(COALESCE(m.min_value, m.value)), MAX(COALESCE(m.max_value, m.value))`
+ params = append(params, stepSecs, stepSecs, stepSecs)
+ }
+ for i, tier := range tiers {
+ branch := `SELECT ` + projection + `
+ FROM metrics AS m INDEXED BY ` + index + ` WHERE ` + scope
+ appendScopeParams(tier)
for _, preferred := range tiers[:i] {
+ branch += ` AND NOT EXISTS (`
// Retention intervals nest on UTC minute/hour/day boundaries. Using the
// larger interval makes both finer and coarser overlap probes indexable.
bucket := max(tierBucketSeconds(tier), tierBucketSeconds(preferred))
- branch += fmt.Sprintf(` AND NOT EXISTS (
+ branch += fmt.Sprintf(`
SELECT 1 FROM metrics AS h
WHERE h.resource_type = m.resource_type AND h.resource_id = m.resource_id
AND h.metric_type = m.metric_type AND h.tier = ?
@@ -1374,7 +1408,23 @@ func retainedQuerySQL(resourceType string, resourceIDs, metricTypes []string, st
}
branches = append(branches, branch)
}
- return strings.Join(branches, " UNION ALL ") + " ORDER BY resource_id, metric_type, timestamp ASC", params
+ query := strings.Join(branches, " UNION ALL ")
+ if directAggregate {
+ query += " GROUP BY " + identityColumns + "bucket_ts ORDER BY " + identityColumns + "bucket_ts ASC"
+ } else if stepSecs > 1 {
+ // Reconcile first, then aggregate inside SQLite so a bounded chart does
+ // not allocate and scan every retained observation in Go.
+ query = `SELECT ` + identityColumns + `
+ (timestamp / ?) * ? + (? / 2) AS bucket_ts,
+ AVG(value), MIN(min_value), MAX(max_value)
+ FROM (` + query + `)
+ GROUP BY ` + identityColumns + `bucket_ts
+ ORDER BY ` + identityColumns + `bucket_ts ASC`
+ params = append([]interface{}{stepSecs, stepSecs, stepSecs}, params...)
+ } else {
+ query += " ORDER BY " + identityColumns + "timestamp ASC"
+ }
+ return query, params
}
func tierBucketSeconds(tier Tier) int64 {
@@ -1390,14 +1440,110 @@ func tierBucketSeconds(tier Tier) int64 {
}
}
+// The database owns statement closure. Once the bounded set is full, uncommon
+// parameter-count shapes run uncached. No result, tier presence or time window
+// is cached, and no eviction can close a statement another reader is binding.
+const maxRetainedPresenceStatements = 32
+
+func (s *Store) retainedPresenceStatement(query string) (*pdb.InstrumentedStmt, error) {
+ s.presenceMu.Lock()
+ defer s.presenceMu.Unlock()
+ if stmt := s.presenceStatements[query]; stmt != nil {
+ return stmt, nil
+ }
+ if len(s.presenceStatements) >= maxRetainedPresenceStatements {
+ return nil, nil
+ }
+ stmt, err := s.db.Prepare(query)
+ if err != nil {
+ return nil, err
+ }
+ if s.presenceStatements == nil {
+ s.presenceStatements = make(map[string]*pdb.InstrumentedStmt)
+ }
+ s.presenceStatements[query] = stmt
+ return stmt, nil
+}
+
func (s *Store) queryRetainedChunk(resourceType string, resourceIDs []string, metricTypes []string, start, end time.Time, stepSecs int64, tiers []Tier) (map[string]map[string][]MetricPoint, error) {
- sqlQuery, params := retainedQuerySQL(resourceType, resourceIDs, metricTypes, start, end, tiers)
+ // Determine which tiers exist in the same read snapshot used below. An
+ // all-raw series should cost a direct range read, not a UNION and an empty
+ // overlap probe for every observation. Presence does not imply coverage:
+ // every tier with any matching observations still participates.
+ // EXISTS needs only indexed identity and time columns. Do not build or
+ // order the value projection for each presence probe.
+ idSlots := strings.TrimSuffix(strings.Repeat("?,", len(resourceIDs)), ",")
+ presenceScope := "resource_type = ? AND resource_id IN (" + idSlots + ")"
+ index := "idx_metrics_query_all"
+ if len(metricTypes) > 0 {
+ presenceScope += " AND metric_type IN (" + strings.TrimSuffix(strings.Repeat("?,", len(metricTypes)), ",") + ")"
+ index = "idx_metrics_lookup"
+ }
+ presenceScope += " AND tier = ? AND timestamp >= ? AND timestamp <= ?"
+ check := "EXISTS (SELECT 1 FROM metrics INDEXED BY " + index + " WHERE " + presenceScope + ")"
+ checks := make([]string, len(tiers))
+ checkParams := make([]interface{}, 0, len(tiers)*(len(resourceIDs)+len(metricTypes)+4))
+ present := make([]bool, len(tiers))
+ checkDestinations := make([]interface{}, len(tiers))
+ for i, tier := range tiers {
+ checks[i] = check
+ checkParams = append(checkParams, resourceType)
+ for _, id := range resourceIDs {
+ checkParams = append(checkParams, id)
+ }
+ for _, metric := range metricTypes {
+ checkParams = append(checkParams, metric)
+ }
+ checkParams = append(checkParams, string(tier), start.Unix(), end.Unix())
+ checkDestinations[i] = &present[i]
+ }
+ presenceQuery := "SELECT " + strings.Join(checks, ", ")
+ statement, err := s.retainedPresenceStatement(presenceQuery)
+ if err != nil {
+ return nil, fmt.Errorf("prepare retained metrics presence: %w", err)
+ }
+ // Prepare before acquiring the transaction so a one-connection pool never
+ // waits for itself while preparing a database-level statement.
+ tx, err := s.db.BeginTx(context.Background(), &sql.TxOptions{ReadOnly: true})
+ if err != nil {
+ return nil, fmt.Errorf("begin retained metrics snapshot: %w", err)
+ }
+ defer tx.Rollback()
+ var presenceRow *sql.Row
+ if statement != nil {
+ bound := tx.Stmt(statement)
+ defer bound.Close()
+ presenceRow = bound.QueryRow(checkParams...)
+ } else {
+ presenceRow = tx.QueryRow(presenceQuery, checkParams...)
+ }
+ if err := presenceRow.Scan(checkDestinations...); err != nil {
+ return nil, fmt.Errorf("inspect retained metrics tiers: %w", err)
+ }
+ available := make([]Tier, 0, len(tiers))
+ for i, tier := range tiers {
+ if present[i] {
+ available = append(available, tier)
+ }
+ }
+ if len(available) == 0 {
+ return make(map[string]map[string][]MetricPoint), nil
+ }
+ tiers = available
+ // Fleet results are already ordered by series and time. Stream their
+ // display buckets to avoid SQLite's fleet-wide temporary GROUP BY tree.
+ // Single-resource charts aggregate in SQLite to bound rows crossing Go.
+ streamBuckets := stepSecs > 1 && len(resourceIDs) > 1
+ queryStep := stepSecs
+ if streamBuckets {
+ queryStep = 0
+ }
+ sqlQuery, params := retainedQuerySQL(resourceType, resourceIDs, metricTypes, start, end, queryStep, tiers)
// Retry on SQLITE_BUSY
var rows *sql.Rows
- var err error
for i := 0; i < 5; i++ {
- rows, err = s.db.Query(sqlQuery, params...)
+ rows, err = tx.Query(sqlQuery, params...)
if err == nil {
break
}
@@ -1411,121 +1557,74 @@ func (s *Store) queryRetainedChunk(resourceType string, resourceIDs []string, me
result := make(map[string]map[string][]MetricPoint, len(resourceIDs))
seriesCapacity := estimateQueryAllBatchSeriesCapacity(start, end, stepSecs)
- if stepSecs > 1 {
- type bucketAggregate struct {
- active bool
- resourceID string
- metricType string
- bucketStart int64
- sum float64
- min float64
- max float64
- count int
+
+ appendPoint := func(resourceID, metricType string, point MetricPoint) {
+ if result[resourceID] == nil {
+ result[resourceID] = make(map[string][]MetricPoint, 8)
}
-
- appendBucketPoint := func(resourceID, metricType string, point MetricPoint) {
- if _, exists := result[resourceID]; !exists {
- result[resourceID] = make(map[string][]MetricPoint, 8)
- }
- if _, exists := result[resourceID][metricType]; !exists {
- result[resourceID][metricType] = make([]MetricPoint, 0, seriesCapacity)
- }
- result[resourceID][metricType] = append(result[resourceID][metricType], point)
+ if result[resourceID][metricType] == nil {
+ result[resourceID][metricType] = make([]MetricPoint, 0, seriesCapacity)
}
-
- var aggregate bucketAggregate
- flushAggregate := func() {
- if !aggregate.active || aggregate.count == 0 {
- return
- }
- appendBucketPoint(aggregate.resourceID, aggregate.metricType, MetricPoint{
- Timestamp: time.Unix(aggregate.bucketStart+(stepSecs/2), 0),
- Value: aggregate.sum / float64(aggregate.count),
- Min: aggregate.min,
- Max: aggregate.max,
- })
- aggregate = bucketAggregate{}
- }
-
- for rows.Next() {
- var resourceID, metricType string
- var ts int64
- var value, minVal, maxVal float64
- if err := rows.Scan(&resourceID, &metricType, &ts, &value, &minVal, &maxVal); err != nil {
- log.Warn().Err(err).Msg("Failed to scan batch metric row")
- continue
- }
-
- bucketStart := (ts / stepSecs) * stepSecs
- if !aggregate.active {
- aggregate = bucketAggregate{
- active: true,
- resourceID: resourceID,
- metricType: metricType,
- bucketStart: bucketStart,
- sum: value,
- min: minVal,
- max: maxVal,
- count: 1,
- }
- continue
- }
-
- if aggregate.resourceID != resourceID || aggregate.metricType != metricType || aggregate.bucketStart != bucketStart {
- flushAggregate()
- aggregate = bucketAggregate{
- active: true,
- resourceID: resourceID,
- metricType: metricType,
- bucketStart: bucketStart,
- sum: value,
- min: minVal,
- max: maxVal,
- count: 1,
- }
- continue
- }
-
- aggregate.sum += value
- if minVal < aggregate.min {
- aggregate.min = minVal
- }
- if maxVal > aggregate.max {
- aggregate.max = maxVal
- }
- aggregate.count++
- }
-
- if err := rows.Err(); err != nil {
- return nil, err
- }
- flushAggregate()
- return result, nil
+ result[resourceID][metricType] = append(result[resourceID][metricType], point)
}
-
+ var bucketResource, bucketMetric string
+ var bucketStart int64
+ var bucketPoint MetricPoint
+ var bucketCount int
+ flushBucket := func() {
+ if bucketCount == 0 {
+ return
+ }
+ bucketPoint.Value /= float64(bucketCount)
+ bucketPoint.Timestamp = time.Unix(bucketStart+stepSecs/2, 0)
+ appendPoint(bucketResource, bucketMetric, bucketPoint)
+ bucketCount = 0
+ }
+ var resourceID, metricType string
+ var ts int64
+ var p MetricPoint
+ destinations := make([]interface{}, 0, 6)
+ if len(resourceIDs) == 1 {
+ resourceID = resourceIDs[0]
+ } else {
+ destinations = append(destinations, &resourceID)
+ }
+ if len(metricTypes) == 1 {
+ metricType = metricTypes[0]
+ } else {
+ destinations = append(destinations, &metricType)
+ }
+ destinations = append(destinations, &ts, &p.Value, &p.Min, &p.Max)
for rows.Next() {
- var resourceID, metricType string
- var ts int64
- var p MetricPoint
- if err := rows.Scan(&resourceID, &metricType, &ts, &p.Value, &p.Min, &p.Max); err != nil {
+ if err := rows.Scan(destinations...); err != nil {
log.Warn().Err(err).Msg("Failed to scan batch metric row")
continue
}
- p.Timestamp = time.Unix(ts, 0)
-
- if _, exists := result[resourceID]; !exists {
- result[resourceID] = make(map[string][]MetricPoint, 8)
+ if !streamBuckets {
+ p.Timestamp = time.Unix(ts, 0)
+ appendPoint(resourceID, metricType, p)
+ continue
}
- if _, exists := result[resourceID][metricType]; !exists {
- result[resourceID][metricType] = make([]MetricPoint, 0, seriesCapacity)
+ start := (ts / stepSecs) * stepSecs
+ if bucketCount > 0 && (bucketResource != resourceID || bucketMetric != metricType || bucketStart != start) {
+ flushBucket()
}
- result[resourceID][metricType] = append(result[resourceID][metricType], p)
+ if bucketCount == 0 {
+ bucketResource, bucketMetric, bucketStart = resourceID, metricType, start
+ bucketPoint = p
+ } else {
+ bucketPoint.Value += p.Value
+ bucketPoint.Min = min(bucketPoint.Min, p.Min)
+ bucketPoint.Max = max(bucketPoint.Max, p.Max)
+ }
+ bucketCount++
}
if err := rows.Err(); err != nil {
return nil, err
}
+ flushBucket()
return result, nil
}
diff --git a/pkg/metrics/store_query_plan_test.go b/pkg/metrics/store_query_plan_test.go
index 49524fcf3..76e7d0e1d 100644
--- a/pkg/metrics/store_query_plan_test.go
+++ b/pkg/metrics/store_query_plan_test.go
@@ -165,32 +165,66 @@ func TestQueryPlansUseIndexes(t *testing.T) {
}
// Query, QueryAll and both batch variants share this exact runtime SQL. The
-// display step is applied by the streaming reader after tier reconciliation.
+// display step is applied in SQLite after tier reconciliation.
func TestRetainedQueryPlansUseIndexes(t *testing.T) {
db := newPlanTestDB(t)
store := &Store{}
for _, window := range []time.Duration{time.Hour, 24 * time.Hour, 7 * 24 * time.Hour, 30 * 24 * time.Hour} {
for _, filtered := range []bool{false, true} {
- t.Run(fmt.Sprintf("%s/filtered_%v", window, filtered), func(t *testing.T) {
+ for _, step := range []int64{0, 60} {
+ t.Run(fmt.Sprintf("%s/filtered_%v/step_%d", window, filtered, step), func(t *testing.T) {
+ var metrics []string
+ if filtered {
+ metrics = []string{"cpu", "memory"}
+ }
+ end := time.Unix(2000000000, 0)
+ sql, args := retainedQuerySQL("vm", []string{"vm-1", "vm-2", "vm-3"}, metrics, end.Add(-window), end, step, store.tierFallbacks(window))
+ plan := explainQueryPlan(t, db, sql, args)
+ searches := 0
+ for _, line := range strings.Split(plan, "\n") {
+ if strings.Contains(line, "SCAN m ") || strings.Contains(line, "SCAN h ") {
+ t.Fatalf("unbounded metrics read: %s", plan)
+ }
+ if strings.Contains(line, "SEARCH m ") || strings.Contains(line, "SEARCH h ") {
+ searches++
+ for _, constraint := range []string{"resource_type=?", "resource_id=?", "tier=?", "timestamp>?", "timestamp"} {
+ if !strings.Contains(line, constraint) {
+ t.Fatalf("retained search missing indexed %s: %s", constraint, line)
+ }
+ }
+ }
+ }
+ n := len(store.tierFallbacks(window))
+ if searches != n*(n+1)/2 {
+ t.Fatalf("expected every tier and overlap probe indexed, got %d: %s", searches, plan)
+ }
+ })
+ }
+ }
+ }
+}
+
+// A single present tier uses direct aggregation. It must retain the same
+// bounded identity/tier/time lookup as the reconciled multi-tier query.
+func TestRetainedSingleTierQueryPlansUseIndexes(t *testing.T) {
+ db := newPlanTestDB(t)
+ end := time.Unix(2000000000, 0)
+ for _, tier := range []Tier{TierRaw, TierMinute, TierHourly, TierDaily} {
+ for _, filtered := range []bool{false, true} {
+ t.Run(fmt.Sprintf("%s/filtered_%v", tier, filtered), func(t *testing.T) {
var metrics []string
if filtered {
- metrics = []string{"cpu", "memory"}
+ metrics = []string{"cpu"}
}
- end := time.Unix(2000000000, 0)
- sql, args := retainedQuerySQL("vm", []string{"vm-1", "vm-2", "vm-3"}, metrics, end.Add(-window), end, store.tierFallbacks(window))
- plan := explainQueryPlan(t, db, sql, args)
- searches := 0
- for _, line := range strings.Split(plan, "\n") {
- if strings.Contains(line, "SCAN m ") || strings.Contains(line, "SCAN h ") {
- t.Fatalf("unbounded metrics read: %s", plan)
- }
- if strings.Contains(line, "SEARCH m ") || strings.Contains(line, "SEARCH h ") {
- searches++
- }
+ query, args := retainedQuerySQL("vm", []string{"vm-1"}, metrics, end.Add(-time.Hour), end, 60, []Tier{tier})
+ plan := explainQueryPlan(t, db, query, args)
+ if strings.Contains(plan, "SCAN m ") || !strings.Contains(plan, "SEARCH m ") {
+ t.Fatalf("single tier must use a bounded lookup: %s", plan)
}
- n := len(store.tierFallbacks(window))
- if searches != n*(n+1)/2 {
- t.Fatalf("expected every tier and overlap probe indexed, got %d: %s", searches, plan)
+ for _, constraint := range []string{"resource_type=?", "resource_id=?", "tier=?", "timestamp>?", "timestamp"} {
+ if !strings.Contains(plan, constraint) {
+ t.Fatalf("missing indexed %s: %s", constraint, plan)
+ }
}
})
}
diff --git a/pkg/metrics/store_tier_coverage_test.go b/pkg/metrics/store_tier_coverage_test.go
index c966bb132..871d30e1a 100644
--- a/pkg/metrics/store_tier_coverage_test.go
+++ b/pkg/metrics/store_tier_coverage_test.go
@@ -171,3 +171,52 @@ func TestStoreRollupPreservesRetainedExtrema(t *testing.T) {
})
}
}
+
+func TestRetainedPresenceStatementsReadCurrentSnapshot(t *testing.T) {
+ db := newPlanTestDB(t)
+ db.SetMaxOpenConns(1)
+ store := &Store{db: pdb.Wrap(db, "presence-snapshot")}
+ end := time.Unix(2000000040, 0)
+ start := end.Add(-time.Hour)
+ query := func(id string) []MetricPoint {
+ t.Helper()
+ points, err := store.Query("node", id, "cpu", start, end, 60)
+ if err != nil {
+ t.Fatal(err)
+ }
+ return points
+ }
+ if points := query("a"); len(points) != 0 {
+ t.Fatalf("empty inventory: %+v", points)
+ }
+ if _, err := db.Exec(`INSERT INTO metrics(resource_type,resource_id,metric_type,tier,timestamp,value) VALUES ('node','a','cpu','minute',?,17)`, end.Add(-time.Minute).Unix()); err != nil {
+ t.Fatal(err)
+ }
+ if points := query("a"); len(points) != 1 || points[0].Value != 17 {
+ t.Fatalf("newly present tier was hidden: %+v", points)
+ }
+ if points := query("b"); len(points) != 0 {
+ t.Fatalf("statement reused another resource's result: %+v", points)
+ }
+ if _, err := db.Exec("DELETE FROM metrics"); err != nil {
+ t.Fatal(err)
+ }
+ if points := query("a"); len(points) != 0 {
+ t.Fatalf("removed history remained present: %+v", points)
+ }
+ // Vary the parameter-count shape beyond the bound. Uncached shapes must
+ // remain functional without expanding the retained statement set.
+ for n := 1; n <= maxRetainedPresenceStatements+2; n++ {
+ ids := make([]string, n)
+ for i := range ids {
+ ids[i] = fmt.Sprintf("missing-%d", i)
+ }
+ got, err := store.QueryAllBatch("node", ids, start, end, 60)
+ if err != nil || len(got) != 0 {
+ t.Fatalf("shape %d: %+v, %v", n, got, err)
+ }
+ }
+ if len(store.presenceStatements) > maxRetainedPresenceStatements {
+ t.Fatalf("unbounded compiled SQL retention: %d", len(store.presenceStatements))
+ }
+}