605bfa1ec9
Implements AgentService gRPC server with all 6 RPCs (Enroll/Register/Heartbeat/ Subscribe/Ack/ReportUsage), Hub command routing with Redis ZSET at-least-once persistence and cross-instance pub/sub delivery, LoadCache for node:load metrics, NodeStore SQL interface + MySQL implementation, and full mTLS gRPC listener in main. Integration tests: 18 tests covering full Enroll→Register→Heartbeat→Subscribe→Ack flow, reconnect resume with last_command_id, and cross-instance pub/sub delivery via bufconn + miniredis + mockNodeStore + real mTLS certificates. Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
61 lines
1.5 KiB
Go
61 lines
1.5 KiB
Go
package nodes
|
|
|
|
import (
|
|
"github.com/redis/go-redis/v9"
|
|
|
|
"github.com/wangjia/pangolin/server/internal/mtls"
|
|
agentv1 "github.com/wangjia/pangolin/server/internal/pb/agentv1"
|
|
)
|
|
|
|
// Service bundles all nodes-domain components and exposes assembly helpers.
|
|
// Construct via NewService after initialising each dependency independently.
|
|
type Service struct {
|
|
handler *Handler
|
|
hub *Hub
|
|
store NodeStore
|
|
load *LoadCache
|
|
}
|
|
|
|
// NewService wires up the full nodes service from its dependencies.
|
|
//
|
|
// - ca / tokens / crl: from the mtls package (task 5b)
|
|
// - rdb: Redis client (for hub + load cache)
|
|
// - store: NodeStore implementation (SQLNodeStore in production; mock in tests)
|
|
func NewService(
|
|
ca *mtls.CA,
|
|
tokens *mtls.BootstrapTokenManager,
|
|
rdb *redis.Client,
|
|
store NodeStore,
|
|
) *Service {
|
|
hub := NewHub(rdb)
|
|
load := NewLoadCache(rdb)
|
|
handler := NewHandler(ca, tokens, hub, store, load)
|
|
return &Service{
|
|
handler: handler,
|
|
hub: hub,
|
|
store: store,
|
|
load: load,
|
|
}
|
|
}
|
|
|
|
// Handler returns the AgentServiceServer implementation for gRPC registration.
|
|
func (s *Service) Handler() agentv1.AgentServiceServer {
|
|
return s.handler
|
|
}
|
|
|
|
// Hub exposes the command routing hub for callers (e.g. task 5d/5e) that need to
|
|
// Push or Broadcast commands.
|
|
func (s *Service) Hub() *Hub {
|
|
return s.hub
|
|
}
|
|
|
|
// Store exposes the NodeStore for callers that need direct DB access.
|
|
func (s *Service) Store() NodeStore {
|
|
return s.store
|
|
}
|
|
|
|
// Load exposes the LoadCache for callers that display per-node load metrics.
|
|
func (s *Service) Load() *LoadCache {
|
|
return s.load
|
|
}
|