From 66954aead00b94525550b0e312557e755c8ceb7c Mon Sep 17 00:00:00 2001 From: Zach <3724288+zedkipp@users.noreply.github.com> Date: Tue, 3 Mar 2026 09:13:11 -0700 Subject: [PATCH] feat: add TagV2 BoundaryMessage envelope protocol (#22520) Extend the wire protocol for the boundary <-> agent unix socket with a message envelope. The envelope creates a boundary <-> agent data path that is separate from the agent <-> coderd path. This lets boundary send operational metadata (drop counts, configuration like jail type, capabilities) that the agent can act on locally (e.g. Prometheus metrics) or use to enrich outbound requests, without polluting the coderd-facing proto with fields coderd never consumes. Co-Authored-By: Claude Opus 4.6 --- Makefile | 8 + agent/boundarylogproxy/codec/boundary.pb.go | 184 ++++++++++++++++++++ agent/boundarylogproxy/codec/boundary.proto | 17 ++ agent/boundarylogproxy/codec/codec.go | 87 +++++++-- agent/boundarylogproxy/codec/codec_test.go | 4 +- agent/boundarylogproxy/proxy.go | 54 +++--- agent/boundarylogproxy/proxy_test.go | 116 ++++++++++-- 7 files changed, 416 insertions(+), 54 deletions(-) create mode 100644 agent/boundarylogproxy/codec/boundary.pb.go create mode 100644 agent/boundarylogproxy/codec/boundary.proto diff --git a/Makefile b/Makefile index ab4dc273f5..0dbd0e3c9e 100644 --- a/Makefile +++ b/Makefile @@ -654,6 +654,7 @@ GEN_FILES := \ tailnet/proto/tailnet.pb.go \ agent/proto/agent.pb.go \ agent/agentsocket/proto/agentsocket.pb.go \ + agent/boundarylogproxy/codec/boundary.pb.go \ provisionersdk/proto/provisioner.pb.go \ provisionerd/proto/provisionerd.pb.go \ vpn/vpn.pb.go \ @@ -709,6 +710,7 @@ gen/mark-fresh: provisionersdk/proto/provisioner.pb.go \ provisionerd/proto/provisionerd.pb.go \ agent/agentsocket/proto/agentsocket.pb.go \ + agent/boundarylogproxy/codec/boundary.pb.go \ vpn/vpn.pb.go \ enterprise/aibridged/proto/aibridged.pb.go \ coderd/database/dump.sql \ @@ -843,6 +845,12 @@ vpn/vpn.pb.go: vpn/vpn.proto --go_opt=paths=source_relative \ ./vpn/vpn.proto +agent/boundarylogproxy/codec/boundary.pb.go: agent/boundarylogproxy/codec/boundary.proto agent/proto/agent.proto + protoc \ + --go_out=. \ + --go_opt=paths=source_relative \ + ./agent/boundarylogproxy/codec/boundary.proto + enterprise/aibridged/proto/aibridged.pb.go: enterprise/aibridged/proto/aibridged.proto protoc \ --go_out=. \ diff --git a/agent/boundarylogproxy/codec/boundary.pb.go b/agent/boundarylogproxy/codec/boundary.pb.go new file mode 100644 index 0000000000..86b18361b7 --- /dev/null +++ b/agent/boundarylogproxy/codec/boundary.pb.go @@ -0,0 +1,184 @@ +// Code generated by protoc-gen-go. DO NOT EDIT. +// versions: +// protoc-gen-go v1.30.0 +// protoc v4.23.4 +// source: agent/boundarylogproxy/codec/boundary.proto + +package codec + +import ( + proto "github.com/coder/coder/v2/agent/proto" + protoreflect "google.golang.org/protobuf/reflect/protoreflect" + protoimpl "google.golang.org/protobuf/runtime/protoimpl" + reflect "reflect" + sync "sync" +) + +const ( + // Verify that this generated code is sufficiently up-to-date. + _ = protoimpl.EnforceVersion(20 - protoimpl.MinVersion) + // Verify that runtime/protoimpl is sufficiently up-to-date. + _ = protoimpl.EnforceVersion(protoimpl.MaxVersion - 20) +) + +// BoundaryMessage is the envelope for all TagV2 messages sent over the +// boundary <-> agent unix socket. TagV1 carries a bare +// ReportBoundaryLogsRequest for backwards compatibility; TagV2 wraps +// everything in this envelope so the protocol can be extended with new +// message types without adding more tags. +type BoundaryMessage struct { + state protoimpl.MessageState + sizeCache protoimpl.SizeCache + unknownFields protoimpl.UnknownFields + + // Types that are assignable to Msg: + // + // *BoundaryMessage_Logs + Msg isBoundaryMessage_Msg `protobuf_oneof:"msg"` +} + +func (x *BoundaryMessage) Reset() { + *x = BoundaryMessage{} + if protoimpl.UnsafeEnabled { + mi := &file_agent_boundarylogproxy_codec_boundary_proto_msgTypes[0] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) + } +} + +func (x *BoundaryMessage) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*BoundaryMessage) ProtoMessage() {} + +func (x *BoundaryMessage) ProtoReflect() protoreflect.Message { + mi := &file_agent_boundarylogproxy_codec_boundary_proto_msgTypes[0] + if protoimpl.UnsafeEnabled && x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use BoundaryMessage.ProtoReflect.Descriptor instead. +func (*BoundaryMessage) Descriptor() ([]byte, []int) { + return file_agent_boundarylogproxy_codec_boundary_proto_rawDescGZIP(), []int{0} +} + +func (m *BoundaryMessage) GetMsg() isBoundaryMessage_Msg { + if m != nil { + return m.Msg + } + return nil +} + +func (x *BoundaryMessage) GetLogs() *proto.ReportBoundaryLogsRequest { + if x, ok := x.GetMsg().(*BoundaryMessage_Logs); ok { + return x.Logs + } + return nil +} + +type isBoundaryMessage_Msg interface { + isBoundaryMessage_Msg() +} + +type BoundaryMessage_Logs struct { + Logs *proto.ReportBoundaryLogsRequest `protobuf:"bytes,1,opt,name=logs,proto3,oneof"` +} + +func (*BoundaryMessage_Logs) isBoundaryMessage_Msg() {} + +var File_agent_boundarylogproxy_codec_boundary_proto protoreflect.FileDescriptor + +var file_agent_boundarylogproxy_codec_boundary_proto_rawDesc = []byte{ + 0x0a, 0x2b, 0x61, 0x67, 0x65, 0x6e, 0x74, 0x2f, 0x62, 0x6f, 0x75, 0x6e, 0x64, 0x61, 0x72, 0x79, + 0x6c, 0x6f, 0x67, 0x70, 0x72, 0x6f, 0x78, 0x79, 0x2f, 0x63, 0x6f, 0x64, 0x65, 0x63, 0x2f, 0x62, + 0x6f, 0x75, 0x6e, 0x64, 0x61, 0x72, 0x79, 0x2e, 0x70, 0x72, 0x6f, 0x74, 0x6f, 0x12, 0x1f, 0x63, + 0x6f, 0x64, 0x65, 0x72, 0x2e, 0x62, 0x6f, 0x75, 0x6e, 0x64, 0x61, 0x72, 0x79, 0x6c, 0x6f, 0x67, + 0x70, 0x72, 0x6f, 0x78, 0x79, 0x2e, 0x63, 0x6f, 0x64, 0x65, 0x63, 0x2e, 0x76, 0x31, 0x1a, 0x17, + 0x61, 0x67, 0x65, 0x6e, 0x74, 0x2f, 0x70, 0x72, 0x6f, 0x74, 0x6f, 0x2f, 0x61, 0x67, 0x65, 0x6e, + 0x74, 0x2e, 0x70, 0x72, 0x6f, 0x74, 0x6f, 0x22, 0x59, 0x0a, 0x0f, 0x42, 0x6f, 0x75, 0x6e, 0x64, + 0x61, 0x72, 0x79, 0x4d, 0x65, 0x73, 0x73, 0x61, 0x67, 0x65, 0x12, 0x3f, 0x0a, 0x04, 0x6c, 0x6f, + 0x67, 0x73, 0x18, 0x01, 0x20, 0x01, 0x28, 0x0b, 0x32, 0x29, 0x2e, 0x63, 0x6f, 0x64, 0x65, 0x72, + 0x2e, 0x61, 0x67, 0x65, 0x6e, 0x74, 0x2e, 0x76, 0x32, 0x2e, 0x52, 0x65, 0x70, 0x6f, 0x72, 0x74, + 0x42, 0x6f, 0x75, 0x6e, 0x64, 0x61, 0x72, 0x79, 0x4c, 0x6f, 0x67, 0x73, 0x52, 0x65, 0x71, 0x75, + 0x65, 0x73, 0x74, 0x48, 0x00, 0x52, 0x04, 0x6c, 0x6f, 0x67, 0x73, 0x42, 0x05, 0x0a, 0x03, 0x6d, + 0x73, 0x67, 0x42, 0x38, 0x5a, 0x36, 0x67, 0x69, 0x74, 0x68, 0x75, 0x62, 0x2e, 0x63, 0x6f, 0x6d, + 0x2f, 0x63, 0x6f, 0x64, 0x65, 0x72, 0x2f, 0x63, 0x6f, 0x64, 0x65, 0x72, 0x2f, 0x76, 0x32, 0x2f, + 0x61, 0x67, 0x65, 0x6e, 0x74, 0x2f, 0x62, 0x6f, 0x75, 0x6e, 0x64, 0x61, 0x72, 0x79, 0x6c, 0x6f, + 0x67, 0x70, 0x72, 0x6f, 0x78, 0x79, 0x2f, 0x63, 0x6f, 0x64, 0x65, 0x63, 0x62, 0x06, 0x70, 0x72, + 0x6f, 0x74, 0x6f, 0x33, +} + +var ( + file_agent_boundarylogproxy_codec_boundary_proto_rawDescOnce sync.Once + file_agent_boundarylogproxy_codec_boundary_proto_rawDescData = file_agent_boundarylogproxy_codec_boundary_proto_rawDesc +) + +func file_agent_boundarylogproxy_codec_boundary_proto_rawDescGZIP() []byte { + file_agent_boundarylogproxy_codec_boundary_proto_rawDescOnce.Do(func() { + file_agent_boundarylogproxy_codec_boundary_proto_rawDescData = protoimpl.X.CompressGZIP(file_agent_boundarylogproxy_codec_boundary_proto_rawDescData) + }) + return file_agent_boundarylogproxy_codec_boundary_proto_rawDescData +} + +var file_agent_boundarylogproxy_codec_boundary_proto_msgTypes = make([]protoimpl.MessageInfo, 1) +var file_agent_boundarylogproxy_codec_boundary_proto_goTypes = []interface{}{ + (*BoundaryMessage)(nil), // 0: coder.boundarylogproxy.codec.v1.BoundaryMessage + (*proto.ReportBoundaryLogsRequest)(nil), // 1: coder.agent.v2.ReportBoundaryLogsRequest +} +var file_agent_boundarylogproxy_codec_boundary_proto_depIdxs = []int32{ + 1, // 0: coder.boundarylogproxy.codec.v1.BoundaryMessage.logs:type_name -> coder.agent.v2.ReportBoundaryLogsRequest + 1, // [1:1] is the sub-list for method output_type + 1, // [1:1] is the sub-list for method input_type + 1, // [1:1] is the sub-list for extension type_name + 1, // [1:1] is the sub-list for extension extendee + 0, // [0:1] is the sub-list for field type_name +} + +func init() { file_agent_boundarylogproxy_codec_boundary_proto_init() } +func file_agent_boundarylogproxy_codec_boundary_proto_init() { + if File_agent_boundarylogproxy_codec_boundary_proto != nil { + return + } + if !protoimpl.UnsafeEnabled { + file_agent_boundarylogproxy_codec_boundary_proto_msgTypes[0].Exporter = func(v interface{}, i int) interface{} { + switch v := v.(*BoundaryMessage); i { + case 0: + return &v.state + case 1: + return &v.sizeCache + case 2: + return &v.unknownFields + default: + return nil + } + } + } + file_agent_boundarylogproxy_codec_boundary_proto_msgTypes[0].OneofWrappers = []interface{}{ + (*BoundaryMessage_Logs)(nil), + } + type x struct{} + out := protoimpl.TypeBuilder{ + File: protoimpl.DescBuilder{ + GoPackagePath: reflect.TypeOf(x{}).PkgPath(), + RawDescriptor: file_agent_boundarylogproxy_codec_boundary_proto_rawDesc, + NumEnums: 0, + NumMessages: 1, + NumExtensions: 0, + NumServices: 0, + }, + GoTypes: file_agent_boundarylogproxy_codec_boundary_proto_goTypes, + DependencyIndexes: file_agent_boundarylogproxy_codec_boundary_proto_depIdxs, + MessageInfos: file_agent_boundarylogproxy_codec_boundary_proto_msgTypes, + }.Build() + File_agent_boundarylogproxy_codec_boundary_proto = out.File + file_agent_boundarylogproxy_codec_boundary_proto_rawDesc = nil + file_agent_boundarylogproxy_codec_boundary_proto_goTypes = nil + file_agent_boundarylogproxy_codec_boundary_proto_depIdxs = nil +} diff --git a/agent/boundarylogproxy/codec/boundary.proto b/agent/boundarylogproxy/codec/boundary.proto new file mode 100644 index 0000000000..ed13160c74 --- /dev/null +++ b/agent/boundarylogproxy/codec/boundary.proto @@ -0,0 +1,17 @@ +syntax = "proto3"; +option go_package = "github.com/coder/coder/v2/agent/boundarylogproxy/codec"; + +package coder.boundarylogproxy.codec.v1; + +import "agent/proto/agent.proto"; + +// BoundaryMessage is the envelope for all TagV2 messages sent over the +// boundary <-> agent unix socket. TagV1 carries a bare +// ReportBoundaryLogsRequest for backwards compatibility; TagV2 wraps +// everything in this envelope so the protocol can be extended with new +// message types without adding more tags. +message BoundaryMessage { + oneof msg { + coder.agent.v2.ReportBoundaryLogsRequest logs = 1; + } +} diff --git a/agent/boundarylogproxy/codec/codec.go b/agent/boundarylogproxy/codec/codec.go index cda876c64d..dd4c023bae 100644 --- a/agent/boundarylogproxy/codec/codec.go +++ b/agent/boundarylogproxy/codec/codec.go @@ -14,14 +14,23 @@ import ( "io" "golang.org/x/xerrors" + "google.golang.org/protobuf/proto" + + agentproto "github.com/coder/coder/v2/agent/proto" ) type Tag uint8 const ( - // TagV1 identifies the first revision of the protocol. This version has a maximum - // data length of MaxMessageSizeV1. + // TagV1 identifies the first revision of the protocol. The payload is a + // bare ReportBoundaryLogsRequest. This version has a maximum data length + // of MaxMessageSizeV1. TagV1 Tag = 1 + + // TagV2 identifies the second revision of the protocol. The payload is + // a BoundaryMessage envelope. This version has a maximum data length of + // MaxMessageSizeV2. + TagV2 Tag = 2 ) const ( @@ -35,6 +44,9 @@ const ( // over the wire for the TagV1 tag. While the wire format allows 24 bits for // length, TagV1 only uses 15 bits. MaxMessageSizeV1 uint32 = 1 << 15 + + // MaxMessageSizeV2 is the maximum data length for TagV2. + MaxMessageSizeV2 = MaxMessageSizeV1 ) var ( @@ -48,12 +60,9 @@ var ( // WriteFrame writes a framed message with the given tag and data. The data // must not exceed 2^DataLength in length. func WriteFrame(w io.Writer, tag Tag, data []byte) error { - var maxSize uint32 - switch tag { - case TagV1: - maxSize = MaxMessageSizeV1 - default: - return xerrors.Errorf("%w: %d", ErrUnsupportedTag, tag) + maxSize, err := maxSizeForTag(tag) + if err != nil { + return err } if len(data) > int(maxSize) { @@ -101,12 +110,9 @@ func ReadFrame(r io.Reader, buf []byte) (Tag, []byte, error) { } tag := Tag(shifted) - var maxSize uint32 - switch tag { - case TagV1: - maxSize = MaxMessageSizeV1 - default: - return 0, nil, xerrors.Errorf("%w: %d", ErrUnsupportedTag, tag) + maxSize, err := maxSizeForTag(tag) + if err != nil { + return 0, nil, err } if length > maxSize { @@ -125,3 +131,56 @@ func ReadFrame(r io.Reader, buf []byte) (Tag, []byte, error) { return tag, buf[:length], nil } + +// maxSizeForTag returns the maximum payload size for the given tag. +func maxSizeForTag(tag Tag) (uint32, error) { + switch tag { + case TagV1: + return MaxMessageSizeV1, nil + case TagV2: + return MaxMessageSizeV2, nil + default: + return 0, xerrors.Errorf("%w: %d", ErrUnsupportedTag, tag) + } +} + +// ReadMessage reads a framed message and unmarshals it based on tag. The +// returned buf should be passed back on the next call for buffer reuse. +func ReadMessage(r io.Reader, buf []byte) (proto.Message, []byte, error) { + tag, data, err := ReadFrame(r, buf) + if err != nil { + return nil, data, err + } + + var msg proto.Message + switch tag { + case TagV1: + var req agentproto.ReportBoundaryLogsRequest + if err := proto.Unmarshal(data, &req); err != nil { + return nil, data, xerrors.Errorf("unmarshal TagV1: %w", err) + } + msg = &req + case TagV2: + var envelope BoundaryMessage + if err := proto.Unmarshal(data, &envelope); err != nil { + return nil, data, xerrors.Errorf("unmarshal TagV2: %w", err) + } + msg = &envelope + default: + // maxSizeForTag already rejects unknown tags during ReadFrame, + // but handle it here for safety. + return nil, data, xerrors.Errorf("%w: %d", ErrUnsupportedTag, tag) + } + + return msg, data, nil +} + +// WriteMessage marshals a proto message and writes it as a framed message +// with the given tag. +func WriteMessage(w io.Writer, tag Tag, msg proto.Message) error { + data, err := proto.Marshal(msg) + if err != nil { + return xerrors.Errorf("marshal: %w", err) + } + return WriteFrame(w, tag, data) +} diff --git a/agent/boundarylogproxy/codec/codec_test.go b/agent/boundarylogproxy/codec/codec_test.go index 4ca719f2d0..1bda4a8f7c 100644 --- a/agent/boundarylogproxy/codec/codec_test.go +++ b/agent/boundarylogproxy/codec/codec_test.go @@ -89,7 +89,7 @@ func TestReadFrameInvalidTag(t *testing.T) { // reading the invalid tag. const ( dataLength uint32 = 10 - bogusTag uint32 = 2 + bogusTag uint32 = 222 ) header := bogusTag<