From bc79daab48a3f91d65849eafc79e2df02c421d47 Mon Sep 17 00:00:00 2001 From: mrhid6 Date: Wed, 29 Jul 2026 12:37:37 +0100 Subject: [PATCH] feat: add ProxyStream wire types for agent-relayed console --- agent/internal/grpc/client.go | 4 + agent/internal/grpc/pb/vantage.pb.go | 82 ++++++++++++++++++ proto/vantage/v1/vantage.proto | 32 ++++++++ server/internal/grpc/pb/vantage.pb.go | 96 ++++++++++++++++++++++ server/internal/grpc/pb/vantage_pb_test.go | 39 +++++++++ server/internal/grpc/server.go | 5 ++ 6 files changed, 258 insertions(+) create mode 100644 server/internal/grpc/pb/vantage_pb_test.go diff --git a/agent/internal/grpc/client.go b/agent/internal/grpc/client.go index 10171e4..fa6daff 100644 --- a/agent/internal/grpc/client.go +++ b/agent/internal/grpc/client.go @@ -151,3 +151,7 @@ func (c *Client) ReportChecks(serverID, agentToken string, results []pb.CheckRes func (c *Client) CommandStream(ctx context.Context) (pb.Vantage_CommandStreamClient, error) { return c.client.CommandStream(ctx) } + +func (c *Client) ProxyStream(ctx context.Context) (pb.Vantage_ProxyStreamClient, error) { + return c.client.ProxyStream(ctx) +} diff --git a/agent/internal/grpc/pb/vantage.pb.go b/agent/internal/grpc/pb/vantage.pb.go index 11f9cdb..d96cd29 100644 --- a/agent/internal/grpc/pb/vantage.pb.go +++ b/agent/internal/grpc/pb/vantage.pb.go @@ -131,6 +131,32 @@ type ReportChecksResponse struct{} type ApplyUpdatesCmd struct{} +type OpenProxyCmd struct { + ProxyId string `json:"proxy_id"` + Port uint32 `json:"port"` +} + +type ProxyOpen struct { + ServerId string `json:"server_id"` + AgentToken string `json:"agent_token"` + ProxyId string `json:"proxy_id"` +} + +type ProxyClose struct { + Reason string `json:"reason,omitempty"` +} + +type ProxyClientMsg struct { + Open *ProxyOpen `json:"open,omitempty"` + Data []byte `json:"data,omitempty"` + Close *ProxyClose `json:"close,omitempty"` +} + +type ProxyServerMsg struct { + Data []byte `json:"data,omitempty"` + Close *ProxyClose `json:"close,omitempty"` +} + type ServerCommand struct { CommandId string `json:"command_id"` GenerateKey *GenerateKeyCmd `json:"generate_key,omitempty"` @@ -139,6 +165,7 @@ type ServerCommand struct { ApplyUpdates *ApplyUpdatesCmd `json:"apply_updates,omitempty"` RunStep *RunStepCmd `json:"run_step,omitempty"` CleanupWorkspace *CleanupWorkspaceCmd `json:"cleanup_workspace,omitempty"` + OpenProxy *OpenProxyCmd `json:"open_proxy,omitempty"` } @@ -254,6 +281,51 @@ func (s *keyManagerCommandStreamServer) Recv() (*AgentMessage, error) { return m, nil } +type Vantage_ProxyStreamServer interface { + Send(*ProxyServerMsg) error + Recv() (*ProxyClientMsg, error) + grpc.ServerStream +} + +type vantageProxyStreamServer struct { + grpc.ServerStream +} + +func (s *vantageProxyStreamServer) Send(m *ProxyServerMsg) error { + return s.ServerStream.SendMsg(m) +} + +func (s *vantageProxyStreamServer) Recv() (*ProxyClientMsg, error) { + m := new(ProxyClientMsg) + if err := s.ServerStream.RecvMsg(m); err != nil { + return nil, err + } + return m, nil +} + +type Vantage_ProxyStreamClient interface { + Send(*ProxyClientMsg) error + Recv() (*ProxyServerMsg, error) + CloseSend() error + grpc.ClientStream +} + +type vantageProxyStreamClient struct { + grpc.ClientStream +} + +func (c *vantageProxyStreamClient) Send(m *ProxyClientMsg) error { + return c.ClientStream.SendMsg(m) +} + +func (c *vantageProxyStreamClient) Recv() (*ProxyServerMsg, error) { + m := new(ProxyServerMsg) + if err := c.ClientStream.RecvMsg(m); err != nil { + return nil, err + } + return m, nil +} + type VantageClient interface { Register(ctx context.Context, in *RegisterRequest, opts ...grpc.CallOption) (*RegisterResponse, error) SyncKeys(ctx context.Context, in *SyncRequest, opts ...grpc.CallOption) (*SyncResponse, error) @@ -263,6 +335,7 @@ type VantageClient interface { 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) + ProxyStream(ctx context.Context, opts ...grpc.CallOption) (Vantage_ProxyStreamClient, error) } type UnimplementedVantageServer struct{} @@ -349,3 +422,12 @@ func (c *keyManagerClient) CommandStream(ctx context.Context, opts ...grpc.CallO } return &vantageCommandStreamClient{stream}, nil } + +func (c *keyManagerClient) ProxyStream(ctx context.Context, opts ...grpc.CallOption) (Vantage_ProxyStreamClient, error) { + desc := &grpc.StreamDesc{StreamName: "ProxyStream", ServerStreams: true, ClientStreams: true} + stream, err := c.cc.NewStream(ctx, desc, "/vantage.v1.Vantage/ProxyStream", opts...) + if err != nil { + return nil, err + } + return &vantageProxyStreamClient{stream}, nil +} diff --git a/proto/vantage/v1/vantage.proto b/proto/vantage/v1/vantage.proto index d0d2de3..5754a39 100644 --- a/proto/vantage/v1/vantage.proto +++ b/proto/vantage/v1/vantage.proto @@ -14,6 +14,7 @@ service Vantage { rpc ReportChecks(ReportChecksRequest) returns (ReportChecksResponse); // Bidirectional stream: agent sends auth once, server pushes commands. rpc CommandStream(stream AgentMessage) returns (stream ServerCommand); + rpc ProxyStream(stream ProxyClientMsg) returns (stream ProxyServerMsg); } message RegisterRequest { @@ -180,6 +181,7 @@ message ServerCommand { ApplyUpdatesCmd apply_updates = 5; RunStepCmd run_step = 6; CleanupWorkspaceCmd cleanup_workspace = 7; + OpenProxyCmd open_proxy = 8; } } @@ -228,3 +230,33 @@ message StepOutputChunk { bytes data = 3; bool eof = 4; } + +// OpenProxyCmd tells the agent to dial 127.0.0.1:port locally and relay that +// connection back over a fresh ProxyStream identified by proxy_id. +message OpenProxyCmd { + string proxy_id = 1; + uint32 port = 2; +} + +message ProxyOpen { + string server_id = 1; + string agent_token = 2; + string proxy_id = 3; +} + +message ProxyClose { string reason = 1; } + +message ProxyClientMsg { + oneof payload { + ProxyOpen open = 1; // first message only + bytes data = 2; + ProxyClose close = 3; + } +} + +message ProxyServerMsg { + oneof payload { + bytes data = 1; + ProxyClose close = 2; + } +} diff --git a/server/internal/grpc/pb/vantage.pb.go b/server/internal/grpc/pb/vantage.pb.go index 24a0053..648f9c2 100644 --- a/server/internal/grpc/pb/vantage.pb.go +++ b/server/internal/grpc/pb/vantage.pb.go @@ -123,6 +123,32 @@ type ReportChecksResponse struct{} type ApplyUpdatesCmd struct{} +type OpenProxyCmd struct { + ProxyId string `json:"proxy_id"` + Port uint32 `json:"port"` +} + +type ProxyOpen struct { + ServerId string `json:"server_id"` + AgentToken string `json:"agent_token"` + ProxyId string `json:"proxy_id"` +} + +type ProxyClose struct { + Reason string `json:"reason,omitempty"` +} + +type ProxyClientMsg struct { + Open *ProxyOpen `json:"open,omitempty"` + Data []byte `json:"data,omitempty"` + Close *ProxyClose `json:"close,omitempty"` +} + +type ProxyServerMsg struct { + Data []byte `json:"data,omitempty"` + Close *ProxyClose `json:"close,omitempty"` +} + type ServerCommand struct { CommandId string `json:"command_id"` GenerateKey *GenerateKeyCmd `json:"generate_key,omitempty"` @@ -131,6 +157,7 @@ type ServerCommand struct { ApplyUpdates *ApplyUpdatesCmd `json:"apply_updates,omitempty"` RunStep *RunStepCmd `json:"run_step,omitempty"` CleanupWorkspace *CleanupWorkspaceCmd `json:"cleanup_workspace,omitempty"` + OpenProxy *OpenProxyCmd `json:"open_proxy,omitempty"` } type CleanupWorkspaceCmd struct { @@ -239,6 +266,55 @@ func (c *vantageCommandStreamClient) Recv() (*ServerCommand, error) { return m, nil } +type Vantage_ProxyStreamServer interface { + Send(*ProxyServerMsg) error + Recv() (*ProxyClientMsg, error) + grpc.ServerStream +} + +type vantageProxyStreamServer struct { + grpc.ServerStream +} + +func (s *vantageProxyStreamServer) Send(m *ProxyServerMsg) error { + return s.ServerStream.SendMsg(m) +} + +func (s *vantageProxyStreamServer) Recv() (*ProxyClientMsg, error) { + m := new(ProxyClientMsg) + if err := s.ServerStream.RecvMsg(m); err != nil { + return nil, err + } + return m, nil +} + +type Vantage_ProxyStreamClient interface { + Send(*ProxyClientMsg) error + Recv() (*ProxyServerMsg, error) + CloseSend() error + grpc.ClientStream +} + +type vantageProxyStreamClient struct { + grpc.ClientStream +} + +func (c *vantageProxyStreamClient) Send(m *ProxyClientMsg) error { + return c.ClientStream.SendMsg(m) +} + +func (c *vantageProxyStreamClient) Recv() (*ProxyServerMsg, error) { + m := new(ProxyServerMsg) + if err := c.ClientStream.RecvMsg(m); err != nil { + return nil, err + } + return m, nil +} + +func _Vantage_ProxyStream_Handler(srv interface{}, stream grpc.ServerStream) error { + return srv.(VantageServer).ProxyStream(&vantageProxyStreamServer{stream}) +} + type VantageServer interface { Register(context.Context, *RegisterRequest) (*RegisterResponse, error) SyncKeys(context.Context, *SyncRequest) (*SyncResponse, error) @@ -248,6 +324,7 @@ type VantageServer interface { SyncMonitors(context.Context, *SyncMonitorsRequest) (*SyncMonitorsResponse, error) ReportChecks(context.Context, *ReportChecksRequest) (*ReportChecksResponse, error) CommandStream(Vantage_CommandStreamServer) error + ProxyStream(Vantage_ProxyStreamServer) error } type UnimplementedVantageServer struct{} @@ -284,6 +361,10 @@ func (UnimplementedVantageServer) CommandStream(Vantage_CommandStreamServer) err return status.Errorf(codes.Unimplemented, "method CommandStream not implemented") } +func (UnimplementedVantageServer) ProxyStream(Vantage_ProxyStreamServer) error { + return status.Errorf(codes.Unimplemented, "method ProxyStream not implemented") +} + type VantageClient interface { Register(ctx context.Context, in *RegisterRequest, opts ...grpc.CallOption) (*RegisterResponse, error) SyncKeys(ctx context.Context, in *SyncRequest, opts ...grpc.CallOption) (*SyncResponse, error) @@ -293,6 +374,7 @@ type VantageClient interface { 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) + ProxyStream(ctx context.Context, opts ...grpc.CallOption) (Vantage_ProxyStreamClient, error) } type keyManagerClient struct { @@ -367,6 +449,14 @@ func (c *keyManagerClient) CommandStream(ctx context.Context, opts ...grpc.CallO return &vantageCommandStreamClient{stream}, nil } +func (c *keyManagerClient) ProxyStream(ctx context.Context, opts ...grpc.CallOption) (Vantage_ProxyStreamClient, error) { + stream, err := c.cc.NewStream(ctx, &Vantage_ServiceDesc.Streams[1], "/vantage.v1.Vantage/ProxyStream", opts...) + if err != nil { + return nil, err + } + return &vantageProxyStreamClient{stream}, nil +} + func RegisterVantageServer(s grpc.ServiceRegistrar, srv VantageServer) { s.RegisterService(&Vantage_ServiceDesc, srv) } @@ -390,6 +480,12 @@ var Vantage_ServiceDesc = grpc.ServiceDesc{ ServerStreams: true, ClientStreams: true, }, + { + StreamName: "ProxyStream", + Handler: _Vantage_ProxyStream_Handler, + ServerStreams: true, + ClientStreams: true, + }, }, Metadata: "vantage/v1/vantage.proto", } diff --git a/server/internal/grpc/pb/vantage_pb_test.go b/server/internal/grpc/pb/vantage_pb_test.go new file mode 100644 index 0000000..56cc75b --- /dev/null +++ b/server/internal/grpc/pb/vantage_pb_test.go @@ -0,0 +1,39 @@ +package pb + +import ( + "encoding/json" + "testing" +) + +func TestProxyClientMsgRoundTrip(t *testing.T) { + in := &ProxyClientMsg{Data: []byte{0x00, 0xff, 0x10}} + raw, err := json.Marshal(in) + if err != nil { + t.Fatalf("marshal: %v", err) + } + var out ProxyClientMsg + if err := json.Unmarshal(raw, &out); err != nil { + t.Fatalf("unmarshal: %v", err) + } + if string(out.Data) != string(in.Data) { + t.Fatalf("data mismatch: got %v want %v", out.Data, in.Data) + } + if out.Open != nil || out.Close != nil { + t.Fatalf("empty oneof fields should stay nil, got open=%v close=%v", out.Open, out.Close) + } +} + +func TestOpenProxyCmdOnServerCommand(t *testing.T) { + cmd := &ServerCommand{CommandId: "c1", OpenProxy: &OpenProxyCmd{ProxyId: "p1", Port: 22}} + raw, err := json.Marshal(cmd) + if err != nil { + t.Fatalf("marshal: %v", err) + } + var out ServerCommand + if err := json.Unmarshal(raw, &out); err != nil { + t.Fatalf("unmarshal: %v", err) + } + if out.OpenProxy == nil || out.OpenProxy.ProxyId != "p1" || out.OpenProxy.Port != 22 { + t.Fatalf("open_proxy did not round-trip: %+v", out.OpenProxy) + } +} diff --git a/server/internal/grpc/server.go b/server/internal/grpc/server.go index fd11bfa..313b2bc 100644 --- a/server/internal/grpc/server.go +++ b/server/internal/grpc/server.go @@ -223,6 +223,11 @@ func (s *vantageServer) CommandStream(stream pb.Vantage_CommandStreamServer) err } } +// ProxyStream is stubbed here; Task 4 implements the real relay handler. +func (s *vantageServer) ProxyStream(stream pb.Vantage_ProxyStreamServer) error { + return status.Errorf(codes.Unimplemented, "method ProxyStream not implemented") +} + func StartGRPC(port int) error { lis, err := net.Listen("tcp", fmt.Sprintf(":%d", port)) if err != nil {