pb: flip WithBrokerGrpcClient to ZAP via brokerPool
dialBrokerZapAddr (PQ-TLS when grpc.msg_broker.cert/.key, else plaintext) + a transport.NewPool, mirroring filerPool. WithBrokerGrpcClient now pools one ZAP conn per broker address and runs fn with mqzap.New(conn) — the broker fully cut over so its address IS the ZAP endpoint (no port offset). DialOption now inert.
This commit is contained in:
@@ -19,6 +19,7 @@ import (
|
||||
"google.golang.org/grpc/metadata"
|
||||
|
||||
"github.com/hanzoai/s3/s3/glog"
|
||||
"github.com/hanzoai/s3/s3/mqzap"
|
||||
"github.com/hanzoai/s3/s3/pb/volume_server_pb"
|
||||
"github.com/hanzoai/s3/s3/security"
|
||||
"github.com/hanzoai/s3/s3/util"
|
||||
@@ -567,13 +568,36 @@ func WithOneOfGrpcMasterClients(streamingMode bool, masterGrpcAddresses map[stri
|
||||
return err
|
||||
}
|
||||
|
||||
// dialBrokerZapAddr opens a ZAP connection to a broker address. When
|
||||
// grpc.msg_broker.cert/.key is configured it is PQ-secured TLS (transport.PQTLSConfig
|
||||
// pins X25519MLKEM768, the PQ X-Wing curve) presenting the client cert and trusting
|
||||
// grpc.ca — the same mTLS the legacy gRPC broker client used. Otherwise plaintext
|
||||
// (loopback / dev), matching the broker server's gating in command/mq_broker.go. The
|
||||
// returned *Conn drives both unary Call and OpenStream. The broker fully cut over to
|
||||
// ZAP, so its address IS the ZAP endpoint (no port offset, unlike master/IAM).
|
||||
func dialBrokerZapAddr(addr string) (transport.Conn, error) {
|
||||
if cfg := security.ClientTLSConfig(util.GetViper(), "grpc.msg_broker"); cfg != nil {
|
||||
return transport.DialTLS("tcp", addr, transport.PQTLSConfig(cfg))
|
||||
}
|
||||
return transport.Dial("tcp", addr)
|
||||
}
|
||||
|
||||
// brokerPool reuses one ZAP connection per broker address across calls — a Conn is
|
||||
// concurrency-safe, so this avoids a fresh TCP (and, under grpc.msg_broker mTLS, a
|
||||
// fresh X25519MLKEM768 handshake) on every broker RPC. Generic pooling lives in the
|
||||
// transport (transport.Pool); only the dial choice is ours. Mirrors filerPool.
|
||||
var brokerPool = transport.NewPool(dialBrokerZapAddr)
|
||||
|
||||
// WithBrokerGrpcClient dials the broker over the native ZAP transport and runs fn
|
||||
// with a mq_pb.HanzoMessagingClient backed by that connection (mqzap.New). The
|
||||
// streamingMode/grpcDialOption parameters are retained for caller compatibility; the
|
||||
// ZAP path needs neither a streaming flag (every stream is a transport stream) nor a
|
||||
// dial option. The pooled connection outlives fn.
|
||||
func WithBrokerGrpcClient(streamingMode bool, brokerGrpcAddress string, grpcDialOption grpc.DialOption, fn func(client mq_pb.HanzoMessagingClient) error) error {
|
||||
|
||||
return WithGrpcClient(context.Background(), streamingMode, 0, func(grpcConnection *grpc.ClientConn) error {
|
||||
client := mq_pb.NewHanzoMessagingClient(grpcConnection)
|
||||
return fn(client)
|
||||
}, brokerGrpcAddress, false, grpcDialOption)
|
||||
|
||||
_, _ = streamingMode, grpcDialOption
|
||||
return brokerPool.With(brokerGrpcAddress, func(conn transport.Conn) error {
|
||||
return fn(mqzap.New(conn, nil))
|
||||
})
|
||||
}
|
||||
|
||||
func WithFilerClient(streamingMode bool, signature int32, filer ServerAddress, grpcDialOption grpc.DialOption, fn func(client filer_pb.HanzoFilerClient) error) error {
|
||||
|
||||
Reference in New Issue
Block a user