From 00dbbec62b127fe254804a365f58650821c94acc Mon Sep 17 00:00:00 2001 From: Pavel Sviderski Date: Wed, 4 Dec 2024 14:47:31 +1000 Subject: [PATCH] refactor grpc-proxy backends to enhance responses from local backends with metadata in one2many mode --- internal/machine/api/proxy/backend.go | 144 +++++++++++++++++++++++++ internal/machine/api/proxy/local.go | 18 ++-- internal/machine/api/proxy/remote.go | 145 ++------------------------ 3 files changed, 159 insertions(+), 148 deletions(-) create mode 100644 internal/machine/api/proxy/backend.go diff --git a/internal/machine/api/proxy/backend.go b/internal/machine/api/proxy/backend.go new file mode 100644 index 00000000..1fc017c1 --- /dev/null +++ b/internal/machine/api/proxy/backend.go @@ -0,0 +1,144 @@ +package proxy + +import ( + "fmt" + "google.golang.org/grpc/status" + "google.golang.org/protobuf/encoding/protowire" + "google.golang.org/protobuf/proto" + "uncloud/internal/machine/api/pb" +) + +// One2ManyResponder converts upstream responses into messages from upstreams, so that multiple +// successful and failure responses might be returned in One2Many mode. +type One2ManyResponder struct { + machine string +} + +// AppendInfo is called to enhance response from the backend with additional data. +// +// AppendInfo enhances upstream response with machine metadata (target). +// +// This method depends on grpc protobuf response structure, each response should +// look like: +// +// message SomeResponse { +// repeated SomeReply messages = 1; // please note field ID == 1 +// } +// +// message SomeReply { +// common.Metadata metadata = 1; +// +// } +// +// As 'SomeReply' is repeated in 'SomeResponse', if we concatenate protobuf representation +// of several 'SomeResponse' messages, we still get valid 'SomeResponse' representation but with more +// entries (feature of protobuf binary representation). +// +// If we look at binary representation of any unary 'SomeResponse' message, it will always contain one +// protobuf field with field ID 1 (see above) and type 2 (embedded message SomeReply is encoded +// as string with length). So if we want to add fields to 'SomeReply', we can simply read field +// header, adjust length for new 'SomeReply' representation, and prepend new field header. +// +// At the same time, we can add 'common.Metadata' structure to 'SomeReply' by simply +// appending or prepending 'common.Metadata' as a single field. This requires 'metadata' +// field to be not defined in original response. (This is due to the fact that protobuf message +// representation is concatenation of each field representation). +// +// To build only single field (Metadata) we use helper message which contains exactly this +// field with same field ID as in every other 'SomeReply': +// +// message Empty { +// common.Metadata metadata = 1; +// } +// +// As streaming replies are not wrapped into 'SomeResponse' with 'repeated', handling is simpler: we just +// need to append Empty with details. +// +// So AppendInfo does the following: validates that response contains field ID 1 encoded as string, +// cuts field header, rest is representation of some reply. Marshal 'Empty' as protobuf, +// which builds 'common.Metadata' field, append it to original response message, build new header +// for new length of some response, and add back new field header. +func (b *One2ManyResponder) AppendInfo(streaming bool, resp []byte) ([]byte, error) { + payload, err := proto.Marshal(&pb.Empty{ + Metadata: &pb.Metadata{ + Machine: b.machine, + }, + }) + + if streaming { + return append(resp, payload...), err + } + + const ( + metadataField = 1 // field number in proto definition for repeated response + metadataType = 2 // "string" for embedded messages + ) + + // decode protobuf embedded header + + typ, n1 := protowire.ConsumeVarint(resp) + if n1 < 0 { + return nil, protowire.ParseError(n1) + } + + _, n2 := protowire.ConsumeVarint(resp[n1:]) // length + if n2 < 0 { + return nil, protowire.ParseError(n2) + } + + if typ != (metadataField<<3)|metadataType { + return nil, fmt.Errorf("unexpected message format: %d", typ) + } + + if n1+n2 > len(resp) { + return nil, fmt.Errorf("unexpected message size: %d", len(resp)) + } + + // cut off embedded message header + resp = resp[n1+n2:] + // build new embedded message header + prefix := protowire.AppendVarint( + protowire.AppendVarint(nil, (metadataField<<3)|metadataType), + uint64(len(resp)+len(payload)), + ) + resp = append(prefix, resp...) + + return append(resp, payload...), err +} + +// BuildError converts upstream error into message from upstream, so that multiple +// successful and failure responses might be returned. +// +// This simply relies on the fact that any response contains 'Empty' message. +// So if 'Empty' is unmarshalled into any other reply message, all the fields +// are undefined but 'Metadata': +// +// message Empty { +// common.Metadata metadata = 1; +// } +// +// message EmptyResponse { +// repeated Empty messages = 1; +// } +// +// Streaming responses are not wrapped into Empty, so we simply marshall EmptyResponse +// message. +func (b *One2ManyResponder) BuildError(streaming bool, err error) ([]byte, error) { + var resp proto.Message = &pb.Empty{ + Metadata: &pb.Metadata{ + Machine: b.machine, + Error: err.Error(), + Status: status.Convert(err).Proto(), + }, + } + + if !streaming { + resp = &pb.EmptyResponse{ + Messages: []*pb.Empty{ + resp.(*pb.Empty), + }, + } + } + + return proto.Marshal(resp) +} diff --git a/internal/machine/api/proxy/local.go b/internal/machine/api/proxy/local.go index 805de6e5..d5638d6c 100644 --- a/internal/machine/api/proxy/local.go +++ b/internal/machine/api/proxy/local.go @@ -9,8 +9,9 @@ import ( "sync" ) -// LocalBackend is a proxy.Backend implementation that proxies to a local gRPC server listening on a Unix socket. +// LocalBackend is a proxy.One2ManyResponder implementation that proxies to a local gRPC server listening on a Unix socket. type LocalBackend struct { + One2ManyResponder sockPath string mu sync.RWMutex @@ -22,12 +23,15 @@ var _ proxy.Backend = (*LocalBackend)(nil) // NewLocalBackend returns a new LocalBackend for the given Unix socket path. func NewLocalBackend(sockPath string) *LocalBackend { return &LocalBackend{ + One2ManyResponder: One2ManyResponder{ + machine: "local", + }, sockPath: sockPath, } } func (b *LocalBackend) String() string { - return "local" + return b.machine } // GetConnection returns a gRPC connection to the local server listening on the Unix socket. @@ -57,16 +61,6 @@ func (b *LocalBackend) GetConnection(ctx context.Context, _ string) (context.Con return outCtx, b.conn, err } -// AppendInfo is called to enhance response from the backend with additional data. -func (b *LocalBackend) AppendInfo(_ bool, resp []byte) ([]byte, error) { - return resp, nil -} - -// BuildError is called to convert error from upstream into response field. -func (b *LocalBackend) BuildError(bool, error) ([]byte, error) { - return nil, nil -} - // Close closes the upstream gRPC connection. func (b *LocalBackend) Close() { b.mu.Lock() diff --git a/internal/machine/api/proxy/remote.go b/internal/machine/api/proxy/remote.go index 4c0dd302..0b47fbb0 100644 --- a/internal/machine/api/proxy/remote.go +++ b/internal/machine/api/proxy/remote.go @@ -8,22 +8,19 @@ import ( "google.golang.org/grpc/backoff" "google.golang.org/grpc/credentials/insecure" "google.golang.org/grpc/metadata" - "google.golang.org/grpc/status" - "google.golang.org/protobuf/encoding/protowire" - "google.golang.org/protobuf/proto" "net" "net/netip" "sync" "time" - "uncloud/internal/machine/api/pb" ) -// RemoteBackend is a proxy.Backend implementation that proxies to a remote gRPC server, injecting machine metadata +// RemoteBackend is a proxy.One2ManyResponder implementation that proxies to a remote gRPC server, injecting machine metadata // into the response. // // Based on the Talos apid implementation: // https://github.com/siderolabs/talos/blob/59a78da42cdea8fbccc35d0851f9b0eef928261b/internal/app/apid/pkg/backend/apid.go type RemoteBackend struct { + One2ManyResponder target string mu sync.RWMutex @@ -43,11 +40,16 @@ func NewRemoteBackend(target string) (*RemoteBackend, error) { return nil, fmt.Errorf("target host must be a valid IPv6 address: %s", host) } - return &RemoteBackend{target: target}, nil + return &RemoteBackend{ + One2ManyResponder: One2ManyResponder{ + machine: target, + }, + target: target, + }, nil } func (b *RemoteBackend) String() string { - return b.target + return b.machine } // GetConnection returns a gRPC connection to the remote server. @@ -100,135 +102,6 @@ func (b *RemoteBackend) GetConnection(ctx context.Context, _ string) (context.Co return outCtx, b.conn, err } -// AppendInfo is called to enhance response from the backend with additional data. -// -// AppendInfo enhances upstream response with machine metadata (target). -// -// This method depends on grpc protobuf response structure, each response should -// look like: -// -// message SomeResponse { -// repeated SomeReply messages = 1; // please note field ID == 1 -// } -// -// message SomeReply { -// common.Metadata metadata = 1; -// -// } -// -// As 'SomeReply' is repeated in 'SomeResponse', if we concatenate protobuf representation -// of several 'SomeResponse' messages, we still get valid 'SomeResponse' representation but with more -// entries (feature of protobuf binary representation). -// -// If we look at binary representation of any unary 'SomeResponse' message, it will always contain one -// protobuf field with field ID 1 (see above) and type 2 (embedded message SomeReply is encoded -// as string with length). So if we want to add fields to 'SomeReply', we can simply read field -// header, adjust length for new 'SomeReply' representation, and prepend new field header. -// -// At the same time, we can add 'common.Metadata' structure to 'SomeReply' by simply -// appending or prepending 'common.Metadata' as a single field. This requires 'metadata' -// field to be not defined in original response. (This is due to the fact that protobuf message -// representation is concatenation of each field representation). -// -// To build only single field (Metadata) we use helper message which contains exactly this -// field with same field ID as in every other 'SomeReply': -// -// message Empty { -// common.Metadata metadata = 1; -// } -// -// As streaming replies are not wrapped into 'SomeResponse' with 'repeated', handling is simpler: we just -// need to append Empty with details. -// -// So AppendInfo does the following: validates that response contains field ID 1 encoded as string, -// cuts field header, rest is representation of some reply. Marshal 'Empty' as protobuf, -// which builds 'common.Metadata' field, append it to original response message, build new header -// for new length of some response, and add back new field header. -func (b *RemoteBackend) AppendInfo(streaming bool, resp []byte) ([]byte, error) { - payload, err := proto.Marshal(&pb.Empty{ - Metadata: &pb.Metadata{ - Machine: b.target, - }, - }) - - if streaming { - return append(resp, payload...), err - } - - const ( - metadataField = 1 // field number in proto definition for repeated response - metadataType = 2 // "string" for embedded messages - ) - - // decode protobuf embedded header - - typ, n1 := protowire.ConsumeVarint(resp) - if n1 < 0 { - return nil, protowire.ParseError(n1) - } - - _, n2 := protowire.ConsumeVarint(resp[n1:]) // length - if n2 < 0 { - return nil, protowire.ParseError(n2) - } - - if typ != (metadataField<<3)|metadataType { - return nil, fmt.Errorf("unexpected message format: %d", typ) - } - - if n1+n2 > len(resp) { - return nil, fmt.Errorf("unexpected message size: %d", len(resp)) - } - - // cut off embedded message header - resp = resp[n1+n2:] - // build new embedded message header - prefix := protowire.AppendVarint( - protowire.AppendVarint(nil, (metadataField<<3)|metadataType), - uint64(len(resp)+len(payload)), - ) - resp = append(prefix, resp...) - - return append(resp, payload...), err -} - -// BuildError converts upstream error into message from upstream, so that multiple -// successful and failure responses might be returned. -// -// This simply relies on the fact that any response contains 'Empty' message. -// So if 'Empty' is unmarshalled into any other reply message, all the fields -// are undefined but 'Metadata': -// -// message Empty { -// common.Metadata metadata = 1; -// } -// -// message EmptyResponse { -// repeated Empty messages = 1; -// } -// -// Streaming responses are not wrapped into Empty, so we simply marshall EmptyResponse -// message. -func (b *RemoteBackend) BuildError(streaming bool, err error) ([]byte, error) { - var resp proto.Message = &pb.Empty{ - Metadata: &pb.Metadata{ - Machine: b.target, - Error: err.Error(), - Status: status.Convert(err).Proto(), - }, - } - - if !streaming { - resp = &pb.EmptyResponse{ - Messages: []*pb.Empty{ - resp.(*pb.Empty), - }, - } - } - - return proto.Marshal(resp) -} - // Close closes the upstream gRPC connection. func (b *RemoteBackend) Close() { b.mu.Lock()