fix(ai): preserve diagnostic evidence and proposal boundaries

Keep canonical disk risk, source freshness and retained history intact when
Assistant and Patrol gather evidence. Proposal acceptance validates an action
contract and must not rewrite uncertain conclusions as established root cause.

Preserve complete subscription tool batches without exposing routing envelopes
as answers. Keep wide answer tables readable and keyboard-scrollable on mobile.
Optimize retained tier reconciliation without discarding gaps or newer samples.

Record failed real-model diagnoses and outstanding autonomous qualification
separately from passing data-path and interface checks.
This commit is contained in:
rcourtman 2026-09-06 01:54:28 +01:00
parent 3a189f31d4
commit c5d2f56dda
46 changed files with 1387 additions and 708 deletions

View file

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

View file

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

View file

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

View file

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

View file

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

View file

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

View file

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

View file

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

View file

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

View file

@ -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(
'<table class="fixed inset-0" style="position:fixed" onclick="alert(1)"><tr><th>Metric</th><th>Value</th></tr><tr><td>Memory</td><td>88%</td></tr></table>',
);
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');

View file

@ -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 <a>; 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) => {

View file

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

View file

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

View file

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

View file

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

View file

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

View file

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

View file

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

View file

@ -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"}}]}}`,

View file

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

View file

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

View file

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

View file

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

View file

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

View file

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

View file

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

View file

@ -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.",
}
}

View file

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

View file

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

View file

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

View file

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

View file

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

View file

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

View file

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

View file

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

View file

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

View file

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

View file

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

View file

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

View file

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

View file

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

View file

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

View file

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

View file

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

View file

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

View file

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