diff --git a/agent/internal/grpc/pb/vantage.pb.go b/agent/internal/grpc/pb/vantage.pb.go index bfd3548..a914114 100644 --- a/agent/internal/grpc/pb/vantage.pb.go +++ b/agent/internal/grpc/pb/vantage.pb.go @@ -206,6 +206,10 @@ type ServerCommand struct { CleanupWorkspace *CleanupWorkspaceCmd `json:"cleanup_workspace,omitempty"` OpenProxy *OpenProxyCmd `json:"open_proxy,omitempty"` Ping *PingCmd `json:"ping,omitempty"` + + RefreshWorkloads *RefreshWorkloadsCmd `json:"refresh_workloads,omitempty"` + ControlWorkload *ControlWorkloadCmd `json:"control_workload,omitempty"` + WorkloadLogs *WorkloadLogsCmd `json:"workload_logs,omitempty"` } // PingCmd is a server-originated liveness beat. It carries nothing and expects @@ -243,6 +247,8 @@ type AgentMessage struct { Result *CommandResult `json:"result,omitempty"` StepResult *StepResult `json:"step_result,omitempty"` StepOutput *StepOutputChunk `json:"step_output,omitempty"` + + WorkloadLogsResult *WorkloadLogsResult `json:"workload_logs_result,omitempty"` } type AgentReady struct{} @@ -377,6 +383,7 @@ 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) ReportPackages(ctx context.Context, in *ReportPackagesRequest, opts ...grpc.CallOption) (*ReportPackagesResponse, error) + ReportWorkloads(ctx context.Context, in *ReportWorkloadsRequest, opts ...grpc.CallOption) (*ReportWorkloadsResponse, 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) diff --git a/agent/internal/grpc/pb/workloads.pb.go b/agent/internal/grpc/pb/workloads.pb.go new file mode 100644 index 0000000..3222e41 --- /dev/null +++ b/agent/internal/grpc/pb/workloads.pb.go @@ -0,0 +1,77 @@ +package pb + +import ( + "context" + + "google.golang.org/grpc" +) + +// Workload registry messages. Hand-written like the rest of this package: the +// .proto is the contract, this file is the Go side of it, and the two must be +// changed together. + +// Workload is one container or one systemd unit. +type Workload struct { + Kind string `json:"kind"` + Id string `json:"id"` + Name string `json:"name"` + State string `json:"state"` + Health string `json:"health,omitempty"` + Image string `json:"image,omitempty"` + Stack string `json:"stack,omitempty"` + Ports []string `json:"ports,omitempty"` + Restarts int32 `json:"restarts,omitempty"` + StartedAt string `json:"started_at,omitempty"` // RFC3339, empty when not running + Protected bool `json:"protected,omitempty"` +} + +// ReportWorkloadsRequest carries what a server is running. +// +// Offer-then-send, the same handshake as ReportPackages: the agent calls once +// with Workloads empty, and resends with the body only if NeedFull is set. +type ReportWorkloadsRequest struct { + ServerId string `json:"server_id"` + AgentToken string `json:"agent_token"` + Hash string `json:"hash"` + DockerOk bool `json:"docker_ok"` + DockerError string `json:"docker_error,omitempty"` + SystemdOk bool `json:"systemd_ok"` + SystemdError string `json:"systemd_error,omitempty"` + Workloads []Workload `json:"workloads,omitempty"` // empty on the offer call +} + +type ReportWorkloadsResponse struct { + NeedFull bool `json:"need_full"` +} + +// RefreshWorkloadsCmd carries no payload back. It makes the agent report +// immediately through ReportWorkloads, so there is exactly one writer for the +// server_workloads collection rather than two arriving by different routes. +type RefreshWorkloadsCmd struct{} + +type ControlWorkloadCmd struct { + Kind string `json:"kind"` + Id string `json:"id"` + Action string `json:"action"` // start | stop | restart +} + +type WorkloadLogsCmd struct { + Kind string `json:"kind"` + Id string `json:"id"` + Tail int32 `json:"tail,omitempty"` +} + +type WorkloadLogsResult struct { + CommandId string `json:"command_id"` + Text string `json:"text,omitempty"` + Truncated bool `json:"truncated,omitempty"` + Error string `json:"error,omitempty"` +} + +func (c *keyManagerClient) ReportWorkloads(ctx context.Context, in *ReportWorkloadsRequest, opts ...grpc.CallOption) (*ReportWorkloadsResponse, error) { + out := new(ReportWorkloadsResponse) + if err := c.cc.Invoke(ctx, "/vantage.v1.Vantage/ReportWorkloads", in, out, opts...); err != nil { + return nil, err + } + return out, nil +} diff --git a/proto/vantage/v1/vantage.proto b/proto/vantage/v1/vantage.proto index 245e7bf..b753dc2 100644 --- a/proto/vantage/v1/vantage.proto +++ b/proto/vantage/v1/vantage.proto @@ -10,6 +10,7 @@ service Vantage { rpc UploadGeneratedKey(UploadKeyRequest) returns (UploadKeyResponse); rpc ReportUpdates(ReportUpdatesRequest) returns (ReportUpdatesResponse); rpc ReportPackages(ReportPackagesRequest) returns (ReportPackagesResponse); + rpc ReportWorkloads(ReportWorkloadsRequest) returns (ReportWorkloadsResponse); rpc ReportInventory(InventoryReport) returns (InventoryReportResponse); rpc SyncMonitors(SyncMonitorsRequest) returns (SyncMonitorsResponse); rpc ReportChecks(ReportChecksRequest) returns (ReportChecksResponse); @@ -107,6 +108,7 @@ message AgentMessage { CommandResult result = 4; StepResult step_result = 5; StepOutputChunk step_output = 6; + WorkloadLogsResult workload_logs_result = 7; } } @@ -229,6 +231,9 @@ message ServerCommand { CleanupWorkspaceCmd cleanup_workspace = 7; OpenProxyCmd open_proxy = 8; PingCmd ping = 9; + RefreshWorkloadsCmd refresh_workloads = 10; + ControlWorkloadCmd control_workload = 11; + WorkloadLogsCmd workload_logs = 12; } } @@ -319,3 +324,63 @@ message ProxyServerMsg { ProxyClose close = 2; } } + +// --------------------------------------------------------------------------- +// Workload registry + +// ReportWorkloads carries what a server is running. +// +// Offer-then-send, the same handshake as ReportPackages: the agent calls once +// with workloads empty, and resends with the body only if need_full is set. +message ReportWorkloadsRequest { + string server_id = 1; + string agent_token = 2; + string hash = 3; + bool docker_ok = 4; + string docker_error = 5; + bool systemd_ok = 6; + string systemd_error = 7; + repeated Workload workloads = 8; // empty on the offer call +} + +message ReportWorkloadsResponse { + bool need_full = 1; +} + +message Workload { + string kind = 1; // "container" | "unit" + string id = 2; + string name = 3; + string state = 4; + string health = 5; + string image = 6; + string stack = 7; + repeated string ports = 8; + int32 restarts = 9; + string started_at = 10; // RFC3339, empty when not running + bool protected = 11; +} + +// RefreshWorkloadsCmd carries no payload back. It makes the agent report +// immediately through ReportWorkloads, so there is exactly one writer for the +// server_workloads collection rather than two arriving by different routes. +message RefreshWorkloadsCmd {} + +message ControlWorkloadCmd { + string kind = 1; + string id = 2; + string action = 3; // "start" | "stop" | "restart" +} + +message WorkloadLogsCmd { + string kind = 1; + string id = 2; + int32 tail = 3; +} + +message WorkloadLogsResult { + string command_id = 1; + string text = 2; + bool truncated = 3; + string error = 4; +} diff --git a/server/internal/grpc/pb/vantage.pb.go b/server/internal/grpc/pb/vantage.pb.go index 0cd42d9..f2b6936 100644 --- a/server/internal/grpc/pb/vantage.pb.go +++ b/server/internal/grpc/pb/vantage.pb.go @@ -198,6 +198,10 @@ type ServerCommand struct { CleanupWorkspace *CleanupWorkspaceCmd `json:"cleanup_workspace,omitempty"` OpenProxy *OpenProxyCmd `json:"open_proxy,omitempty"` Ping *PingCmd `json:"ping,omitempty"` + + RefreshWorkloads *RefreshWorkloadsCmd `json:"refresh_workloads,omitempty"` + ControlWorkload *ControlWorkloadCmd `json:"control_workload,omitempty"` + WorkloadLogs *WorkloadLogsCmd `json:"workload_logs,omitempty"` } // PingCmd is a server-originated liveness beat. It carries nothing and expects @@ -233,6 +237,8 @@ type AgentMessage struct { Result *CommandResult `json:"result,omitempty"` StepResult *StepResult `json:"step_result,omitempty"` StepOutput *StepOutputChunk `json:"step_output,omitempty"` + + WorkloadLogsResult *WorkloadLogsResult `json:"workload_logs_result,omitempty"` } type AgentReady struct{} @@ -366,6 +372,7 @@ type VantageServer interface { UploadGeneratedKey(context.Context, *UploadKeyRequest) (*UploadKeyResponse, error) ReportUpdates(context.Context, *ReportUpdatesRequest) (*ReportUpdatesResponse, error) ReportPackages(context.Context, *ReportPackagesRequest) (*ReportPackagesResponse, error) + ReportWorkloads(context.Context, *ReportWorkloadsRequest) (*ReportWorkloadsResponse, error) ReportInventory(context.Context, *InventoryReport) (*InventoryReportResponse, error) SyncMonitors(context.Context, *SyncMonitorsRequest) (*SyncMonitorsResponse, error) ReportChecks(context.Context, *ReportChecksRequest) (*ReportChecksResponse, error) @@ -421,6 +428,7 @@ 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) ReportPackages(ctx context.Context, in *ReportPackagesRequest, opts ...grpc.CallOption) (*ReportPackagesResponse, error) + ReportWorkloads(ctx context.Context, in *ReportWorkloadsRequest, opts ...grpc.CallOption) (*ReportWorkloadsResponse, 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) @@ -529,6 +537,7 @@ var Vantage_ServiceDesc = grpc.ServiceDesc{ {MethodName: "UploadGeneratedKey", Handler: _Vantage_UploadGeneratedKey_Handler}, {MethodName: "ReportUpdates", Handler: _Vantage_ReportUpdates_Handler}, {MethodName: "ReportPackages", Handler: _Vantage_ReportPackages_Handler}, + {MethodName: "ReportWorkloads", Handler: _Vantage_ReportWorkloads_Handler}, {MethodName: "ReportInventory", Handler: _Vantage_ReportInventory_Handler}, {MethodName: "SyncMonitors", Handler: _Vantage_SyncMonitors_Handler}, {MethodName: "ReportChecks", Handler: _Vantage_ReportChecks_Handler}, diff --git a/server/internal/grpc/pb/workloads.pb.go b/server/internal/grpc/pb/workloads.pb.go new file mode 100644 index 0000000..176030a --- /dev/null +++ b/server/internal/grpc/pb/workloads.pb.go @@ -0,0 +1,98 @@ +package pb + +import ( + "context" + + "google.golang.org/grpc" + "google.golang.org/grpc/codes" + "google.golang.org/grpc/status" +) + +// Workload registry messages. Hand-written like the rest of this package: the +// .proto is the contract, this file is the Go side of it, and the two must be +// changed together. The agent module carries the same declarations. + +// Workload is one container or one systemd unit. +type Workload struct { + Kind string `json:"kind"` + Id string `json:"id"` + Name string `json:"name"` + State string `json:"state"` + Health string `json:"health,omitempty"` + Image string `json:"image,omitempty"` + Stack string `json:"stack,omitempty"` + Ports []string `json:"ports,omitempty"` + Restarts int32 `json:"restarts,omitempty"` + StartedAt string `json:"started_at,omitempty"` // RFC3339, empty when not running + Protected bool `json:"protected,omitempty"` +} + +// ReportWorkloadsRequest carries what a server is running. +// +// Offer-then-send, the same handshake as ReportPackages: the agent calls once +// with Workloads empty, and resends with the body only if NeedFull is set. +type ReportWorkloadsRequest struct { + ServerId string `json:"server_id"` + AgentToken string `json:"agent_token"` + Hash string `json:"hash"` + DockerOk bool `json:"docker_ok"` + DockerError string `json:"docker_error,omitempty"` + SystemdOk bool `json:"systemd_ok"` + SystemdError string `json:"systemd_error,omitempty"` + Workloads []Workload `json:"workloads,omitempty"` // empty on the offer call +} + +type ReportWorkloadsResponse struct { + NeedFull bool `json:"need_full"` +} + +// RefreshWorkloadsCmd carries no payload back. It makes the agent report +// immediately through ReportWorkloads, so there is exactly one writer for the +// server_workloads collection rather than two arriving by different routes. +type RefreshWorkloadsCmd struct{} + +type ControlWorkloadCmd struct { + Kind string `json:"kind"` + Id string `json:"id"` + Action string `json:"action"` // start | stop | restart +} + +type WorkloadLogsCmd struct { + Kind string `json:"kind"` + Id string `json:"id"` + Tail int32 `json:"tail,omitempty"` +} + +type WorkloadLogsResult struct { + CommandId string `json:"command_id"` + Text string `json:"text,omitempty"` + Truncated bool `json:"truncated,omitempty"` + Error string `json:"error,omitempty"` +} + +func (UnimplementedVantageServer) ReportWorkloads(context.Context, *ReportWorkloadsRequest) (*ReportWorkloadsResponse, error) { + return nil, status.Errorf(codes.Unimplemented, "method ReportWorkloads not implemented") +} + +func (c *keyManagerClient) ReportWorkloads(ctx context.Context, in *ReportWorkloadsRequest, opts ...grpc.CallOption) (*ReportWorkloadsResponse, error) { + out := new(ReportWorkloadsResponse) + if err := c.cc.Invoke(ctx, "/vantage.v1.Vantage/ReportWorkloads", in, out, opts...); err != nil { + return nil, err + } + return out, nil +} + +func _Vantage_ReportWorkloads_Handler(srv interface{}, ctx context.Context, dec func(interface{}) error, interceptor grpc.UnaryServerInterceptor) (interface{}, error) { + in := new(ReportWorkloadsRequest) + if err := dec(in); err != nil { + return nil, err + } + if interceptor == nil { + return srv.(VantageServer).ReportWorkloads(ctx, in) + } + info := &grpc.UnaryServerInfo{Server: srv, FullMethod: "/vantage.v1.Vantage/ReportWorkloads"} + handler := func(ctx context.Context, req interface{}) (interface{}, error) { + return srv.(VantageServer).ReportWorkloads(ctx, req.(*ReportWorkloadsRequest)) + } + return interceptor(ctx, in, info, handler) +}