mirror of
https://github.com/rcourtman/Pulse.git
synced 2026-10-03 04:38:48 +00:00
Remove redundant connected-dashboard projection and snapshot copies
List the continuity-aware registry once, decorate one owned host projection, and sort final frontend rows rather than a second estate-sized conversion array. Encode concrete frontend resources directly into immutable snapshot buffers, decoding every authoritative ID from those bytes; leave generic marshalers and all other fields on the existing path. Keep live freshness, alert and lifecycle semantics without a cache. Change-source: pulse-maintainer
This commit is contained in:
parent
db9e089691
commit
4b01e25cdf
9 changed files with 520 additions and 65 deletions
|
|
@ -15,6 +15,18 @@
|
|||
|
||||
## Purpose
|
||||
|
||||
### Continuity-aware broadcast read ownership
|
||||
|
||||
Broadcast state takes its single resource list from the current continuity-aware
|
||||
read view, not a preliminary registry clone discarded before that view is
|
||||
resolved. The resulting presentation slice is caller-owned; URL/health
|
||||
decoration cannot mutate nested registry data. Ignored/re-enrollment surfaces,
|
||||
parent identities and agent action targets remain unchanged, and a prior client
|
||||
baseline must survive later mutable updates. No freshness cache substitutes for
|
||||
live reads. `TestBroadcastProjectionListsRegistryOnceAndKeepsLiveChanges` checks
|
||||
same-freshness changes and ignored-host inventory; the previous-pipeline JSON
|
||||
oracle checks split-host identity/action/infrastructure composition.
|
||||
|
||||
### Import preview lifetime and setup authority
|
||||
|
||||
The node credential editor binds each monitored-system impact preview to the
|
||||
|
|
|
|||
|
|
@ -17,6 +17,27 @@
|
|||
|
||||
## Purpose
|
||||
|
||||
### Single-pass live broadcast projection — issue #2199
|
||||
|
||||
A frontend broadcast lists the continuity-aware read store once, coalesces host
|
||||
presentation once, and decorates its owned outer resource slice with current
|
||||
persisted URLs and freshly evaluated health. Generic conversion callers still
|
||||
coalesce their own input; an already-coalesced broadcast uses the prepared
|
||||
projection converter. Conversion sorts final frontend rows beside precomputed
|
||||
keys instead of retaining an estate-sized array of conversion inputs. This
|
||||
must preserve canonical identity, parent/action targets, catalogs, ordering,
|
||||
ignore surfaces and source isolation. No timestamp/generation cache is added:
|
||||
metrics, labels, alert state, metadata clears and time-sensitive health continue
|
||||
to be evaluated on each requested broadcast.
|
||||
|
||||
`TestBroadcastProjectionMatchesPreviousPipeline` compares full JSON with the
|
||||
pre-repair conversion/sort/coalesce pipeline on split hosts and mixed workloads.
|
||||
`TestBroadcastProjectionListsRegistryOnceAndKeepsLiveChanges` pins one registry
|
||||
list and verifies same-freshness metric/label changes, live alerts, ignored
|
||||
hosts and immutable prior projections. The 1,000-resource broadcast benchmark
|
||||
includes real registry cloning, projection, encoding, per-client deltas and
|
||||
queues with zero/one/four viewers, not persistence, transport or installed CPU.
|
||||
|
||||
### Synthetic machine identity survives restart
|
||||
|
||||
Linked mock hosts and Docker hosts derive their machine ID from their stable
|
||||
|
|
|
|||
|
|
@ -15,6 +15,23 @@
|
|||
|
||||
## Purpose
|
||||
|
||||
### Connected-dashboard snapshot ownership — issue #2199
|
||||
|
||||
Concrete frontend snapshots encode resources individually into immutable owned
|
||||
buffers, avoiding a whole-state resources encode/decode/copy round trip. Each
|
||||
resource identity is decoded from its encoded bytes, including every tail entry;
|
||||
source ID hints are never authoritative. Unknown shapes and custom marshalers
|
||||
retain full generic encoding. The remaining top-level fields use that same
|
||||
generic path, preserving future fields, nil/empty/omitempty semantics and keyed
|
||||
alert/infrastructure fallback. Reconnect baselines, removals, ordering, queue
|
||||
failure and REST-hydration frame limits keep their existing ownership.
|
||||
|
||||
Full-wire differential tests and fuzzing compare typed value/pointer snapshots,
|
||||
errors and encoded identities with generic encoding; mutation tests protect
|
||||
retained buffers. Complete broadcast benchmarks include projection and delta
|
||||
queues rather than treating an ID-only or component allocation result as field
|
||||
CPU relief. Installed connected/closed-dashboard CPU remains separate evidence.
|
||||
|
||||
The resource adapter fast delta path must remain content-equivalent to the full merge for explicit Proxmox memory withdrawal. An absent canonical metric plus incoming Proxmox usageUnavailable clears the old display value; store patch operations must emit that clear even when only the raw facet key changed. This bounded per-changed-row check must not introduce an estate-wide scan or defeat untouched-row identity preservation. Adapter tests cover both delta paths and store writes, including trusted-zero recovery and ordinary partial omission.
|
||||
|
||||
The PR #1935 log-level parser benchmark remains an unresolved environment-bound
|
||||
|
|
|
|||
92
internal/monitoring/broadcast_projection_reference_test.go
Normal file
92
internal/monitoring/broadcast_projection_reference_test.go
Normal file
|
|
@ -0,0 +1,92 @@
|
|||
package monitoring
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"github.com/rcourtman/pulse-go-rewrite/internal/models"
|
||||
"github.com/rcourtman/pulse-go-rewrite/internal/unifiedresources"
|
||||
"sort"
|
||||
"strings"
|
||||
)
|
||||
|
||||
// Retain the pre-repair projection as a content oracle, including its second
|
||||
// coalesce and input-first sorting. This is never production or persistence code.
|
||||
func convertResourcesForBroadcastReference(
|
||||
allResources []unifiedresources.Resource,
|
||||
metricsTargetResolvers ...MetricsTargetResourceStore,
|
||||
) ([]models.ResourceFrontend, broadcastResourceCatalogs) {
|
||||
if len(allResources) == 0 {
|
||||
return []models.ResourceFrontend{}, broadcastResourceCatalogs{}
|
||||
}
|
||||
allResources = attachBroadcastMetricsTargets(
|
||||
allResources,
|
||||
firstBroadcastMetricsTargetResolver(metricsTargetResolvers),
|
||||
)
|
||||
allResources = unifiedresources.CoalescePresentationHostResources(allResources)
|
||||
type broadcastResource struct {
|
||||
input models.ResourceConvertInput
|
||||
sortKey string
|
||||
resourceID string
|
||||
}
|
||||
|
||||
converted := make([]broadcastResource, 0, len(allResources))
|
||||
catalogs := broadcastResourceCatalogs{
|
||||
capabilities: make(map[string]json.RawMessage),
|
||||
policies: make(map[string]json.RawMessage),
|
||||
aiSafeSummaries: make(map[string]string),
|
||||
}
|
||||
for _, r := range allResources {
|
||||
input := monitorResourceToConvertInput(r)
|
||||
if len(input.Capabilities) > 0 {
|
||||
id := capabilityCatalogID(input.Capabilities)
|
||||
catalogs.capabilities[id] = input.Capabilities
|
||||
input.CapabilitiesRef = id
|
||||
input.Capabilities = nil
|
||||
}
|
||||
// Non-default policies and AI-safe summaries dedupe the same way:
|
||||
// estates carry a handful of distinct postures and templated summary
|
||||
// strings, so refs replace per-resource inline duplication.
|
||||
if len(input.Policy) > 0 {
|
||||
id := capabilityCatalogID(input.Policy)
|
||||
catalogs.policies[id] = input.Policy
|
||||
input.PolicyRef = id
|
||||
input.Policy = nil
|
||||
}
|
||||
if input.AISafeSummary != "" {
|
||||
id := capabilityCatalogID(json.RawMessage(input.AISafeSummary))
|
||||
catalogs.aiSafeSummaries[id] = input.AISafeSummary
|
||||
input.AISafeSummaryRef = id
|
||||
input.AISafeSummary = ""
|
||||
}
|
||||
sortKey := strings.ToLower(input.DisplayName)
|
||||
if sortKey == "" {
|
||||
sortKey = strings.ToLower(input.Name)
|
||||
}
|
||||
converted = append(converted, broadcastResource{
|
||||
input: input,
|
||||
sortKey: sortKey,
|
||||
resourceID: input.ID,
|
||||
})
|
||||
}
|
||||
|
||||
sort.Slice(converted, func(i, j int) bool {
|
||||
if converted[i].sortKey == converted[j].sortKey {
|
||||
return converted[i].resourceID < converted[j].resourceID
|
||||
}
|
||||
return converted[i].sortKey < converted[j].sortKey
|
||||
})
|
||||
|
||||
result := make([]models.ResourceFrontend, len(converted))
|
||||
for i, resource := range converted {
|
||||
result[i] = models.ConvertResourceToFrontend(resource.input)
|
||||
}
|
||||
if len(catalogs.capabilities) == 0 {
|
||||
catalogs.capabilities = nil
|
||||
}
|
||||
if len(catalogs.policies) == 0 {
|
||||
catalogs.policies = nil
|
||||
}
|
||||
if len(catalogs.aiSafeSummaries) == 0 {
|
||||
catalogs.aiSafeSummaries = nil
|
||||
}
|
||||
return result, catalogs
|
||||
}
|
||||
|
|
@ -1,6 +1,8 @@
|
|||
package monitoring
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"math"
|
||||
"os"
|
||||
"path/filepath"
|
||||
|
|
@ -2902,3 +2904,67 @@ func TestMockGuestChartHistorySkipsMemoryUsedForNonProxmoxGuests(t *testing.T) {
|
|||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestBroadcastProjectionMatchesPreviousPipeline(t *testing.T) {
|
||||
now := time.Now().UTC()
|
||||
resources := []unifiedresources.Resource{
|
||||
{ID: "agent-api", Type: unifiedresources.ResourceTypeAgent, Name: "tower", Status: unifiedresources.StatusOnline, LastSeen: now,
|
||||
Identity: unifiedresources.ResourceIdentity{Hostnames: []string{"tower.local"}},
|
||||
Sources: []unifiedresources.DataSource{unifiedresources.SourceProxmox},
|
||||
Proxmox: &unifiedresources.ProxmoxData{NodeName: "tower.local", ClusterName: "lab", PVEVersion: "9"}},
|
||||
{ID: "agent-runtime", Type: unifiedresources.ResourceTypeAgent, Name: "tower", Status: unifiedresources.StatusOnline, LastSeen: now,
|
||||
Identity: unifiedresources.ResourceIdentity{MachineID: "tower-machine", Hostnames: []string{"tower.local"}},
|
||||
Sources: []unifiedresources.DataSource{unifiedresources.SourceAgent},
|
||||
Agent: &unifiedresources.AgentData{AgentID: "tower-machine", Hostname: "tower.local", Platform: "linux", AgentVersion: "6.4.5"},
|
||||
Metrics: &unifiedresources.ResourceMetrics{CPU: &unifiedresources.MetricValue{Value: 25, Unit: "percent"}}},
|
||||
}
|
||||
for i, name := range []string{"zebra", "Alpha", "alpha", "", "尾"} {
|
||||
parent := "agent-api"
|
||||
resources = append(resources, unifiedresources.Resource{
|
||||
ID: fmt.Sprintf("vm-%d", i), Type: unifiedresources.ResourceTypeVM, Name: name, DisplayName: name,
|
||||
ParentID: &parent, Status: unifiedresources.StatusOnline, LastSeen: now,
|
||||
Proxmox: &unifiedresources.ProxmoxData{NodeName: "tower.local", VMID: i + 100},
|
||||
Sources: []unifiedresources.DataSource{unifiedresources.SourceProxmox},
|
||||
Metrics: &unifiedresources.ResourceMetrics{CPU: &unifiedresources.MetricValue{Value: 10, Unit: "percent"}},
|
||||
})
|
||||
}
|
||||
registry := unifiedresources.NewRegistry(nil)
|
||||
registry.IngestResources(resources)
|
||||
store := &broadcastProjectionCountingStore{MonitorAdapter: unifiedresources.NewMonitorAdapter(registry)}
|
||||
m := &Monitor{resourceStore: store}
|
||||
snapshot := models.EmptyStateSnapshot()
|
||||
snapshot.ActiveAlerts = []models.Alert{{ID: "warning", ResourceID: "agent-runtime", Level: "warning", Type: "cpu"}}
|
||||
snapshot.RemovedHostAgents = []models.RemovedHostAgent{{ID: "ignored-host", Hostname: "ignored.local", RemovedAt: now}}
|
||||
snapshot.RemovedDockerHosts = []models.RemovedDockerHost{{ID: "ignored-docker", Hostname: "docker.local", RemovedAt: now}}
|
||||
sourceBefore, _ := json.Marshal(store.GetAll())
|
||||
view := m.currentUnifiedStateView()
|
||||
previousResources := unifiedresources.CoalescePresentationHostResources(view.resources)
|
||||
previousResources = m.applyPersistedMetadataToUnifiedResources(previousResources)
|
||||
previousResources = unifiedresources.AttachResourceHealth(previousResources, resourceHealthAlerts(snapshot.ActiveAlerts), now)
|
||||
want := snapshot.ToFrontend()
|
||||
projected, catalogs := convertResourcesForBroadcastReference(previousResources, broadcastMetricsTargetResolver(view.readState))
|
||||
want.Resources = projected
|
||||
want.CapabilityCatalog = catalogs.capabilities
|
||||
want.PolicyCatalog = catalogs.policies
|
||||
want.AISafeSummaryCatalog = catalogs.aiSafeSummaries
|
||||
want.ConnectedInfrastructure = buildConnectedInfrastructure(previousResources, snapshot)
|
||||
if !view.freshness.IsZero() {
|
||||
want.LastUpdate = view.freshness.UnixMilli()
|
||||
}
|
||||
got := m.buildBroadcastFrontendStateFromSnapshot(snapshot)
|
||||
wantJSON, err := json.Marshal(want)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
gotJSON, err := json.Marshal(got)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if string(gotJSON) != string(wantJSON) {
|
||||
t.Fatalf("projection differs from previous pipeline\nwant=%s\ngot=%s", wantJSON, gotJSON)
|
||||
}
|
||||
sourceAfter, _ := json.Marshal(store.GetAll())
|
||||
if string(sourceAfter) != string(sourceBefore) {
|
||||
t.Fatal("broadcast decoration mutated registry")
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -4340,13 +4340,18 @@ func (m *Monitor) buildBroadcastFrontendStateFromSnapshot(snapshot models.StateS
|
|||
unifiedView := m.currentUnifiedStateView()
|
||||
metricsTargetResolver := broadcastMetricsTargetResolver(unifiedView.readState)
|
||||
broadcastResources := unifiedresources.CoalescePresentationHostResources(unifiedView.resources)
|
||||
broadcastResources = m.applyPersistedMetadataToUnifiedResources(broadcastResources)
|
||||
broadcastResources = unifiedresources.AttachResourceHealth(
|
||||
broadcastResources,
|
||||
resourceHealthAlerts(frontendState.ActiveAlerts),
|
||||
time.Now().UTC(),
|
||||
// Coalescing owns the outer slice. Decorate that one projection in place,
|
||||
// not three full-resource copies; nested store data is still read-only.
|
||||
healthAlerts := resourceHealthAlerts(frontendState.ActiveAlerts)
|
||||
now := time.Now().UTC()
|
||||
for i := range broadcastResources {
|
||||
m.applyPersistedMetadataToUnifiedResource(&broadcastResources[i])
|
||||
health := unifiedresources.EvaluateResourceHealth(broadcastResources[i], healthAlerts, now)
|
||||
broadcastResources[i].Health = &health
|
||||
}
|
||||
broadcastFrontendResources, broadcastCatalogs := convertPresentationResourcesForBroadcast(
|
||||
attachBroadcastMetricsTargets(broadcastResources, metricsTargetResolver),
|
||||
)
|
||||
broadcastFrontendResources, broadcastCatalogs := convertResourcesForBroadcast(broadcastResources, metricsTargetResolver)
|
||||
frontendState.Resources = broadcastFrontendResources
|
||||
frontendState.CapabilityCatalog = broadcastCatalogs.capabilities
|
||||
frontendState.PolicyCatalog = broadcastCatalogs.policies
|
||||
|
|
@ -4968,17 +4973,18 @@ func (m *Monitor) currentUnifiedStateView() monitorUnifiedStateView {
|
|||
return m.unifiedStateViewWithStandaloneHostContinuity(monitorUnifiedStateViewFromSnapshot(m.GetState()))
|
||||
}
|
||||
|
||||
resources := store.GetAll()
|
||||
freshness := unifiedResourceFreshness(store, state)
|
||||
|
||||
if readState, ok := store.(unifiedresources.ReadState); ok {
|
||||
// The continuity view lists this same store (or its overlay) below.
|
||||
// Cloning here too discarded a complete registry on every broadcast.
|
||||
return m.unifiedStateViewWithStandaloneHostContinuity(monitorUnifiedStateView{
|
||||
resources: resources,
|
||||
readState: readState,
|
||||
freshness: freshness,
|
||||
})
|
||||
}
|
||||
|
||||
resources := store.GetAll()
|
||||
if len(resources) > 0 || state == nil {
|
||||
return m.unifiedStateViewWithStandaloneHostContinuity(monitorUnifiedStateViewFromResources(resources, freshness))
|
||||
}
|
||||
|
|
@ -6132,44 +6138,48 @@ func (m *Monitor) applyPersistedMetadataToUnifiedResources(resources []unifiedre
|
|||
out := make([]unifiedresources.Resource, len(resources))
|
||||
copy(out, resources)
|
||||
for i := range out {
|
||||
resource := &out[i]
|
||||
|
||||
switch unifiedresources.ContractResourceType(*resource) {
|
||||
case unifiedresources.ResourceTypeAppContainer:
|
||||
if resource.Docker == nil {
|
||||
continue
|
||||
}
|
||||
hostID := strings.TrimSpace(resource.Docker.HostSourceID)
|
||||
containerID := strings.TrimSpace(resource.Docker.ContainerID)
|
||||
if hostID == "" {
|
||||
continue
|
||||
}
|
||||
if customURL, ok := m.dockerAppContainerCustomURL(*resource, hostID, containerID); ok {
|
||||
// The metadata record is authoritative even when empty: an
|
||||
// explicit clear must remove a stale URL still carried by an
|
||||
// older unified-resource snapshot.
|
||||
resource.CustomURL = strings.TrimSpace(customURL)
|
||||
}
|
||||
case unifiedresources.ResourceTypePod,
|
||||
unifiedresources.ResourceTypeK8sDeployment,
|
||||
unifiedresources.ResourceTypeK8sService:
|
||||
if customURL, ok := m.kubernetesWorkloadCustomURL(*resource); ok {
|
||||
resource.CustomURL = strings.TrimSpace(customURL)
|
||||
}
|
||||
case unifiedresources.ResourceTypeAgent,
|
||||
unifiedresources.ResourceType("docker-host"),
|
||||
unifiedresources.ResourceTypePBS,
|
||||
unifiedresources.ResourceTypePMG,
|
||||
unifiedresources.ResourceTypeK8sCluster,
|
||||
unifiedresources.ResourceTypeK8sNode:
|
||||
if customURL, ok := m.hostResourceCustomURL(*resource); ok {
|
||||
resource.CustomURL = customURL
|
||||
}
|
||||
}
|
||||
m.applyPersistedMetadataToUnifiedResource(&out[i])
|
||||
}
|
||||
return out
|
||||
}
|
||||
|
||||
// applyPersistedMetadataToUnifiedResource changes only a caller-owned resource
|
||||
// value, never its nested registry state. Empty persisted values clear stale URLs.
|
||||
func (m *Monitor) applyPersistedMetadataToUnifiedResource(resource *unifiedresources.Resource) {
|
||||
switch unifiedresources.ContractResourceType(*resource) {
|
||||
case unifiedresources.ResourceTypeAppContainer:
|
||||
if resource.Docker == nil {
|
||||
return
|
||||
}
|
||||
hostID := strings.TrimSpace(resource.Docker.HostSourceID)
|
||||
containerID := strings.TrimSpace(resource.Docker.ContainerID)
|
||||
if hostID == "" {
|
||||
return
|
||||
}
|
||||
if customURL, ok := m.dockerAppContainerCustomURL(*resource, hostID, containerID); ok {
|
||||
// The metadata record is authoritative even when empty: an
|
||||
// explicit clear must remove a stale URL still carried by an
|
||||
// older unified-resource snapshot.
|
||||
resource.CustomURL = strings.TrimSpace(customURL)
|
||||
}
|
||||
case unifiedresources.ResourceTypePod,
|
||||
unifiedresources.ResourceTypeK8sDeployment,
|
||||
unifiedresources.ResourceTypeK8sService:
|
||||
if customURL, ok := m.kubernetesWorkloadCustomURL(*resource); ok {
|
||||
resource.CustomURL = strings.TrimSpace(customURL)
|
||||
}
|
||||
case unifiedresources.ResourceTypeAgent,
|
||||
unifiedresources.ResourceType("docker-host"),
|
||||
unifiedresources.ResourceTypePBS,
|
||||
unifiedresources.ResourceTypePMG,
|
||||
unifiedresources.ResourceTypeK8sCluster,
|
||||
unifiedresources.ResourceTypeK8sNode:
|
||||
if customURL, ok := m.hostResourceCustomURL(*resource); ok {
|
||||
resource.CustomURL = customURL
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func appendUniqueMetadataCandidate(candidates []string, seen map[string]struct{}, value string) []string {
|
||||
value = strings.TrimSpace(value)
|
||||
if value == "" {
|
||||
|
|
@ -6346,19 +6356,25 @@ func convertResourcesForBroadcast(
|
|||
firstBroadcastMetricsTargetResolver(metricsTargetResolvers),
|
||||
)
|
||||
allResources = unifiedresources.CoalescePresentationHostResources(allResources)
|
||||
type broadcastResource struct {
|
||||
input models.ResourceConvertInput
|
||||
sortKey string
|
||||
resourceID string
|
||||
return convertPresentationResourcesForBroadcast(allResources)
|
||||
}
|
||||
|
||||
// convertPresentationResourcesForBroadcast consumes an already-coalesced
|
||||
// projection. Sorting the final rows avoids retaining a second estate-sized
|
||||
// array of conversion inputs and coalescing the same hosts a second time.
|
||||
func convertPresentationResourcesForBroadcast(allResources []unifiedresources.Resource) ([]models.ResourceFrontend, broadcastResourceCatalogs) {
|
||||
if len(allResources) == 0 {
|
||||
return []models.ResourceFrontend{}, broadcastResourceCatalogs{}
|
||||
}
|
||||
|
||||
converted := make([]broadcastResource, 0, len(allResources))
|
||||
result := make([]models.ResourceFrontend, len(allResources))
|
||||
sortKeys := make([]string, len(allResources))
|
||||
catalogs := broadcastResourceCatalogs{
|
||||
capabilities: make(map[string]json.RawMessage),
|
||||
policies: make(map[string]json.RawMessage),
|
||||
aiSafeSummaries: make(map[string]string),
|
||||
}
|
||||
for _, r := range allResources {
|
||||
for i, r := range allResources {
|
||||
input := monitorResourceToConvertInput(r)
|
||||
if len(input.Capabilities) > 0 {
|
||||
id := capabilityCatalogID(input.Capabilities)
|
||||
|
|
@ -6385,24 +6401,11 @@ func convertResourcesForBroadcast(
|
|||
if sortKey == "" {
|
||||
sortKey = strings.ToLower(input.Name)
|
||||
}
|
||||
converted = append(converted, broadcastResource{
|
||||
input: input,
|
||||
sortKey: sortKey,
|
||||
resourceID: input.ID,
|
||||
})
|
||||
sortKeys[i] = sortKey
|
||||
result[i] = models.ConvertResourceToFrontend(input)
|
||||
}
|
||||
|
||||
sort.Slice(converted, func(i, j int) bool {
|
||||
if converted[i].sortKey == converted[j].sortKey {
|
||||
return converted[i].resourceID < converted[j].resourceID
|
||||
}
|
||||
return converted[i].sortKey < converted[j].sortKey
|
||||
})
|
||||
|
||||
result := make([]models.ResourceFrontend, len(converted))
|
||||
for i, resource := range converted {
|
||||
result[i] = models.ConvertResourceToFrontend(resource.input)
|
||||
}
|
||||
sort.Sort(broadcastFrontendSort{resources: result, keys: sortKeys})
|
||||
if len(catalogs.capabilities) == 0 {
|
||||
catalogs.capabilities = nil
|
||||
}
|
||||
|
|
@ -6415,6 +6418,24 @@ func convertResourcesForBroadcast(
|
|||
return result, catalogs
|
||||
}
|
||||
|
||||
// Keep each precomputed sort key beside its row while sorting in place.
|
||||
type broadcastFrontendSort struct {
|
||||
resources []models.ResourceFrontend
|
||||
keys []string
|
||||
}
|
||||
|
||||
func (s broadcastFrontendSort) Len() int { return len(s.resources) }
|
||||
func (s broadcastFrontendSort) Less(i, j int) bool {
|
||||
if s.keys[i] == s.keys[j] {
|
||||
return s.resources[i].ID < s.resources[j].ID
|
||||
}
|
||||
return s.keys[i] < s.keys[j]
|
||||
}
|
||||
func (s broadcastFrontendSort) Swap(i, j int) {
|
||||
s.resources[i], s.resources[j] = s.resources[j], s.resources[i]
|
||||
s.keys[i], s.keys[j] = s.keys[j], s.keys[i]
|
||||
}
|
||||
|
||||
func broadcastMetricsTargetResolver(source interface{}) MetricsTargetResourceStore {
|
||||
resolver, ok := source.(MetricsTargetResourceStore)
|
||||
if !ok {
|
||||
|
|
|
|||
|
|
@ -6630,3 +6630,72 @@ func TestGetLiveHostsSnapshotCopiesOnlyHosts(t *testing.T) {
|
|||
t.Fatalf("GetLiveHostsSnapshot allocated %.0f times with 1 host and 1000 guests; it is copying guests", allocs)
|
||||
}
|
||||
}
|
||||
|
||||
// Reads model the completed-ingest boundary without persistence/background work.
|
||||
// The real adapter still owns clone isolation, mutable facets and identity.
|
||||
type broadcastProjectionCountingStore struct {
|
||||
*unifiedresources.MonitorAdapter
|
||||
reads int
|
||||
}
|
||||
|
||||
func (s *broadcastProjectionCountingStore) GetAll() []unifiedresources.Resource {
|
||||
s.reads++
|
||||
return s.MonitorAdapter.GetAll()
|
||||
}
|
||||
func (*broadcastProjectionCountingStore) TryReplaceRegistryForRead(models.StateSnapshot, time.Duration, func() map[unifiedresources.DataSource][]unifiedresources.IngestRecord) bool {
|
||||
return false
|
||||
}
|
||||
|
||||
func TestBroadcastProjectionListsRegistryOnceAndKeepsLiveChanges(t *testing.T) {
|
||||
m, adapter, _ := newReadStateCloneTestMonitor(t, 4)
|
||||
store := &broadcastProjectionCountingStore{MonitorAdapter: adapter}
|
||||
m.resourceStore = store
|
||||
first := m.BuildFrontendState()
|
||||
if store.reads != 1 {
|
||||
t.Fatalf("GetAll calls=%d, want one owned continuity-aware clone", store.reads)
|
||||
}
|
||||
if len(first.Resources) != 4 {
|
||||
t.Fatalf("resources=%d", len(first.Resources))
|
||||
}
|
||||
oldCPU := first.Resources[0].CPU.Current
|
||||
changed := store.GetAll()
|
||||
// Keep LastSeen/overall freshness unchanged: a timestamp-only cache would
|
||||
// miss these metric/status/label changes.
|
||||
changed[0].Metrics.CPU.Value = 77
|
||||
changed[0].Labels = map[string]string{"new": "label"}
|
||||
changed[0].Status = unifiedresources.StatusOffline
|
||||
registry := unifiedresources.NewRegistry(nil)
|
||||
registry.IngestResources(changed)
|
||||
store.MonitorAdapter = unifiedresources.NewMonitorAdapter(registry)
|
||||
m.state.UpdateActiveAlerts([]models.Alert{{ID: "live-alert", ResourceID: changed[0].ID, Level: "critical", Type: "cpu"}})
|
||||
m.state.RemovedHostAgents = []models.RemovedHostAgent{{ID: "ignored", Hostname: "ignored.local", RemovedAt: time.Now()}}
|
||||
store.reads = 0
|
||||
second := m.BuildFrontendState()
|
||||
if store.reads != 1 {
|
||||
t.Fatalf("next GetAll calls=%d", store.reads)
|
||||
}
|
||||
var row *models.ResourceFrontend
|
||||
for i := range second.Resources {
|
||||
if second.Resources[i].ID == changed[0].ID {
|
||||
row = &second.Resources[i]
|
||||
}
|
||||
}
|
||||
if row == nil || row.CPU.Current != 77 || row.Labels["new"] != "label" {
|
||||
t.Fatalf("mutable row was stale: %#v", row)
|
||||
}
|
||||
if first.Resources[0].CPU.Current != oldCPU {
|
||||
t.Fatal("later projection mutated an accepted baseline")
|
||||
}
|
||||
if len(second.ActiveAlerts) != 1 || second.ActiveAlerts[0].ID != "live-alert" {
|
||||
t.Fatal("live alert disappeared")
|
||||
}
|
||||
foundIgnored := false
|
||||
for _, item := range second.ConnectedInfrastructure {
|
||||
if item.Status == "ignored" && item.Name == "ignored.local" {
|
||||
foundIgnored = true
|
||||
}
|
||||
}
|
||||
if !foundIgnored {
|
||||
t.Fatal("removed host lifecycle surface disappeared")
|
||||
}
|
||||
}
|
||||
|
|
|
|||
100
internal/websocket/frontend_snapshot_test.go
Normal file
100
internal/websocket/frontend_snapshot_test.go
Normal file
|
|
@ -0,0 +1,100 @@
|
|||
package websocket
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"math"
|
||||
"reflect"
|
||||
"testing"
|
||||
|
||||
"github.com/rcourtman/pulse-go-rewrite/internal/models"
|
||||
)
|
||||
|
||||
func assertFrontendSnapshotEquivalent(t testing.TB, state models.StateFrontend) {
|
||||
t.Helper()
|
||||
want, wantErr := buildGenericClientStateSnapshot(state)
|
||||
for _, value := range []interface{}{state, &state} {
|
||||
got, gotErr := buildClientStateSnapshot(value)
|
||||
if (wantErr == nil) != (gotErr == nil) {
|
||||
t.Fatalf("%T error differs: generic=%v typed=%v", value, wantErr, gotErr)
|
||||
}
|
||||
if wantErr == nil && !reflect.DeepEqual(got, want) {
|
||||
t.Fatalf("%T snapshot differs from full wire encoding", value)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestFrontendSnapshotMatchesFullWireEncoding(t *testing.T) {
|
||||
state := benchmarkFrontendState(20)
|
||||
state.Resources[19].ID = "尾<&\"\\\n"
|
||||
state.Resources[19].PlatformData = json.RawMessage(`{"id":"not-the-key","extra":[1,true,null]}`)
|
||||
state.ConnectedInfrastructure = []models.ConnectedInfrastructureItemFrontend{{ID: "i-1"}, {ID: "i-2"}}
|
||||
state.ActiveAlerts = []models.Alert{{ID: "a-1", Level: "critical"}, {ID: "a-2"}}
|
||||
state.CapabilityCatalog = map[string]json.RawMessage{"ref": json.RawMessage(`{"allowed":true}`)}
|
||||
state.PolicyCatalog = map[string]json.RawMessage{"ref": json.RawMessage(`{"denied":true}`)}
|
||||
state.AISafeSummaryCatalog = map[string]string{"ref": "summary<&>"}
|
||||
state.ConnectionHealth = map[string]bool{"lab": false}
|
||||
state.PVETagColors = map[string]string{"tag": "#abcdef"}
|
||||
state.LastUpdate = 1234
|
||||
for name, value := range map[string]models.StateFrontend{
|
||||
"populated": state, "zero": {}, "empty": models.EmptyStateFrontend(),
|
||||
} {
|
||||
t.Run(name, func(t *testing.T) { assertFrontendSnapshotEquivalent(t, value) })
|
||||
}
|
||||
// Missing and duplicate IDs, bad raw data and unencodable metrics must fail
|
||||
// rather than quietly adopt a different resource baseline.
|
||||
state.Resources[1].ID = state.Resources[0].ID
|
||||
assertFrontendSnapshotEquivalent(t, state)
|
||||
state.Resources[1].ID = ""
|
||||
assertFrontendSnapshotEquivalent(t, state)
|
||||
state.Resources[1].ID = "fixed"
|
||||
state.Resources[0].CPU.Current = math.NaN()
|
||||
assertFrontendSnapshotEquivalent(t, state)
|
||||
state.Resources[0].CPU.Current = 0
|
||||
state.Resources[19].PlatformData = json.RawMessage(`{"invalid":`)
|
||||
assertFrontendSnapshotEquivalent(t, state)
|
||||
}
|
||||
|
||||
func TestFrontendSnapshotOwnsMutableSource(t *testing.T) {
|
||||
state := benchmarkFrontendState(2)
|
||||
state.Resources[0].Labels = map[string]string{"label": "original"}
|
||||
state.Resources[0].Tags = []string{"original"}
|
||||
state.Resources[0].PlatformData = json.RawMessage(`{"value":"original"}`)
|
||||
snapshot, err := buildClientStateSnapshot(state)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
before, err := json.Marshal(snapshot.resources)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
state.Resources[0].ID = "different"
|
||||
state.Resources[0].CPU.Current++
|
||||
state.Resources[0].Labels["label"] = "changed"
|
||||
state.Resources[0].Tags[0] = "changed"
|
||||
state.Resources[0].PlatformData[10] = 'X'
|
||||
after, err := json.Marshal(snapshot.resources)
|
||||
if err != nil || string(before) != string(after) {
|
||||
t.Fatal("snapshot aliases mutable source")
|
||||
}
|
||||
}
|
||||
|
||||
func FuzzFrontendSnapshotMatchesFullWireEncoding(f *testing.F) {
|
||||
f.Add("id-1", "payload<&>", false, false)
|
||||
f.Add("尾\n\\\"", "ID", true, true)
|
||||
f.Add("", "", true, false)
|
||||
f.Fuzz(func(t *testing.T, id, payload string, duplicate, nilCollections bool) {
|
||||
state := benchmarkFrontendState(3)
|
||||
state.Resources[2].ID = id
|
||||
state.Resources[2].Name = payload
|
||||
if duplicate {
|
||||
state.Resources[1].ID = id
|
||||
}
|
||||
if nilCollections {
|
||||
state.ActiveAlerts = nil
|
||||
state.ConnectedInfrastructure = nil
|
||||
state.Metrics = nil
|
||||
state.ConnectionHealth = nil
|
||||
}
|
||||
assertFrontendSnapshotEquivalent(t, state)
|
||||
})
|
||||
}
|
||||
|
|
@ -6,6 +6,8 @@ import (
|
|||
"fmt"
|
||||
"reflect"
|
||||
"sort"
|
||||
|
||||
"github.com/rcourtman/pulse-go-rewrite/internal/models"
|
||||
)
|
||||
|
||||
const resourceDeltaField = "resourceDelta"
|
||||
|
|
@ -86,6 +88,61 @@ func extractKeyedEntries(
|
|||
}
|
||||
|
||||
func buildClientStateSnapshot(state interface{}) (*clientStateSnapshot, error) {
|
||||
// Unknown shapes and custom marshalers retain the generic wire-authoritative
|
||||
// path. Only the concrete frontend projection can avoid the whole resources
|
||||
// array's encode/decode/copy round trip.
|
||||
if _, custom := state.(json.Marshaler); !custom {
|
||||
switch frontend := state.(type) {
|
||||
case models.StateFrontend:
|
||||
return buildFrontendStateSnapshot(frontend)
|
||||
case *models.StateFrontend:
|
||||
if frontend != nil {
|
||||
return buildFrontendStateSnapshot(*frontend)
|
||||
}
|
||||
}
|
||||
}
|
||||
return buildGenericClientStateSnapshot(state)
|
||||
}
|
||||
|
||||
func buildFrontendStateSnapshot(state models.StateFrontend) (*clientStateSnapshot, error) {
|
||||
resources := state.Resources
|
||||
// Encode the rest through the existing generic path so future top-level
|
||||
// fields, omitempty, nil/empty arrays and keyed-field fallback stay identical.
|
||||
// This is a local value copy; no source slices/maps are modified or cached.
|
||||
state.Resources = nil
|
||||
snapshot, err := buildGenericClientStateSnapshot(state)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
snapshot.resources = make(map[string]json.RawMessage, len(resources))
|
||||
snapshot.resourceOrder = make([]string, 0, len(resources))
|
||||
for i := range resources {
|
||||
encoded, err := json.Marshal(&resources[i])
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("marshal state resource: %w", err)
|
||||
}
|
||||
// Identity still comes from EVERY encoded entry, never a source-ID hint.
|
||||
// json.Marshal owns this buffer; it can be retained directly as immutable
|
||||
// snapshot data without whole-state and RawMessage decoding copies.
|
||||
var identity struct {
|
||||
ID string `json:"id"`
|
||||
}
|
||||
if err := json.Unmarshal(encoded, &identity); err != nil {
|
||||
return nil, fmt.Errorf("decode state resource identity: %w", err)
|
||||
}
|
||||
if identity.ID == "" {
|
||||
return nil, fmt.Errorf("state resource entry is missing id")
|
||||
}
|
||||
if _, exists := snapshot.resources[identity.ID]; exists {
|
||||
return nil, fmt.Errorf("state resource id %q is duplicated", identity.ID)
|
||||
}
|
||||
snapshot.resources[identity.ID] = encoded
|
||||
snapshot.resourceOrder = append(snapshot.resourceOrder, identity.ID)
|
||||
}
|
||||
return snapshot, nil
|
||||
}
|
||||
|
||||
func buildGenericClientStateSnapshot(state interface{}) (*clientStateSnapshot, error) {
|
||||
encoded, err := json.Marshal(state)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("marshal state snapshot: %w", err)
|
||||
|
|
@ -125,7 +182,7 @@ func buildClientStateSnapshot(state interface{}) (*clientStateSnapshot, error) {
|
|||
snapshot.keyed[keyedField.field] = &keyedFieldSnapshot{
|
||||
entries: entries,
|
||||
order: order,
|
||||
raw: append(json.RawMessage(nil), encodedField...),
|
||||
raw: encodedField,
|
||||
}
|
||||
delete(fields, keyedField.field)
|
||||
}
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue