mirror of
https://github.com/henrygd/beszel.git
synced 2026-09-30 13:27:46 +02:00
refactor(hub): decode each agent response into a new struct
The WebSocket path reused sys.data, and CBOR leaves omitted fields untouched, so stale values carried over (e.g. Details re-saved every poll, SystemdServicesUpdated stuck true, old map entries kept). sys.data is now set only by the regular update for alerts. Removes the Wi-Fi-only workaround in UnmarshalResponse.
This commit is contained in:
@@ -153,6 +153,7 @@ func (sys *System) update() error {
|
|||||||
|
|
||||||
// ensure deprecated fields from older agents are migrated to current fields
|
// ensure deprecated fields from older agents are migrated to current fields
|
||||||
migrateDeprecatedFields(data, !sys.detailsFetched.Load())
|
migrateDeprecatedFields(data, !sys.detailsFetched.Load())
|
||||||
|
sys.data = data
|
||||||
|
|
||||||
// create system records
|
// create system records
|
||||||
_, err = sys.createRecords(data)
|
_, err = sys.createRecords(data)
|
||||||
@@ -704,11 +705,9 @@ func (sys *System) ensureSSHTransport() error {
|
|||||||
}
|
}
|
||||||
|
|
||||||
// fetchDataFromAgent attempts to fetch data from the agent, prioritizing WebSocket if available.
|
// fetchDataFromAgent attempts to fetch data from the agent, prioritizing WebSocket if available.
|
||||||
|
// Each fetch decodes into a new struct: CBOR leaves fields the agent omits
|
||||||
|
// untouched, and real-time and regular updates may fetch concurrently.
|
||||||
func (sys *System) fetchDataFromAgent(options common.DataRequestOptions) (*system.CombinedData, error) {
|
func (sys *System) fetchDataFromAgent(options common.DataRequestOptions) (*system.CombinedData, error) {
|
||||||
if sys.data == nil {
|
|
||||||
sys.data = &system.CombinedData{}
|
|
||||||
}
|
|
||||||
|
|
||||||
if sys.WsConn != nil && sys.WsConn.IsConnected() {
|
if sys.WsConn != nil && sys.WsConn.IsConnected() {
|
||||||
wsData, err := sys.fetchDataViaWebSocket(options)
|
wsData, err := sys.fetchDataViaWebSocket(options)
|
||||||
if err == nil {
|
if err == nil {
|
||||||
@@ -744,11 +743,11 @@ func (sys *System) fetchDataViaWebSocket(options common.DataRequestOptions) (*sy
|
|||||||
ctx, cancel := context.WithTimeout(context.Background(), wsDataRequestTimeout)
|
ctx, cancel := context.WithTimeout(context.Background(), wsDataRequestTimeout)
|
||||||
defer cancel()
|
defer cancel()
|
||||||
wsTransport := transport.NewWebSocketTransport(sys.WsConn)
|
wsTransport := transport.NewWebSocketTransport(sys.WsConn)
|
||||||
err := wsTransport.Request(ctx, common.GetData, options, sys.data)
|
data := &system.CombinedData{}
|
||||||
if err != nil {
|
if err := wsTransport.Request(ctx, common.GetData, options, data); err != nil {
|
||||||
return nil, err
|
return nil, err
|
||||||
}
|
}
|
||||||
return sys.data, nil
|
return data, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
// FetchContainerInfoFromAgent fetches container info from the agent
|
// FetchContainerInfoFromAgent fetches container info from the agent
|
||||||
@@ -810,9 +809,8 @@ func MakeStableHashId(strings ...string) string {
|
|||||||
}
|
}
|
||||||
|
|
||||||
// fetchDataViaSSH handles fetching data using SSH.
|
// fetchDataViaSSH handles fetching data using SSH.
|
||||||
// This function encapsulates the original SSH logic.
|
|
||||||
// It updates sys.data directly upon successful fetch.
|
|
||||||
func (sys *System) fetchDataViaSSH(options common.DataRequestOptions) (*system.CombinedData, error) {
|
func (sys *System) fetchDataViaSSH(options common.DataRequestOptions) (*system.CombinedData, error) {
|
||||||
|
data := &system.CombinedData{}
|
||||||
err := sys.runSSHOperation(4*time.Second, 1, func(session *ssh.Session) (bool, error) {
|
err := sys.runSSHOperation(4*time.Second, 1, func(session *ssh.Session) (bool, error) {
|
||||||
stdout, err := session.StdoutPipe()
|
stdout, err := session.StdoutPipe()
|
||||||
if err != nil {
|
if err != nil {
|
||||||
@@ -823,7 +821,8 @@ func (sys *System) fetchDataViaSSH(options common.DataRequestOptions) (*system.C
|
|||||||
return false, err
|
return false, err
|
||||||
}
|
}
|
||||||
|
|
||||||
*sys.data = system.CombinedData{}
|
// reset in case of retry after a partial decode
|
||||||
|
*data = system.CombinedData{}
|
||||||
|
|
||||||
if sys.agentVersion.GTE(beszel.MinVersionAgentResponse) && stdinErr == nil {
|
if sys.agentVersion.GTE(beszel.MinVersionAgentResponse) && stdinErr == nil {
|
||||||
req := common.HubRequest[any]{Action: common.GetData, Data: options}
|
req := common.HubRequest[any]{Action: common.GetData, Data: options}
|
||||||
@@ -832,7 +831,7 @@ func (sys *System) fetchDataViaSSH(options common.DataRequestOptions) (*system.C
|
|||||||
|
|
||||||
var resp common.AgentResponse
|
var resp common.AgentResponse
|
||||||
if decErr := cbor.NewDecoder(stdout).Decode(&resp); decErr == nil && resp.SystemData != nil {
|
if decErr := cbor.NewDecoder(stdout).Decode(&resp); decErr == nil && resp.SystemData != nil {
|
||||||
*sys.data = *resp.SystemData
|
*data = *resp.SystemData
|
||||||
if err := session.Wait(); err != nil {
|
if err := session.Wait(); err != nil {
|
||||||
return false, err
|
return false, err
|
||||||
}
|
}
|
||||||
@@ -842,9 +841,9 @@ func (sys *System) fetchDataViaSSH(options common.DataRequestOptions) (*system.C
|
|||||||
|
|
||||||
var decodeErr error
|
var decodeErr error
|
||||||
if sys.agentVersion.GTE(beszel.MinVersionCbor) {
|
if sys.agentVersion.GTE(beszel.MinVersionCbor) {
|
||||||
decodeErr = cbor.NewDecoder(stdout).Decode(sys.data)
|
decodeErr = cbor.NewDecoder(stdout).Decode(data)
|
||||||
} else {
|
} else {
|
||||||
decodeErr = json.NewDecoder(stdout).Decode(sys.data)
|
decodeErr = json.NewDecoder(stdout).Decode(data)
|
||||||
}
|
}
|
||||||
|
|
||||||
if decodeErr != nil {
|
if decodeErr != nil {
|
||||||
@@ -861,7 +860,7 @@ func (sys *System) fetchDataViaSSH(options common.DataRequestOptions) (*system.C
|
|||||||
return nil, err
|
return nil, err
|
||||||
}
|
}
|
||||||
|
|
||||||
return sys.data, nil
|
return data, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
// runSSHOperation establishes an SSH session and executes the provided operation.
|
// runSSHOperation establishes an SSH session and executes the provided operation.
|
||||||
|
|||||||
92
internal/hub/systems/ws_fresh_data_test.go
Normal file
92
internal/hub/systems/ws_fresh_data_test.go
Normal file
@@ -0,0 +1,92 @@
|
|||||||
|
//go:build testing
|
||||||
|
|
||||||
|
package systems
|
||||||
|
|
||||||
|
import (
|
||||||
|
"net/http"
|
||||||
|
"net/http/httptest"
|
||||||
|
"strings"
|
||||||
|
"testing"
|
||||||
|
"time"
|
||||||
|
|
||||||
|
"github.com/blang/semver"
|
||||||
|
"github.com/fxamacker/cbor/v2"
|
||||||
|
"github.com/henrygd/beszel/internal/common"
|
||||||
|
esystem "github.com/henrygd/beszel/internal/entities/system"
|
||||||
|
"github.com/henrygd/beszel/internal/hub/ws"
|
||||||
|
"github.com/lxzan/gws"
|
||||||
|
"github.com/stretchr/testify/require"
|
||||||
|
)
|
||||||
|
|
||||||
|
// sequenceDataClient answers each GetData request with the next queued payload.
|
||||||
|
type sequenceDataClient struct {
|
||||||
|
gws.BuiltinEventHandler
|
||||||
|
responses chan esystem.CombinedData
|
||||||
|
}
|
||||||
|
|
||||||
|
func (c *sequenceDataClient) OnMessage(conn *gws.Conn, message *gws.Message) {
|
||||||
|
defer message.Close()
|
||||||
|
var req common.HubRequest[cbor.RawMessage]
|
||||||
|
if err := cbor.Unmarshal(message.Bytes(), &req); err != nil || req.Action != common.GetData {
|
||||||
|
return
|
||||||
|
}
|
||||||
|
data, _ := cbor.Marshal(<-c.responses)
|
||||||
|
response, _ := cbor.Marshal(common.AgentResponse{Id: req.Id, Data: data})
|
||||||
|
_ = conn.WriteMessage(gws.OpcodeBinary, response)
|
||||||
|
}
|
||||||
|
|
||||||
|
// Fields the agent omits must not carry over from a previous response.
|
||||||
|
func TestFetchDataViaWebSocketDoesNotRetainOmittedFields(t *testing.T) {
|
||||||
|
connections := make(chan *ws.WsConn, 1)
|
||||||
|
upgrader := gws.NewUpgrader(&monitorSyncServer{}, nil)
|
||||||
|
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||||
|
conn, err := upgrader.Upgrade(w, r)
|
||||||
|
if err != nil {
|
||||||
|
t.Error(err)
|
||||||
|
return
|
||||||
|
}
|
||||||
|
wsConn := ws.NewWsConnection(conn, semver.MustParse("0.20.0"))
|
||||||
|
conn.Session().Store("wsConn", wsConn)
|
||||||
|
connections <- wsConn
|
||||||
|
conn.ReadLoop()
|
||||||
|
}))
|
||||||
|
t.Cleanup(server.Close)
|
||||||
|
|
||||||
|
client := &sequenceDataClient{responses: make(chan esystem.CombinedData, 2)}
|
||||||
|
client.responses <- esystem.CombinedData{
|
||||||
|
Details: &esystem.Details{Hostname: "host"},
|
||||||
|
SystemdServicesUpdated: true,
|
||||||
|
Stats: esystem.Stats{Batteries: map[string]uint8{"BAT0": 50, "BAT1": 60}},
|
||||||
|
Info: esystem.Info{WiFi: map[string]esystem.WiFi{"wlan0": {SSID: "home"}}},
|
||||||
|
}
|
||||||
|
client.responses <- esystem.CombinedData{
|
||||||
|
Stats: esystem.Stats{Batteries: map[string]uint8{"BAT0": 40}},
|
||||||
|
}
|
||||||
|
conn, _, err := gws.NewClient(client, &gws.ClientOption{Addr: "ws" + strings.TrimPrefix(server.URL, "http")})
|
||||||
|
require.NoError(t, err)
|
||||||
|
t.Cleanup(func() { _ = conn.NetConn().Close() })
|
||||||
|
go conn.ReadLoop()
|
||||||
|
|
||||||
|
var sys *System
|
||||||
|
select {
|
||||||
|
case wsConn := <-connections:
|
||||||
|
sys = &System{WsConn: wsConn}
|
||||||
|
case <-time.After(3 * time.Second):
|
||||||
|
t.Fatal("websocket connection was not established")
|
||||||
|
}
|
||||||
|
|
||||||
|
first, err := sys.fetchDataFromAgent(common.DataRequestOptions{})
|
||||||
|
require.NoError(t, err)
|
||||||
|
require.NotNil(t, first.Details)
|
||||||
|
require.Len(t, first.Stats.Batteries, 2)
|
||||||
|
|
||||||
|
second, err := sys.fetchDataFromAgent(common.DataRequestOptions{})
|
||||||
|
require.NoError(t, err)
|
||||||
|
require.Nil(t, second.Details)
|
||||||
|
require.False(t, second.SystemdServicesUpdated)
|
||||||
|
require.Equal(t, map[string]uint8{"BAT0": 40}, second.Stats.Batteries)
|
||||||
|
require.Empty(t, second.Info.WiFi)
|
||||||
|
|
||||||
|
// the first result is not mutated by the second fetch
|
||||||
|
require.Len(t, first.Stats.Batteries, 2)
|
||||||
|
}
|
||||||
@@ -35,14 +35,6 @@ func UnmarshalResponse(resp common.AgentResponse, action common.WebSocketAction,
|
|||||||
}
|
}
|
||||||
// Try generic Data field first (0.19+)
|
// Try generic Data field first (0.19+)
|
||||||
if len(resp.Data) > 0 {
|
if len(resp.Data) > 0 {
|
||||||
// Wi-Fi maps are complete snapshots. CBOR otherwise merges entries into
|
|
||||||
// reused destinations, retaining disconnected interfaces and old RSSI.
|
|
||||||
if action == common.GetData {
|
|
||||||
if data, ok := dest.(*system.CombinedData); ok {
|
|
||||||
data.Info.WiFi = nil
|
|
||||||
data.Stats.WiFi = nil
|
|
||||||
}
|
|
||||||
}
|
|
||||||
if err := cbor.Unmarshal(resp.Data, dest); err != nil {
|
if err := cbor.Unmarshal(resp.Data, dest); err != nil {
|
||||||
return fmt.Errorf("failed to unmarshal generic response data: %w", err)
|
return fmt.Errorf("failed to unmarshal generic response data: %w", err)
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -1,43 +0,0 @@
|
|||||||
package transport
|
|
||||||
|
|
||||||
import (
|
|
||||||
"github.com/fxamacker/cbor/v2"
|
|
||||||
"github.com/henrygd/beszel/internal/common"
|
|
||||||
"github.com/henrygd/beszel/internal/entities/system"
|
|
||||||
"github.com/stretchr/testify/require"
|
|
||||||
"testing"
|
|
||||||
)
|
|
||||||
|
|
||||||
func TestWiFiSequentialResponseSnapshots(t *testing.T) {
|
|
||||||
signal := -50.0
|
|
||||||
var decoded system.CombinedData
|
|
||||||
for _, snapshot := range []map[string]system.WiFi{
|
|
||||||
{"wlan0": {SSID: "home", Signal: &signal}, "wlan1": {Signal: &signal}},
|
|
||||||
{"wlan0": {SSID: "home"}}, {}, nil,
|
|
||||||
{"wlan1": {SSID: "new", Signal: &signal}},
|
|
||||||
} {
|
|
||||||
signals := make(map[string]int8)
|
|
||||||
for id, reading := range snapshot {
|
|
||||||
if reading.Signal != nil {
|
|
||||||
signals[id] = int8(*reading.Signal)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
payload, err := cbor.Marshal(system.CombinedData{Info: system.Info{WiFi: snapshot}, Stats: system.Stats{WiFi: signals}})
|
|
||||||
require.NoError(t, err)
|
|
||||||
require.NoError(t, UnmarshalResponse(common.AgentResponse{Data: payload}, common.GetData, &decoded))
|
|
||||||
require.Len(t, decoded.Info.WiFi, len(snapshot))
|
|
||||||
require.Len(t, decoded.Stats.WiFi, len(signals))
|
|
||||||
for id, want := range snapshot {
|
|
||||||
require.Equal(t, want, decoded.Info.WiFi[id])
|
|
||||||
}
|
|
||||||
for id, want := range signals {
|
|
||||||
require.Equal(t, want, decoded.Stats.WiFi[id])
|
|
||||||
}
|
|
||||||
}
|
|
||||||
// An older generic-response agent may omit both fields entirely.
|
|
||||||
payload, err := cbor.Marshal(map[int]any{0: map[int]any{}, 1: map[int]any{}})
|
|
||||||
require.NoError(t, err)
|
|
||||||
require.NoError(t, UnmarshalResponse(common.AgentResponse{Data: payload}, common.GetData, &decoded))
|
|
||||||
require.Empty(t, decoded.Info.WiFi)
|
|
||||||
require.Empty(t, decoded.Stats.WiFi)
|
|
||||||
}
|
|
||||||
Reference in New Issue
Block a user