gRPC dynamicProxyController startup (#1035)

This commit is contained in:
Michael Quigley
2025-09-10 14:22:09 -04:00
parent 6092178682
commit 058b308859
11 changed files with 541 additions and 34 deletions
+1 -1
View File
@@ -1,6 +1,6 @@
// Code generated by protoc-gen-go. DO NOT EDIT.
// versions:
// protoc-gen-go v1.36.6
// protoc-gen-go v1.36.9
// protoc v6.31.1
// source: agent/agentGrpc/agent.proto
+5 -4
View File
@@ -287,6 +287,7 @@ func local_request_Agent_Version_0(ctx context.Context, marshaler runtime.Marsha
// UnaryRPC :call AgentServer directly.
// StreamingRPC :currently unsupported pending https://github.com/grpc/grpc-go/issues/906.
// Note that using this registration option will cause many gRPC library features to stop working. Consider using RegisterAgentHandlerFromEndpoint instead.
// GRPC interceptors will not work for this type of registration. To use interceptors, you must use the "runtime.WithMiddlewares" option in the "runtime.NewServeMux" call.
func RegisterAgentHandlerServer(ctx context.Context, mux *runtime.ServeMux, server AgentServer) error {
mux.Handle("POST", pattern_Agent_AccessPrivate_0, func(w http.ResponseWriter, req *http.Request, pathParams map[string]string) {
@@ -495,21 +496,21 @@ func RegisterAgentHandlerServer(ctx context.Context, mux *runtime.ServeMux, serv
// RegisterAgentHandlerFromEndpoint is same as RegisterAgentHandler but
// automatically dials to "endpoint" and closes the connection when "ctx" gets done.
func RegisterAgentHandlerFromEndpoint(ctx context.Context, mux *runtime.ServeMux, endpoint string, opts []grpc.DialOption) (err error) {
conn, err := grpc.DialContext(ctx, endpoint, opts...)
conn, err := grpc.NewClient(endpoint, opts...)
if err != nil {
return err
}
defer func() {
if err != nil {
if cerr := conn.Close(); cerr != nil {
grpclog.Infof("Failed to close conn to %s: %v", endpoint, cerr)
grpclog.Errorf("Failed to close conn to %s: %v", endpoint, cerr)
}
return
}
go func() {
<-ctx.Done()
if cerr := conn.Close(); cerr != nil {
grpclog.Infof("Failed to close conn to %s: %v", endpoint, cerr)
grpclog.Errorf("Failed to close conn to %s: %v", endpoint, cerr)
}
}()
}()
@@ -527,7 +528,7 @@ func RegisterAgentHandler(ctx context.Context, mux *runtime.ServeMux, conn *grpc
// to "mux". The handlers forward requests to the grpc endpoint over the given implementation of "AgentClient".
// Note: the gRPC framework executes interceptors within the gRPC handler. If the passed in "AgentClient"
// doesn't go through the normal gRPC flow (creating a gRPC client etc.) then it will be up to the passed in
// "AgentClient" to call the correct interceptors.
// "AgentClient" to call the correct interceptors. This client ignores the HTTP middlewares.
func RegisterAgentHandlerClient(ctx context.Context, mux *runtime.ServeMux, client AgentClient) error {
mux.Handle("POST", pattern_Agent_AccessPrivate_0, func(w http.ResponseWriter, req *http.Request, pathParams map[string]string) {
+18 -18
View File
@@ -29,7 +29,7 @@
"default": {
"description": "An unexpected error response.",
"schema": {
"$ref": "#/definitions/rpcStatus"
"$ref": "#/definitions/googlerpcStatus"
}
}
},
@@ -101,7 +101,7 @@
"default": {
"description": "An unexpected error response.",
"schema": {
"$ref": "#/definitions/rpcStatus"
"$ref": "#/definitions/googlerpcStatus"
}
}
},
@@ -157,7 +157,7 @@
"default": {
"description": "An unexpected error response.",
"schema": {
"$ref": "#/definitions/rpcStatus"
"$ref": "#/definitions/googlerpcStatus"
}
}
},
@@ -187,7 +187,7 @@
"default": {
"description": "An unexpected error response.",
"schema": {
"$ref": "#/definitions/rpcStatus"
"$ref": "#/definitions/googlerpcStatus"
}
}
},
@@ -217,7 +217,7 @@
"default": {
"description": "An unexpected error response.",
"schema": {
"$ref": "#/definitions/rpcStatus"
"$ref": "#/definitions/googlerpcStatus"
}
}
},
@@ -275,7 +275,7 @@
"default": {
"description": "An unexpected error response.",
"schema": {
"$ref": "#/definitions/rpcStatus"
"$ref": "#/definitions/googlerpcStatus"
}
}
},
@@ -375,7 +375,7 @@
"default": {
"description": "An unexpected error response.",
"schema": {
"$ref": "#/definitions/rpcStatus"
"$ref": "#/definitions/googlerpcStatus"
}
}
},
@@ -397,7 +397,7 @@
"default": {
"description": "An unexpected error response.",
"schema": {
"$ref": "#/definitions/rpcStatus"
"$ref": "#/definitions/googlerpcStatus"
}
}
},
@@ -560,16 +560,7 @@
}
}
},
"protobufAny": {
"type": "object",
"properties": {
"@type": {
"type": "string"
}
},
"additionalProperties": {}
},
"rpcStatus": {
"googlerpcStatus": {
"type": "object",
"properties": {
"code": {
@@ -587,6 +578,15 @@
}
}
}
},
"protobufAny": {
"type": "object",
"properties": {
"@type": {
"type": "string"
}
},
"additionalProperties": {}
}
}
}
+3
View File
@@ -12,3 +12,6 @@ protoc --go_out=. --go_opt=paths=source_relative \
--openapiv2_out=. \
agent/agentGrpc/agent.proto
protoc --go_out=. --go_opt=paths=source_relative \
--go-grpc_out=. --go-grpc_opt=paths=source_relative \
controller/dynamicProxyController/dynamicProxyController.proto
+8 -8
View File
@@ -135,14 +135,6 @@ func Run(inCfg *config.Config) error {
api.ShareUpdateAccessHandler = newUpdateAccessHandler()
api.ShareUpdateShareHandler = newUpdateShareHandler()
if cfg.DynamicProxyController != nil {
dPCtrl, err = dynamicProxyController.NewController(cfg.DynamicProxyController)
if err != nil {
return err
}
logrus.Infof("started dynamic proxy controller")
}
if err := controllerStartup(); err != nil {
return err
}
@@ -153,6 +145,14 @@ func Run(inCfg *config.Config) error {
return errors.Wrap(err, "error opening store")
}
if cfg.DynamicProxyController != nil {
dPCtrl, err = dynamicProxyController.NewController(cfg.DynamicProxyController, str)
if err != nil {
return err
}
logrus.Infof("started dynamic proxy controller")
}
if cfg.Metrics != nil && cfg.Metrics.Influx != nil {
idb = influxdb2.NewClient(cfg.Metrics.Influx.Url, cfg.Metrics.Influx.Token)
} else {
+3 -1
View File
@@ -1,5 +1,7 @@
package dynamicProxyController
type Config struct {
AmqpPublisher *AmqpPublisherConfig
IdentityPath string `df:"required"`
ServiceName string `df:"required"`
AmqpPublisher *AmqpPublisherConfig `df:"required"`
}
@@ -3,20 +3,56 @@ package dynamicProxyController
import (
"context"
"github.com/openziti/sdk-golang/ziti"
"github.com/openziti/zrok/controller/store"
"github.com/openziti/zrok/dynamicProxyModel"
"github.com/sirupsen/logrus"
"google.golang.org/grpc"
)
type Controller struct {
UnimplementedDynamicProxyControllerServer
str *store.Store
publisher *AmqpPublisher
zCfg *ziti.Config
zCtx ziti.Context
}
func NewController(cfg *Config) (*Controller, error) {
func NewController(cfg *Config, str *store.Store) (*Controller, error) {
publisher, err := NewAmqpPublisher(cfg.AmqpPublisher)
if err != nil {
return nil, err
}
return &Controller{publisher: publisher}, nil
zCfg, err := ziti.NewConfigFromFile(cfg.IdentityPath)
if err != nil {
return nil, err
}
zCtx, err := ziti.NewContext(zCfg)
if err != nil {
return nil, err
}
srv := grpc.NewServer()
ctrl := &Controller{
str: str,
publisher: publisher,
zCfg: zCfg,
zCtx: zCtx,
}
RegisterDynamicProxyControllerServer(srv, ctrl)
l, err := zCtx.Listen(cfg.ServiceName)
if err != nil {
return nil, err
}
go func() {
if err := srv.Serve(l); err != nil {
logrus.Errorf("error serving dynamic proxy controller: %v", err)
return
}
}()
logrus.Infof("started dynamic proxy controller server")
return ctrl, nil
}
func (c *Controller) SendMappingUpdate(frontendToken string, m dynamicProxyModel.Mapping) error {
@@ -26,3 +62,32 @@ func (c *Controller) SendMappingUpdate(frontendToken string, m dynamicProxyModel
logrus.Infof("sent mapping update '%+v' -> '%s'", m, frontendToken)
return nil
}
func (c *Controller) FrontendMappings(_ context.Context, req *FrontendMappingsRequest) (*FrontendMappingsResponse, error) {
trx, err := c.str.Begin()
if err != nil {
return nil, err
}
defer trx.Rollback()
var mappings []*store.FrontendMapping
if req.GetName() == "" {
mappings, err = c.str.FindFrontendMappingsByFrontendTokenWithVersionOrHigher(req.GetFrontendToken(), req.GetVersion(), trx)
} else {
mappings, err = c.str.FindFrontendMappingsWithVersionOrHigher(req.GetFrontendToken(), req.GetName(), req.GetVersion(), trx)
}
if err != nil {
return nil, err
}
out := make([]*FrontendMapping, len(mappings))
for i, storeMapping := range mappings {
out[i] = &FrontendMapping{
Name: storeMapping.Name,
Version: storeMapping.Version,
ShareToken: storeMapping.ShareToken,
}
}
return &FrontendMappingsResponse{FrontendMappings: out}, nil
}
@@ -0,0 +1,259 @@
// Code generated by protoc-gen-go. DO NOT EDIT.
// versions:
// protoc-gen-go v1.36.9
// protoc v6.31.1
// source: controller/dynamicProxyController/dynamicProxyController.proto
package dynamicProxyController
import (
protoreflect "google.golang.org/protobuf/reflect/protoreflect"
protoimpl "google.golang.org/protobuf/runtime/protoimpl"
reflect "reflect"
sync "sync"
unsafe "unsafe"
)
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)
)
type FrontendMappingsRequest struct {
state protoimpl.MessageState `protogen:"open.v1"`
FrontendToken string `protobuf:"bytes,1,opt,name=frontendToken,proto3" json:"frontendToken,omitempty"`
Name string `protobuf:"bytes,2,opt,name=name,proto3" json:"name,omitempty"`
Version int64 `protobuf:"varint,3,opt,name=version,proto3" json:"version,omitempty"`
unknownFields protoimpl.UnknownFields
sizeCache protoimpl.SizeCache
}
func (x *FrontendMappingsRequest) Reset() {
*x = FrontendMappingsRequest{}
mi := &file_controller_dynamicProxyController_dynamicProxyController_proto_msgTypes[0]
ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
ms.StoreMessageInfo(mi)
}
func (x *FrontendMappingsRequest) String() string {
return protoimpl.X.MessageStringOf(x)
}
func (*FrontendMappingsRequest) ProtoMessage() {}
func (x *FrontendMappingsRequest) ProtoReflect() protoreflect.Message {
mi := &file_controller_dynamicProxyController_dynamicProxyController_proto_msgTypes[0]
if x != nil {
ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
if ms.LoadMessageInfo() == nil {
ms.StoreMessageInfo(mi)
}
return ms
}
return mi.MessageOf(x)
}
// Deprecated: Use FrontendMappingsRequest.ProtoReflect.Descriptor instead.
func (*FrontendMappingsRequest) Descriptor() ([]byte, []int) {
return file_controller_dynamicProxyController_dynamicProxyController_proto_rawDescGZIP(), []int{0}
}
func (x *FrontendMappingsRequest) GetFrontendToken() string {
if x != nil {
return x.FrontendToken
}
return ""
}
func (x *FrontendMappingsRequest) GetName() string {
if x != nil {
return x.Name
}
return ""
}
func (x *FrontendMappingsRequest) GetVersion() int64 {
if x != nil {
return x.Version
}
return 0
}
type FrontendMappingsResponse struct {
state protoimpl.MessageState `protogen:"open.v1"`
FrontendMappings []*FrontendMapping `protobuf:"bytes,1,rep,name=frontendMappings,proto3" json:"frontendMappings,omitempty"`
unknownFields protoimpl.UnknownFields
sizeCache protoimpl.SizeCache
}
func (x *FrontendMappingsResponse) Reset() {
*x = FrontendMappingsResponse{}
mi := &file_controller_dynamicProxyController_dynamicProxyController_proto_msgTypes[1]
ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
ms.StoreMessageInfo(mi)
}
func (x *FrontendMappingsResponse) String() string {
return protoimpl.X.MessageStringOf(x)
}
func (*FrontendMappingsResponse) ProtoMessage() {}
func (x *FrontendMappingsResponse) ProtoReflect() protoreflect.Message {
mi := &file_controller_dynamicProxyController_dynamicProxyController_proto_msgTypes[1]
if x != nil {
ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
if ms.LoadMessageInfo() == nil {
ms.StoreMessageInfo(mi)
}
return ms
}
return mi.MessageOf(x)
}
// Deprecated: Use FrontendMappingsResponse.ProtoReflect.Descriptor instead.
func (*FrontendMappingsResponse) Descriptor() ([]byte, []int) {
return file_controller_dynamicProxyController_dynamicProxyController_proto_rawDescGZIP(), []int{1}
}
func (x *FrontendMappingsResponse) GetFrontendMappings() []*FrontendMapping {
if x != nil {
return x.FrontendMappings
}
return nil
}
type FrontendMapping struct {
state protoimpl.MessageState `protogen:"open.v1"`
Name string `protobuf:"bytes,1,opt,name=name,proto3" json:"name,omitempty"`
Version int64 `protobuf:"varint,2,opt,name=version,proto3" json:"version,omitempty"`
ShareToken string `protobuf:"bytes,3,opt,name=shareToken,proto3" json:"shareToken,omitempty"`
unknownFields protoimpl.UnknownFields
sizeCache protoimpl.SizeCache
}
func (x *FrontendMapping) Reset() {
*x = FrontendMapping{}
mi := &file_controller_dynamicProxyController_dynamicProxyController_proto_msgTypes[2]
ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
ms.StoreMessageInfo(mi)
}
func (x *FrontendMapping) String() string {
return protoimpl.X.MessageStringOf(x)
}
func (*FrontendMapping) ProtoMessage() {}
func (x *FrontendMapping) ProtoReflect() protoreflect.Message {
mi := &file_controller_dynamicProxyController_dynamicProxyController_proto_msgTypes[2]
if x != nil {
ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
if ms.LoadMessageInfo() == nil {
ms.StoreMessageInfo(mi)
}
return ms
}
return mi.MessageOf(x)
}
// Deprecated: Use FrontendMapping.ProtoReflect.Descriptor instead.
func (*FrontendMapping) Descriptor() ([]byte, []int) {
return file_controller_dynamicProxyController_dynamicProxyController_proto_rawDescGZIP(), []int{2}
}
func (x *FrontendMapping) GetName() string {
if x != nil {
return x.Name
}
return ""
}
func (x *FrontendMapping) GetVersion() int64 {
if x != nil {
return x.Version
}
return 0
}
func (x *FrontendMapping) GetShareToken() string {
if x != nil {
return x.ShareToken
}
return ""
}
var File_controller_dynamicProxyController_dynamicProxyController_proto protoreflect.FileDescriptor
const file_controller_dynamicProxyController_dynamicProxyController_proto_rawDesc = "" +
"\n" +
">controller/dynamicProxyController/dynamicProxyController.proto\"m\n" +
"\x17FrontendMappingsRequest\x12$\n" +
"\rfrontendToken\x18\x01 \x01(\tR\rfrontendToken\x12\x12\n" +
"\x04name\x18\x02 \x01(\tR\x04name\x12\x18\n" +
"\aversion\x18\x03 \x01(\x03R\aversion\"X\n" +
"\x18FrontendMappingsResponse\x12<\n" +
"\x10frontendMappings\x18\x01 \x03(\v2\x10.FrontendMappingR\x10frontendMappings\"_\n" +
"\x0fFrontendMapping\x12\x12\n" +
"\x04name\x18\x01 \x01(\tR\x04name\x12\x18\n" +
"\aversion\x18\x02 \x01(\x03R\aversion\x12\x1e\n" +
"\n" +
"shareToken\x18\x03 \x01(\tR\n" +
"shareToken2a\n" +
"\x16DynamicProxyController\x12G\n" +
"\x10FrontendMappings\x12\x18.FrontendMappingsRequest\x1a\x19.FrontendMappingsResponseB<Z:github.com/openziti/zrok/controller/dynamicProxyControllerb\x06proto3"
var (
file_controller_dynamicProxyController_dynamicProxyController_proto_rawDescOnce sync.Once
file_controller_dynamicProxyController_dynamicProxyController_proto_rawDescData []byte
)
func file_controller_dynamicProxyController_dynamicProxyController_proto_rawDescGZIP() []byte {
file_controller_dynamicProxyController_dynamicProxyController_proto_rawDescOnce.Do(func() {
file_controller_dynamicProxyController_dynamicProxyController_proto_rawDescData = protoimpl.X.CompressGZIP(unsafe.Slice(unsafe.StringData(file_controller_dynamicProxyController_dynamicProxyController_proto_rawDesc), len(file_controller_dynamicProxyController_dynamicProxyController_proto_rawDesc)))
})
return file_controller_dynamicProxyController_dynamicProxyController_proto_rawDescData
}
var file_controller_dynamicProxyController_dynamicProxyController_proto_msgTypes = make([]protoimpl.MessageInfo, 3)
var file_controller_dynamicProxyController_dynamicProxyController_proto_goTypes = []any{
(*FrontendMappingsRequest)(nil), // 0: FrontendMappingsRequest
(*FrontendMappingsResponse)(nil), // 1: FrontendMappingsResponse
(*FrontendMapping)(nil), // 2: FrontendMapping
}
var file_controller_dynamicProxyController_dynamicProxyController_proto_depIdxs = []int32{
2, // 0: FrontendMappingsResponse.frontendMappings:type_name -> FrontendMapping
0, // 1: DynamicProxyController.FrontendMappings:input_type -> FrontendMappingsRequest
1, // 2: DynamicProxyController.FrontendMappings:output_type -> FrontendMappingsResponse
2, // [2:3] is the sub-list for method output_type
1, // [1:2] 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_controller_dynamicProxyController_dynamicProxyController_proto_init() }
func file_controller_dynamicProxyController_dynamicProxyController_proto_init() {
if File_controller_dynamicProxyController_dynamicProxyController_proto != nil {
return
}
type x struct{}
out := protoimpl.TypeBuilder{
File: protoimpl.DescBuilder{
GoPackagePath: reflect.TypeOf(x{}).PkgPath(),
RawDescriptor: unsafe.Slice(unsafe.StringData(file_controller_dynamicProxyController_dynamicProxyController_proto_rawDesc), len(file_controller_dynamicProxyController_dynamicProxyController_proto_rawDesc)),
NumEnums: 0,
NumMessages: 3,
NumExtensions: 0,
NumServices: 1,
},
GoTypes: file_controller_dynamicProxyController_dynamicProxyController_proto_goTypes,
DependencyIndexes: file_controller_dynamicProxyController_dynamicProxyController_proto_depIdxs,
MessageInfos: file_controller_dynamicProxyController_dynamicProxyController_proto_msgTypes,
}.Build()
File_controller_dynamicProxyController_dynamicProxyController_proto = out.File
file_controller_dynamicProxyController_dynamicProxyController_proto_goTypes = nil
file_controller_dynamicProxyController_dynamicProxyController_proto_depIdxs = nil
}
@@ -0,0 +1,23 @@
syntax = "proto3";
option go_package = "github.com/openziti/zrok/controller/dynamicProxyController";
service DynamicProxyController {
rpc FrontendMappings(FrontendMappingsRequest) returns (FrontendMappingsResponse);
}
message FrontendMappingsRequest {
string frontendToken = 1;
string name = 2;
int64 version = 3;
}
message FrontendMappingsResponse {
repeated FrontendMapping frontendMappings = 1;
}
message FrontendMapping {
string name = 1;
int64 version = 2;
string shareToken = 3;
}
@@ -0,0 +1,122 @@
// Code generated by protoc-gen-go-grpc. DO NOT EDIT.
// versions:
// - protoc-gen-go-grpc v1.5.1
// - protoc v6.31.1
// source: controller/dynamicProxyController/dynamicProxyController.proto
package dynamicProxyController
import (
context "context"
grpc "google.golang.org/grpc"
codes "google.golang.org/grpc/codes"
status "google.golang.org/grpc/status"
)
// This is a compile-time assertion to ensure that this generated file
// is compatible with the grpc package it is being compiled against.
// Requires gRPC-Go v1.64.0 or later.
const _ = grpc.SupportPackageIsVersion9
const (
DynamicProxyController_FrontendMappings_FullMethodName = "/DynamicProxyController/FrontendMappings"
)
// DynamicProxyControllerClient is the client API for DynamicProxyController service.
//
// For semantics around ctx use and closing/ending streaming RPCs, please refer to https://pkg.go.dev/google.golang.org/grpc/?tab=doc#ClientConn.NewStream.
type DynamicProxyControllerClient interface {
FrontendMappings(ctx context.Context, in *FrontendMappingsRequest, opts ...grpc.CallOption) (*FrontendMappingsResponse, error)
}
type dynamicProxyControllerClient struct {
cc grpc.ClientConnInterface
}
func NewDynamicProxyControllerClient(cc grpc.ClientConnInterface) DynamicProxyControllerClient {
return &dynamicProxyControllerClient{cc}
}
func (c *dynamicProxyControllerClient) FrontendMappings(ctx context.Context, in *FrontendMappingsRequest, opts ...grpc.CallOption) (*FrontendMappingsResponse, error) {
cOpts := append([]grpc.CallOption{grpc.StaticMethod()}, opts...)
out := new(FrontendMappingsResponse)
err := c.cc.Invoke(ctx, DynamicProxyController_FrontendMappings_FullMethodName, in, out, cOpts...)
if err != nil {
return nil, err
}
return out, nil
}
// DynamicProxyControllerServer is the server API for DynamicProxyController service.
// All implementations must embed UnimplementedDynamicProxyControllerServer
// for forward compatibility.
type DynamicProxyControllerServer interface {
FrontendMappings(context.Context, *FrontendMappingsRequest) (*FrontendMappingsResponse, error)
mustEmbedUnimplementedDynamicProxyControllerServer()
}
// UnimplementedDynamicProxyControllerServer must be embedded to have
// forward compatible implementations.
//
// NOTE: this should be embedded by value instead of pointer to avoid a nil
// pointer dereference when methods are called.
type UnimplementedDynamicProxyControllerServer struct{}
func (UnimplementedDynamicProxyControllerServer) FrontendMappings(context.Context, *FrontendMappingsRequest) (*FrontendMappingsResponse, error) {
return nil, status.Errorf(codes.Unimplemented, "method FrontendMappings not implemented")
}
func (UnimplementedDynamicProxyControllerServer) mustEmbedUnimplementedDynamicProxyControllerServer() {
}
func (UnimplementedDynamicProxyControllerServer) testEmbeddedByValue() {}
// UnsafeDynamicProxyControllerServer may be embedded to opt out of forward compatibility for this service.
// Use of this interface is not recommended, as added methods to DynamicProxyControllerServer will
// result in compilation errors.
type UnsafeDynamicProxyControllerServer interface {
mustEmbedUnimplementedDynamicProxyControllerServer()
}
func RegisterDynamicProxyControllerServer(s grpc.ServiceRegistrar, srv DynamicProxyControllerServer) {
// If the following call pancis, it indicates UnimplementedDynamicProxyControllerServer was
// embedded by pointer and is nil. This will cause panics if an
// unimplemented method is ever invoked, so we test this at initialization
// time to prevent it from happening at runtime later due to I/O.
if t, ok := srv.(interface{ testEmbeddedByValue() }); ok {
t.testEmbeddedByValue()
}
s.RegisterService(&DynamicProxyController_ServiceDesc, srv)
}
func _DynamicProxyController_FrontendMappings_Handler(srv interface{}, ctx context.Context, dec func(interface{}) error, interceptor grpc.UnaryServerInterceptor) (interface{}, error) {
in := new(FrontendMappingsRequest)
if err := dec(in); err != nil {
return nil, err
}
if interceptor == nil {
return srv.(DynamicProxyControllerServer).FrontendMappings(ctx, in)
}
info := &grpc.UnaryServerInfo{
Server: srv,
FullMethod: DynamicProxyController_FrontendMappings_FullMethodName,
}
handler := func(ctx context.Context, req interface{}) (interface{}, error) {
return srv.(DynamicProxyControllerServer).FrontendMappings(ctx, req.(*FrontendMappingsRequest))
}
return interceptor(ctx, in, info, handler)
}
// DynamicProxyController_ServiceDesc is the grpc.ServiceDesc for DynamicProxyController service.
// It's only intended for direct use with grpc.RegisterService,
// and not to be introspected or modified (even as a copy)
var DynamicProxyController_ServiceDesc = grpc.ServiceDesc{
ServiceName: "DynamicProxyController",
HandlerType: (*DynamicProxyControllerServer)(nil),
Methods: []grpc.MethodDesc{
{
MethodName: "FrontendMappings",
Handler: _DynamicProxyController_FrontendMappings_Handler,
},
},
Streams: []grpc.StreamDesc{},
Metadata: "controller/dynamicProxyController/dynamicProxyController.proto",
}
+32
View File
@@ -105,4 +105,36 @@ func (str *Store) DeleteFrontendMappingsByFrontendToken(frontendToken string, tx
return errors.Wrap(err, "error executing frontend_mappings delete by frontend_token statement")
}
return nil
}
func (str *Store) FindFrontendMappingsWithVersionOrHigher(frontendToken, name string, version int64, tx *sqlx.Tx) ([]*FrontendMapping, error) {
rows, err := tx.Queryx("select * from frontend_mappings where frontend_token = $1 and name = $2 and version >= $3 order by version asc", frontendToken, name, version)
if err != nil {
return nil, errors.Wrap(err, "error selecting frontend mappings with version or higher")
}
var mappings []*FrontendMapping
for rows.Next() {
fm := &FrontendMapping{}
if err := rows.StructScan(fm); err != nil {
return nil, errors.Wrap(err, "error scanning frontend mapping")
}
mappings = append(mappings, fm)
}
return mappings, nil
}
func (str *Store) FindFrontendMappingsByFrontendTokenWithVersionOrHigher(frontendToken string, version int64, tx *sqlx.Tx) ([]*FrontendMapping, error) {
rows, err := tx.Queryx("select * from frontend_mappings where frontend_token = $1 and version >= $2 order by name, version asc", frontendToken, version)
if err != nil {
return nil, errors.Wrap(err, "error selecting frontend mappings by frontend_token with version or higher")
}
var mappings []*FrontendMapping
for rows.Next() {
fm := &FrontendMapping{}
if err := rows.StructScan(fm); err != nil {
return nil, errors.Wrap(err, "error scanning frontend mapping")
}
mappings = append(mappings, fm)
}
return mappings, nil
}