From a77ed345d76daf7676dbf2d7b68873ffdd240537 Mon Sep 17 00:00:00 2001 From: user01010111 <12504630+user01010111@users.noreply.github.com> Date: Wed, 30 Sep 2026 02:41:48 +1300 Subject: [PATCH] refactor(hub): scope monitor result writes to the reporting system (#2449) Co-authored-by: hank --- .../hub/systems/network_monitor_stats_test.go | 107 ++++++++++++++++++ internal/hub/systems/system.go | 27 ++++- 2 files changed, 131 insertions(+), 3 deletions(-) diff --git a/internal/hub/systems/network_monitor_stats_test.go b/internal/hub/systems/network_monitor_stats_test.go index f3b1572d7..a7e8e6b73 100644 --- a/internal/hub/systems/network_monitor_stats_test.go +++ b/internal/hub/systems/network_monitor_stats_test.go @@ -3,6 +3,7 @@ package systems import ( + "encoding/json" "fmt" "testing" "time" @@ -16,6 +17,112 @@ import ( "github.com/stretchr/testify/require" ) +func TestNetworkMonitorResultOwnership(t *testing.T) { + for _, realtime := range []bool{false, true} { + t.Run(fmt.Sprintf("realtime=%t", realtime), func(t *testing.T) { + sys, app := newTestSystemWithHub(t) + if realtime { + client := subscriptions.NewDefaultClient() + client.Subscribe("network_monitors/*") + app.SubscriptionsBroker().Register(client) + t.Cleanup(func() { app.SubscriptionsBroker().Unregister(client.Id()) }) + } + systems, err := app.FindCachedCollectionByNameOrId("systems") + require.NoError(t, err) + foreignSystem := core.NewRecord(systems) + require.NoError(t, app.SaveNoValidate(foreignSystem)) + other := &System{Id: foreignSystem.Id, manager: sys.manager} + collection, err := app.FindCachedCollectionByNameOrId("network_monitors") + require.NoError(t, err) + for _, cfg := range []struct { + id, system string + enabled bool + }{ + {"owned", sys.Id, true}, + {"disabled", sys.Id, false}, + {"foreign", other.Id, true}, + } { + record := core.NewRecord(collection) + record.Id = cfg.id + record.Load(map[string]any{"system": cfg.system, "enabled": cfg.enabled}) + require.NoError(t, app.SaveNoValidate(record)) + } + // Establish a legitimate result and history for the other system. + _, err = other.createRecords(&system.CombinedData{Monitors: map[string]monitor.Result{ + "foreign": {LastProbeAt: 1000, AvgResponse: 77, TotalCount: 1, SuccessCount: 1, ResponseSum: 77}, + }}) + require.NoError(t, err) + snapshot := func(id string) []byte { + t.Helper() + record, err := app.FindRecordById("network_monitors", id) + require.NoError(t, err) + data, err := json.Marshal(record) + require.NoError(t, err) + return data + } + before := snapshot("foreign") + result := monitor.Result{ + LastProbeAt: 2000, AvgResponse: 22, AvgResponse1h: 33, + MinResponse1h: 11, MaxResponse1h: 44, PacketLoss1h: 50, + TotalCount: 2, SuccessCount: 1, ResponseSum: 22, + Cert: &monitor.CertInfo{Expires: 1_800_000_000_000, Issuer: "Test CA"}, + } + data := &system.CombinedData{Monitors: map[string]monitor.Result{ + "owned": result, "disabled": result, "foreign": result, "nonexistent": result, + }} + for range 2 { // Repeated results must remain deduplicated. + _, err = sys.createRecords(data) + require.NoError(t, err) + } + assert.JSONEq(t, string(before), string(snapshot("foreign"))) + assert.Equal(t, map[string]int64{"owned": 2000, "disabled": 2000}, sys.lastSavedMonitorProbe) + assert.Len(t, data.Monitors, 4, "ingestion must not mutate the shared telemetry payload") + for _, id := range []string{"owned", "disabled"} { + record, err := app.FindRecordById("network_monitors", id) + require.NoError(t, err) + assert.EqualValues(t, result.AvgResponse, record.GetInt("res")) + assert.EqualValues(t, result.AvgResponse1h, record.GetInt("resAvg1h")) + var cert monitor.CertInfo + require.NoError(t, record.UnmarshalJSONField("certInfo", &cert)) + assert.Equal(t, *result.Cert, cert) + } + stats, err := app.FindAllRecords("network_monitor_stats") + require.NoError(t, err) + assert.Len(t, stats, 3) + for _, stat := range stats { + record, err := app.FindRecordById("network_monitors", stat.GetString("monitor")) + require.NoError(t, err) + assert.Equal(t, record.GetString("system"), stat.GetString("system")) + if record.Id == "foreign" { + assert.Equal(t, 77, stat.GetInt("res_sum")) + } else { + assert.Equal(t, 22, stat.GetInt("res_sum")) + } + } + + // Ownership must be read afresh, even for IDs with saved probe markers. + moved, err := app.FindRecordById("network_monitors", "owned") + require.NoError(t, err) + moved.Set("system", other.Id) + require.NoError(t, app.SaveNoValidate(moved)) + movedBefore := snapshot("owned") + for id, result := range data.Monitors { + result.LastProbeAt = 3000 + result.AvgResponse = 999 + data.Monitors[id] = result + } + _, err = sys.createRecords(data) + require.NoError(t, err) + assert.JSONEq(t, string(movedBefore), string(snapshot("owned"))) + assert.JSONEq(t, string(before), string(snapshot("foreign"))) + count, err := app.CountRecords("network_monitor_stats") + require.NoError(t, err) + assert.EqualValues(t, 4, count, "only the still-owned disabled monitor gets another sample") + assert.Equal(t, map[string]int64{"owned": 2000, "disabled": 3000}, sys.lastSavedMonitorProbe) + }) + } +} + func TestNetworkMonitorProbePruning(t *testing.T) { for _, tc := range []struct { name string diff --git a/internal/hub/systems/system.go b/internal/hub/systems/system.go index 49d1f5d63..55a5878f0 100644 --- a/internal/hub/systems/system.go +++ b/internal/hub/systems/system.go @@ -2,6 +2,7 @@ package systems import ( "context" + "database/sql" "encoding/json" "errors" "fmt" @@ -433,7 +434,6 @@ func (sys *System) updateNetworkMonitorsRecords(app core.App, monitorResults map if len(monitorResults) == 0 { return nil } - var err error systemId := sys.Id const monitorCollectionName = "network_monitors" @@ -457,14 +457,19 @@ func (sys *System) updateNetworkMonitorsRecords(app core.App, monitorResults map } // Results omit certInfo unless it changed, so keep the stored value. setClauses = append(setClauses, "certInfo=COALESCE({:certInfo}, certInfo)") - queryString := fmt.Sprintf("UPDATE %s SET %s WHERE id={:id}", monitorCollectionName, strings.Join(setClauses, ", ")) + queryString := fmt.Sprintf("UPDATE %s SET %s WHERE id={:id} AND system={:system}", monitorCollectionName, strings.Join(setClauses, ", ")) updateQuery = db.NewQuery(queryString) } + // Results are keyed by agent-supplied IDs. Record history only for monitors + // this system owns, as confirmed by the update below + owned := make(map[string]struct{}, len(monitorResults)) + // update network_monitors records for id, result := range monitorResults { monitorData := map[string]any{ "id": id, + "system": systemId, "res": result.AvgResponse, "resAvg1h": result.AvgResponse1h, "resMin1h": result.MinResponse1h, @@ -472,11 +477,15 @@ func (sys *System) updateNetworkMonitorsRecords(app core.App, monitorResults map "loss1h": result.PacketLoss1h, "updated": nowString, } + var err error switch realtimeActive { case true: var record *core.Record record, err = app.FindRecordById(monitorCollectionName, id) if err == nil { + if record.GetString("system") != systemId { + continue + } if result.Cert != nil { monitorData["certInfo"] = result.Cert } @@ -492,12 +501,20 @@ func (sys *System) updateNetworkMonitorsRecords(app core.App, monitorResults map } } if err == nil { - _, err = updateQuery.Bind(dbx.Params(monitorData)).Execute() + var res sql.Result + // Zero rows means the monitor is foreign or no longer exists. + if res, err = updateQuery.Bind(dbx.Params(monitorData)).Execute(); err == nil { + if n, _ := res.RowsAffected(); n == 0 { + continue + } + } } } if err != nil { app.Logger().Warn("Failed to update monitor", "system", systemId, "monitor", id, "err", err) + continue } + owned[id] = struct{}{} } // handle stats collection — one record per monitor @@ -509,6 +526,9 @@ func (sys *System) updateNetworkMonitorsRecords(app core.App, monitorResults map } for monitorId, result := range monitorResults { + if _, ok := owned[monitorId]; !ok { + continue + } // Compare identity, not ordering, so agent clock changes don't stall writes. if result.LastProbeAt == sys.lastSavedMonitorProbe[monitorId] { continue @@ -524,6 +544,7 @@ func (sys *System) updateNetworkMonitorsRecords(app core.App, monitorResults map "success_count": result.SuccessCount, "res_sum": result.ResponseSum, } + var err error switch realtimeActive { case true: record := core.NewRecord(statsCollection)