From e18fbcc0fa40233db212a9fb2f1dc8a0ebd8813a Mon Sep 17 00:00:00 2001 From: Paul Wells Date: Sat, 5 Sep 2026 06:58:59 -0700 Subject: [PATCH] wire psrpc bus compression into the message bus constructor (#4844) getMessageBus built the bus from the redis client alone, so there was no way to reach the gzip settings psrpc v0.7.6 added at the bus boundary. Take rpc.PSRPCConfig, which the wire graph already provides, and pass its bus options to both the redis and the local bus. Compression is off by default. A peer on an older psrpc cannot decode a compressed payload and drops it silently, so egress, ingress, SIP and agent workers all have to be upgraded before quality is raised. --- config-sample.yaml | 9 +++++++ go.mod | 4 +-- go.sum | 8 +++--- pkg/config/config_test.go | 21 +++++++++++++++ pkg/service/wire.go | 7 ++--- pkg/service/wire_gen.go | 55 ++++++++++++++++++++------------------- 6 files changed, 68 insertions(+), 36 deletions(-) diff --git a/config-sample.yaml b/config-sample.yaml index f07299972..459b693f6 100644 --- a/config-sample.yaml +++ b/config-sample.yaml @@ -268,6 +268,15 @@ keys: # backoff: 500ms # # number of messages to buffer before dropping # buffer_size: 1000 +# # optional gzip compression of bus payloads +# compression: +# # gzip level 1-9; 0 disables. every node on the bus must support +# # compression before enabling it +# quality: 0 +# # payload bytes below which compression is skipped +# threshold: 1024 +# # cap on an inbound payload after decompression, 0 for unlimited +# max_decompressed_size: 0 # customize audio level sensitivity # audio: diff --git a/go.mod b/go.mod index 488bf8d0c..4f13448d4 100644 --- a/go.mod +++ b/go.mod @@ -21,8 +21,8 @@ require ( github.com/jxskiss/base62 v1.1.0 github.com/livekit/mageutil v0.0.0-20250511045019-0f1ff63f7731 github.com/livekit/mediatransportutil v0.0.0-20260821083140-f234b534b095 - github.com/livekit/protocol v1.51.1-0.20260903060125-0cf5ba018b8e - github.com/livekit/psrpc v0.7.5 + github.com/livekit/protocol v1.51.1-0.20260905133529-a4f4b5c0c23f + github.com/livekit/psrpc v0.7.6 github.com/mackerelio/go-osstat v0.2.8 github.com/magefile/mage v1.17.2 github.com/mitchellh/go-homedir v1.1.0 diff --git a/go.sum b/go.sum index 458914a7b..083396f13 100644 --- a/go.sum +++ b/go.sum @@ -164,10 +164,10 @@ github.com/livekit/mageutil v0.0.0-20250511045019-0f1ff63f7731 h1:9x+U2HGLrSw5AT github.com/livekit/mageutil v0.0.0-20250511045019-0f1ff63f7731/go.mod h1:Rs3MhFwutWhGwmY1VQsygw28z5bWcnEYmS1OG9OxjOQ= github.com/livekit/mediatransportutil v0.0.0-20260821083140-f234b534b095 h1:BcliKAXoMhl/nWmzQweQ5kmh4Qqagxl4s3Z5pvM/7AY= github.com/livekit/mediatransportutil v0.0.0-20260821083140-f234b534b095/go.mod h1:o8CFmAdrVwzJNOCsQCLUzXRjokkufNshnQHOe4fRaqU= -github.com/livekit/protocol v1.51.1-0.20260903060125-0cf5ba018b8e h1:QAfYvElm7Jr4ZHs1NiX4OTKpVnRO2d9hzSfCPzGMyMI= -github.com/livekit/protocol v1.51.1-0.20260903060125-0cf5ba018b8e/go.mod h1:qt6LOKs6pzIwJF+P/dSUKHJE1SPowzFWDK6byqktdfU= -github.com/livekit/psrpc v0.7.5 h1:WxfJIQ41X1b+48A1uzc8Gy9FhYEMxBJNCpCYqPdO/Ds= -github.com/livekit/psrpc v0.7.5/go.mod h1:rAI+m2+/cb4x9RXhLRtUx5ZwdfjjXOl4zi46IjEetaw= +github.com/livekit/protocol v1.51.1-0.20260905133529-a4f4b5c0c23f h1:+48IWNrsoTgbB0JGv+xJl4umx6e/vh1BX1Tcokuacwc= +github.com/livekit/protocol v1.51.1-0.20260905133529-a4f4b5c0c23f/go.mod h1:zxowkRnQlJ2VMn6ZyinXMDi985wcKXuWNeXmEERqFAs= +github.com/livekit/psrpc v0.7.6 h1:YG07lUMTtf+eaYI2goT9zcVZ0kGJNWN1K6ETNFtv1HQ= +github.com/livekit/psrpc v0.7.6/go.mod h1:DMw15RO7x5XmcgfwzWJYk2In605kx+wu1QRVbPfzf8M= github.com/livekit/webrtc-pion/v4 v4.2.18-warp.1 h1:fH+v4W+NFp9FfPzON6FaUFNmazGcctaAhb2P+Ksf+1s= github.com/livekit/webrtc-pion/v4 v4.2.18-warp.1/go.mod h1:rbKGHo2OpNUImWTvRIV776/3xjjq/t47H3IZiTtwluc= github.com/mackerelio/go-osstat v0.2.8 h1:I2duicTaCGWoM53XwAwA9OIe1inu0xnVs8/pqOWWVr4= diff --git a/pkg/config/config_test.go b/pkg/config/config_test.go index 3e9cfbf48..0b6439945 100644 --- a/pkg/config/config_test.go +++ b/pkg/config/config_test.go @@ -61,6 +61,27 @@ func TestConfig_SignalMessageSizeLimitOverride(t *testing.T) { require.Equal(t, int64(0), conf.Limit.AgentSignalMessageSizeLimit) } +func TestConfig_PSRPCCompressionDefaults(t *testing.T) { + conf, err := NewConfig("", true, nil, nil) + require.NoError(t, err) + require.Equal(t, 0, conf.PSRPC.Compression.Quality) + require.Equal(t, 1024, conf.PSRPC.Compression.Threshold) + require.Equal(t, 0, conf.PSRPC.Compression.MaxDecompressedSize) +} + +func TestConfig_PSRPCCompressionOverride(t *testing.T) { + const content = `psrpc: + compression: + quality: 6 + max_decompressed_size: 4096` + conf, err := NewConfig(content, true, nil, nil) + require.NoError(t, err) + require.Equal(t, 6, conf.PSRPC.Compression.Quality) + require.Equal(t, 4096, conf.PSRPC.Compression.MaxDecompressedSize) + require.Equal(t, 1024, conf.PSRPC.Compression.Threshold) + require.Equal(t, 3, conf.PSRPC.MaxAttempts) +} + func TestConfig_UnknownKeys(t *testing.T) { const content = `unknown: 10 room: diff --git a/pkg/service/wire.go b/pkg/service/wire.go index a2d60fd32..9d1842b4b 100644 --- a/pkg/service/wire.go +++ b/pkg/service/wire.go @@ -193,11 +193,12 @@ func createStore(rc redis.UniversalClient) ObjectStore { return NewLocalStore() } -func getMessageBus(rc redis.UniversalClient) psrpc.MessageBus { +func getMessageBus(rc redis.UniversalClient, psrpcConf rpc.PSRPCConfig) psrpc.MessageBus { + opts := psrpcConf.BusOptions() if rc == nil { - return psrpc.NewLocalMessageBus() + return psrpc.NewLocalMessageBus(opts...) } - return psrpc.NewRedisMessageBus(rc) + return psrpc.NewRedisMessageBus(rc, opts...) } func getEgressStore(s ObjectStore) EgressStore { diff --git a/pkg/service/wire_gen.go b/pkg/service/wire_gen.go index 0da441132..f7c5a6267 100644 --- a/pkg/service/wire_gen.go +++ b/pkg/service/wire_gen.go @@ -39,14 +39,14 @@ func InitializeServer(conf *config.Config, currentNode routing.LocalNode) (*Live return nil, err } nodeID := getNodeID(currentNode) - messageBus := getMessageBus(universalClient) + psrpcConfig := getPSRPCConfig(conf) + v := getMessageBus(universalClient, psrpcConfig) signalRelayConfig := getSignalRelayConfig(conf) - signalClient, err := routing.NewSignalClient(nodeID, messageBus, signalRelayConfig) + signalClient, err := routing.NewSignalClient(nodeID, v, signalRelayConfig) if err != nil { return nil, err } - psrpcConfig := getPSRPCConfig(conf) - clientParams := getPSRPCClientParams(psrpcConfig, messageBus) + clientParams := getPSRPCClientParams(psrpcConfig, v) roomConfig := getRoomConfig(conf) roomManagerClient, err := routing.NewRoomManagerClient(clientParams, roomConfig) if err != nil { @@ -80,57 +80,57 @@ func InitializeServer(conf *config.Config, currentNode routing.LocalNode) (*Live } analyticsService := telemetry.NewAnalyticsService(conf, currentNode) telemetryService := createTelemetryService(queuedNotifier, analyticsService) - ioInfoService, err := NewIOInfoService(messageBus, egressStore, ingressStore, sipStore, telemetryService) + ioInfoService, err := NewIOInfoService(v, egressStore, ingressStore, sipStore, telemetryService) if err != nil { return nil, err } rtcEgressLauncher := NewEgressLauncher(egressClient, ioInfoService, objectStore) topicFormatter := rpc.NewTopicFormatter() - v, err := rpc.NewTypedRoomClient(clientParams) + v2, err := rpc.NewTypedRoomClient(clientParams) if err != nil { return nil, err } - v2, err := rpc.NewTypedParticipantClient(clientParams) + v3, err := rpc.NewTypedParticipantClient(clientParams) if err != nil { return nil, err } - roomService, err := NewRoomService(limitConfig, apiConfig, router, roomAllocator, objectStore, rtcEgressLauncher, topicFormatter, v, v2) + roomService, err := NewRoomService(limitConfig, apiConfig, router, roomAllocator, objectStore, rtcEgressLauncher, topicFormatter, v2, v3) if err != nil { return nil, err } - v3, err := rpc.NewTypedAgentDispatchInternalClient(clientParams) + v4, err := rpc.NewTypedAgentDispatchInternalClient(clientParams) if err != nil { return nil, err } - agentDispatchService := NewAgentDispatchService(limitConfig, v3, topicFormatter, roomAllocator, router) + agentDispatchService := NewAgentDispatchService(limitConfig, v4, topicFormatter, roomAllocator, router) egressService := NewEgressService(egressClient, rtcEgressLauncher, ioInfoService, roomService) ingressConfig := getIngressConfig(conf) ingressClient, err := rpc.NewIngressClient(clientParams) if err != nil { return nil, err } - ingressService := NewIngressService(ingressConfig, nodeID, messageBus, ingressClient, ingressStore, ioInfoService, telemetryService) + ingressService := NewIngressService(ingressConfig, nodeID, v, ingressClient, ingressStore, ioInfoService, telemetryService) sipConfig := getSIPConfig(conf) sipClient, err := newSIPClient(clientParams) if err != nil { return nil, err } - sipService := NewSIPService(sipConfig, nodeID, messageBus, sipClient, sipStore, roomService, telemetryService) + sipService := NewSIPService(sipConfig, nodeID, v, sipClient, sipStore, roomService, telemetryService) rtcService := NewRTCService(conf, roomAllocator, router, telemetryService) - v4, err := rpc.NewTypedWHIPParticipantClient(clientParams) + v5, err := rpc.NewTypedWHIPParticipantClient(clientParams) if err != nil { return nil, err } - serviceWHIPService, err := NewWHIPService(conf, router, roomAllocator, clientParams, topicFormatter, v4) + serviceWHIPService, err := NewWHIPService(conf, router, roomAllocator, clientParams, topicFormatter, v5) if err != nil { return nil, err } - agentService, err := NewAgentService(conf, currentNode, messageBus, keyProvider) + agentService, err := NewAgentService(conf, currentNode, v, keyProvider) if err != nil { return nil, err } agentConfig := getAgentConfig(conf) - client, err := agent.NewAgentClient(messageBus, agentConfig) + client, err := agent.NewAgentClient(v, agentConfig) if err != nil { return nil, err } @@ -138,16 +138,16 @@ func InitializeServer(conf *config.Config, currentNode routing.LocalNode) (*Live timedVersionGenerator := utils.NewDefaultTimedVersionGenerator() turnAuthHandler := NewTURNAuthHandler(keyProvider) forwardStats := createForwardStats(conf) - roomManager, err := NewLocalRoomManager(conf, objectStore, currentNode, router, roomAllocator, telemetryService, client, agentStore, rtcEgressLauncher, timedVersionGenerator, turnAuthHandler, messageBus, forwardStats) + roomManager, err := NewLocalRoomManager(conf, objectStore, currentNode, router, roomAllocator, telemetryService, client, agentStore, rtcEgressLauncher, timedVersionGenerator, turnAuthHandler, v, forwardStats) if err != nil { return nil, err } - signalServer, err := NewDefaultSignalServer(currentNode, messageBus, signalRelayConfig, router, roomManager) + signalServer, err := NewDefaultSignalServer(currentNode, v, signalRelayConfig, router, roomManager) if err != nil { return nil, err } - v5 := getTURNAuthHandlerFunc(turnAuthHandler) - server, err := newInProcessTurnServer(conf, v5) + v6 := getTURNAuthHandlerFunc(turnAuthHandler) + server, err := newInProcessTurnServer(conf, v6) if err != nil { return nil, err } @@ -164,14 +164,14 @@ func InitializeRouter(conf *config.Config, currentNode routing.LocalNode) (routi return nil, err } nodeID := getNodeID(currentNode) - messageBus := getMessageBus(universalClient) + psrpcConfig := getPSRPCConfig(conf) + v := getMessageBus(universalClient, psrpcConfig) signalRelayConfig := getSignalRelayConfig(conf) - signalClient, err := routing.NewSignalClient(nodeID, messageBus, signalRelayConfig) + signalClient, err := routing.NewSignalClient(nodeID, v, signalRelayConfig) if err != nil { return nil, err } - psrpcConfig := getPSRPCConfig(conf) - clientParams := getPSRPCClientParams(psrpcConfig, messageBus) + clientParams := getPSRPCClientParams(psrpcConfig, v) roomConfig := getRoomConfig(conf) roomManagerClient, err := routing.NewRoomManagerClient(clientParams, roomConfig) if err != nil { @@ -254,11 +254,12 @@ func createStore(rc redis.UniversalClient) ObjectStore { return NewLocalStore() } -func getMessageBus(rc redis.UniversalClient) psrpc.MessageBus { +func getMessageBus(rc redis.UniversalClient, psrpcConf rpc.PSRPCConfig) psrpc.MessageBus { + opts := psrpcConf.BusOptions() if rc == nil { - return psrpc.NewLocalMessageBus() + return psrpc.NewLocalMessageBus(opts...) } - return psrpc.NewRedisMessageBus(rc) + return psrpc.NewRedisMessageBus(rc, opts...) } func getEgressStore(s ObjectStore) EgressStore {