Enhance API and services with catalog update notifications and improved node management

- Added response structure for node details in API documentation.
- Updated maintainer workflow to include dashboard sync instructions.
- Refactored node service to start a worker for catalog synchronization.
- Implemented onChange hook for catalog updates to trigger node evaluations.
- Enhanced node service to manage sensor updates based on active revisions.
- Added WebSocket event handling for catalog updates in the UI.
- Updated fleet store to fetch data on catalog updates.
- Introduced tests for catalog refresh deduplication and node details handling.
This commit is contained in:
AndReicscs
2026-06-16 14:51:32 +00:00
parent c721dabd3d
commit fcad057363
13 changed files with 483 additions and 28 deletions
+22
View File
@@ -153,6 +153,28 @@ Creates a new node entry and returns the generated node credentials.
Returns all registered nodes and their installed sensors. Each node includes current status, heartbeat metadata, and pending config state.
**Response:**
```json
[
{
"nodeId": "node-alpha",
"alias": "Production Web Server",
"activeRevision": "rev_e4f2a1",
"desiredRevision": "rev_b7c9d0",
"hasPendingConfig": true,
"hasUpdateAvailable": true,
"status": "online",
"installedSensors": [
{
"sensorId": "hw-tcp-tarpit",
"customName": "SSH Decoy",
"updateAvailable": true
}
]
}
]
```
#### GET /api/v1/nodes/{nodeId}
Returns details for a single node, including installed sensors and per-sensor event counts.
+3
View File
@@ -124,3 +124,6 @@ Check the `registry-pages` branch to confirm:
- The versioned JSON file exists
- `index.json` lists the new version
- The `latest` field points to the new version
### Step 5: Dashboard Sync
Refresh your HoneyWire dashboard or wait for the automatic UI sync. The Hub's event-driven catalog hook will instantly detect the registry mutation, generate the new cryptographic hashes for your nodes, lock them into 'Pending Sync' state, and immediately broadcast the update to all active UI sessions without latency. You do not need to wait for a 5-minute background loop.
+1 -1
View File
@@ -57,7 +57,6 @@ func main() {
catalogService := catalog.NewService(dbStore, wsService)
nodeSvc := node.NewService(dbStore, wsService, catalogService)
nodeSvc.StartAutoEvaluator()
sensorSvc := sensor.NewService(dbStore, wsService)
siemService := siem.NewService(nodeSvc)
notifyService := notify.NewService(nodeSvc)
@@ -95,6 +94,7 @@ func main() {
go authService.StartWorkers(rootCtx)
go siemService.StartWorker(rootCtx)
go notifyService.StartWorker(rootCtx)
go nodeSvc.StartWorker(rootCtx)
// 2. Load Runtime Configurations Safely
isArmed := loadConfigSafe(dbStore, "is_armed", "false") == "true"
+7 -6
View File
@@ -107,12 +107,13 @@ func (h *NodeHandler) GetNodeDetails(w http.ResponseWriter, r *http.Request) {
"activeRevision": node.ActiveRevision,
"desiredRevision": node.DesiredRevision,
"publicIp": node.PublicIP,
"privateIp": node.PrivateIP,
"tags": node.Tags,
"hasPendingConfig": node.HasPendingConfig,
"lastHeartbeat": node.LastHeartbeat,
"status": node.Status,
"installedSensors": node.InstalledSensors,
"privateIp": node.PrivateIP,
"tags": node.Tags,
"hasPendingConfig": node.HasPendingConfig,
"hasUpdateAvailable": node.HasUpdateAvailable,
"lastHeartbeat": node.LastHeartbeat,
"status": node.Status,
"installedSensors": node.InstalledSensors,
})
}
+26 -2
View File
@@ -35,12 +35,19 @@ type Service struct {
broadcaster Broadcaster
indexCache *RegistryIndex
mu sync.RWMutex
onChange func()
}
func NewService(store Store, broadcaster Broadcaster) *Service {
return &Service{store: store, broadcaster: broadcaster}
}
func (s *Service) SetOnChangeHook(hook func()) {
s.mu.Lock()
s.onChange = hook
s.mu.Unlock()
}
func (s *Service) RefreshIndex() error {
registryURL, err := s.store.GetConfigValue("registry_url")
if err != nil || registryURL == "" {
@@ -61,11 +68,28 @@ func (s *Service) RefreshIndex() error {
}
s.mu.Lock()
var oldJSON, newJSON []byte
if s.indexCache != nil {
oldJSON, _ = json.Marshal(s.indexCache)
}
newJSON, _ = json.Marshal(&idx)
changed := string(oldJSON) != string(newJSON)
s.indexCache = &idx
var hook func()
if changed {
hook = s.onChange
}
s.mu.Unlock()
if s.broadcaster != nil {
s.broadcaster.Broadcast("CATALOG_UPDATED", nil)
if changed {
if s.broadcaster != nil {
s.broadcaster.Broadcast("CATALOG_UPDATED", nil)
}
if hook != nil {
go hook()
}
}
return nil
+153
View File
@@ -0,0 +1,153 @@
package catalog_test
import (
"encoding/json"
"net/http"
"net/http/httptest"
"sync"
"testing"
"time"
"github.com/honeywire/hub/internal/catalog"
)
type MockStore struct {
RegistryURL string
}
func (m *MockStore) GetConfigValue(key string) (string, error) {
if key == "registry_url" {
return m.RegistryURL, nil
}
return "", nil
}
type MockBroadcaster struct {
BroadcastCount int
}
func (m *MockBroadcaster) Broadcast(eventType string, payload interface{}) {
if eventType == "CATALOG_UPDATED" {
m.BroadcastCount++
}
}
func TestRefreshIndexDeduplication(t *testing.T) {
responseIndex := map[string]interface{}{
"sensors": []map[string]interface{}{
{
"id": "hw-sensor-test",
"latest": "v1.0.0",
"versions": []map[string]interface{}{
{"v": "v1.0.0", "min_hub_version": "v1.0.0"},
},
},
},
}
ts := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
json.NewEncoder(w).Encode(responseIndex)
}))
defer ts.Close()
store := &MockStore{RegistryURL: ts.URL}
broadcaster := &MockBroadcaster{}
svc := catalog.NewService(store, broadcaster)
// First fetch should broadcast
if err := svc.RefreshIndex(); err != nil {
t.Fatalf("First RefreshIndex failed: %v", err)
}
if broadcaster.BroadcastCount != 1 {
t.Errorf("Expected exactly 1 broadcast on initial pull, got %d", broadcaster.BroadcastCount)
}
// Second fetch with IDENTICAL data should NOT broadcast
if err := svc.RefreshIndex(); err != nil {
t.Fatalf("Second RefreshIndex failed: %v", err)
}
if broadcaster.BroadcastCount != 1 {
t.Errorf("Expected broadcast count to remain 1 on identical pull, got %d", broadcaster.BroadcastCount)
}
// Third fetch with CHANGED data SHOULD broadcast
responseIndex["sensors"].([]map[string]interface{})[0]["latest"] = "v2.0.0"
if err := svc.RefreshIndex(); err != nil {
t.Fatalf("Third RefreshIndex failed: %v", err)
}
if broadcaster.BroadcastCount != 2 {
t.Errorf("Expected broadcast count to increment to 2 on changed pull, got %d", broadcaster.BroadcastCount)
}
}
func TestRefreshIndexTriggersHook(t *testing.T) {
responseIndex := map[string]interface{}{
"sensors": []map[string]interface{}{
{
"id": "hw-sensor-test",
"latest": "v1.0.0",
},
},
}
ts := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
json.NewEncoder(w).Encode(responseIndex)
}))
defer ts.Close()
store := &MockStore{RegistryURL: ts.URL}
broadcaster := &MockBroadcaster{}
svc := catalog.NewService(store, broadcaster)
var mu sync.Mutex
hookFiredCount := 0
svc.SetOnChangeHook(func() {
mu.Lock()
hookFiredCount++
mu.Unlock()
})
if err := svc.RefreshIndex(); err != nil {
t.Fatalf("First RefreshIndex failed: %v", err)
}
time.Sleep(50 * time.Millisecond)
mu.Lock()
count := hookFiredCount
mu.Unlock()
if count != 1 {
t.Errorf("Expected hook to fire exactly 1 time on initial pull, got %d", count)
}
// Second fetch with IDENTICAL data should NOT trigger hook
if err := svc.RefreshIndex(); err != nil {
t.Fatalf("Second RefreshIndex failed: %v", err)
}
time.Sleep(50 * time.Millisecond)
mu.Lock()
count = hookFiredCount
mu.Unlock()
if count != 1 {
t.Errorf("Expected hook count to remain 1 on identical pull, got %d", count)
}
// Third fetch with CHANGED data SHOULD trigger hook
responseIndex["sensors"].([]map[string]interface{})[0]["latest"] = "v2.0.0"
if err := svc.RefreshIndex(); err != nil {
t.Fatalf("Third RefreshIndex failed: %v", err)
}
time.Sleep(50 * time.Millisecond)
mu.Lock()
count = hookFiredCount
mu.Unlock()
if count != 2 {
t.Errorf("Expected hook count to increment to 2 on changed pull, got %d", count)
}
}
+2 -9
View File
@@ -149,19 +149,12 @@ func (s *Service) GetNodeCompose(token, hostFallback string, currentHubVersion s
hubEndpoint = hostFallback
}
effectiveRevision := nodeDetails.DesiredRevision
if effectiveRevision == "" && nodeDetails.HasPendingConfig {
effectiveRevision = node.GenerateRevisionHash(nodeDetails.InstalledSensors, s.catalog, currentHubVersion)
effectiveRevision := node.GenerateRevisionHash(nodeDetails.InstalledSensors, s.catalog, currentHubVersion)
if effectiveRevision != nodeDetails.DesiredRevision {
if err := s.store.SetNodeDesiredRevision(nodeID, effectiveRevision); err != nil {
return nil, fmt.Errorf("failed_to_allocate")
}
}
if effectiveRevision == "" {
effectiveRevision = nodeDetails.ActiveRevision
}
if effectiveRevision == "" {
effectiveRevision = node.GenerateRevisionHash(nodeDetails.InstalledSensors, s.catalog, currentHubVersion)
}
// Auto-reconcile empty nodes since they have no sensors to report back via heartbeats
if len(nodeDetails.InstalledSensors) == 0 && nodeDetails.HasPendingConfig {
@@ -2,6 +2,7 @@ package composesvc_test
import (
"encoding/json"
"fmt"
"net/http"
"net/http/httptest"
"strings"
@@ -14,6 +15,8 @@ import (
type MockStore struct {
RegistryURL string
AssignedRevision string
ShouldFailSetRevision bool
}
func (m *MockStore) GetConfigValue(key string) (string, error) {
@@ -36,10 +39,15 @@ func (m *MockStore) GetNodeDetails(nodeID string) (*models.Node, error) {
InstalledSensors: []models.NodeSensor{
{ID: "hw-sensor-test"},
},
DesiredRevision: "old_revision",
}, nil
}
func (m *MockStore) SetNodeDesiredRevision(nodeID, rev string) error {
if m.ShouldFailSetRevision {
return fmt.Errorf("simulated database lock timeout")
}
m.AssignedRevision = rev
return nil
}
@@ -165,3 +173,88 @@ func TestComposeSmartVersionSelection(t *testing.T) {
func contains(b []byte, s string) bool {
return strings.Contains(string(b), s)
}
func TestGetNodeComposeSynchronousSync(t *testing.T) {
ts := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
if r.URL.Path == "/index.json" {
json.NewEncoder(w).Encode(map[string]interface{}{
"sensors": []map[string]interface{}{
{
"id": "hw-sensor-test",
"version": "v1.0.0",
},
},
})
return
}
if strings.HasPrefix(r.URL.Path, "/hw-sensor-test-v1.0.0.json") {
json.NewEncoder(w).Encode(map[string]interface{}{
"id": "hw-sensor-test",
"version": "1.0.0",
"deployment": map[string]interface{}{
"image": "test-image",
},
})
return
}
w.WriteHeader(http.StatusNotFound)
}))
defer ts.Close()
store := &MockStore{RegistryURL: ts.URL}
catSvc := catalog.NewService(store, nil)
svc := composesvc.NewService(store, catSvc)
// Refresh index to simulate worker background state
catSvc.RefreshIndex()
_, err := svc.GetNodeCompose("node-1", "http://fallback.com", "v1.0.0")
if err != nil {
t.Fatalf("GetNodeCompose failed: %v", err)
}
// Verify that GetNodeCompose calculated a new hash and explicitly synchronized it!
if store.AssignedRevision == "" || store.AssignedRevision == "old_revision" {
t.Fatalf("Expected DesiredRevision to be dynamically assigned new hash, got: %v", store.AssignedRevision)
}
}
func TestComposeSetDesiredRevisionErrorHandling(t *testing.T) {
ts := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
if r.URL.Path == "/index.json" {
json.NewEncoder(w).Encode(map[string]interface{}{
"sensors": []map[string]interface{}{
{
"id": "hw-sensor-test",
"latest": "v1.0.0",
},
},
})
return
}
if strings.HasPrefix(r.URL.Path, "/hw-sensor-test-v1.0.0.json") {
json.NewEncoder(w).Encode(map[string]interface{}{
"id": "hw-sensor-test",
"version": "1.0.0",
"deployment": map[string]interface{}{
"image": "test-image",
},
})
return
}
w.WriteHeader(http.StatusNotFound)
}))
defer ts.Close()
// Simulate database failure
store := &MockStore{RegistryURL: ts.URL, ShouldFailSetRevision: true}
catSvc := catalog.NewService(store, nil)
svc := composesvc.NewService(store, catSvc)
catSvc.RefreshIndex()
_, err := svc.GetNodeCompose("node-1", "http://fallback.com", "v1.0.0")
if err == nil || !strings.Contains(err.Error(), "failed_to_allocate") {
t.Fatalf("Expected failed_to_allocate error to cascade up, got: %v", err)
}
}
+40 -9
View File
@@ -1,9 +1,11 @@
package node
import (
"context"
"crypto/sha256"
"encoding/hex"
"encoding/json"
"log"
"sort"
"time"
@@ -45,16 +47,34 @@ func NewService(store Store, broadcaster Broadcaster, cat *catalog.Service) *Ser
}
}
// StartAutoEvaluator runs a background thread that periodically refreshes the catalog
// StartWorker runs a background thread that periodically refreshes the catalog
// and recalculates the node sync states to instantly flag updates natively.
func (s *Service) StartAutoEvaluator() {
if s.catalog == nil {
return
func (s *Service) StartWorker(ctx context.Context) {
log.Println("[INFO] Starting node sync background worker...")
if s.catalog != nil {
s.catalog.SetOnChangeHook(func() {
nodes, err := s.store.GetNodes()
if err == nil {
for _, n := range nodes {
s.evaluateNodeSyncState(n.ID)
}
}
})
}
go func() {
for {
time.Sleep(5 * time.Minute)
s.catalog.RefreshIndex()
ticker := time.NewTicker(5 * time.Minute)
defer ticker.Stop()
for {
select {
case <-ctx.Done():
log.Println("[INFO] Node sync worker stopped")
return
case <-ticker.C:
if s.catalog != nil {
s.catalog.RefreshIndex()
}
nodes, err := s.store.GetNodes()
if err == nil {
for _, n := range nodes {
@@ -62,7 +82,7 @@ func (s *Service) StartAutoEvaluator() {
}
}
}
}()
}
}
func (s *Service) CreateNode(alias string, tags []string) (string, string, error) {
@@ -107,6 +127,11 @@ func (s *Service) GetNodes() ([]models.Node, error) {
if err == nil && latest != "" {
deployedVer := sensor.DeployedVersion
if deployedVer != latest {
newHash := GenerateRevisionHash(nodes[i].InstalledSensors, s.catalog, models.HubVersion)
if nodes[i].ActiveRevision != "" && nodes[i].ActiveRevision == newHash {
_ = s.store.SetNodeSensorDeployedVersion(nodes[i].ID, sensor.ID, latest)
continue
}
nodes[i].HasUpdateAvailable = true
break
}
@@ -131,6 +156,12 @@ func (s *Service) GetNodeDetails(nodeID string) (*models.Node, error) {
if err == nil && latest != "" {
deployedVer := sensor.DeployedVersion
if deployedVer != latest {
newHash := GenerateRevisionHash(node.InstalledSensors, s.catalog, models.HubVersion)
if node.ActiveRevision != "" && node.ActiveRevision == newHash {
_ = s.store.SetNodeSensorDeployedVersion(node.ID, sensor.ID, latest)
node.InstalledSensors[i].DeployedVersion = latest
continue
}
node.InstalledSensors[i].UpdateAvailable = true
node.HasUpdateAvailable = true
}
+132
View File
@@ -0,0 +1,132 @@
package node_test
import (
"encoding/json"
"net/http"
"net/http/httptest"
"testing"
"github.com/honeywire/hub/internal/catalog"
"github.com/honeywire/hub/internal/models"
"github.com/honeywire/hub/internal/services/node"
)
type MockNodeStore struct {
nodes map[string]*models.Node
DeployedUpdates map[string]string // Tracks bumped DeployedVersions for assertion
}
func NewMockNodeStore() *MockNodeStore {
return &MockNodeStore{
nodes: make(map[string]*models.Node),
DeployedUpdates: make(map[string]string),
}
}
func (m *MockNodeStore) GetNodes() ([]models.Node, error) {
var result []models.Node
for _, n := range m.nodes {
result = append(result, *n)
}
return result, nil
}
func (m *MockNodeStore) GetNodeDetails(nodeID string) (*models.Node, error) {
if n, ok := m.nodes[nodeID]; ok {
return n, nil
}
return nil, nil
}
func (m *MockNodeStore) SetNodeSensorDeployedVersion(nodeID, sensorID, version string) error {
m.DeployedUpdates[nodeID+":"+sensorID] = version
return nil
}
func (m *MockNodeStore) GetConfigValue(key string) (string, error) {
return "", nil
}
// Unused methods required by interface
func (m *MockNodeStore) CreateNode(alias, tags string) (string, string, error) { return "", "", nil }
func (m *MockNodeStore) UpdateNodeMeta(nodeID, alias, tags string, publicIP, privateIP *string) error { return nil }
func (m *MockNodeStore) AddSensorToNode(nodeID, sensorID, customName string, configValues map[string]interface{}) error { return nil }
func (m *MockNodeStore) UpdateNodeSensor(nodeID, sensorID, customName string, configValues map[string]interface{}) error { return nil }
func (m *MockNodeStore) RemoveNodeSensor(nodeID, sensorID string) error { return nil }
func (m *MockNodeStore) SetNodePendingConfig(nodeID string) error { return nil }
func (m *MockNodeStore) ClearNodePendingConfig(nodeID string) error { return nil }
func (m *MockNodeStore) DeleteNode(nodeID string) error { return nil }
type MockNodeBroadcaster struct{}
func (m *MockNodeBroadcaster) Broadcast(eventType string, payload interface{}) {}
func TestGetNodeDetailsStrictHashMatch(t *testing.T) {
responseIndex := map[string]interface{}{
"sensors": []map[string]interface{}{
{
"id": "hw-sensor-test",
"latest": "v2.0.0", // Latest version available is v2.0.0
"versions": []map[string]interface{}{
{"v": "v2.0.0", "min_hub_version": "v1.0.0"},
},
},
},
}
ts := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
json.NewEncoder(w).Encode(responseIndex)
}))
defer ts.Close()
catStore := &CatalogStore{url: ts.URL}
catSvc := catalog.NewService(catStore, nil)
catSvc.RefreshIndex() // Hydrate catalog
store := NewMockNodeStore()
store.nodes["node-1"] = &models.Node{
ID: "node-1",
ActiveRevision: "old_revision_hash",
InstalledSensors: []models.NodeSensor{
{
ID: "hw-sensor-test",
DeployedVersion: "v1.0.0", // Node currently has v1.0.0
},
},
}
svc := node.NewService(store, &MockNodeBroadcaster{}, catSvc)
// Call GetNodeDetails - ActiveRevision ("old_revision_hash") DOES NOT match newly generated hash!
_, err := svc.GetNodeDetails("node-1")
if err != nil {
t.Fatalf("GetNodeDetails failed: %v", err)
}
// Because hashes did not match, it MUST NOT auto-bump the DeployedVersion in the database
if val, ok := store.DeployedUpdates["node-1:hw-sensor-test"]; ok {
t.Fatalf("Expected DeployedVersion to REMAIN v1.0.0, but it aggressively auto-bumped to %s", val)
}
// Now we spoof a successful Edge Node heartbeat that updates ActiveRevision to perfectly match the newest hash
newHash := node.GenerateRevisionHash(store.nodes["node-1"].InstalledSensors, catSvc, models.HubVersion)
store.nodes["node-1"].ActiveRevision = newHash
_, err = svc.GetNodeDetails("node-1")
if err != nil {
t.Fatalf("GetNodeDetails failed: %v", err)
}
// Because hashes NOW MATCH perfectly, it MUST auto-bump the DeployedVersion
if val, ok := store.DeployedUpdates["node-1:hw-sensor-test"]; !ok || val != "v2.0.0" {
t.Fatalf("Expected DeployedVersion to auto-bump to v2.0.0 upon valid hash match, got: %v", val)
}
}
type CatalogStore struct{ url string }
func (c *CatalogStore) GetConfigValue(key string) (string, error) {
if key == "registry_url" {
return c.url, nil
}
return "", nil
}
+1
View File
@@ -70,6 +70,7 @@ const loadAppData = async () => {
wsService.on('onUpdateNode', (payload: any) => fleetStore.handleWsUpdate('UPDATE_NODE', payload))
wsService.on('onDeleteNode', (payload: any) => fleetStore.handleWsUpdate('DELETE_NODE', payload))
wsService.on('onNodeSynced', (payload: any) => fleetStore.handleWsUpdate('NODE_SYNCED', payload))
wsService.on('onCatalogUpdated', (payload: any) => fleetStore.handleWsUpdate('CATALOG_UPDATED', payload))
wsService.on('onReconnect', async () => {
console.log("WebSocket Reconnected: Syncing missed data...")
+2 -1
View File
@@ -52,7 +52,8 @@ export class HoneyWireWS {
'NEW_NODE': 'onNewNode',
'UPDATE_NODE': 'onUpdateNode',
'DELETE_NODE': 'onDeleteNode',
'NODE_SYNCED': 'onNodeSynced'
'NODE_SYNCED': 'onNodeSynced',
'CATALOG_UPDATED': 'onCatalogUpdated'
};
// Callback registry
+1
View File
@@ -458,6 +458,7 @@ export const useFleetStore = defineStore('fleet', () => {
}
if (type === 'CATALOG_UPDATED') {
fetchFleet()
fetchManifests()
return
}