mirror of
https://github.com/henrygd/beszel.git
synced 2026-10-01 05:47:48 +02:00
refactor(hub): scope monitor result writes to the reporting system (#2449)
Co-authored-by: hank <hank@henrygd.me>
This commit is contained in:
@@ -3,6 +3,7 @@
|
|||||||
package systems
|
package systems
|
||||||
|
|
||||||
import (
|
import (
|
||||||
|
"encoding/json"
|
||||||
"fmt"
|
"fmt"
|
||||||
"testing"
|
"testing"
|
||||||
"time"
|
"time"
|
||||||
@@ -16,6 +17,112 @@ import (
|
|||||||
"github.com/stretchr/testify/require"
|
"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) {
|
func TestNetworkMonitorProbePruning(t *testing.T) {
|
||||||
for _, tc := range []struct {
|
for _, tc := range []struct {
|
||||||
name string
|
name string
|
||||||
|
|||||||
@@ -2,6 +2,7 @@ package systems
|
|||||||
|
|
||||||
import (
|
import (
|
||||||
"context"
|
"context"
|
||||||
|
"database/sql"
|
||||||
"encoding/json"
|
"encoding/json"
|
||||||
"errors"
|
"errors"
|
||||||
"fmt"
|
"fmt"
|
||||||
@@ -433,7 +434,6 @@ func (sys *System) updateNetworkMonitorsRecords(app core.App, monitorResults map
|
|||||||
if len(monitorResults) == 0 {
|
if len(monitorResults) == 0 {
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
var err error
|
|
||||||
systemId := sys.Id
|
systemId := sys.Id
|
||||||
const monitorCollectionName = "network_monitors"
|
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.
|
// Results omit certInfo unless it changed, so keep the stored value.
|
||||||
setClauses = append(setClauses, "certInfo=COALESCE({:certInfo}, certInfo)")
|
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)
|
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
|
// update network_monitors records
|
||||||
for id, result := range monitorResults {
|
for id, result := range monitorResults {
|
||||||
monitorData := map[string]any{
|
monitorData := map[string]any{
|
||||||
"id": id,
|
"id": id,
|
||||||
|
"system": systemId,
|
||||||
"res": result.AvgResponse,
|
"res": result.AvgResponse,
|
||||||
"resAvg1h": result.AvgResponse1h,
|
"resAvg1h": result.AvgResponse1h,
|
||||||
"resMin1h": result.MinResponse1h,
|
"resMin1h": result.MinResponse1h,
|
||||||
@@ -472,11 +477,15 @@ func (sys *System) updateNetworkMonitorsRecords(app core.App, monitorResults map
|
|||||||
"loss1h": result.PacketLoss1h,
|
"loss1h": result.PacketLoss1h,
|
||||||
"updated": nowString,
|
"updated": nowString,
|
||||||
}
|
}
|
||||||
|
var err error
|
||||||
switch realtimeActive {
|
switch realtimeActive {
|
||||||
case true:
|
case true:
|
||||||
var record *core.Record
|
var record *core.Record
|
||||||
record, err = app.FindRecordById(monitorCollectionName, id)
|
record, err = app.FindRecordById(monitorCollectionName, id)
|
||||||
if err == nil {
|
if err == nil {
|
||||||
|
if record.GetString("system") != systemId {
|
||||||
|
continue
|
||||||
|
}
|
||||||
if result.Cert != nil {
|
if result.Cert != nil {
|
||||||
monitorData["certInfo"] = result.Cert
|
monitorData["certInfo"] = result.Cert
|
||||||
}
|
}
|
||||||
@@ -492,12 +501,20 @@ func (sys *System) updateNetworkMonitorsRecords(app core.App, monitorResults map
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
if err == nil {
|
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 {
|
if err != nil {
|
||||||
app.Logger().Warn("Failed to update monitor", "system", systemId, "monitor", id, "err", err)
|
app.Logger().Warn("Failed to update monitor", "system", systemId, "monitor", id, "err", err)
|
||||||
|
continue
|
||||||
}
|
}
|
||||||
|
owned[id] = struct{}{}
|
||||||
}
|
}
|
||||||
|
|
||||||
// handle stats collection — one record per monitor
|
// 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 {
|
for monitorId, result := range monitorResults {
|
||||||
|
if _, ok := owned[monitorId]; !ok {
|
||||||
|
continue
|
||||||
|
}
|
||||||
// Compare identity, not ordering, so agent clock changes don't stall writes.
|
// Compare identity, not ordering, so agent clock changes don't stall writes.
|
||||||
if result.LastProbeAt == sys.lastSavedMonitorProbe[monitorId] {
|
if result.LastProbeAt == sys.lastSavedMonitorProbe[monitorId] {
|
||||||
continue
|
continue
|
||||||
@@ -524,6 +544,7 @@ func (sys *System) updateNetworkMonitorsRecords(app core.App, monitorResults map
|
|||||||
"success_count": result.SuccessCount,
|
"success_count": result.SuccessCount,
|
||||||
"res_sum": result.ResponseSum,
|
"res_sum": result.ResponseSum,
|
||||||
}
|
}
|
||||||
|
var err error
|
||||||
switch realtimeActive {
|
switch realtimeActive {
|
||||||
case true:
|
case true:
|
||||||
record := core.NewRecord(statsCollection)
|
record := core.NewRecord(statsCollection)
|
||||||
|
|||||||
Reference in New Issue
Block a user