From a2bfa98a2d75f91b199117f4c5c6425e452c16c0 Mon Sep 17 00:00:00 2001 From: mrhid6 Date: Tue, 21 Jul 2026 14:18:49 +0100 Subject: [PATCH] feat(proto): SyncMonitors + ReportChecks RPCs --- agent/internal/grpc/client.go | 17 +++++ agent/internal/grpc/pb/vantage.pb.go | 54 +++++++++++++++ proto/vantage/v1/vantage.proto | 41 ++++++++++++ server/internal/grpc/pb/vantage.pb.go | 96 +++++++++++++++++++++++++++ server/internal/grpc/server.go | 46 +++++++++++++ 5 files changed, 254 insertions(+) diff --git a/agent/internal/grpc/client.go b/agent/internal/grpc/client.go index f5e4746..9f07404 100644 --- a/agent/internal/grpc/client.go +++ b/agent/internal/grpc/client.go @@ -133,6 +133,23 @@ func (c *Client) ReportInventory(report *pb.InventoryReport) error { return err } +func (c *Client) SyncMonitors(serverID, agentToken string) ([]pb.MonitorSpec, error) { + ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second) + defer cancel() + resp, err := c.client.SyncMonitors(ctx, &pb.SyncMonitorsRequest{ServerId: serverID, AgentToken: agentToken}) + if err != nil { + return nil, err + } + return resp.Monitors, nil +} + +func (c *Client) ReportChecks(serverID, agentToken string, results []pb.CheckResult) error { + ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second) + defer cancel() + _, err := c.client.ReportChecks(ctx, &pb.ReportChecksRequest{ServerId: serverID, AgentToken: agentToken, Results: results}) + return err +} + // CommandStream opens a long-lived bidirectional stream for server-pushed commands. // The caller controls the stream lifetime via ctx. func (c *Client) CommandStream(ctx context.Context) (pb.Vantage_CommandStreamClient, error) { diff --git a/agent/internal/grpc/pb/vantage.pb.go b/agent/internal/grpc/pb/vantage.pb.go index bcbd19b..723b59e 100644 --- a/agent/internal/grpc/pb/vantage.pb.go +++ b/agent/internal/grpc/pb/vantage.pb.go @@ -92,6 +92,42 @@ type InventoryReport struct { } type InventoryReportResponse struct{} +// Monitor sync / check report message types + +type MonitorSpec struct { + MonitorId string `json:"monitor_id"` + Type string `json:"type"` + URL string `json:"url,omitempty"` + Host string `json:"host,omitempty"` + Port int `json:"port,omitempty"` + Method string `json:"method,omitempty"` + ExpectedStatus int `json:"expected_status,omitempty"` + Keyword string `json:"keyword,omitempty"` + TLSWarnDays int `json:"tls_warn_days,omitempty"` + IntervalSec int `json:"interval_sec"` + Retries int `json:"retries"` +} +type SyncMonitorsRequest struct { + ServerId string `json:"server_id"` + AgentToken string `json:"agent_token"` +} +type SyncMonitorsResponse struct { + Monitors []MonitorSpec `json:"monitors,omitempty"` +} +type CheckResult struct { + MonitorId string `json:"monitor_id"` + Up bool `json:"up"` + LatencyMs int `json:"latency_ms"` + Message string `json:"message,omitempty"` + CertExpiryUnix int64 `json:"cert_expiry_unix,omitempty"` +} +type ReportChecksRequest struct { + ServerId string `json:"server_id"` + AgentToken string `json:"agent_token"` + Results []CheckResult `json:"results,omitempty"` +} +type ReportChecksResponse struct{} + type ApplyUpdatesCmd struct{} type ServerCommand struct { @@ -223,6 +259,8 @@ type VantageClient interface { UploadGeneratedKey(ctx context.Context, in *UploadKeyRequest, opts ...grpc.CallOption) (*UploadKeyResponse, error) ReportUpdates(ctx context.Context, in *ReportUpdatesRequest, opts ...grpc.CallOption) (*ReportUpdatesResponse, error) ReportInventory(ctx context.Context, in *InventoryReport, opts ...grpc.CallOption) (*InventoryReportResponse, error) + SyncMonitors(ctx context.Context, in *SyncMonitorsRequest, opts ...grpc.CallOption) (*SyncMonitorsResponse, error) + ReportChecks(ctx context.Context, in *ReportChecksRequest, opts ...grpc.CallOption) (*ReportChecksResponse, error) CommandStream(ctx context.Context, opts ...grpc.CallOption) (Vantage_CommandStreamClient, error) } @@ -286,6 +324,22 @@ func (c *keyManagerClient) ReportInventory(ctx context.Context, in *InventoryRep return out, nil } +func (c *keyManagerClient) SyncMonitors(ctx context.Context, in *SyncMonitorsRequest, opts ...grpc.CallOption) (*SyncMonitorsResponse, error) { + out := new(SyncMonitorsResponse) + if err := c.cc.Invoke(ctx, "/vantage.v1.Vantage/SyncMonitors", in, out, opts...); err != nil { + return nil, err + } + return out, nil +} + +func (c *keyManagerClient) ReportChecks(ctx context.Context, in *ReportChecksRequest, opts ...grpc.CallOption) (*ReportChecksResponse, error) { + out := new(ReportChecksResponse) + if err := c.cc.Invoke(ctx, "/vantage.v1.Vantage/ReportChecks", in, out, opts...); err != nil { + return nil, err + } + return out, nil +} + func (c *keyManagerClient) CommandStream(ctx context.Context, opts ...grpc.CallOption) (Vantage_CommandStreamClient, error) { desc := &grpc.StreamDesc{StreamName: "CommandStream", ServerStreams: true, ClientStreams: true} stream, err := c.cc.NewStream(ctx, desc, "/vantage.v1.Vantage/CommandStream", opts...) diff --git a/proto/vantage/v1/vantage.proto b/proto/vantage/v1/vantage.proto index a4fc92a..e5fda50 100644 --- a/proto/vantage/v1/vantage.proto +++ b/proto/vantage/v1/vantage.proto @@ -10,6 +10,8 @@ service Vantage { rpc UploadGeneratedKey(UploadKeyRequest) returns (UploadKeyResponse); rpc ReportUpdates(ReportUpdatesRequest) returns (ReportUpdatesResponse); rpc ReportInventory(InventoryReport) returns (InventoryReportResponse); + rpc SyncMonitors(SyncMonitorsRequest) returns (SyncMonitorsResponse); + rpc ReportChecks(ReportChecksRequest) returns (ReportChecksResponse); // Bidirectional stream: agent sends auth once, server pushes commands. rpc CommandStream(stream AgentMessage) returns (stream ServerCommand); } @@ -117,6 +119,45 @@ message InventoryReport { message InventoryReportResponse {} +message MonitorSpec { + string monitor_id = 1; + string type = 2; + string url = 3; + string host = 4; + int32 port = 5; + string method = 6; + int32 expected_status = 7; + string keyword = 8; + int32 tls_warn_days = 9; + int32 interval_sec = 10; + int32 retries = 11; +} + +message SyncMonitorsRequest { + string server_id = 1; + string agent_token = 2; +} + +message SyncMonitorsResponse { + repeated MonitorSpec monitors = 1; +} + +message CheckResult { + string monitor_id = 1; + bool up = 2; + int32 latency_ms = 3; + string message = 4; + int64 cert_expiry_unix = 5; +} + +message ReportChecksRequest { + string server_id = 1; + string agent_token = 2; + repeated CheckResult results = 3; +} + +message ReportChecksResponse {} + message ApplyUpdatesCmd {} message ServerCommand { diff --git a/server/internal/grpc/pb/vantage.pb.go b/server/internal/grpc/pb/vantage.pb.go index b835d65..3cf3882 100644 --- a/server/internal/grpc/pb/vantage.pb.go +++ b/server/internal/grpc/pb/vantage.pb.go @@ -95,6 +95,42 @@ type InventoryReport struct { } type InventoryReportResponse struct{} +// Monitor sync / check report message types + +type MonitorSpec struct { + MonitorId string `json:"monitor_id"` + Type string `json:"type"` + URL string `json:"url,omitempty"` + Host string `json:"host,omitempty"` + Port int `json:"port,omitempty"` + Method string `json:"method,omitempty"` + ExpectedStatus int `json:"expected_status,omitempty"` + Keyword string `json:"keyword,omitempty"` + TLSWarnDays int `json:"tls_warn_days,omitempty"` + IntervalSec int `json:"interval_sec"` + Retries int `json:"retries"` +} +type SyncMonitorsRequest struct { + ServerId string `json:"server_id"` + AgentToken string `json:"agent_token"` +} +type SyncMonitorsResponse struct { + Monitors []MonitorSpec `json:"monitors,omitempty"` +} +type CheckResult struct { + MonitorId string `json:"monitor_id"` + Up bool `json:"up"` + LatencyMs int `json:"latency_ms"` + Message string `json:"message,omitempty"` + CertExpiryUnix int64 `json:"cert_expiry_unix,omitempty"` +} +type ReportChecksRequest struct { + ServerId string `json:"server_id"` + AgentToken string `json:"agent_token"` + Results []CheckResult `json:"results,omitempty"` +} +type ReportChecksResponse struct{} + type ApplyUpdatesCmd struct{} type ServerCommand struct { @@ -228,6 +264,8 @@ type VantageServer interface { UploadGeneratedKey(context.Context, *UploadKeyRequest) (*UploadKeyResponse, error) ReportUpdates(context.Context, *ReportUpdatesRequest) (*ReportUpdatesResponse, error) ReportInventory(context.Context, *InventoryReport) (*InventoryReportResponse, error) + SyncMonitors(context.Context, *SyncMonitorsRequest) (*SyncMonitorsResponse, error) + ReportChecks(context.Context, *ReportChecksRequest) (*ReportChecksResponse, error) CommandStream(Vantage_CommandStreamServer) error } @@ -253,6 +291,14 @@ func (UnimplementedVantageServer) ReportInventory(context.Context, *InventoryRep return nil, status.Errorf(codes.Unimplemented, "method ReportInventory not implemented") } +func (UnimplementedVantageServer) SyncMonitors(context.Context, *SyncMonitorsRequest) (*SyncMonitorsResponse, error) { + return nil, status.Errorf(codes.Unimplemented, "method SyncMonitors not implemented") +} + +func (UnimplementedVantageServer) ReportChecks(context.Context, *ReportChecksRequest) (*ReportChecksResponse, error) { + return nil, status.Errorf(codes.Unimplemented, "method ReportChecks not implemented") +} + func (UnimplementedVantageServer) CommandStream(Vantage_CommandStreamServer) error { return status.Errorf(codes.Unimplemented, "method CommandStream not implemented") } @@ -265,6 +311,8 @@ type VantageClient interface { UploadGeneratedKey(ctx context.Context, in *UploadKeyRequest, opts ...grpc.CallOption) (*UploadKeyResponse, error) ReportUpdates(ctx context.Context, in *ReportUpdatesRequest, opts ...grpc.CallOption) (*ReportUpdatesResponse, error) ReportInventory(ctx context.Context, in *InventoryReport, opts ...grpc.CallOption) (*InventoryReportResponse, error) + SyncMonitors(ctx context.Context, in *SyncMonitorsRequest, opts ...grpc.CallOption) (*SyncMonitorsResponse, error) + ReportChecks(ctx context.Context, in *ReportChecksRequest, opts ...grpc.CallOption) (*ReportChecksResponse, error) CommandStream(ctx context.Context, opts ...grpc.CallOption) (Vantage_CommandStreamClient, error) } @@ -316,6 +364,22 @@ func (c *keyManagerClient) ReportInventory(ctx context.Context, in *InventoryRep return out, nil } +func (c *keyManagerClient) SyncMonitors(ctx context.Context, in *SyncMonitorsRequest, opts ...grpc.CallOption) (*SyncMonitorsResponse, error) { + out := new(SyncMonitorsResponse) + if err := c.cc.Invoke(ctx, "/vantage.v1.Vantage/SyncMonitors", in, out, opts...); err != nil { + return nil, err + } + return out, nil +} + +func (c *keyManagerClient) ReportChecks(ctx context.Context, in *ReportChecksRequest, opts ...grpc.CallOption) (*ReportChecksResponse, error) { + out := new(ReportChecksResponse) + if err := c.cc.Invoke(ctx, "/vantage.v1.Vantage/ReportChecks", in, out, opts...); err != nil { + return nil, err + } + return out, nil +} + func (c *keyManagerClient) CommandStream(ctx context.Context, opts ...grpc.CallOption) (Vantage_CommandStreamClient, error) { stream, err := c.cc.NewStream(ctx, &Vantage_ServiceDesc.Streams[0], "/vantage.v1.Vantage/CommandStream", opts...) if err != nil { @@ -339,6 +403,8 @@ var Vantage_ServiceDesc = grpc.ServiceDesc{ {MethodName: "UploadGeneratedKey", Handler: _Vantage_UploadGeneratedKey_Handler}, {MethodName: "ReportUpdates", Handler: _Vantage_ReportUpdates_Handler}, {MethodName: "ReportInventory", Handler: _Vantage_ReportInventory_Handler}, + {MethodName: "SyncMonitors", Handler: _Vantage_SyncMonitors_Handler}, + {MethodName: "ReportChecks", Handler: _Vantage_ReportChecks_Handler}, }, Streams: []grpc.StreamDesc{ { @@ -426,6 +492,36 @@ func _Vantage_ReportInventory_Handler(srv interface{}, ctx context.Context, dec return interceptor(ctx, in, info, handler) } +func _Vantage_SyncMonitors_Handler(srv interface{}, ctx context.Context, dec func(interface{}) error, interceptor grpc.UnaryServerInterceptor) (interface{}, error) { + in := new(SyncMonitorsRequest) + if err := dec(in); err != nil { + return nil, err + } + if interceptor == nil { + return srv.(VantageServer).SyncMonitors(ctx, in) + } + info := &grpc.UnaryServerInfo{Server: srv, FullMethod: "/vantage.v1.Vantage/SyncMonitors"} + handler := func(ctx context.Context, req interface{}) (interface{}, error) { + return srv.(VantageServer).SyncMonitors(ctx, req.(*SyncMonitorsRequest)) + } + return interceptor(ctx, in, info, handler) +} + +func _Vantage_ReportChecks_Handler(srv interface{}, ctx context.Context, dec func(interface{}) error, interceptor grpc.UnaryServerInterceptor) (interface{}, error) { + in := new(ReportChecksRequest) + if err := dec(in); err != nil { + return nil, err + } + if interceptor == nil { + return srv.(VantageServer).ReportChecks(ctx, in) + } + info := &grpc.UnaryServerInfo{Server: srv, FullMethod: "/vantage.v1.Vantage/ReportChecks"} + handler := func(ctx context.Context, req interface{}) (interface{}, error) { + return srv.(VantageServer).ReportChecks(ctx, req.(*ReportChecksRequest)) + } + return interceptor(ctx, in, info, handler) +} + func _Vantage_CommandStream_Handler(srv interface{}, stream grpc.ServerStream) error { return srv.(VantageServer).CommandStream(&keyManagerCommandStreamServer{stream}) } diff --git a/server/internal/grpc/server.go b/server/internal/grpc/server.go index 2a0e7f8..e7888b6 100644 --- a/server/internal/grpc/server.go +++ b/server/internal/grpc/server.go @@ -7,6 +7,7 @@ import ( "net" "time" + "github.com/mrhid6/vantage/server/internal/checker" "github.com/mrhid6/vantage/server/internal/grpc/pb" "github.com/mrhid6/vantage/server/internal/models" "github.com/mrhid6/vantage/server/internal/services" @@ -106,6 +107,51 @@ func (s *vantageServer) ReportInventory(ctx context.Context, req *pb.InventoryRe return &pb.InventoryReportResponse{}, nil } +func (s *vantageServer) SyncMonitors(ctx context.Context, req *pb.SyncMonitorsRequest) (*pb.SyncMonitorsResponse, error) { + srv, err := services.ValidateAgentToken(req.ServerId, req.AgentToken) + if err != nil { + return nil, status.Errorf(codes.Unauthenticated, "invalid agent token") + } + monitors, err := services.ListMonitorsForRunner(srv.ServerID) + if err != nil { + return nil, status.Errorf(codes.Internal, "list monitors") + } + specs := make([]pb.MonitorSpec, 0, len(monitors)) + for _, m := range monitors { + specs = append(specs, pb.MonitorSpec{ + MonitorId: m.MonitorID, + Type: m.Type, + URL: m.Target.URL, + Host: m.Target.Host, + Port: m.Target.Port, + Method: m.Target.Method, + ExpectedStatus: m.Target.ExpectedStatus, + Keyword: m.Target.Keyword, + TLSWarnDays: m.Target.TLSWarnDays, + IntervalSec: m.IntervalSec, + Retries: m.Retries, + }) + } + return &pb.SyncMonitorsResponse{Monitors: specs}, nil +} + +func (s *vantageServer) ReportChecks(ctx context.Context, req *pb.ReportChecksRequest) (*pb.ReportChecksResponse, error) { + if _, err := services.ValidateAgentToken(req.ServerId, req.AgentToken); err != nil { + return nil, status.Errorf(codes.Unauthenticated, "invalid agent token") + } + for _, r := range req.Results { + res := checker.Result{Up: r.Up, LatencyMs: r.LatencyMs, Message: r.Message} + if r.CertExpiryUnix > 0 { + t := time.Unix(r.CertExpiryUnix, 0) + res.CertExpiry = &t + } + if err := services.IngestResult(r.MonitorId, res); err != nil { + log.Printf("ingest check %s: %v", r.MonitorId, err) + } + } + return &pb.ReportChecksResponse{}, nil +} + func (s *vantageServer) CommandStream(stream pb.Vantage_CommandStreamServer) error { // First message authenticates the agent and signals readiness. msg, err := stream.Recv()