From 492422cd5359005b5e3b00a1465b7b14b793c0be Mon Sep 17 00:00:00 2001 From: boks1971 Date: Sun, 14 Jun 2026 01:06:25 +0530 Subject: [PATCH] data blob --- go.mod | 4 +- go.sum | 37 +- pkg/config/config.go | 34 +- pkg/rtc/participant.go | 6 +- pkg/rtc/participant_async_attributes.go | 107 ----- .../participant_async_attributes_handler.go | 127 ------ ...rticipant_async_attributes_handler_test.go | 324 --------------- pkg/rtc/participant_async_attributes_test.go | 218 ---------- pkg/rtc/participant_data_blob.go | 91 +++++ pkg/rtc/participant_data_blob_handler.go | 111 +++++ pkg/rtc/participant_data_blob_handler_test.go | 291 +++++++++++++ pkg/rtc/participant_data_blob_test.go | 175 ++++++++ pkg/rtc/participant_signal.go | 6 +- pkg/rtc/room.go | 10 +- pkg/rtc/signalling/interfaces.go | 2 +- pkg/rtc/signalling/signalhandler.go | 8 +- pkg/rtc/signalling/signalling.go | 6 +- pkg/rtc/signalling/signallingunimplemented.go | 2 +- pkg/rtc/types/interfaces.go | 20 +- .../typesfakes/fake_local_participant.go | 382 +++++++++--------- .../fake_local_participant_listener.go | 126 +++--- pkg/rtc/types/typesfakes/fake_participant.go | 136 +++---- pkg/service/roommanager.go | 6 +- test/multinode_test.go | 61 +-- test/singlenode_test.go | 145 ++++--- 25 files changed, 1149 insertions(+), 1286 deletions(-) delete mode 100644 pkg/rtc/participant_async_attributes.go delete mode 100644 pkg/rtc/participant_async_attributes_handler.go delete mode 100644 pkg/rtc/participant_async_attributes_handler_test.go delete mode 100644 pkg/rtc/participant_async_attributes_test.go create mode 100644 pkg/rtc/participant_data_blob.go create mode 100644 pkg/rtc/participant_data_blob_handler.go create mode 100644 pkg/rtc/participant_data_blob_handler_test.go create mode 100644 pkg/rtc/participant_data_blob_test.go diff --git a/go.mod b/go.mod index 24632413c..020701222 100644 --- a/go.mod +++ b/go.mod @@ -70,9 +70,9 @@ require ( github.com/distribution/reference v0.6.0 // indirect github.com/fatih/color v1.19.0 // indirect github.com/felixge/httpsnoop v1.0.4 // indirect - github.com/go-jose/go-jose/v3 v3.0.5 // indirect github.com/go-logr/stdr v1.2.2 // indirect github.com/goccy/go-json v0.10.6 // indirect + github.com/golang-jwt/jwt/v5 v5.3.1 // indirect github.com/grpc-ecosystem/grpc-gateway/v2 v2.29.0 // indirect github.com/mattn/go-colorable v0.1.15 // indirect github.com/mattn/go-isatty v0.0.22 // indirect @@ -81,7 +81,7 @@ require ( github.com/olekukonko/cat v0.0.0-20250911104152-50322a0618f6 // indirect github.com/olekukonko/errors v1.3.0 // indirect github.com/olekukonko/ll v0.1.8 // indirect - github.com/puzpuzpuz/xsync/v3 v3.5.1 // indirect + github.com/puzpuzpuz/xsync/v4 v4.5.0 // indirect go.opentelemetry.io/auto/sdk v1.2.1 // indirect go.opentelemetry.io/contrib/instrumentation/net/http/otelhttp v0.69.0 // indirect go.opentelemetry.io/otel v1.44.0 // indirect diff --git a/go.sum b/go.sum index 66a2eb201..fd3b05834 100644 --- a/go.sum +++ b/go.sum @@ -74,8 +74,6 @@ github.com/gammazero/deque v1.2.1 h1:9fnQVFCCZ9/NOc7ccTNqzoKd1tCWOqeI05/lPqFPMGQ github.com/gammazero/deque v1.2.1/go.mod h1:5nSFkzVm+afG9+gy0VIowlqVAW4N8zNcMne+CMQVD2g= github.com/gammazero/workerpool v1.2.1 h1:MEDvUJsNYGuCvl1RwIXNKu2YtQtHqCSF9XWF04N7lqs= github.com/gammazero/workerpool v1.2.1/go.mod h1:E32GVRUanF4d6QtRmdss3AScgaDkIyrvPtgRQUWgmx4= -github.com/go-jose/go-jose/v3 v3.0.5 h1:BLLJWbC4nMZOfuPVxoZIxeYsn6Nl2r1fITaJ78UQlVQ= -github.com/go-jose/go-jose/v3 v3.0.5/go.mod h1:5b+7YgP7ZICgJDBdfjZaIt+H/9L9T/YQrVfLAMboGkQ= github.com/go-logr/logr v1.2.2/go.mod h1:jdQByPbusPIv2/zmleS9BjJVeZ6kBagPoEUsqbVz/1A= github.com/go-logr/logr v1.4.3 h1:CjnDlHq8ikf6E492q6eKboGOC0T8CDaOvkHCIg8idEI= github.com/go-logr/logr v1.4.3/go.mod h1:9T104GzyrTigFIr8wt5mBrctHMim0Nb2HLGrmQ40KvY= @@ -83,6 +81,8 @@ github.com/go-logr/stdr v1.2.2 h1:hSWxHoqTgW2S2qGc0LTAI563KZ5YKYRhT3MFKZMbjag= github.com/go-logr/stdr v1.2.2/go.mod h1:mMo/vtBO5dYbehREoey6XUKy/eSumjCCveDpRre4VKE= github.com/goccy/go-json v0.10.6 h1:p8HrPJzOakx/mn/bQtjgNjdTcN+/S6FcG2CTtQOrHVU= github.com/goccy/go-json v0.10.6/go.mod h1:oq7eo15ShAhp70Anwd5lgX2pLfOS3QCiwU/PULtXL6M= +github.com/golang-jwt/jwt/v5 v5.3.1 h1:kYf81DTWFe7t+1VvL7eS+jKFVWaUnK9cB1qbwn63YCY= +github.com/golang-jwt/jwt/v5 v5.3.1/go.mod h1:fxCRLWMO43lRc8nhHWY6LGqRcf+1gQWArsqaEUEa5bE= github.com/golang/protobuf v1.5.4 h1:i7eJL8qZTpSEXOPTxNKhASYpMn+8e5Q6AdndVa1dWek= github.com/golang/protobuf v1.5.4/go.mod h1:lnTiLA8Wa4RWRcIUkrtSVa5nRhsEGBg48fD6rSs7xps= github.com/google/cel-go v0.28.1 h1:YWIwi77J4xIsYUwAF/iIuS6haffzIHS8yWI8glSbLWM= @@ -276,8 +276,8 @@ github.com/prometheus/common v0.68.1 h1:omjRRl4QP4komogpXuhfeOiisQg7xdy8VM1UY+pS github.com/prometheus/common v0.68.1/go.mod h1:ZzL3f6u94qUxh9p+tJTrF+FvBS1XXbbRAZCQkytAL0Y= github.com/prometheus/procfs v0.20.1 h1:XwbrGOIplXW/AU3YhIhLODXMJYyC1isLFfYCsTEycfc= github.com/prometheus/procfs v0.20.1/go.mod h1:o9EMBZGRyvDrSPH1RqdxhojkuXstoe4UlK79eF5TGGo= -github.com/puzpuzpuz/xsync/v3 v3.5.1 h1:GJYJZwO6IdxN/IKbneznS6yPkVC+c3zyY/j19c++5Fg= -github.com/puzpuzpuz/xsync/v3 v3.5.1/go.mod h1:VjzYrABPabuM4KyBh1Ftq6u8nhwY5tBPKP9jpmh0nnA= +github.com/puzpuzpuz/xsync/v4 v4.5.0 h1:vOSWu6b57/emh+L/Cw0BeQfvxa/cogFywXHeGUxQxAg= +github.com/puzpuzpuz/xsync/v4 v4.5.0/go.mod h1:VJDmTCJMBt8igNxnkQd86r+8KUeN1quSfNKu5bLYFQo= github.com/redis/go-redis/v9 v9.20.0 h1:WnQYxLkgO2xiXTCJY0ldIiI8dNqCDlQAG+AtaH7a2a0= github.com/redis/go-redis/v9 v9.20.0/go.mod h1:v/M13XI1PVCDcm01VtPFOADfZtHf8YW3baQf57KlIkA= github.com/rodaine/protogofakeit v0.1.1 h1:ZKouljuRM3A+TArppfBqnH8tGZHOwM/pjvtXe9DaXH8= @@ -293,7 +293,6 @@ github.com/shoenig/test v1.7.0 h1:eWcHtTXa6QLnBvm0jgEabMRN/uJ4DMV3M8xUGgRkZmk= github.com/shoenig/test v1.7.0/go.mod h1:UxJ6u/x2v/TNs/LoLxBNJRV9DiwBBKYxXSyczsBHFoI= github.com/stretchr/objx v0.1.0/go.mod h1:HFkY916IF+rwdDfMAkV7OtwuqBVzrE8GR6GFx+wExME= github.com/stretchr/testify v1.4.0/go.mod h1:j7eGeouHqKxXV5pUuKE4zz7dFj8WfuZ+81PSLYec5m4= -github.com/stretchr/testify v1.7.0/go.mod h1:6Fq8oRcR53rry900zMqJjRRixrwX3KX962/h/Wwjteg= github.com/stretchr/testify v1.11.1 h1:7s2iGBzp5EwR7/aIZr8ao5+dra3wiQyKjjFuvgVKu7U= github.com/stretchr/testify v1.11.1/go.mod h1:wZwfW3scLgRK+23gO65QZefKpKQRnfz6sD981Nm4B6U= github.com/thoas/go-funk v0.9.3 h1:7+nAEx3kn5ZJcnDm2Bh23N2yOtweO14bi//dvRtgLpw= @@ -310,7 +309,6 @@ github.com/urfave/negroni/v3 v3.1.1 h1:6MS4nG9Jk/UuCACaUlNXCbiKa0ywF9LXz5dGu09v8 github.com/urfave/negroni/v3 v3.1.1/go.mod h1:jWvnX03kcSjDBl/ShB0iHvx5uOs7mAzZXW+JvJ5XYAs= github.com/wlynxg/anet v0.0.5 h1:J3VJGi1gvo0JwZ/P1/Yc/8p63SoW98B5dHkYDmpgvvU= github.com/wlynxg/anet v0.0.5/go.mod h1:eay5PRQr7fIVAMbTbchTnO9gG65Hg/uYGdc7mguHxoA= -github.com/yuin/goldmark v1.4.13/go.mod h1:6yULJ656Px+3vBD8DxQVa3kxgyrAnzto9xy5taEt/CY= github.com/zeebo/assert v1.3.0 h1:g7C04CbJuIDKNPFHmsk4hwZDO5O+kntRxzaUoNXj+IQ= github.com/zeebo/assert v1.3.0/go.mod h1:Pq9JiuJQpG8JLJdtkwrJESF0Foym2/D9XMU5ciN/wJ0= github.com/zeebo/xxh3 v1.1.0 h1:s7DLGDK45Dyfg7++yxI0khrfwq9661w9EN78eP/UZVs= @@ -351,20 +349,15 @@ go.yaml.in/yaml/v3 v3.0.4 h1:tfq32ie2Jv2UxXFdLJdh3jXuOzWiL1fo0bu/FbuKpbc= go.yaml.in/yaml/v3 v3.0.4/go.mod h1:DhzuOOF2ATzADvBadXxruRBLzYTpT36CKvDb3+aBEFg= golang.org/x/crypto v0.0.0-20190308221718-c2843e01d9a2/go.mod h1:djNgcEr1/C05ACkg1iLfiJU5Ep61QUkGW8qpdssI0+w= golang.org/x/crypto v0.0.0-20200622213623-75b288015ac9/go.mod h1:LzIPMQfyMNhhGPhUkYOs5KpL4U8rLKemX1yGLhDgUto= -golang.org/x/crypto v0.0.0-20210921155107-089bfa567519/go.mod h1:GvvjBRRGRdwPK5ydBHafDWAxML/pGHZbMvKqRZ5+Abc= -golang.org/x/crypto v0.19.0/go.mod h1:Iy9bg/ha4yyC70EfRS8jz+B6ybOBKMaSxLj6P6oBDfU= golang.org/x/crypto v0.52.0 h1:RMs7fP2rXdep0CftQlK8Uf+kibLm7qkCcradZWYz988= golang.org/x/crypto v0.52.0/go.mod h1:1QgfPxDqh0T2M/elOJtp9RvuR95kVjir0e6/BvEmGbc= golang.org/x/exp v0.0.0-20260603202125-055de637280b h1:v1uXiEBHo8QA0LiGCo7UgHMzHT4Kdfpl2zmtH5vaP1Q= golang.org/x/exp v0.0.0-20260603202125-055de637280b/go.mod h1:d2fgXJLVs4dYDHUk5lwMIfzRzSrWCfGZb0ZqeLa/Vcw= -golang.org/x/mod v0.6.0-dev.0.20220419223038-86c51ed26bb4/go.mod h1:jJ57K6gSWd91VN4djpZkiMVwK6gcyfeH4XE8wZrZaV4= -golang.org/x/mod v0.8.0/go.mod h1:iBbtSCu2XBx23ZKBPSOrRkjjQPZFPuis4dIYUhu/chs= golang.org/x/mod v0.36.0 h1:JJjpVx6myfUsUdAzZuOSTTmRE0PfZeNWzzvKrP7amb4= golang.org/x/mod v0.36.0/go.mod h1:moc6ELqsWcOw5Ef3xVprK5ul/MvtVvkIXLziUOICjUQ= golang.org/x/net v0.0.0-20190311183353-d8887717615a/go.mod h1:t9HGtf8HONx5eT2rtn7q6eTqICYqUVnKs3thJo3Qplg= golang.org/x/net v0.0.0-20190404232315-eb5bcb51f2a3/go.mod h1:t9HGtf8HONx5eT2rtn7q6eTqICYqUVnKs3thJo3Qplg= golang.org/x/net v0.0.0-20190503192946-f4e77d36d62c/go.mod h1:t9HGtf8HONx5eT2rtn7q6eTqICYqUVnKs3thJo3Qplg= -golang.org/x/net v0.0.0-20190620200207-3b0461eec859/go.mod h1:z5CRVTTTmAJ677TzLLGU+0bjPO0LkuOLi4/5GtJWs/s= golang.org/x/net v0.0.0-20190827160401-ba9fcec4b297/go.mod h1:z5CRVTTTmAJ677TzLLGU+0bjPO0LkuOLi4/5GtJWs/s= golang.org/x/net v0.0.0-20191007182048-72f939374954/go.mod h1:z5CRVTTTmAJ677TzLLGU+0bjPO0LkuOLi4/5GtJWs/s= golang.org/x/net v0.0.0-20200202094626-16171245cfb2/go.mod h1:z5CRVTTTmAJ677TzLLGU+0bjPO0LkuOLi4/5GtJWs/s= @@ -375,17 +368,11 @@ golang.org/x/net v0.0.0-20201224014010-6772e930b67b/go.mod h1:m0MpNAwzfU5UDzcl9v golang.org/x/net v0.0.0-20210119194325-5f4716e94777/go.mod h1:m0MpNAwzfU5UDzcl9v0D8zg8gWTRqZa9RBIspLL5mdg= golang.org/x/net v0.0.0-20210226172049-e18ecbb05110/go.mod h1:m0MpNAwzfU5UDzcl9v0D8zg8gWTRqZa9RBIspLL5mdg= golang.org/x/net v0.0.0-20210525063256-abc453219eb5/go.mod h1:9nx3DQGgdP8bBQD5qxJ1jj9UTztislL4KSBs9R2vV5Y= -golang.org/x/net v0.0.0-20220722155237-a158d28d115b/go.mod h1:XRhObCWvk6IyKnWLug+ECip1KBveYUHfp+8e9klMJ9c= golang.org/x/net v0.0.0-20220923203811-8be639271d50/go.mod h1:YDH+HFinaLZZlnHAfSS6ZXJJ9M9t4Dl22yv3iI2vPwk= -golang.org/x/net v0.6.0/go.mod h1:2Tu9+aMcznHK/AK1HMvgo6xiTLG5rD5rZLDS+rp2Bjs= -golang.org/x/net v0.10.0/go.mod h1:0qNGK6F8kojg2nk9dLZ2mShWaEBan6FAoqfSigmmuDg= golang.org/x/net v0.55.0 h1:bcvxaJn3e1U6InsFWt1JUq1aSjnRxLzT2rtD2KfkDF8= golang.org/x/net v0.55.0/go.mod h1:L5U2KuzuOe1lY7Z+aWVIKK6qEeJXnXV9yzGA+WCHJww= -golang.org/x/sync v0.0.0-20190423024810-112230192c58/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM= golang.org/x/sync v0.0.0-20210220032951-036812b2e83c/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM= -golang.org/x/sync v0.0.0-20220722155255-886fb9371eb4/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM= golang.org/x/sync v0.0.0-20220923202941-7f9b1623fab7/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM= -golang.org/x/sync v0.1.0/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM= golang.org/x/sync v0.20.0 h1:e0PTpb7pjO8GAtTs2dQ6jYa5BWYlMuX047Dco/pItO4= golang.org/x/sync v0.20.0/go.mod h1:9xrNwdLfx4jkKbNva9FpL6vEN7evnE43NNNJQ2LF3+0= golang.org/x/sys v0.0.0-20190215142949-d0b11bdaac8a/go.mod h1:STP8DvDyc/dI5b8T5hshtkjS+E42TnysNCUPdjciGhY= @@ -410,37 +397,22 @@ golang.org/x/sys v0.0.0-20210525143221-35b2ab0089ea/go.mod h1:oPkhp1MJrh7nUepCBc golang.org/x/sys v0.0.0-20210615035016-665e8c7367d1/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= golang.org/x/sys v0.0.0-20210906170528-6f6e22806c34/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= golang.org/x/sys v0.0.0-20220319134239-a9b59b0215f8/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= -golang.org/x/sys v0.0.0-20220520151302-bc2c85ada10a/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= -golang.org/x/sys v0.0.0-20220722155257-8c9f86f7a55f/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= golang.org/x/sys v0.0.0-20220728004956-3c1f35247d10/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= -golang.org/x/sys v0.5.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= -golang.org/x/sys v0.8.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= -golang.org/x/sys v0.17.0/go.mod h1:/VUhepiaJMQUp4+oa/7Zr1D23ma6VTLIYjOOTFZPUcA= golang.org/x/sys v0.45.0 h1:dO4czNzziLiiXplLQgBCEpCvXQ3dnkn0SdaZSYdQ+FY= golang.org/x/sys v0.45.0/go.mod h1:4GL1E5IUh+htKOUEOaiffhrAeqysfVGipDYzABqnCmw= golang.org/x/term v0.0.0-20201126162022-7de9c90e9dd1/go.mod h1:bj7SfCRtBDWHUb9snDiAeCFNEtKQo2Wmx5Cou7ajbmo= golang.org/x/term v0.0.0-20210927222741-03fcf44c2211/go.mod h1:jbD1KX2456YbFQfuXm/mYQcufACuNUgVhRMnK/tPxf8= -golang.org/x/term v0.5.0/go.mod h1:jMB1sMXY+tzblOD4FWmEbocvup2/aLOaQEp7JmGp78k= -golang.org/x/term v0.8.0/go.mod h1:xPskH00ivmX89bAKVGSKKtLOWNx2+17Eiy94tnKShWo= -golang.org/x/term v0.17.0/go.mod h1:lLRBjIVuehSbZlaOtGMbcMncT+aqLLLmKrsjNrUguwk= golang.org/x/text v0.3.0/go.mod h1:NqM8EUOU14njkJ3fqMW+pc6Ldnwhi/IjpwHt7yyuwOQ= golang.org/x/text v0.3.3/go.mod h1:5Zoc/QRtKVWzQhOtBMvqHzDpF6irO9z98xDceosuGiQ= golang.org/x/text v0.3.6/go.mod h1:5Zoc/QRtKVWzQhOtBMvqHzDpF6irO9z98xDceosuGiQ= golang.org/x/text v0.3.7/go.mod h1:u+2+/6zg+i71rQMx5EYifcz6MCKuco9NR6JIITiCfzQ= -golang.org/x/text v0.7.0/go.mod h1:mrYo+phRRbMaCq/xk9113O4dZlRixOauAjOtrjsXDZ8= -golang.org/x/text v0.9.0/go.mod h1:e1OnstbJyHTd6l/uOt8jFFHp6TRDWZR/bV3emEE/zU8= -golang.org/x/text v0.14.0/go.mod h1:18ZOQIKpY8NJVqYksKHtTdi31H5itFRjB5/qKTNYzSU= golang.org/x/text v0.37.0 h1:Cqjiwd9eSg8e0QAkyCaQTNHFIIzWtidPahFWR83rTrc= golang.org/x/text v0.37.0/go.mod h1:a5sjxXGs9hsn/AJVwuElvCAo9v8QYLzvavO5z2PiM38= golang.org/x/time v0.15.0 h1:bbrp8t3bGUeFOx08pvsMYRTCVSMk89u4tKbNOZbp88U= golang.org/x/time v0.15.0/go.mod h1:Y4YMaQmXwGQZoFaVFk4YpCt4FLQMYKZe9oeV/f4MSno= golang.org/x/tools v0.0.0-20180917221912-90fa682c2a6e/go.mod h1:n7NCudcB/nEzxVGmLbDWY5pfWTLqBcC2KZ6jyYvM4mQ= -golang.org/x/tools v0.0.0-20191119224855-298f0cb1881e/go.mod h1:b+2E5dAYhXwXZwtnZ6UAqBI28+e2cm9otk0dWdXHAEo= -golang.org/x/tools v0.1.12/go.mod h1:hNGJHUnrk76NpqgfD5Aqm5Crs+Hm0VOH/i9J2+nxYbc= -golang.org/x/tools v0.6.0/go.mod h1:Xwgl3UAJ/d3gWutnCtw505GrjyAbvKui8lOU390QaIU= golang.org/x/tools v0.45.0 h1:18qN3FAooORvApf5XjCXgsuayZOEtXf6JK18I3+ONa8= golang.org/x/tools v0.45.0/go.mod h1:LuUGqqaXcXMEFEruIVJVm5mgDD8vww/z/SR1gQ4uE/0= -golang.org/x/xerrors v0.0.0-20190717185122-a985d3407aa7/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0= golang.org/x/xerrors v0.0.0-20191204190536-9bdfabe68543/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0= golang.org/x/xerrors v0.0.0-20200804184101-5ec99f83aff1/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0= gonum.org/v1/gonum v0.17.0 h1:VbpOemQlsSMrYmn7T2OUvQ4dqxQXU+ouZFQsZOx50z4= @@ -459,7 +431,6 @@ gopkg.in/check.v1 v1.0.0-20201130134442-10cb98267c6c h1:Hei/4ADfdWqJk1ZMxUNpqntN gopkg.in/check.v1 v1.0.0-20201130134442-10cb98267c6c/go.mod h1:JHkPIbrfpd72SG/EVd6muEfDQjcINNoR0C8j2r3qZ4Q= gopkg.in/errgo.v2 v2.1.0/go.mod h1:hNsd1EY+bozCKY1Ytp96fpM3vjJbqLJn88ws8XvfDNI= gopkg.in/yaml.v2 v2.2.2/go.mod h1:hI93XBmqTisBFMUTm0b8Fm+jr3Dg1NNxqwp+5A1VGuI= -gopkg.in/yaml.v3 v3.0.0-20200313102051-9f266ea9e77c/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA= gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= gotest.tools/v3 v3.5.2 h1:7koQfIKdy+I8UTetycgUqXWSDwpgv193Ka+qRsmBY8Q= diff --git a/pkg/config/config.go b/pkg/config/config.go index 05d1dbe30..80a5ece55 100644 --- a/pkg/config/config.go +++ b/pkg/config/config.go @@ -91,7 +91,7 @@ type Config struct { EnableDataTracks bool `yaml:"enable_data_tracks,omitempty"` - EnableParticipantAsyncAttributes bool `yaml:"enable_participant_async_attributes,omitempty"` + EnableParticipantDataBlob bool `yaml:"enable_participant_data_blob,omitempty"` API APIConfig `yaml:"api,omitempty"` } @@ -297,8 +297,8 @@ type LimitConfig struct { MaxParticipantIdentityLength int `yaml:"max_participant_identity_length,omitempty"` MaxParticipantNameLength int `yaml:"max_participant_name_length,omitempty"` - MaxAsyncAttributeNameLength int `yaml:"max_async_attributes_name_length,omitempty"` - MaxAsyncAttributesSize uint32 `yaml:"max_async_attributes_size,omitempty"` + MaxDataBlobKeyLength int `yaml:"max_data_blob_key_length,omitempty"` + MaxDataBlobSize uint32 `yaml:"max_data_blobs_size,omitempty"` } func (l LimitConfig) CheckRoomNameLength(name string) bool { @@ -329,32 +329,32 @@ func (l LimitConfig) CheckAttributesSize(attributes map[string]string) bool { return uint32(total) <= l.MaxAttributesSize } -func (l LimitConfig) CheckAsyncAttributeNameLength(name string) bool { - return l.MaxAsyncAttributeNameLength == 0 || len(name) <= l.MaxAsyncAttributeNameLength +func (l LimitConfig) CheckDataBlobKeyLength(key string) bool { + return l.MaxDataBlobKeyLength == 0 || len(key) <= l.MaxDataBlobKeyLength } -func (l LimitConfig) CheckAsyncAttributesSize(asyncAttributes []*livekit.DataTrackSchemaDefinition) bool { - if l.MaxAsyncAttributesSize == 0 { +func (l LimitConfig) CheckDataBlobsSize(dataBlobs []*livekit.DataBlob) bool { + if l.MaxDataBlobSize == 0 { return true } total := 0 - for _, asyncAttribute := range asyncAttributes { - total += len(asyncAttribute.Id.Name) + len(asyncAttribute.Definition) + for _, dataBlob := range dataBlobs { + total += len(dataBlob.GetKey().String()) + len(dataBlob.Contents) } - return uint32(total) <= l.MaxAsyncAttributesSize + return uint32(total) <= l.MaxDataBlobSize } -func (l LimitConfig) CanAddAsyncAttribute(asyncAttributes []*livekit.DataTrackSchemaDefinition, toAdd *livekit.DataTrackSchemaDefinition) bool { - if l.MaxAsyncAttributesSize == 0 { +func (l LimitConfig) CanAddDataBlob(dataBlobs []*livekit.DataBlob, toAdd *livekit.DataBlob) bool { + if l.MaxDataBlobSize == 0 { return true } total := 0 - for _, asyncAttribute := range asyncAttributes { - total += len(asyncAttribute.Id.Name) + len(asyncAttribute.Definition) + for _, dataBlob := range dataBlobs { + total += len(dataBlob.Key.String()) + len(dataBlob.Contents) } - return uint32(total+len(toAdd.Id.Name)+len(toAdd.Definition)) <= l.MaxAsyncAttributesSize + return uint32(total+len(toAdd.GetKey().String())+len(toAdd.Contents)) <= l.MaxDataBlobSize } // --------------------------------- @@ -482,8 +482,8 @@ var DefaultConfig = Config{ MaxRoomNameLength: 256, MaxParticipantIdentityLength: 256, MaxParticipantNameLength: 256, - MaxAsyncAttributeNameLength: 256, - MaxAsyncAttributesSize: 256000, + MaxDataBlobKeyLength: 256, + MaxDataBlobSize: 64000, }, Logging: LoggingConfig{ PionLevel: "error", diff --git a/pkg/rtc/participant.go b/pkg/rtc/participant.go index 57312c034..32035dfb5 100644 --- a/pkg/rtc/participant.go +++ b/pkg/rtc/participant.go @@ -225,7 +225,7 @@ type ParticipantParams struct { EnableRTPStreamRestartDetection bool ForceBackupCodecPolicySimulcast bool DisableTransceiverReuseForE2EE bool - EnableParticipantDataBlobs bool + EnableParticipantDataBlob bool } type ParticipantImpl struct { @@ -335,7 +335,7 @@ type ParticipantImpl struct { rpcPendingAcks map[string]*utils.DataChannelRpcPendingAckHandler rpcPendingResponses map[string]*utils.DataChannelRpcPendingResponseHandler - asyncAttributes *ParticipantAsyncAttributes + dataBlob *ParticipantDataBlob } func NewParticipant(params ParticipantParams) (*ParticipantImpl, error) { @@ -374,7 +374,7 @@ func NewParticipant(params ParticipantParams) (*ParticipantImpl, error) { telemetryGuard: &telemetry.ReferenceGuard{}, nextSubscribedDataTrackHandle: uint16(rand.Intn(256)), requireBroadcast: params.Grants.Metadata != "" || len(params.Grants.Attributes) != 0, - asyncAttributes: NewParticipantAsyncAttributes(ParticipantAsyncAttributesParams{ + dataBlob: NewParticipantDataBlob(ParticipantDataBlobParams{ Logger: params.Logger, }), } diff --git a/pkg/rtc/participant_async_attributes.go b/pkg/rtc/participant_async_attributes.go deleted file mode 100644 index 46e7bdfe9..000000000 --- a/pkg/rtc/participant_async_attributes.go +++ /dev/null @@ -1,107 +0,0 @@ -// Copyright 2026 LiveKit, Inc. -// -// Licensed under the Apache License, Version 2.0 (the "License"); -// you may not use this file except in compliance with the License. -// You may obtain a copy of the License at -// -// http://www.apache.org/licenses/LICENSE-2.0 -// -// Unless required by applicable law or agreed to in writing, software -// distributed under the License is distributed on an "AS IS" BASIS, -// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. -// See the License for the specific language governing permissions and -// limitations under the License. - -package rtc - -import ( - "fmt" - "strconv" - "strings" - "sync" - - "github.com/livekit/protocol/livekit" - "github.com/livekit/protocol/logger" - "github.com/livekit/protocol/utils" -) - -type ParticipantAsyncAttributesParams struct { - Logger logger.Logger -} - -type ParticipantAsyncAttributes struct { - params ParticipantAsyncAttributesParams - lock sync.Mutex - attributes map[string]*livekit.DataTrackSchemaDefinition -} - -func NewParticipantAsyncAttributes(params ParticipantAsyncAttributesParams) *ParticipantAsyncAttributes { - return &ParticipantAsyncAttributes{ - params: params, - attributes: make(map[string]*livekit.DataTrackSchemaDefinition), - } -} - -func (p *ParticipantAsyncAttributes) Add(aa *livekit.DataTrackSchemaDefinition) { - p.lock.Lock() - defer p.lock.Unlock() - - if aa.Id == nil { - return - } - - p.attributes[ToParticipantAsyncAttributeKey(aa.Id)] = utils.CloneProto(aa) -} - -func (p *ParticipantAsyncAttributes) Delete(id *livekit.DataTrackSchemaId) { - p.lock.Lock() - defer p.lock.Unlock() - - if id == nil { - return - } - - delete(p.attributes, ToParticipantAsyncAttributeKey(id)) -} - -func (p *ParticipantAsyncAttributes) Get(id *livekit.DataTrackSchemaId) *livekit.DataTrackSchemaDefinition { - p.lock.Lock() - defer p.lock.Unlock() - - if id == nil { - return nil - } - - aa, ok := p.attributes[ToParticipantAsyncAttributeKey(id)] - if !ok { - return nil - } - - return aa -} - -func (p *ParticipantAsyncAttributes) GetAll() []*livekit.DataTrackSchemaDefinition { - p.lock.Lock() - defer p.lock.Unlock() - - all := make([]*livekit.DataTrackSchemaDefinition, 0, len(p.attributes)) - for _, aa := range p.attributes { - all = append(all, utils.CloneProto(aa)) - } - return all -} - -// ------------------------------- - -func ToParticipantAsyncAttributeKey(id *livekit.DataTrackSchemaId) string { - return fmt.Sprintf("%s/%d", id.Name, id.Encoding) -} - -func fromParticipantAsyncAttributeKey(key string) *livekit.DataTrackSchemaId { - parts := strings.Split(key, "/") - encoding, _ := strconv.Atoi(parts[1]) - return &livekit.DataTrackSchemaId{ - Name: parts[0], - Encoding: livekit.DataTrackSchemaEncoding(encoding), - } -} diff --git a/pkg/rtc/participant_async_attributes_handler.go b/pkg/rtc/participant_async_attributes_handler.go deleted file mode 100644 index c0ce16113..000000000 --- a/pkg/rtc/participant_async_attributes_handler.go +++ /dev/null @@ -1,127 +0,0 @@ -// Copyright 2026 LiveKit, Inc. -// -// Licensed under the Apache License, Version 2.0 (the "License"); -// you may not use this file except in compliance with the License. -// You may obtain a copy of the License at -// -// http://www.apache.org/licenses/LICENSE-2.0 -// -// Unless required by applicable law or agreed to in writing, software -// distributed under the License is distributed on an "AS IS" BASIS, -// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. -// See the License for the specific language governing permissions and -// limitations under the License. - -package rtc - -import ( - "github.com/livekit/livekit-server/pkg/rtc/types" - "github.com/livekit/protocol/livekit" - "github.com/livekit/protocol/logger" - "github.com/livekit/protocol/utils" -) - -func (p *ParticipantImpl) HandleDefineDataTrackSchemaRequest(req *livekit.DefineDataTrackSchemaRequest) { - if !p.params.EnableParticipantAsyncAttributes { - p.pubLogger.Warnw("async attributes not enabled", nil, "req", logger.Proto(req)) - p.sendRequestResponse(&livekit.RequestResponse{ - Reason: livekit.RequestResponse_NOT_ALLOWED, - Message: "async attributes not enabled", - Request: &livekit.RequestResponse_DefineDataTrackSchema{ - DefineDataTrackSchema: utils.CloneProto(req), - }, - }) - return - } - - if req.SchemaDefinition == nil || req.SchemaDefinition.Id == nil || len(req.SchemaDefinition.Id.Name) == 0 || !p.params.LimitConfig.CheckAsyncAttributeNameLength(req.SchemaDefinition.Id.Name) { - p.pubLogger.Warnw("aync attribute definition is invalid", nil, "req", logger.Proto(req)) - p.sendRequestResponse(&livekit.RequestResponse{ - Reason: livekit.RequestResponse_INVALID_REQUEST, - Message: "async attribute definition is invalid", - Request: &livekit.RequestResponse_DefineDataTrackSchema{ - DefineDataTrackSchema: utils.CloneProto(req), - }, - }) - return - } - - if len(req.SchemaDefinition.Definition) == 0 { - p.sendRequestResponse(&livekit.RequestResponse{ - Reason: livekit.RequestResponse_INVALID_REQUEST, - Message: "async attribute definition is empty", - Request: &livekit.RequestResponse_DefineDataTrackSchema{ - DefineDataTrackSchema: utils.CloneProto(req), - }, - }) - return - } - - if !p.params.LimitConfig.CanAddAsyncAttribute(p.asyncAttributes.GetAll(), req.SchemaDefinition) { - p.sendRequestResponse(&livekit.RequestResponse{ - Reason: livekit.RequestResponse_LIMIT_EXCEEDED, - Message: "async attribute definition exceeds limit", - Request: &livekit.RequestResponse_DefineDataTrackSchema{ - DefineDataTrackSchema: utils.CloneProto(req), - }, - }) - return - } - - p.AddDataTrackSchema(req.SchemaDefinition) - p.listener().OnDefineDataTrackSchema(p, req.SchemaDefinition) -} - -func (p *ParticipantImpl) HandleGetDataTrackSchemaRequest(req *livekit.GetDataTrackSchemaRequest) { - if req.SchemaId == nil { - p.sendRequestResponse(&livekit.RequestResponse{ - Reason: livekit.RequestResponse_INVALID_REQUEST, - Message: "async attribute id is required", - Request: &livekit.RequestResponse_GetDataTrackSchema{ - GetDataTrackSchema: utils.CloneProto(req), - }, - }) - return - } - - p.listener().OnGetDataTrackSchema(p, req) -} - -func (p *ParticipantImpl) AddDataTrackSchema(definition *livekit.DataTrackSchemaDefinition) { - p.asyncAttributes.Add(definition) -} - -func (p *ParticipantImpl) GetDataTrackSchema(id *livekit.DataTrackSchemaId) *livekit.DataTrackSchemaDefinition { - return p.asyncAttributes.Get(id) -} - -func (p *ParticipantImpl) ProcessGetDataTrackSchemaRequest(req *livekit.GetDataTrackSchemaRequest, publisher types.Participant) { - if publisher == nil { - p.sendRequestResponse(&livekit.RequestResponse{ - Reason: livekit.RequestResponse_NOT_FOUND, - Message: "participant not found", - Request: &livekit.RequestResponse_GetDataTrackSchema{ - GetDataTrackSchema: utils.CloneProto(req), - }, - }) - return - } - - asyncAttribute := publisher.GetDataTrackSchema(req.SchemaId) - if asyncAttribute == nil { - p.sendRequestResponse(&livekit.RequestResponse{ - Reason: livekit.RequestResponse_NOT_FOUND, - Message: "async attribute not found", - Request: &livekit.RequestResponse_GetDataTrackSchema{ - GetDataTrackSchema: utils.CloneProto(req), - }, - }) - return - } - - p.sendGetDataTrackSchemaResponse(asyncAttribute) -} - -func (p *ParticipantImpl) GetAllAsyncAttributes() []*livekit.DataTrackSchemaDefinition { - return p.asyncAttributes.GetAll() -} diff --git a/pkg/rtc/participant_async_attributes_handler_test.go b/pkg/rtc/participant_async_attributes_handler_test.go deleted file mode 100644 index 173d98eb4..000000000 --- a/pkg/rtc/participant_async_attributes_handler_test.go +++ /dev/null @@ -1,324 +0,0 @@ -// Copyright 2026 LiveKit, Inc. -// -// Licensed under the Apache License, Version 2.0 (the "License"); -// you may not use this file except in compliance with the License. -// You may obtain a copy of the License at -// -// http://www.apache.org/licenses/LICENSE-2.0 -// -// Unless required by applicable law or agreed to in writing, software -// distributed under the License is distributed on an "AS IS" BASIS, -// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. -// See the License for the specific language governing permissions and -// limitations under the License. - -package rtc - -import ( - "strings" - "testing" - - "github.com/stretchr/testify/require" - - "github.com/livekit/protocol/livekit" - - "github.com/livekit/livekit-server/pkg/config" - "github.com/livekit/livekit-server/pkg/routing/routingfakes" - "github.com/livekit/livekit-server/pkg/rtc/types/typesfakes" -) - -func newParticipantWithAsyncAttributes(t *testing.T, enabled bool, maxNameLength int, maxSize uint32) *ParticipantImpl { - t.Helper() - p := newParticipantForTest("test") - p.params.EnableParticipantAsyncAttributes = enabled - p.params.LimitConfig = config.LimitConfig{ - MaxAsyncAttributeNameLength: maxNameLength, - MaxAsyncAttributesSize: maxSize, - } - return p -} - -func lastRequestResponse(t *testing.T, sink *routingfakes.FakeMessageSink, idx int) *livekit.RequestResponse { - t.Helper() - msg := sink.WriteMessageArgsForCall(idx).(*livekit.SignalResponse) - rr, ok := msg.Message.(*livekit.SignalResponse_RequestResponse) - require.True(t, ok, "expected SignalResponse_RequestResponse, got %T", msg.Message) - return rr.RequestResponse -} - -func TestHandleDefineDataTrackSchemaRequest(t *testing.T) { - t.Run("returns NOT_ALLOWED when feature not enabled", func(t *testing.T) { - p := newParticipantWithAsyncAttributes(t, false, 0, 0) - sink := p.params.Sink.(*routingfakes.FakeMessageSink) - - req := &livekit.DefineDataTrackSchemaRequest{ - SchemaDefinition: &livekit.DataTrackSchemaDefinition{ - Id: &livekit.DataTrackSchemaId{ - Name: "schema-1", - Encoding: livekit.DataTrackSchemaEncoding_DATA_TRACK_SCHEMA_ENCODING_PROTOBUF, - }, - Definition: []byte("def"), - }, - } - p.HandleDefineDataTrackSchemaRequest(req) - - require.Equal(t, 1, sink.WriteMessageCallCount()) - rr := lastRequestResponse(t, sink, 0) - require.Equal(t, livekit.RequestResponse_NOT_ALLOWED, rr.Reason) - require.NotNil(t, rr.GetDefineDataTrackSchema()) - require.Empty(t, p.asyncAttributes.GetAll()) - }) - - t.Run("returns INVALID_REQUEST when schema definition is nil", func(t *testing.T) { - p := newParticipantWithAsyncAttributes(t, true, 0, 0) - sink := p.params.Sink.(*routingfakes.FakeMessageSink) - - p.HandleDefineDataTrackSchemaRequest(&livekit.DefineDataTrackSchemaRequest{}) - - require.Equal(t, 1, sink.WriteMessageCallCount()) - rr := lastRequestResponse(t, sink, 0) - require.Equal(t, livekit.RequestResponse_INVALID_REQUEST, rr.Reason) - require.Empty(t, p.asyncAttributes.GetAll()) - }) - - t.Run("returns INVALID_REQUEST when id is nil", func(t *testing.T) { - p := newParticipantWithAsyncAttributes(t, true, 0, 0) - sink := p.params.Sink.(*routingfakes.FakeMessageSink) - - p.HandleDefineDataTrackSchemaRequest(&livekit.DefineDataTrackSchemaRequest{ - SchemaDefinition: &livekit.DataTrackSchemaDefinition{ - Definition: []byte("def"), - }, - }) - - require.Equal(t, 1, sink.WriteMessageCallCount()) - rr := lastRequestResponse(t, sink, 0) - require.Equal(t, livekit.RequestResponse_INVALID_REQUEST, rr.Reason) - }) - - t.Run("returns INVALID_REQUEST when name is empty", func(t *testing.T) { - p := newParticipantWithAsyncAttributes(t, true, 0, 0) - sink := p.params.Sink.(*routingfakes.FakeMessageSink) - - p.HandleDefineDataTrackSchemaRequest(&livekit.DefineDataTrackSchemaRequest{ - SchemaDefinition: &livekit.DataTrackSchemaDefinition{ - Id: &livekit.DataTrackSchemaId{ - Name: "", - Encoding: livekit.DataTrackSchemaEncoding_DATA_TRACK_SCHEMA_ENCODING_PROTOBUF, - }, - Definition: []byte("def"), - }, - }) - - require.Equal(t, 1, sink.WriteMessageCallCount()) - rr := lastRequestResponse(t, sink, 0) - require.Equal(t, livekit.RequestResponse_INVALID_REQUEST, rr.Reason) - }) - - t.Run("returns INVALID_REQUEST when name exceeds 256 chars", func(t *testing.T) { - p := newParticipantWithAsyncAttributes(t, true, 256, 0) - sink := p.params.Sink.(*routingfakes.FakeMessageSink) - - p.HandleDefineDataTrackSchemaRequest(&livekit.DefineDataTrackSchemaRequest{ - SchemaDefinition: &livekit.DataTrackSchemaDefinition{ - Id: &livekit.DataTrackSchemaId{ - Name: strings.Repeat("a", 257), - Encoding: livekit.DataTrackSchemaEncoding_DATA_TRACK_SCHEMA_ENCODING_PROTOBUF, - }, - Definition: []byte("def"), - }, - }) - - require.Equal(t, 1, sink.WriteMessageCallCount()) - rr := lastRequestResponse(t, sink, 0) - require.Equal(t, livekit.RequestResponse_INVALID_REQUEST, rr.Reason) - }) - - t.Run("returns INVALID_REQUEST when definition is empty", func(t *testing.T) { - p := newParticipantWithAsyncAttributes(t, true, 0, 0) - sink := p.params.Sink.(*routingfakes.FakeMessageSink) - - p.HandleDefineDataTrackSchemaRequest(&livekit.DefineDataTrackSchemaRequest{ - SchemaDefinition: &livekit.DataTrackSchemaDefinition{ - Id: &livekit.DataTrackSchemaId{ - Name: "schema-1", - Encoding: livekit.DataTrackSchemaEncoding_DATA_TRACK_SCHEMA_ENCODING_PROTOBUF, - }, - Definition: nil, - }, - }) - - require.Equal(t, 1, sink.WriteMessageCallCount()) - rr := lastRequestResponse(t, sink, 0) - require.Equal(t, livekit.RequestResponse_INVALID_REQUEST, rr.Reason) - require.Empty(t, p.asyncAttributes.GetAll()) - }) - - t.Run("returns LIMIT_EXCEEDED when adding would breach the limit", func(t *testing.T) { - p := newParticipantWithAsyncAttributes(t, true, 0, 16) - sink := p.params.Sink.(*routingfakes.FakeMessageSink) - - p.HandleDefineDataTrackSchemaRequest(&livekit.DefineDataTrackSchemaRequest{ - SchemaDefinition: &livekit.DataTrackSchemaDefinition{ - Id: &livekit.DataTrackSchemaId{ - Name: "schema-1", - Encoding: livekit.DataTrackSchemaEncoding_DATA_TRACK_SCHEMA_ENCODING_PROTOBUF, - }, - Definition: []byte(strings.Repeat("x", 32)), - }, - }) - - require.Equal(t, 1, sink.WriteMessageCallCount()) - rr := lastRequestResponse(t, sink, 0) - require.Equal(t, livekit.RequestResponse_LIMIT_EXCEEDED, rr.Reason) - require.Empty(t, p.asyncAttributes.GetAll()) - }) - - t.Run("stores a valid definition and sends no response", func(t *testing.T) { - p := newParticipantWithAsyncAttributes(t, true, 0, 0) - sink := p.params.Sink.(*routingfakes.FakeMessageSink) - - id := &livekit.DataTrackSchemaId{ - Name: "schema-1", - Encoding: livekit.DataTrackSchemaEncoding_DATA_TRACK_SCHEMA_ENCODING_PROTOBUF, - } - def := []byte("definition-bytes") - - p.HandleDefineDataTrackSchemaRequest(&livekit.DefineDataTrackSchemaRequest{ - SchemaDefinition: &livekit.DataTrackSchemaDefinition{ - Id: id, - Definition: def, - }, - }) - - // on success no response is sent to the client - require.Equal(t, 0, sink.WriteMessageCallCount()) - - stored := p.asyncAttributes.Get(id) - require.NotNil(t, stored) - require.Equal(t, def, stored.Definition) - }) -} - -func TestHandleGetDataTrackSchemaRequest(t *testing.T) { - t.Run("returns INVALID_REQUEST when schema id is missing", func(t *testing.T) { - p := newParticipantWithAsyncAttributes(t, true, 0, 0) - sink := p.params.Sink.(*routingfakes.FakeMessageSink) - - p.HandleGetDataTrackSchemaRequest(&livekit.GetDataTrackSchemaRequest{ - ParticipantIdentity: "other", - }) - - require.Equal(t, 1, sink.WriteMessageCallCount()) - rr := lastRequestResponse(t, sink, 0) - require.Equal(t, livekit.RequestResponse_INVALID_REQUEST, rr.Reason) - }) - - t.Run("forwards request to listener when schema id is provided", func(t *testing.T) { - p := newParticipantWithAsyncAttributes(t, true, 0, 0) - listener := p.params.ParticipantListener.(*typesfakes.FakeLocalParticipantListener) - - req := &livekit.GetDataTrackSchemaRequest{ - ParticipantIdentity: "other", - SchemaId: &livekit.DataTrackSchemaId{ - Name: "schema-1", - Encoding: livekit.DataTrackSchemaEncoding_DATA_TRACK_SCHEMA_ENCODING_PROTOBUF, - }, - } - p.HandleGetDataTrackSchemaRequest(req) - - require.Equal(t, 1, listener.OnGetDataTrackSchemaCallCount()) - gotParticipant, gotReq := listener.OnGetDataTrackSchemaArgsForCall(0) - require.Equal(t, p, gotParticipant) - require.Equal(t, req, gotReq) - }) -} - -func TestGetDataTrackSchema(t *testing.T) { - p := newParticipantWithAsyncAttributes(t, true, 0, 0) - - id := &livekit.DataTrackSchemaId{ - Name: "schema-1", - Encoding: livekit.DataTrackSchemaEncoding_DATA_TRACK_SCHEMA_ENCODING_PROTOBUF, - } - require.Nil(t, p.GetDataTrackSchema(id)) - - definition := &livekit.DataTrackSchemaDefinition{ - Id: id, - Definition: []byte("definition"), - } - p.asyncAttributes.Add(definition) - got := p.GetDataTrackSchema(id) - require.NotNil(t, got) - require.Equal(t, id.Name, got.Id.Name) - require.Equal(t, []byte("definition"), got.Definition) -} - -func TestProcessGetDataTrackSchemaRequest(t *testing.T) { - t.Run("returns NOT_FOUND when publisher is nil", func(t *testing.T) { - p := newParticipantWithAsyncAttributes(t, true, 0, 0) - sink := p.params.Sink.(*routingfakes.FakeMessageSink) - - p.ProcessGetDataTrackSchemaRequest(&livekit.GetDataTrackSchemaRequest{ - SchemaId: &livekit.DataTrackSchemaId{ - Name: "schema-1", - Encoding: livekit.DataTrackSchemaEncoding_DATA_TRACK_SCHEMA_ENCODING_PROTOBUF, - }, - }, nil) - - require.Equal(t, 1, sink.WriteMessageCallCount()) - rr := lastRequestResponse(t, sink, 0) - require.Equal(t, livekit.RequestResponse_NOT_FOUND, rr.Reason) - require.Contains(t, rr.Message, "participant") - }) - - t.Run("returns NOT_FOUND when publisher has no matching schema", func(t *testing.T) { - p := newParticipantWithAsyncAttributes(t, true, 0, 0) - sink := p.params.Sink.(*routingfakes.FakeMessageSink) - - publisher := &typesfakes.FakeParticipant{} - publisher.GetDataTrackSchemaReturns(nil) - - req := &livekit.GetDataTrackSchemaRequest{ - SchemaId: &livekit.DataTrackSchemaId{ - Name: "schema-1", - Encoding: livekit.DataTrackSchemaEncoding_DATA_TRACK_SCHEMA_ENCODING_PROTOBUF, - }, - } - p.ProcessGetDataTrackSchemaRequest(req, publisher) - - require.Equal(t, 1, publisher.GetDataTrackSchemaCallCount()) - require.Equal(t, req.SchemaId, publisher.GetDataTrackSchemaArgsForCall(0)) - - require.Equal(t, 1, sink.WriteMessageCallCount()) - rr := lastRequestResponse(t, sink, 0) - require.Equal(t, livekit.RequestResponse_NOT_FOUND, rr.Reason) - }) - - t.Run("sends schema response when publisher has a matching schema", func(t *testing.T) { - p := newParticipantWithAsyncAttributes(t, true, 0, 0) - sink := p.params.Sink.(*routingfakes.FakeMessageSink) - - id := &livekit.DataTrackSchemaId{ - Name: "schema-1", - Encoding: livekit.DataTrackSchemaEncoding_DATA_TRACK_SCHEMA_ENCODING_PROTOBUF, - } - def := &livekit.DataTrackSchemaDefinition{ - Id: id, - Definition: []byte("definition-bytes"), - } - - publisher := &typesfakes.FakeParticipant{} - publisher.GetDataTrackSchemaReturns(def) - - p.ProcessGetDataTrackSchemaRequest(&livekit.GetDataTrackSchemaRequest{ - SchemaId: id, - }, publisher) - - require.Equal(t, 1, sink.WriteMessageCallCount()) - msg := sink.WriteMessageArgsForCall(0).(*livekit.SignalResponse) - response, ok := msg.Message.(*livekit.SignalResponse_GetDataTrackSchemaResponse) - require.True(t, ok, "expected SignalResponse_GetDataTrackSchemaResponse, got %T", msg.Message) - require.Equal(t, def, response.GetDataTrackSchemaResponse.SchemaDefinition) - }) -} diff --git a/pkg/rtc/participant_async_attributes_test.go b/pkg/rtc/participant_async_attributes_test.go deleted file mode 100644 index 5f68110d0..000000000 --- a/pkg/rtc/participant_async_attributes_test.go +++ /dev/null @@ -1,218 +0,0 @@ -// Copyright 2026 LiveKit, Inc. -// -// Licensed under the Apache License, Version 2.0 (the "License"); -// you may not use this file except in compliance with the License. -// You may obtain a copy of the License at -// -// http://www.apache.org/licenses/LICENSE-2.0 -// -// Unless required by applicable law or agreed to in writing, software -// distributed under the License is distributed on an "AS IS" BASIS, -// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. -// See the License for the specific language governing permissions and -// limitations under the License. - -package rtc - -import ( - "sync" - "testing" - - "github.com/stretchr/testify/require" - - "github.com/livekit/protocol/livekit" - "github.com/livekit/protocol/logger" -) - -func newTestAsyncAttributes() *ParticipantAsyncAttributes { - return NewParticipantAsyncAttributes(ParticipantAsyncAttributesParams{ - Logger: logger.GetLogger(), - }) -} - -func TestParticipantAsyncAttributes_AddAndGet(t *testing.T) { - a := newTestAsyncAttributes() - - id := &livekit.DataTrackSchemaId{ - Name: "schema-1", - Encoding: livekit.DataTrackSchemaEncoding_DATA_TRACK_SCHEMA_ENCODING_PROTOBUF, - } - value := []byte("definition-bytes") - - a.Add(&livekit.DataTrackSchemaDefinition{Id: id, Definition: value}) - - got := a.Get(id) - require.NotNil(t, got) - require.Equal(t, id.Name, got.Id.Name) - require.Equal(t, id.Encoding, got.Id.Encoding) - require.Equal(t, value, got.Definition) -} - -func TestParticipantAsyncAttributes_AddOverwrites(t *testing.T) { - a := newTestAsyncAttributes() - - id := &livekit.DataTrackSchemaId{ - Name: "schema-1", - Encoding: livekit.DataTrackSchemaEncoding_DATA_TRACK_SCHEMA_ENCODING_PROTOBUF, - } - a.Add(&livekit.DataTrackSchemaDefinition{Id: id, Definition: []byte("v1")}) - a.Add(&livekit.DataTrackSchemaDefinition{Id: id, Definition: []byte("v2")}) - - got := a.Get(id) - require.NotNil(t, got) - require.Equal(t, []byte("v2"), got.Definition) - - require.Len(t, a.GetAll(), 1) -} - -func TestParticipantAsyncAttributes_DifferentEncodingsAreDistinct(t *testing.T) { - a := newTestAsyncAttributes() - - idProto := &livekit.DataTrackSchemaId{ - Name: "schema-1", - Encoding: livekit.DataTrackSchemaEncoding_DATA_TRACK_SCHEMA_ENCODING_PROTOBUF, - } - idJSON := &livekit.DataTrackSchemaId{ - Name: "schema-1", - Encoding: livekit.DataTrackSchemaEncoding_DATA_TRACK_SCHEMA_ENCODING_JSON_SCHEMA, - } - - a.Add(&livekit.DataTrackSchemaDefinition{Id: idProto, Definition: []byte("proto-def")}) - a.Add(&livekit.DataTrackSchemaDefinition{Id: idJSON, Definition: []byte("json-def")}) - - gotProto := a.Get(idProto) - require.NotNil(t, gotProto) - require.Equal(t, []byte("proto-def"), gotProto.Definition) - - gotJSON := a.Get(idJSON) - require.NotNil(t, gotJSON) - require.Equal(t, []byte("json-def"), gotJSON.Definition) - - require.Len(t, a.GetAll(), 2) -} - -func TestParticipantAsyncAttributes_Delete(t *testing.T) { - a := newTestAsyncAttributes() - - id := &livekit.DataTrackSchemaId{ - Name: "schema-1", - Encoding: livekit.DataTrackSchemaEncoding_DATA_TRACK_SCHEMA_ENCODING_PROTOBUF, - } - a.Add(&livekit.DataTrackSchemaDefinition{ - Id: id, - Definition: []byte("definition"), - }) - - a.Delete(id) - require.Nil(t, a.Get(id)) - require.Empty(t, a.GetAll()) - - // deleting a non-existent key is a no-op - a.Delete(id) - require.Empty(t, a.GetAll()) -} - -func TestParticipantAsyncAttributes_NilId(t *testing.T) { - a := newTestAsyncAttributes() - - // nil id should be silently ignored, not panic - a.Add(&livekit.DataTrackSchemaDefinition{Definition: []byte("definition")}) - require.Empty(t, a.GetAll()) - - require.Nil(t, a.Get(nil)) - - a.Delete(nil) - require.Empty(t, a.GetAll()) -} - -func TestParticipantAsyncAttributes_GetMissing(t *testing.T) { - a := newTestAsyncAttributes() - - id := &livekit.DataTrackSchemaId{ - Name: "missing", - Encoding: livekit.DataTrackSchemaEncoding_DATA_TRACK_SCHEMA_ENCODING_PROTOBUF, - } - require.Nil(t, a.Get(id)) -} - -func TestParticipantAsyncAttributes_GetAllContents(t *testing.T) { - a := newTestAsyncAttributes() - - id1 := &livekit.DataTrackSchemaId{ - Name: "schema-1", - Encoding: livekit.DataTrackSchemaEncoding_DATA_TRACK_SCHEMA_ENCODING_PROTOBUF, - } - id2 := &livekit.DataTrackSchemaId{ - Name: "schema-2", - Encoding: livekit.DataTrackSchemaEncoding_DATA_TRACK_SCHEMA_ENCODING_FLATBUFFER, - } - - a.Add(&livekit.DataTrackSchemaDefinition{Id: id1, Definition: []byte("def-1")}) - a.Add(&livekit.DataTrackSchemaDefinition{Id: id2, Definition: []byte("def-2")}) - - all := a.GetAll() - require.Len(t, all, 2) - for _, aa := range all { - switch aa.Id.Name { - case "schema-1": - require.Equal(t, []byte("def-1"), aa.Definition) - case "schema-2": - require.Equal(t, []byte("def-2"), aa.Definition) - default: - require.Fail(t, "unexpected name", aa.Id.Name) - } - } -} - -func TestParticipantAsyncAttributes_ConcurrentAccess(t *testing.T) { - a := newTestAsyncAttributes() - - const numGoroutines = 16 - const opsPerGoroutine = 100 - - var wg sync.WaitGroup - wg.Add(numGoroutines) - for g := 0; g < numGoroutines; g++ { - go func(g int) { - defer wg.Done() - for i := 0; i < opsPerGoroutine; i++ { - id := &livekit.DataTrackSchemaId{ - Name: "schema", - Encoding: livekit.DataTrackSchemaEncoding(g % 8), - } - a.Add(&livekit.DataTrackSchemaDefinition{Id: id, Definition: []byte("v")}) - _ = a.Get(id) - _ = a.GetAll() - if i%3 == 0 { - a.Delete(id) - } - } - }(g) - } - wg.Wait() -} - -func TestToKeyFromKey(t *testing.T) { - ids := []*livekit.DataTrackSchemaId{ - { - Name: "schema-1", - Encoding: livekit.DataTrackSchemaEncoding_DATA_TRACK_SCHEMA_ENCODING_PROTOBUF, - }, - { - Name: "another-schema", - Encoding: livekit.DataTrackSchemaEncoding_DATA_TRACK_SCHEMA_ENCODING_JSON_SCHEMA, - }, - { - Name: "x", - Encoding: livekit.DataTrackSchemaEncoding_DATA_TRACK_SCHEMA_ENCODING_UNSPECIFIED, - }, - } - for _, id := range ids { - t.Run(id.Name, func(t *testing.T) { - key := ToParticipantAsyncAttributeKey(id) - roundTripped := fromParticipantAsyncAttributeKey(key) - require.Equal(t, id.Name, roundTripped.Name) - require.Equal(t, id.Encoding, roundTripped.Encoding) - }) - } -} diff --git a/pkg/rtc/participant_data_blob.go b/pkg/rtc/participant_data_blob.go new file mode 100644 index 000000000..20306fe55 --- /dev/null +++ b/pkg/rtc/participant_data_blob.go @@ -0,0 +1,91 @@ +// Copyright 2026 LiveKit, Inc. +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +package rtc + +import ( + "sync" + + "github.com/livekit/protocol/livekit" + "github.com/livekit/protocol/logger" + "github.com/livekit/protocol/utils" +) + +type ParticipantDataBlobParams struct { + Logger logger.Logger +} + +type ParticipantDataBlob struct { + params ParticipantDataBlobParams + lock sync.Mutex + blobs map[string]*livekit.DataBlob +} + +func NewParticipantDataBlob(params ParticipantDataBlobParams) *ParticipantDataBlob { + return &ParticipantDataBlob{ + params: params, + blobs: make(map[string]*livekit.DataBlob), + } +} + +func (p *ParticipantDataBlob) Add(db *livekit.DataBlob) { + p.lock.Lock() + defer p.lock.Unlock() + + if db.Key == nil { + return + } + + p.blobs[db.Key.String()] = utils.CloneProto(db) +} + +func (p *ParticipantDataBlob) Delete(dbKey *livekit.DataBlobKey) { + p.lock.Lock() + defer p.lock.Unlock() + + if dbKey == nil { + return + } + + delete(p.blobs, dbKey.String()) +} + +func (p *ParticipantDataBlob) Get(dbKey *livekit.DataBlobKey) *livekit.DataBlob { + p.lock.Lock() + defer p.lock.Unlock() + + if dbKey == nil { + return nil + } + + db, ok := p.blobs[dbKey.String()] + if !ok { + return nil + } + + return db +} + +func (p *ParticipantDataBlob) GetAll() []*livekit.DataBlob { + p.lock.Lock() + defer p.lock.Unlock() + + all := make([]*livekit.DataBlob, 0, len(p.blobs)) + for _, db := range p.blobs { + all = append(all, utils.CloneProto(db)) + } + return all +} + +// ------------------------------- diff --git a/pkg/rtc/participant_data_blob_handler.go b/pkg/rtc/participant_data_blob_handler.go new file mode 100644 index 000000000..4559d5d90 --- /dev/null +++ b/pkg/rtc/participant_data_blob_handler.go @@ -0,0 +1,111 @@ +// Copyright 2026 LiveKit, Inc. +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +package rtc + +import ( + "github.com/livekit/livekit-server/pkg/rtc/types" + "github.com/livekit/protocol/livekit" + "github.com/livekit/protocol/logger" +) + +func (p *ParticipantImpl) HandleStoreDataBlobRequest(req *livekit.StoreDataBlobRequest) { + if !p.params.EnableParticipantDataBlob { + p.pubLogger.Warnw("data blob not enabled", nil, "req", logger.Proto(req)) + p.sendRequestResponse(&livekit.RequestResponse{ + RequestId: req.RequestId, + Reason: livekit.RequestResponse_NOT_ALLOWED, + Message: "data blob not enabled", + }) + return + } + + if req.Blob == nil || req.Blob.Key == nil || len(req.Blob.Key.String()) == 0 || !p.params.LimitConfig.CheckDataBlobKeyLength(req.Blob.Key.String()) { + p.pubLogger.Warnw("data blob is invalid", nil, "req", logger.Proto(req)) + p.sendRequestResponse(&livekit.RequestResponse{ + RequestId: req.RequestId, + Reason: livekit.RequestResponse_INVALID_REQUEST, + Message: "data blob is invalid", + }) + return + } + + if len(req.Blob.Contents) == 0 { + p.sendRequestResponse(&livekit.RequestResponse{ + RequestId: req.RequestId, + Reason: livekit.RequestResponse_INVALID_REQUEST, + Message: "data blob is empty", + }) + return + } + + if !p.params.LimitConfig.CanAddDataBlob(p.dataBlob.GetAll(), req.Blob) { + p.sendRequestResponse(&livekit.RequestResponse{ + RequestId: req.RequestId, + Reason: livekit.RequestResponse_LIMIT_EXCEEDED, + Message: "async attribute definition exceeds limit", + }) + return + } + + p.AddDataBlob(req.Blob) + p.listener().OnStoreDataBlob(p, req.Blob) +} + +func (p *ParticipantImpl) HandleGetDataBlobRequest(req *livekit.GetDataBlobRequest) { + if req.Key == nil { + p.sendRequestResponse(&livekit.RequestResponse{ + RequestId: req.RequestId, + Reason: livekit.RequestResponse_INVALID_REQUEST, + Message: "data blob key is required", + }) + return + } + + p.listener().OnGetDataBlob(p, req) +} + +func (p *ParticipantImpl) AddDataBlob(dataBlob *livekit.DataBlob) { + p.dataBlob.Add(dataBlob) +} + +func (p *ParticipantImpl) GetDataBlob(key *livekit.DataBlobKey) *livekit.DataBlob { + return p.dataBlob.Get(key) +} + +func (p *ParticipantImpl) ProcessGetDataBlobRequest(req *livekit.GetDataBlobRequest, publisher types.Participant) { + if publisher == nil { + p.sendRequestResponse(&livekit.RequestResponse{ + RequestId: req.RequestId, + Reason: livekit.RequestResponse_NOT_FOUND, + Message: "participant not found", + }) + return + } + + dataBlob := publisher.GetDataBlob(req.Key) + if dataBlob == nil { + p.sendRequestResponse(&livekit.RequestResponse{ + Reason: livekit.RequestResponse_NOT_FOUND, + Message: "data blob not found", + }) + return + } + + p.sendGetDataBlobResponse(dataBlob) +} + +func (p *ParticipantImpl) GetAllDataBlob() []*livekit.DataBlob { + return p.dataBlob.GetAll() +} diff --git a/pkg/rtc/participant_data_blob_handler_test.go b/pkg/rtc/participant_data_blob_handler_test.go new file mode 100644 index 000000000..236b14379 --- /dev/null +++ b/pkg/rtc/participant_data_blob_handler_test.go @@ -0,0 +1,291 @@ +// Copyright 2026 LiveKit, Inc. +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +package rtc + +import ( + "strings" + "testing" + + "github.com/stretchr/testify/require" + + "github.com/livekit/protocol/livekit" + + "github.com/livekit/livekit-server/pkg/config" + "github.com/livekit/livekit-server/pkg/routing/routingfakes" + "github.com/livekit/livekit-server/pkg/rtc/types/typesfakes" +) + +func newParticipantWithDataBlob(t *testing.T, enabled bool, maxKeyLength int, maxSize uint32) *ParticipantImpl { + t.Helper() + p := newParticipantForTest("test") + p.params.EnableParticipantDataBlob = enabled + p.params.LimitConfig = config.LimitConfig{ + MaxDataBlobKeyLength: maxKeyLength, + MaxDataBlobSize: maxSize, + } + return p +} + +func lastRequestResponse(t *testing.T, sink *routingfakes.FakeMessageSink, idx int) *livekit.RequestResponse { + t.Helper() + msg := sink.WriteMessageArgsForCall(idx).(*livekit.SignalResponse) + rr, ok := msg.Message.(*livekit.SignalResponse_RequestResponse) + require.True(t, ok, "expected SignalResponse_RequestResponse, got %T", msg.Message) + return rr.RequestResponse +} + +func TestHandleStoreDataBlobRequest(t *testing.T) { + t.Run("returns NOT_ALLOWED when feature not enabled", func(t *testing.T) { + p := newParticipantWithDataBlob(t, false, 0, 0) + sink := p.params.Sink.(*routingfakes.FakeMessageSink) + + req := &livekit.StoreDataBlobRequest{ + Blob: &livekit.DataBlob{ + Key: genericKey("blob-1"), + Contents: []byte("def"), + }, + } + p.HandleStoreDataBlobRequest(req) + + require.Equal(t, 1, sink.WriteMessageCallCount()) + rr := lastRequestResponse(t, sink, 0) + require.Equal(t, livekit.RequestResponse_NOT_ALLOWED, rr.Reason) + require.Empty(t, p.dataBlob.GetAll()) + }) + + t.Run("returns INVALID_REQUEST when blob is nil", func(t *testing.T) { + p := newParticipantWithDataBlob(t, true, 0, 0) + sink := p.params.Sink.(*routingfakes.FakeMessageSink) + + p.HandleStoreDataBlobRequest(&livekit.StoreDataBlobRequest{}) + + require.Equal(t, 1, sink.WriteMessageCallCount()) + rr := lastRequestResponse(t, sink, 0) + require.Equal(t, livekit.RequestResponse_INVALID_REQUEST, rr.Reason) + require.Empty(t, p.dataBlob.GetAll()) + }) + + t.Run("returns INVALID_REQUEST when key is nil", func(t *testing.T) { + p := newParticipantWithDataBlob(t, true, 0, 0) + sink := p.params.Sink.(*routingfakes.FakeMessageSink) + + p.HandleStoreDataBlobRequest(&livekit.StoreDataBlobRequest{ + Blob: &livekit.DataBlob{ + Contents: []byte("def"), + }, + }) + + require.Equal(t, 1, sink.WriteMessageCallCount()) + rr := lastRequestResponse(t, sink, 0) + require.Equal(t, livekit.RequestResponse_INVALID_REQUEST, rr.Reason) + }) + + t.Run("returns INVALID_REQUEST when key has no oneof set", func(t *testing.T) { + p := newParticipantWithDataBlob(t, true, 0, 0) + sink := p.params.Sink.(*routingfakes.FakeMessageSink) + + p.HandleStoreDataBlobRequest(&livekit.StoreDataBlobRequest{ + Blob: &livekit.DataBlob{ + Key: &livekit.DataBlobKey{}, + Contents: []byte("def"), + }, + }) + + require.Equal(t, 1, sink.WriteMessageCallCount()) + rr := lastRequestResponse(t, sink, 0) + require.Equal(t, livekit.RequestResponse_INVALID_REQUEST, rr.Reason) + }) + + t.Run("returns INVALID_REQUEST when key exceeds length limit", func(t *testing.T) { + p := newParticipantWithDataBlob(t, true, 5, 0) + sink := p.params.Sink.(*routingfakes.FakeMessageSink) + + p.HandleStoreDataBlobRequest(&livekit.StoreDataBlobRequest{ + Blob: &livekit.DataBlob{ + Key: genericKey(strings.Repeat("a", 64)), + Contents: []byte("def"), + }, + }) + + require.Equal(t, 1, sink.WriteMessageCallCount()) + rr := lastRequestResponse(t, sink, 0) + require.Equal(t, livekit.RequestResponse_INVALID_REQUEST, rr.Reason) + }) + + t.Run("returns INVALID_REQUEST when contents is empty", func(t *testing.T) { + p := newParticipantWithDataBlob(t, true, 0, 0) + sink := p.params.Sink.(*routingfakes.FakeMessageSink) + + p.HandleStoreDataBlobRequest(&livekit.StoreDataBlobRequest{ + Blob: &livekit.DataBlob{ + Key: genericKey("blob-1"), + }, + }) + + require.Equal(t, 1, sink.WriteMessageCallCount()) + rr := lastRequestResponse(t, sink, 0) + require.Equal(t, livekit.RequestResponse_INVALID_REQUEST, rr.Reason) + require.Empty(t, p.dataBlob.GetAll()) + }) + + t.Run("returns LIMIT_EXCEEDED when adding would breach the limit", func(t *testing.T) { + p := newParticipantWithDataBlob(t, true, 0, 16) + sink := p.params.Sink.(*routingfakes.FakeMessageSink) + + p.HandleStoreDataBlobRequest(&livekit.StoreDataBlobRequest{ + Blob: &livekit.DataBlob{ + Key: genericKey("blob-1"), + Contents: []byte(strings.Repeat("x", 32)), + }, + }) + + require.Equal(t, 1, sink.WriteMessageCallCount()) + rr := lastRequestResponse(t, sink, 0) + require.Equal(t, livekit.RequestResponse_LIMIT_EXCEEDED, rr.Reason) + require.Empty(t, p.dataBlob.GetAll()) + }) + + t.Run("stores a valid blob, notifies listener, and sends no response", func(t *testing.T) { + p := newParticipantWithDataBlob(t, true, 0, 0) + sink := p.params.Sink.(*routingfakes.FakeMessageSink) + listener := p.params.ParticipantListener.(*typesfakes.FakeLocalParticipantListener) + + key := genericKey("blob-1") + contents := []byte("definition-bytes") + blob := &livekit.DataBlob{Key: key, Contents: contents} + + p.HandleStoreDataBlobRequest(&livekit.StoreDataBlobRequest{Blob: blob}) + + // on success no response is sent to the client + require.Equal(t, 0, sink.WriteMessageCallCount()) + + stored := p.dataBlob.Get(key) + require.NotNil(t, stored) + require.Equal(t, contents, stored.Contents) + + require.Equal(t, 1, listener.OnStoreDataBlobCallCount()) + gotParticipant, gotBlob := listener.OnStoreDataBlobArgsForCall(0) + require.Equal(t, p, gotParticipant) + require.Equal(t, blob, gotBlob) + }) +} + +func TestHandleGetDataBlobRequest(t *testing.T) { + t.Run("returns INVALID_REQUEST when key is missing", func(t *testing.T) { + p := newParticipantWithDataBlob(t, true, 0, 0) + sink := p.params.Sink.(*routingfakes.FakeMessageSink) + + p.HandleGetDataBlobRequest(&livekit.GetDataBlobRequest{ + ParticipantIdentity: "other", + }) + + require.Equal(t, 1, sink.WriteMessageCallCount()) + rr := lastRequestResponse(t, sink, 0) + require.Equal(t, livekit.RequestResponse_INVALID_REQUEST, rr.Reason) + }) + + t.Run("forwards request to listener when key is provided", func(t *testing.T) { + p := newParticipantWithDataBlob(t, true, 0, 0) + listener := p.params.ParticipantListener.(*typesfakes.FakeLocalParticipantListener) + + req := &livekit.GetDataBlobRequest{ + ParticipantIdentity: "other", + Key: genericKey("blob-1"), + } + p.HandleGetDataBlobRequest(req) + + require.Equal(t, 1, listener.OnGetDataBlobCallCount()) + gotParticipant, gotReq := listener.OnGetDataBlobArgsForCall(0) + require.Equal(t, p, gotParticipant) + require.Equal(t, req, gotReq) + }) +} + +func TestGetDataBlob(t *testing.T) { + p := newParticipantWithDataBlob(t, true, 0, 0) + + key := genericKey("blob-1") + require.Nil(t, p.GetDataBlob(key)) + + blob := &livekit.DataBlob{ + Key: key, + Contents: []byte("definition"), + } + p.dataBlob.Add(blob) + got := p.GetDataBlob(key) + require.NotNil(t, got) + require.Equal(t, key.String(), got.Key.String()) + require.Equal(t, []byte("definition"), got.Contents) +} + +func TestProcessGetDataBlobRequest(t *testing.T) { + t.Run("returns NOT_FOUND when publisher is nil", func(t *testing.T) { + p := newParticipantWithDataBlob(t, true, 0, 0) + sink := p.params.Sink.(*routingfakes.FakeMessageSink) + + p.ProcessGetDataBlobRequest(&livekit.GetDataBlobRequest{ + Key: genericKey("blob-1"), + }, nil) + + require.Equal(t, 1, sink.WriteMessageCallCount()) + rr := lastRequestResponse(t, sink, 0) + require.Equal(t, livekit.RequestResponse_NOT_FOUND, rr.Reason) + require.Contains(t, rr.Message, "participant") + }) + + t.Run("returns NOT_FOUND when publisher has no matching blob", func(t *testing.T) { + p := newParticipantWithDataBlob(t, true, 0, 0) + sink := p.params.Sink.(*routingfakes.FakeMessageSink) + + publisher := &typesfakes.FakeParticipant{} + publisher.GetDataBlobReturns(nil) + + req := &livekit.GetDataBlobRequest{ + Key: genericKey("blob-1"), + } + p.ProcessGetDataBlobRequest(req, publisher) + + require.Equal(t, 1, publisher.GetDataBlobCallCount()) + require.Equal(t, req.Key, publisher.GetDataBlobArgsForCall(0)) + + require.Equal(t, 1, sink.WriteMessageCallCount()) + rr := lastRequestResponse(t, sink, 0) + require.Equal(t, livekit.RequestResponse_NOT_FOUND, rr.Reason) + }) + + t.Run("sends blob response when publisher has a matching blob", func(t *testing.T) { + p := newParticipantWithDataBlob(t, true, 0, 0) + sink := p.params.Sink.(*routingfakes.FakeMessageSink) + + key := genericKey("blob-1") + blob := &livekit.DataBlob{ + Key: key, + Contents: []byte("definition-bytes"), + } + + publisher := &typesfakes.FakeParticipant{} + publisher.GetDataBlobReturns(blob) + + p.ProcessGetDataBlobRequest(&livekit.GetDataBlobRequest{ + Key: key, + }, publisher) + + require.Equal(t, 1, sink.WriteMessageCallCount()) + msg := sink.WriteMessageArgsForCall(0).(*livekit.SignalResponse) + response, ok := msg.Message.(*livekit.SignalResponse_GetDataBlobResponse) + require.True(t, ok, "expected SignalResponse_GetDataBlobResponse, got %T", msg.Message) + require.Equal(t, blob, response.GetDataBlobResponse.Blob) + }) +} diff --git a/pkg/rtc/participant_data_blob_test.go b/pkg/rtc/participant_data_blob_test.go new file mode 100644 index 000000000..a6d28fb20 --- /dev/null +++ b/pkg/rtc/participant_data_blob_test.go @@ -0,0 +1,175 @@ +// Copyright 2026 LiveKit, Inc. +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +package rtc + +import ( + "fmt" + "sync" + "testing" + + "github.com/stretchr/testify/require" + + "github.com/livekit/protocol/livekit" + "github.com/livekit/protocol/logger" +) + +func newTestDataBlob() *ParticipantDataBlob { + return NewParticipantDataBlob(ParticipantDataBlobParams{ + Logger: logger.GetLogger(), + }) +} + +func genericKey(name string) *livekit.DataBlobKey { + return &livekit.DataBlobKey{ + Key: &livekit.DataBlobKey_Generic{ + Generic: name, + }, + } +} + +func TestParticipantDataBlob_AddAndGet(t *testing.T) { + a := newTestDataBlob() + + key := genericKey("blob-1") + contents := []byte("definition-bytes") + + a.Add(&livekit.DataBlob{Key: key, Contents: contents}) + + got := a.Get(key) + require.NotNil(t, got) + require.Equal(t, key.String(), got.Key.String()) + require.Equal(t, contents, got.Contents) +} + +func TestParticipantDataBlob_AddOverwrites(t *testing.T) { + a := newTestDataBlob() + + key := genericKey("blob-1") + a.Add(&livekit.DataBlob{Key: key, Contents: []byte("v1")}) + a.Add(&livekit.DataBlob{Key: key, Contents: []byte("v2")}) + + got := a.Get(key) + require.NotNil(t, got) + require.Equal(t, []byte("v2"), got.Contents) + + require.Len(t, a.GetAll(), 1) +} + +func TestParticipantDataBlob_DistinctKeys(t *testing.T) { + a := newTestDataBlob() + + key1 := genericKey("blob-1") + key2 := genericKey("blob-2") + + a.Add(&livekit.DataBlob{Key: key1, Contents: []byte("c1")}) + a.Add(&livekit.DataBlob{Key: key2, Contents: []byte("c2")}) + + got1 := a.Get(key1) + require.NotNil(t, got1) + require.Equal(t, []byte("c1"), got1.Contents) + + got2 := a.Get(key2) + require.NotNil(t, got2) + require.Equal(t, []byte("c2"), got2.Contents) + + require.Len(t, a.GetAll(), 2) +} + +func TestParticipantDataBlob_Delete(t *testing.T) { + a := newTestDataBlob() + + key := genericKey("blob-1") + a.Add(&livekit.DataBlob{Key: key, Contents: []byte("definition")}) + + a.Delete(key) + require.Nil(t, a.Get(key)) + require.Empty(t, a.GetAll()) + + // deleting a non-existent key is a no-op + a.Delete(key) + require.Empty(t, a.GetAll()) +} + +func TestParticipantDataBlob_NilKey(t *testing.T) { + a := newTestDataBlob() + + // nil key should be silently ignored, not panic + a.Add(&livekit.DataBlob{Contents: []byte("definition")}) + require.Empty(t, a.GetAll()) + + require.Nil(t, a.Get(nil)) + + a.Delete(nil) + require.Empty(t, a.GetAll()) +} + +func TestParticipantDataBlob_GetMissing(t *testing.T) { + a := newTestDataBlob() + + require.Nil(t, a.Get(genericKey("missing"))) +} + +func TestParticipantDataBlob_GetAllContents(t *testing.T) { + a := newTestDataBlob() + + key1 := genericKey("blob-1") + key2 := genericKey("blob-2") + + a.Add(&livekit.DataBlob{Key: key1, Contents: []byte("def-1")}) + a.Add(&livekit.DataBlob{Key: key2, Contents: []byte("def-2")}) + + all := a.GetAll() + require.Len(t, all, 2) + for _, db := range all { + switch key := db.Key.Key.(type) { + case *livekit.DataBlobKey_Generic: + switch key.Generic { + case "blob-1": + require.Equal(t, []byte("def-1"), db.Contents) + case "blob-2": + require.Equal(t, []byte("def-2"), db.Contents) + default: + require.Fail(t, "unexpected key", key.Generic) + } + default: + require.Fail(t, "unexpected key type", "Generic") + } + } +} + +func TestParticipantDataBlob_ConcurrentAccess(t *testing.T) { + a := newTestDataBlob() + + const numGoroutines = 16 + const opsPerGoroutine = 100 + + var wg sync.WaitGroup + wg.Add(numGoroutines) + for g := 0; g < numGoroutines; g++ { + go func(g int) { + defer wg.Done() + for i := 0; i < opsPerGoroutine; i++ { + key := genericKey(fmt.Sprintf("blob-%d", g%8)) + a.Add(&livekit.DataBlob{Key: key, Contents: []byte("v")}) + _ = a.Get(key) + _ = a.GetAll() + if i%3 == 0 { + a.Delete(key) + } + } + }(g) + } + wg.Wait() +} diff --git a/pkg/rtc/participant_signal.go b/pkg/rtc/participant_signal.go index 14b3a8a2c..5462129d7 100644 --- a/pkg/rtc/participant_signal.go +++ b/pkg/rtc/participant_signal.go @@ -369,8 +369,8 @@ func (p *ParticipantImpl) SendDataTrackSubscriberHandles(handles map[uint32]*liv })) } -func (p *ParticipantImpl) sendGetDataTrackSchemaResponse(definition *livekit.DataTrackSchemaDefinition) error { - return p.signaller.WriteMessage(p.signalling.SignalGetDataTrackSchemaResponse(&livekit.GetDataTrackSchemaResponse{ - SchemaDefinition: definition, +func (p *ParticipantImpl) sendGetDataBlobResponse(dataBlob *livekit.DataBlob) error { + return p.signaller.WriteMessage(p.signalling.SignalGetDataBlobResponse(&livekit.GetDataBlobResponse{ + Blob: dataBlob, })) } diff --git a/pkg/rtc/room.go b/pkg/rtc/room.go index f9ea860d1..56a7c9370 100644 --- a/pkg/rtc/room.go +++ b/pkg/rtc/room.go @@ -1384,9 +1384,9 @@ func (r *Room) onUpdateDataSubscriptions(participant types.LocalParticipant, req } } -func (r *Room) onGetDataTrackSchema(participant types.LocalParticipant, req *livekit.GetDataTrackSchemaRequest) { +func (r *Room) onGetDataBlob(participant types.LocalParticipant, req *livekit.GetDataBlobRequest) { publisher := r.GetParticipant(livekit.ParticipantIdentity(req.ParticipantIdentity)) - participant.ProcessGetDataTrackSchemaRequest(req, publisher) + participant.ProcessGetDataBlobRequest(req, publisher) } func (r *Room) onLeave(p types.LocalParticipant, reason types.ParticipantCloseReason) { @@ -1998,11 +1998,11 @@ func (l *localParticipantListener) OnUpdateDataSubscriptions(p types.LocalPartic l.room.onUpdateDataSubscriptions(p, req) } -func (l *localParticipantListener) OnDefineDataTrackSchema(_p types.LocalParticipant, _definition *livekit.DataTrackSchemaDefinition) { +func (l *localParticipantListener) OnStoreDataBlob(_p types.LocalParticipant, _dataBlob *livekit.DataBlob) { } -func (l *localParticipantListener) OnGetDataTrackSchema(p types.LocalParticipant, req *livekit.GetDataTrackSchemaRequest) { - l.room.onGetDataTrackSchema(p, req) +func (l *localParticipantListener) OnGetDataBlob(p types.LocalParticipant, req *livekit.GetDataBlobRequest) { + l.room.onGetDataBlob(p, req) } func (l *localParticipantListener) OnSyncState(p types.LocalParticipant, state *livekit.SyncState) error { diff --git a/pkg/rtc/signalling/interfaces.go b/pkg/rtc/signalling/interfaces.go index d6ed03c16..1c8cf5459 100644 --- a/pkg/rtc/signalling/interfaces.go +++ b/pkg/rtc/signalling/interfaces.go @@ -62,5 +62,5 @@ type ParticipantSignalling interface { SignalPublishDataTrackResponse(publishDataTrackResponse *livekit.PublishDataTrackResponse) proto.Message SignalUnpublishDataTrackResponse(unpublishDataTrackResponse *livekit.UnpublishDataTrackResponse) proto.Message SignalDataTrackSubscriberHandles(dataTrackSubscriberHandles *livekit.DataTrackSubscriberHandles) proto.Message - SignalGetDataTrackSchemaResponse(getDataTrackSchemaResponse *livekit.GetDataTrackSchemaResponse) proto.Message + SignalGetDataBlobResponse(getDataBlobResponse *livekit.GetDataBlobResponse) proto.Message } diff --git a/pkg/rtc/signalling/signalhandler.go b/pkg/rtc/signalling/signalhandler.go index 0b3b282a8..740e65899 100644 --- a/pkg/rtc/signalling/signalhandler.go +++ b/pkg/rtc/signalling/signalhandler.go @@ -152,11 +152,11 @@ func (s *signalhandler) HandleMessage(msg proto.Message) error { case *livekit.SignalRequest_UpdateDataSubscription: s.params.Participant.HandleUpdateDataSubscription(msg.UpdateDataSubscription) - case *livekit.SignalRequest_DefineDataTrackSchema: - s.params.Participant.HandleDefineDataTrackSchemaRequest(msg.DefineDataTrackSchema) + case *livekit.SignalRequest_StoreDataBlobRequest: + s.params.Participant.HandleStoreDataBlobRequest(msg.StoreDataBlobRequest) - case *livekit.SignalRequest_GetDataTrackSchema: - s.params.Participant.HandleGetDataTrackSchemaRequest(msg.GetDataTrackSchema) + case *livekit.SignalRequest_GetDataBlobRequest: + s.params.Participant.HandleGetDataBlobRequest(msg.GetDataBlobRequest) } return nil diff --git a/pkg/rtc/signalling/signalling.go b/pkg/rtc/signalling/signalling.go index 946ac9914..9f224c81e 100644 --- a/pkg/rtc/signalling/signalling.go +++ b/pkg/rtc/signalling/signalling.go @@ -259,10 +259,10 @@ func (s *signalling) SignalDataTrackSubscriberHandles(dataTrackSubscriberHandles } } -func (u *signalling) SignalGetDataTrackSchemaResponse(getDataTrackSchemaResponse *livekit.GetDataTrackSchemaResponse) proto.Message { +func (u *signalling) SignalGetDataBlobResponse(getDataBlobResponse *livekit.GetDataBlobResponse) proto.Message { return &livekit.SignalResponse{ - Message: &livekit.SignalResponse_GetDataTrackSchemaResponse{ - GetDataTrackSchemaResponse: getDataTrackSchemaResponse, + Message: &livekit.SignalResponse_GetDataBlobResponse{ + GetDataBlobResponse: getDataBlobResponse, }, } } diff --git a/pkg/rtc/signalling/signallingunimplemented.go b/pkg/rtc/signalling/signallingunimplemented.go index b229dff63..767eeb15a 100644 --- a/pkg/rtc/signalling/signallingunimplemented.go +++ b/pkg/rtc/signalling/signallingunimplemented.go @@ -128,6 +128,6 @@ func (u *signallingUnimplemented) SignalDataTrackSubscriberHandles(dataTrackSubs return nil } -func (u *signallingUnimplemented) SignalGetDataTrackSchemaResponse(getDataTrackSchemaResponse *livekit.GetDataTrackSchemaResponse) proto.Message { +func (u *signallingUnimplemented) SignalGetDataBlobResponse(getDataBlobResponse *livekit.GetDataBlobResponse) proto.Message { return nil } diff --git a/pkg/rtc/types/interfaces.go b/pkg/rtc/types/interfaces.go index f2ac85829..891392ab4 100644 --- a/pkg/rtc/types/interfaces.go +++ b/pkg/rtc/types/interfaces.go @@ -356,8 +356,8 @@ type Participant interface { GetParticipantListener() ParticipantListener - AddDataTrackSchema(definition *livekit.DataTrackSchemaDefinition) - GetDataTrackSchema(id *livekit.DataTrackSchemaId) *livekit.DataTrackSchemaDefinition + AddDataBlob(dataBlob *livekit.DataBlob) + GetDataBlob(key *livekit.DataBlobKey) *livekit.DataBlob } // ------------------------------------------------------- @@ -564,9 +564,9 @@ type LocalParticipant interface { HandlePublishDataTrackRequest(*livekit.PublishDataTrackRequest) HandleUnpublishDataTrackRequest(*livekit.UnpublishDataTrackRequest) HandleUpdateDataSubscription(*livekit.UpdateDataSubscription) - HandleDefineDataTrackSchemaRequest(*livekit.DefineDataTrackSchemaRequest) - HandleGetDataTrackSchemaRequest(*livekit.GetDataTrackSchemaRequest) - ProcessGetDataTrackSchemaRequest(*livekit.GetDataTrackSchemaRequest, Participant) + HandleStoreDataBlobRequest(*livekit.StoreDataBlobRequest) + HandleGetDataBlobRequest(*livekit.GetDataBlobRequest) + ProcessGetDataBlobRequest(*livekit.GetDataBlobRequest, Participant) HandleSignalMessage(msg proto.Message) error @@ -578,7 +578,7 @@ type LocalParticipant interface { GetNextSubscribedDataTrackHandle() uint16 - GetAllAsyncAttributes() []*livekit.DataTrackSchemaDefinition + GetAllDataBlob() []*livekit.DataBlob } // --------------------------------------------- @@ -628,8 +628,8 @@ type LocalParticipantListener interface { ) OnUpdateSubscriptionPermission(LocalParticipant, *livekit.SubscriptionPermission) error OnUpdateDataSubscriptions(LocalParticipant, *livekit.UpdateDataSubscription) - OnDefineDataTrackSchema(LocalParticipant, *livekit.DataTrackSchemaDefinition) - OnGetDataTrackSchema(LocalParticipant, *livekit.GetDataTrackSchemaRequest) + OnStoreDataBlob(LocalParticipant, *livekit.DataBlob) + OnGetDataBlob(LocalParticipant, *livekit.GetDataBlobRequest) OnSyncState(LocalParticipant, *livekit.SyncState) error OnSimulateScenario(LocalParticipant, *livekit.SimulateScenario) error OnLeave(LocalParticipant, ParticipantCloseReason) @@ -661,9 +661,9 @@ func (*NullLocalParticipantListener) OnUpdateSubscriptionPermission(LocalPartici } func (*NullLocalParticipantListener) OnUpdateDataSubscriptions(LocalParticipant, *livekit.UpdateDataSubscription) { } -func (*NullLocalParticipantListener) OnDefineDataTrackSchema(LocalParticipant, *livekit.DataTrackSchemaDefinition) { +func (*NullLocalParticipantListener) OnStoreDataBlob(LocalParticipant, *livekit.DataBlob) { } -func (*NullLocalParticipantListener) OnGetDataTrackSchema(LocalParticipant, *livekit.GetDataTrackSchemaRequest) { +func (*NullLocalParticipantListener) OnGetDataBlob(LocalParticipant, *livekit.GetDataBlobRequest) { } func (*NullLocalParticipantListener) OnSyncState(LocalParticipant, *livekit.SyncState) error { return nil diff --git a/pkg/rtc/types/typesfakes/fake_local_participant.go b/pkg/rtc/types/typesfakes/fake_local_participant.go index fec988305..021d16620 100644 --- a/pkg/rtc/types/typesfakes/fake_local_participant.go +++ b/pkg/rtc/types/typesfakes/fake_local_participant.go @@ -33,10 +33,10 @@ type FakeLocalParticipant struct { activeAtReturnsOnCall map[int]struct { result1 time.Time } - AddDataTrackSchemaStub func(*livekit.DataTrackSchemaDefinition) - addDataTrackSchemaMutex sync.RWMutex - addDataTrackSchemaArgsForCall []struct { - arg1 *livekit.DataTrackSchemaDefinition + AddDataBlobStub func(*livekit.DataBlob) + addDataBlobMutex sync.RWMutex + addDataBlobArgsForCall []struct { + arg1 *livekit.DataBlob } AddOnCloseStub func(string, func(types.LocalParticipant)) addOnCloseMutex sync.RWMutex @@ -221,15 +221,15 @@ type FakeLocalParticipant struct { getAdaptiveStreamReturnsOnCall map[int]struct { result1 bool } - GetAllAsyncAttributesStub func() []*livekit.DataTrackSchemaDefinition - getAllAsyncAttributesMutex sync.RWMutex - getAllAsyncAttributesArgsForCall []struct { + GetAllDataBlobStub func() []*livekit.DataBlob + getAllDataBlobMutex sync.RWMutex + getAllDataBlobArgsForCall []struct { } - getAllAsyncAttributesReturns struct { - result1 []*livekit.DataTrackSchemaDefinition + getAllDataBlobReturns struct { + result1 []*livekit.DataBlob } - getAllAsyncAttributesReturnsOnCall map[int]struct { - result1 []*livekit.DataTrackSchemaDefinition + getAllDataBlobReturnsOnCall map[int]struct { + result1 []*livekit.DataBlob } GetAnswerStub func() (webrtc.SessionDescription, uint32, error) getAnswerMutex sync.RWMutex @@ -320,16 +320,16 @@ type FakeLocalParticipant struct { getCountryReturnsOnCall map[int]struct { result1 string } - GetDataTrackSchemaStub func(*livekit.DataTrackSchemaId) *livekit.DataTrackSchemaDefinition - getDataTrackSchemaMutex sync.RWMutex - getDataTrackSchemaArgsForCall []struct { - arg1 *livekit.DataTrackSchemaId + GetDataBlobStub func(*livekit.DataBlobKey) *livekit.DataBlob + getDataBlobMutex sync.RWMutex + getDataBlobArgsForCall []struct { + arg1 *livekit.DataBlobKey } - getDataTrackSchemaReturns struct { - result1 *livekit.DataTrackSchemaDefinition + getDataBlobReturns struct { + result1 *livekit.DataBlob } - getDataTrackSchemaReturnsOnCall map[int]struct { - result1 *livekit.DataTrackSchemaDefinition + getDataBlobReturnsOnCall map[int]struct { + result1 *livekit.DataBlob } GetDataTrackTransportStub func() types.DataTrackTransport getDataTrackTransportMutex sync.RWMutex @@ -592,15 +592,10 @@ type FakeLocalParticipant struct { handleAnswerArgsForCall []struct { arg1 *livekit.SessionDescription } - HandleDefineDataTrackSchemaRequestStub func(*livekit.DefineDataTrackSchemaRequest) - handleDefineDataTrackSchemaRequestMutex sync.RWMutex - handleDefineDataTrackSchemaRequestArgsForCall []struct { - arg1 *livekit.DefineDataTrackSchemaRequest - } - HandleGetDataTrackSchemaRequestStub func(*livekit.GetDataTrackSchemaRequest) - handleGetDataTrackSchemaRequestMutex sync.RWMutex - handleGetDataTrackSchemaRequestArgsForCall []struct { - arg1 *livekit.GetDataTrackSchemaRequest + HandleGetDataBlobRequestStub func(*livekit.GetDataBlobRequest) + handleGetDataBlobRequestMutex sync.RWMutex + handleGetDataBlobRequestArgsForCall []struct { + arg1 *livekit.GetDataBlobRequest } HandleICERestartSDPFragmentStub func(string) (string, error) handleICERestartSDPFragmentMutex sync.RWMutex @@ -715,6 +710,11 @@ type FakeLocalParticipant struct { handleSimulateScenarioReturnsOnCall map[int]struct { result1 error } + HandleStoreDataBlobRequestStub func(*livekit.StoreDataBlobRequest) + handleStoreDataBlobRequestMutex sync.RWMutex + handleStoreDataBlobRequestArgsForCall []struct { + arg1 *livekit.StoreDataBlobRequest + } HandleSyncStateStub func(*livekit.SyncState) error handleSyncStateMutex sync.RWMutex handleSyncStateArgsForCall []struct { @@ -1012,10 +1012,10 @@ type FakeLocalParticipant struct { arg2 chan string arg3 chan error } - ProcessGetDataTrackSchemaRequestStub func(*livekit.GetDataTrackSchemaRequest, types.Participant) - processGetDataTrackSchemaRequestMutex sync.RWMutex - processGetDataTrackSchemaRequestArgsForCall []struct { - arg1 *livekit.GetDataTrackSchemaRequest + ProcessGetDataBlobRequestStub func(*livekit.GetDataBlobRequest, types.Participant) + processGetDataBlobRequestMutex sync.RWMutex + processGetDataBlobRequestArgsForCall []struct { + arg1 *livekit.GetDataBlobRequest arg2 types.Participant } ProtocolVersionStub func() types.ProtocolVersion @@ -1626,35 +1626,35 @@ func (fake *FakeLocalParticipant) ActiveAtReturnsOnCall(i int, result1 time.Time }{result1} } -func (fake *FakeLocalParticipant) AddDataTrackSchema(arg1 *livekit.DataTrackSchemaDefinition) { - fake.addDataTrackSchemaMutex.Lock() - fake.addDataTrackSchemaArgsForCall = append(fake.addDataTrackSchemaArgsForCall, struct { - arg1 *livekit.DataTrackSchemaDefinition +func (fake *FakeLocalParticipant) AddDataBlob(arg1 *livekit.DataBlob) { + fake.addDataBlobMutex.Lock() + fake.addDataBlobArgsForCall = append(fake.addDataBlobArgsForCall, struct { + arg1 *livekit.DataBlob }{arg1}) - stub := fake.AddDataTrackSchemaStub - fake.recordInvocation("AddDataTrackSchema", []interface{}{arg1}) - fake.addDataTrackSchemaMutex.Unlock() + stub := fake.AddDataBlobStub + fake.recordInvocation("AddDataBlob", []interface{}{arg1}) + fake.addDataBlobMutex.Unlock() if stub != nil { - fake.AddDataTrackSchemaStub(arg1) + fake.AddDataBlobStub(arg1) } } -func (fake *FakeLocalParticipant) AddDataTrackSchemaCallCount() int { - fake.addDataTrackSchemaMutex.RLock() - defer fake.addDataTrackSchemaMutex.RUnlock() - return len(fake.addDataTrackSchemaArgsForCall) +func (fake *FakeLocalParticipant) AddDataBlobCallCount() int { + fake.addDataBlobMutex.RLock() + defer fake.addDataBlobMutex.RUnlock() + return len(fake.addDataBlobArgsForCall) } -func (fake *FakeLocalParticipant) AddDataTrackSchemaCalls(stub func(*livekit.DataTrackSchemaDefinition)) { - fake.addDataTrackSchemaMutex.Lock() - defer fake.addDataTrackSchemaMutex.Unlock() - fake.AddDataTrackSchemaStub = stub +func (fake *FakeLocalParticipant) AddDataBlobCalls(stub func(*livekit.DataBlob)) { + fake.addDataBlobMutex.Lock() + defer fake.addDataBlobMutex.Unlock() + fake.AddDataBlobStub = stub } -func (fake *FakeLocalParticipant) AddDataTrackSchemaArgsForCall(i int) *livekit.DataTrackSchemaDefinition { - fake.addDataTrackSchemaMutex.RLock() - defer fake.addDataTrackSchemaMutex.RUnlock() - argsForCall := fake.addDataTrackSchemaArgsForCall[i] +func (fake *FakeLocalParticipant) AddDataBlobArgsForCall(i int) *livekit.DataBlob { + fake.addDataBlobMutex.RLock() + defer fake.addDataBlobMutex.RUnlock() + argsForCall := fake.addDataBlobArgsForCall[i] return argsForCall.arg1 } @@ -2603,15 +2603,15 @@ func (fake *FakeLocalParticipant) GetAdaptiveStreamReturnsOnCall(i int, result1 }{result1} } -func (fake *FakeLocalParticipant) GetAllAsyncAttributes() []*livekit.DataTrackSchemaDefinition { - fake.getAllAsyncAttributesMutex.Lock() - ret, specificReturn := fake.getAllAsyncAttributesReturnsOnCall[len(fake.getAllAsyncAttributesArgsForCall)] - fake.getAllAsyncAttributesArgsForCall = append(fake.getAllAsyncAttributesArgsForCall, struct { +func (fake *FakeLocalParticipant) GetAllDataBlob() []*livekit.DataBlob { + fake.getAllDataBlobMutex.Lock() + ret, specificReturn := fake.getAllDataBlobReturnsOnCall[len(fake.getAllDataBlobArgsForCall)] + fake.getAllDataBlobArgsForCall = append(fake.getAllDataBlobArgsForCall, struct { }{}) - stub := fake.GetAllAsyncAttributesStub - fakeReturns := fake.getAllAsyncAttributesReturns - fake.recordInvocation("GetAllAsyncAttributes", []interface{}{}) - fake.getAllAsyncAttributesMutex.Unlock() + stub := fake.GetAllDataBlobStub + fakeReturns := fake.getAllDataBlobReturns + fake.recordInvocation("GetAllDataBlob", []interface{}{}) + fake.getAllDataBlobMutex.Unlock() if stub != nil { return stub() } @@ -2621,38 +2621,38 @@ func (fake *FakeLocalParticipant) GetAllAsyncAttributes() []*livekit.DataTrackSc return fakeReturns.result1 } -func (fake *FakeLocalParticipant) GetAllAsyncAttributesCallCount() int { - fake.getAllAsyncAttributesMutex.RLock() - defer fake.getAllAsyncAttributesMutex.RUnlock() - return len(fake.getAllAsyncAttributesArgsForCall) +func (fake *FakeLocalParticipant) GetAllDataBlobCallCount() int { + fake.getAllDataBlobMutex.RLock() + defer fake.getAllDataBlobMutex.RUnlock() + return len(fake.getAllDataBlobArgsForCall) } -func (fake *FakeLocalParticipant) GetAllAsyncAttributesCalls(stub func() []*livekit.DataTrackSchemaDefinition) { - fake.getAllAsyncAttributesMutex.Lock() - defer fake.getAllAsyncAttributesMutex.Unlock() - fake.GetAllAsyncAttributesStub = stub +func (fake *FakeLocalParticipant) GetAllDataBlobCalls(stub func() []*livekit.DataBlob) { + fake.getAllDataBlobMutex.Lock() + defer fake.getAllDataBlobMutex.Unlock() + fake.GetAllDataBlobStub = stub } -func (fake *FakeLocalParticipant) GetAllAsyncAttributesReturns(result1 []*livekit.DataTrackSchemaDefinition) { - fake.getAllAsyncAttributesMutex.Lock() - defer fake.getAllAsyncAttributesMutex.Unlock() - fake.GetAllAsyncAttributesStub = nil - fake.getAllAsyncAttributesReturns = struct { - result1 []*livekit.DataTrackSchemaDefinition +func (fake *FakeLocalParticipant) GetAllDataBlobReturns(result1 []*livekit.DataBlob) { + fake.getAllDataBlobMutex.Lock() + defer fake.getAllDataBlobMutex.Unlock() + fake.GetAllDataBlobStub = nil + fake.getAllDataBlobReturns = struct { + result1 []*livekit.DataBlob }{result1} } -func (fake *FakeLocalParticipant) GetAllAsyncAttributesReturnsOnCall(i int, result1 []*livekit.DataTrackSchemaDefinition) { - fake.getAllAsyncAttributesMutex.Lock() - defer fake.getAllAsyncAttributesMutex.Unlock() - fake.GetAllAsyncAttributesStub = nil - if fake.getAllAsyncAttributesReturnsOnCall == nil { - fake.getAllAsyncAttributesReturnsOnCall = make(map[int]struct { - result1 []*livekit.DataTrackSchemaDefinition +func (fake *FakeLocalParticipant) GetAllDataBlobReturnsOnCall(i int, result1 []*livekit.DataBlob) { + fake.getAllDataBlobMutex.Lock() + defer fake.getAllDataBlobMutex.Unlock() + fake.GetAllDataBlobStub = nil + if fake.getAllDataBlobReturnsOnCall == nil { + fake.getAllDataBlobReturnsOnCall = make(map[int]struct { + result1 []*livekit.DataBlob }) } - fake.getAllAsyncAttributesReturnsOnCall[i] = struct { - result1 []*livekit.DataTrackSchemaDefinition + fake.getAllDataBlobReturnsOnCall[i] = struct { + result1 []*livekit.DataBlob }{result1} } @@ -3100,16 +3100,16 @@ func (fake *FakeLocalParticipant) GetCountryReturnsOnCall(i int, result1 string) }{result1} } -func (fake *FakeLocalParticipant) GetDataTrackSchema(arg1 *livekit.DataTrackSchemaId) *livekit.DataTrackSchemaDefinition { - fake.getDataTrackSchemaMutex.Lock() - ret, specificReturn := fake.getDataTrackSchemaReturnsOnCall[len(fake.getDataTrackSchemaArgsForCall)] - fake.getDataTrackSchemaArgsForCall = append(fake.getDataTrackSchemaArgsForCall, struct { - arg1 *livekit.DataTrackSchemaId +func (fake *FakeLocalParticipant) GetDataBlob(arg1 *livekit.DataBlobKey) *livekit.DataBlob { + fake.getDataBlobMutex.Lock() + ret, specificReturn := fake.getDataBlobReturnsOnCall[len(fake.getDataBlobArgsForCall)] + fake.getDataBlobArgsForCall = append(fake.getDataBlobArgsForCall, struct { + arg1 *livekit.DataBlobKey }{arg1}) - stub := fake.GetDataTrackSchemaStub - fakeReturns := fake.getDataTrackSchemaReturns - fake.recordInvocation("GetDataTrackSchema", []interface{}{arg1}) - fake.getDataTrackSchemaMutex.Unlock() + stub := fake.GetDataBlobStub + fakeReturns := fake.getDataBlobReturns + fake.recordInvocation("GetDataBlob", []interface{}{arg1}) + fake.getDataBlobMutex.Unlock() if stub != nil { return stub(arg1) } @@ -3119,45 +3119,45 @@ func (fake *FakeLocalParticipant) GetDataTrackSchema(arg1 *livekit.DataTrackSche return fakeReturns.result1 } -func (fake *FakeLocalParticipant) GetDataTrackSchemaCallCount() int { - fake.getDataTrackSchemaMutex.RLock() - defer fake.getDataTrackSchemaMutex.RUnlock() - return len(fake.getDataTrackSchemaArgsForCall) +func (fake *FakeLocalParticipant) GetDataBlobCallCount() int { + fake.getDataBlobMutex.RLock() + defer fake.getDataBlobMutex.RUnlock() + return len(fake.getDataBlobArgsForCall) } -func (fake *FakeLocalParticipant) GetDataTrackSchemaCalls(stub func(*livekit.DataTrackSchemaId) *livekit.DataTrackSchemaDefinition) { - fake.getDataTrackSchemaMutex.Lock() - defer fake.getDataTrackSchemaMutex.Unlock() - fake.GetDataTrackSchemaStub = stub +func (fake *FakeLocalParticipant) GetDataBlobCalls(stub func(*livekit.DataBlobKey) *livekit.DataBlob) { + fake.getDataBlobMutex.Lock() + defer fake.getDataBlobMutex.Unlock() + fake.GetDataBlobStub = stub } -func (fake *FakeLocalParticipant) GetDataTrackSchemaArgsForCall(i int) *livekit.DataTrackSchemaId { - fake.getDataTrackSchemaMutex.RLock() - defer fake.getDataTrackSchemaMutex.RUnlock() - argsForCall := fake.getDataTrackSchemaArgsForCall[i] +func (fake *FakeLocalParticipant) GetDataBlobArgsForCall(i int) *livekit.DataBlobKey { + fake.getDataBlobMutex.RLock() + defer fake.getDataBlobMutex.RUnlock() + argsForCall := fake.getDataBlobArgsForCall[i] return argsForCall.arg1 } -func (fake *FakeLocalParticipant) GetDataTrackSchemaReturns(result1 *livekit.DataTrackSchemaDefinition) { - fake.getDataTrackSchemaMutex.Lock() - defer fake.getDataTrackSchemaMutex.Unlock() - fake.GetDataTrackSchemaStub = nil - fake.getDataTrackSchemaReturns = struct { - result1 *livekit.DataTrackSchemaDefinition +func (fake *FakeLocalParticipant) GetDataBlobReturns(result1 *livekit.DataBlob) { + fake.getDataBlobMutex.Lock() + defer fake.getDataBlobMutex.Unlock() + fake.GetDataBlobStub = nil + fake.getDataBlobReturns = struct { + result1 *livekit.DataBlob }{result1} } -func (fake *FakeLocalParticipant) GetDataTrackSchemaReturnsOnCall(i int, result1 *livekit.DataTrackSchemaDefinition) { - fake.getDataTrackSchemaMutex.Lock() - defer fake.getDataTrackSchemaMutex.Unlock() - fake.GetDataTrackSchemaStub = nil - if fake.getDataTrackSchemaReturnsOnCall == nil { - fake.getDataTrackSchemaReturnsOnCall = make(map[int]struct { - result1 *livekit.DataTrackSchemaDefinition +func (fake *FakeLocalParticipant) GetDataBlobReturnsOnCall(i int, result1 *livekit.DataBlob) { + fake.getDataBlobMutex.Lock() + defer fake.getDataBlobMutex.Unlock() + fake.GetDataBlobStub = nil + if fake.getDataBlobReturnsOnCall == nil { + fake.getDataBlobReturnsOnCall = make(map[int]struct { + result1 *livekit.DataBlob }) } - fake.getDataTrackSchemaReturnsOnCall[i] = struct { - result1 *livekit.DataTrackSchemaDefinition + fake.getDataBlobReturnsOnCall[i] = struct { + result1 *livekit.DataBlob }{result1} } @@ -4553,67 +4553,35 @@ func (fake *FakeLocalParticipant) HandleAnswerArgsForCall(i int) *livekit.Sessio return argsForCall.arg1 } -func (fake *FakeLocalParticipant) HandleDefineDataTrackSchemaRequest(arg1 *livekit.DefineDataTrackSchemaRequest) { - fake.handleDefineDataTrackSchemaRequestMutex.Lock() - fake.handleDefineDataTrackSchemaRequestArgsForCall = append(fake.handleDefineDataTrackSchemaRequestArgsForCall, struct { - arg1 *livekit.DefineDataTrackSchemaRequest +func (fake *FakeLocalParticipant) HandleGetDataBlobRequest(arg1 *livekit.GetDataBlobRequest) { + fake.handleGetDataBlobRequestMutex.Lock() + fake.handleGetDataBlobRequestArgsForCall = append(fake.handleGetDataBlobRequestArgsForCall, struct { + arg1 *livekit.GetDataBlobRequest }{arg1}) - stub := fake.HandleDefineDataTrackSchemaRequestStub - fake.recordInvocation("HandleDefineDataTrackSchemaRequest", []interface{}{arg1}) - fake.handleDefineDataTrackSchemaRequestMutex.Unlock() + stub := fake.HandleGetDataBlobRequestStub + fake.recordInvocation("HandleGetDataBlobRequest", []interface{}{arg1}) + fake.handleGetDataBlobRequestMutex.Unlock() if stub != nil { - fake.HandleDefineDataTrackSchemaRequestStub(arg1) + fake.HandleGetDataBlobRequestStub(arg1) } } -func (fake *FakeLocalParticipant) HandleDefineDataTrackSchemaRequestCallCount() int { - fake.handleDefineDataTrackSchemaRequestMutex.RLock() - defer fake.handleDefineDataTrackSchemaRequestMutex.RUnlock() - return len(fake.handleDefineDataTrackSchemaRequestArgsForCall) +func (fake *FakeLocalParticipant) HandleGetDataBlobRequestCallCount() int { + fake.handleGetDataBlobRequestMutex.RLock() + defer fake.handleGetDataBlobRequestMutex.RUnlock() + return len(fake.handleGetDataBlobRequestArgsForCall) } -func (fake *FakeLocalParticipant) HandleDefineDataTrackSchemaRequestCalls(stub func(*livekit.DefineDataTrackSchemaRequest)) { - fake.handleDefineDataTrackSchemaRequestMutex.Lock() - defer fake.handleDefineDataTrackSchemaRequestMutex.Unlock() - fake.HandleDefineDataTrackSchemaRequestStub = stub +func (fake *FakeLocalParticipant) HandleGetDataBlobRequestCalls(stub func(*livekit.GetDataBlobRequest)) { + fake.handleGetDataBlobRequestMutex.Lock() + defer fake.handleGetDataBlobRequestMutex.Unlock() + fake.HandleGetDataBlobRequestStub = stub } -func (fake *FakeLocalParticipant) HandleDefineDataTrackSchemaRequestArgsForCall(i int) *livekit.DefineDataTrackSchemaRequest { - fake.handleDefineDataTrackSchemaRequestMutex.RLock() - defer fake.handleDefineDataTrackSchemaRequestMutex.RUnlock() - argsForCall := fake.handleDefineDataTrackSchemaRequestArgsForCall[i] - return argsForCall.arg1 -} - -func (fake *FakeLocalParticipant) HandleGetDataTrackSchemaRequest(arg1 *livekit.GetDataTrackSchemaRequest) { - fake.handleGetDataTrackSchemaRequestMutex.Lock() - fake.handleGetDataTrackSchemaRequestArgsForCall = append(fake.handleGetDataTrackSchemaRequestArgsForCall, struct { - arg1 *livekit.GetDataTrackSchemaRequest - }{arg1}) - stub := fake.HandleGetDataTrackSchemaRequestStub - fake.recordInvocation("HandleGetDataTrackSchemaRequest", []interface{}{arg1}) - fake.handleGetDataTrackSchemaRequestMutex.Unlock() - if stub != nil { - fake.HandleGetDataTrackSchemaRequestStub(arg1) - } -} - -func (fake *FakeLocalParticipant) HandleGetDataTrackSchemaRequestCallCount() int { - fake.handleGetDataTrackSchemaRequestMutex.RLock() - defer fake.handleGetDataTrackSchemaRequestMutex.RUnlock() - return len(fake.handleGetDataTrackSchemaRequestArgsForCall) -} - -func (fake *FakeLocalParticipant) HandleGetDataTrackSchemaRequestCalls(stub func(*livekit.GetDataTrackSchemaRequest)) { - fake.handleGetDataTrackSchemaRequestMutex.Lock() - defer fake.handleGetDataTrackSchemaRequestMutex.Unlock() - fake.HandleGetDataTrackSchemaRequestStub = stub -} - -func (fake *FakeLocalParticipant) HandleGetDataTrackSchemaRequestArgsForCall(i int) *livekit.GetDataTrackSchemaRequest { - fake.handleGetDataTrackSchemaRequestMutex.RLock() - defer fake.handleGetDataTrackSchemaRequestMutex.RUnlock() - argsForCall := fake.handleGetDataTrackSchemaRequestArgsForCall[i] +func (fake *FakeLocalParticipant) HandleGetDataBlobRequestArgsForCall(i int) *livekit.GetDataBlobRequest { + fake.handleGetDataBlobRequestMutex.RLock() + defer fake.handleGetDataBlobRequestMutex.RUnlock() + argsForCall := fake.handleGetDataBlobRequestArgsForCall[i] return argsForCall.arg1 } @@ -5241,6 +5209,38 @@ func (fake *FakeLocalParticipant) HandleSimulateScenarioReturnsOnCall(i int, res }{result1} } +func (fake *FakeLocalParticipant) HandleStoreDataBlobRequest(arg1 *livekit.StoreDataBlobRequest) { + fake.handleStoreDataBlobRequestMutex.Lock() + fake.handleStoreDataBlobRequestArgsForCall = append(fake.handleStoreDataBlobRequestArgsForCall, struct { + arg1 *livekit.StoreDataBlobRequest + }{arg1}) + stub := fake.HandleStoreDataBlobRequestStub + fake.recordInvocation("HandleStoreDataBlobRequest", []interface{}{arg1}) + fake.handleStoreDataBlobRequestMutex.Unlock() + if stub != nil { + fake.HandleStoreDataBlobRequestStub(arg1) + } +} + +func (fake *FakeLocalParticipant) HandleStoreDataBlobRequestCallCount() int { + fake.handleStoreDataBlobRequestMutex.RLock() + defer fake.handleStoreDataBlobRequestMutex.RUnlock() + return len(fake.handleStoreDataBlobRequestArgsForCall) +} + +func (fake *FakeLocalParticipant) HandleStoreDataBlobRequestCalls(stub func(*livekit.StoreDataBlobRequest)) { + fake.handleStoreDataBlobRequestMutex.Lock() + defer fake.handleStoreDataBlobRequestMutex.Unlock() + fake.HandleStoreDataBlobRequestStub = stub +} + +func (fake *FakeLocalParticipant) HandleStoreDataBlobRequestArgsForCall(i int) *livekit.StoreDataBlobRequest { + fake.handleStoreDataBlobRequestMutex.RLock() + defer fake.handleStoreDataBlobRequestMutex.RUnlock() + argsForCall := fake.handleStoreDataBlobRequestArgsForCall[i] + return argsForCall.arg1 +} + func (fake *FakeLocalParticipant) HandleSyncState(arg1 *livekit.SyncState) error { fake.handleSyncStateMutex.Lock() ret, specificReturn := fake.handleSyncStateReturnsOnCall[len(fake.handleSyncStateArgsForCall)] @@ -6869,36 +6869,36 @@ func (fake *FakeLocalParticipant) PerformRpcArgsForCall(i int) (*livekit.Perform return argsForCall.arg1, argsForCall.arg2, argsForCall.arg3 } -func (fake *FakeLocalParticipant) ProcessGetDataTrackSchemaRequest(arg1 *livekit.GetDataTrackSchemaRequest, arg2 types.Participant) { - fake.processGetDataTrackSchemaRequestMutex.Lock() - fake.processGetDataTrackSchemaRequestArgsForCall = append(fake.processGetDataTrackSchemaRequestArgsForCall, struct { - arg1 *livekit.GetDataTrackSchemaRequest +func (fake *FakeLocalParticipant) ProcessGetDataBlobRequest(arg1 *livekit.GetDataBlobRequest, arg2 types.Participant) { + fake.processGetDataBlobRequestMutex.Lock() + fake.processGetDataBlobRequestArgsForCall = append(fake.processGetDataBlobRequestArgsForCall, struct { + arg1 *livekit.GetDataBlobRequest arg2 types.Participant }{arg1, arg2}) - stub := fake.ProcessGetDataTrackSchemaRequestStub - fake.recordInvocation("ProcessGetDataTrackSchemaRequest", []interface{}{arg1, arg2}) - fake.processGetDataTrackSchemaRequestMutex.Unlock() + stub := fake.ProcessGetDataBlobRequestStub + fake.recordInvocation("ProcessGetDataBlobRequest", []interface{}{arg1, arg2}) + fake.processGetDataBlobRequestMutex.Unlock() if stub != nil { - fake.ProcessGetDataTrackSchemaRequestStub(arg1, arg2) + fake.ProcessGetDataBlobRequestStub(arg1, arg2) } } -func (fake *FakeLocalParticipant) ProcessGetDataTrackSchemaRequestCallCount() int { - fake.processGetDataTrackSchemaRequestMutex.RLock() - defer fake.processGetDataTrackSchemaRequestMutex.RUnlock() - return len(fake.processGetDataTrackSchemaRequestArgsForCall) +func (fake *FakeLocalParticipant) ProcessGetDataBlobRequestCallCount() int { + fake.processGetDataBlobRequestMutex.RLock() + defer fake.processGetDataBlobRequestMutex.RUnlock() + return len(fake.processGetDataBlobRequestArgsForCall) } -func (fake *FakeLocalParticipant) ProcessGetDataTrackSchemaRequestCalls(stub func(*livekit.GetDataTrackSchemaRequest, types.Participant)) { - fake.processGetDataTrackSchemaRequestMutex.Lock() - defer fake.processGetDataTrackSchemaRequestMutex.Unlock() - fake.ProcessGetDataTrackSchemaRequestStub = stub +func (fake *FakeLocalParticipant) ProcessGetDataBlobRequestCalls(stub func(*livekit.GetDataBlobRequest, types.Participant)) { + fake.processGetDataBlobRequestMutex.Lock() + defer fake.processGetDataBlobRequestMutex.Unlock() + fake.ProcessGetDataBlobRequestStub = stub } -func (fake *FakeLocalParticipant) ProcessGetDataTrackSchemaRequestArgsForCall(i int) (*livekit.GetDataTrackSchemaRequest, types.Participant) { - fake.processGetDataTrackSchemaRequestMutex.RLock() - defer fake.processGetDataTrackSchemaRequestMutex.RUnlock() - argsForCall := fake.processGetDataTrackSchemaRequestArgsForCall[i] +func (fake *FakeLocalParticipant) ProcessGetDataBlobRequestArgsForCall(i int) (*livekit.GetDataBlobRequest, types.Participant) { + fake.processGetDataBlobRequestMutex.RLock() + defer fake.processGetDataBlobRequestMutex.RUnlock() + argsForCall := fake.processGetDataBlobRequestArgsForCall[i] return argsForCall.arg1, argsForCall.arg2 } diff --git a/pkg/rtc/types/typesfakes/fake_local_participant_listener.go b/pkg/rtc/types/typesfakes/fake_local_participant_listener.go index d585f2d8c..4c0a3a07a 100644 --- a/pkg/rtc/types/typesfakes/fake_local_participant_listener.go +++ b/pkg/rtc/types/typesfakes/fake_local_participant_listener.go @@ -42,17 +42,11 @@ type FakeLocalParticipantListener struct { arg1 types.Participant arg2 types.DataTrack } - OnDefineDataTrackSchemaStub func(types.LocalParticipant, *livekit.DataTrackSchemaDefinition) - onDefineDataTrackSchemaMutex sync.RWMutex - onDefineDataTrackSchemaArgsForCall []struct { + OnGetDataBlobStub func(types.LocalParticipant, *livekit.GetDataBlobRequest) + onGetDataBlobMutex sync.RWMutex + onGetDataBlobArgsForCall []struct { arg1 types.LocalParticipant - arg2 *livekit.DataTrackSchemaDefinition - } - OnGetDataTrackSchemaStub func(types.LocalParticipant, *livekit.GetDataTrackSchemaRequest) - onGetDataTrackSchemaMutex sync.RWMutex - onGetDataTrackSchemaArgsForCall []struct { - arg1 types.LocalParticipant - arg2 *livekit.GetDataTrackSchemaRequest + arg2 *livekit.GetDataBlobRequest } OnLeaveStub func(types.LocalParticipant, types.ParticipantCloseReason) onLeaveMutex sync.RWMutex @@ -94,6 +88,12 @@ type FakeLocalParticipantListener struct { onStateChangeArgsForCall []struct { arg1 types.LocalParticipant } + OnStoreDataBlobStub func(types.LocalParticipant, *livekit.DataBlob) + onStoreDataBlobMutex sync.RWMutex + onStoreDataBlobArgsForCall []struct { + arg1 types.LocalParticipant + arg2 *livekit.DataBlob + } OnSubscribeStatusChangedStub func(types.LocalParticipant, livekit.ParticipantID, bool) onSubscribeStatusChangedMutex sync.RWMutex onSubscribeStatusChangedArgsForCall []struct { @@ -343,69 +343,36 @@ func (fake *FakeLocalParticipantListener) OnDataTrackUnpublishedArgsForCall(i in return argsForCall.arg1, argsForCall.arg2 } -func (fake *FakeLocalParticipantListener) OnDefineDataTrackSchema(arg1 types.LocalParticipant, arg2 *livekit.DataTrackSchemaDefinition) { - fake.onDefineDataTrackSchemaMutex.Lock() - fake.onDefineDataTrackSchemaArgsForCall = append(fake.onDefineDataTrackSchemaArgsForCall, struct { +func (fake *FakeLocalParticipantListener) OnGetDataBlob(arg1 types.LocalParticipant, arg2 *livekit.GetDataBlobRequest) { + fake.onGetDataBlobMutex.Lock() + fake.onGetDataBlobArgsForCall = append(fake.onGetDataBlobArgsForCall, struct { arg1 types.LocalParticipant - arg2 *livekit.DataTrackSchemaDefinition + arg2 *livekit.GetDataBlobRequest }{arg1, arg2}) - stub := fake.OnDefineDataTrackSchemaStub - fake.recordInvocation("OnDefineDataTrackSchema", []interface{}{arg1, arg2}) - fake.onDefineDataTrackSchemaMutex.Unlock() + stub := fake.OnGetDataBlobStub + fake.recordInvocation("OnGetDataBlob", []interface{}{arg1, arg2}) + fake.onGetDataBlobMutex.Unlock() if stub != nil { - fake.OnDefineDataTrackSchemaStub(arg1, arg2) + fake.OnGetDataBlobStub(arg1, arg2) } } -func (fake *FakeLocalParticipantListener) OnDefineDataTrackSchemaCallCount() int { - fake.onDefineDataTrackSchemaMutex.RLock() - defer fake.onDefineDataTrackSchemaMutex.RUnlock() - return len(fake.onDefineDataTrackSchemaArgsForCall) +func (fake *FakeLocalParticipantListener) OnGetDataBlobCallCount() int { + fake.onGetDataBlobMutex.RLock() + defer fake.onGetDataBlobMutex.RUnlock() + return len(fake.onGetDataBlobArgsForCall) } -func (fake *FakeLocalParticipantListener) OnDefineDataTrackSchemaCalls(stub func(types.LocalParticipant, *livekit.DataTrackSchemaDefinition)) { - fake.onDefineDataTrackSchemaMutex.Lock() - defer fake.onDefineDataTrackSchemaMutex.Unlock() - fake.OnDefineDataTrackSchemaStub = stub +func (fake *FakeLocalParticipantListener) OnGetDataBlobCalls(stub func(types.LocalParticipant, *livekit.GetDataBlobRequest)) { + fake.onGetDataBlobMutex.Lock() + defer fake.onGetDataBlobMutex.Unlock() + fake.OnGetDataBlobStub = stub } -func (fake *FakeLocalParticipantListener) OnDefineDataTrackSchemaArgsForCall(i int) (types.LocalParticipant, *livekit.DataTrackSchemaDefinition) { - fake.onDefineDataTrackSchemaMutex.RLock() - defer fake.onDefineDataTrackSchemaMutex.RUnlock() - argsForCall := fake.onDefineDataTrackSchemaArgsForCall[i] - return argsForCall.arg1, argsForCall.arg2 -} - -func (fake *FakeLocalParticipantListener) OnGetDataTrackSchema(arg1 types.LocalParticipant, arg2 *livekit.GetDataTrackSchemaRequest) { - fake.onGetDataTrackSchemaMutex.Lock() - fake.onGetDataTrackSchemaArgsForCall = append(fake.onGetDataTrackSchemaArgsForCall, struct { - arg1 types.LocalParticipant - arg2 *livekit.GetDataTrackSchemaRequest - }{arg1, arg2}) - stub := fake.OnGetDataTrackSchemaStub - fake.recordInvocation("OnGetDataTrackSchema", []interface{}{arg1, arg2}) - fake.onGetDataTrackSchemaMutex.Unlock() - if stub != nil { - fake.OnGetDataTrackSchemaStub(arg1, arg2) - } -} - -func (fake *FakeLocalParticipantListener) OnGetDataTrackSchemaCallCount() int { - fake.onGetDataTrackSchemaMutex.RLock() - defer fake.onGetDataTrackSchemaMutex.RUnlock() - return len(fake.onGetDataTrackSchemaArgsForCall) -} - -func (fake *FakeLocalParticipantListener) OnGetDataTrackSchemaCalls(stub func(types.LocalParticipant, *livekit.GetDataTrackSchemaRequest)) { - fake.onGetDataTrackSchemaMutex.Lock() - defer fake.onGetDataTrackSchemaMutex.Unlock() - fake.OnGetDataTrackSchemaStub = stub -} - -func (fake *FakeLocalParticipantListener) OnGetDataTrackSchemaArgsForCall(i int) (types.LocalParticipant, *livekit.GetDataTrackSchemaRequest) { - fake.onGetDataTrackSchemaMutex.RLock() - defer fake.onGetDataTrackSchemaMutex.RUnlock() - argsForCall := fake.onGetDataTrackSchemaArgsForCall[i] +func (fake *FakeLocalParticipantListener) OnGetDataBlobArgsForCall(i int) (types.LocalParticipant, *livekit.GetDataBlobRequest) { + fake.onGetDataBlobMutex.RLock() + defer fake.onGetDataBlobMutex.RUnlock() + argsForCall := fake.onGetDataBlobArgsForCall[i] return argsForCall.arg1, argsForCall.arg2 } @@ -634,6 +601,39 @@ func (fake *FakeLocalParticipantListener) OnStateChangeArgsForCall(i int) types. return argsForCall.arg1 } +func (fake *FakeLocalParticipantListener) OnStoreDataBlob(arg1 types.LocalParticipant, arg2 *livekit.DataBlob) { + fake.onStoreDataBlobMutex.Lock() + fake.onStoreDataBlobArgsForCall = append(fake.onStoreDataBlobArgsForCall, struct { + arg1 types.LocalParticipant + arg2 *livekit.DataBlob + }{arg1, arg2}) + stub := fake.OnStoreDataBlobStub + fake.recordInvocation("OnStoreDataBlob", []interface{}{arg1, arg2}) + fake.onStoreDataBlobMutex.Unlock() + if stub != nil { + fake.OnStoreDataBlobStub(arg1, arg2) + } +} + +func (fake *FakeLocalParticipantListener) OnStoreDataBlobCallCount() int { + fake.onStoreDataBlobMutex.RLock() + defer fake.onStoreDataBlobMutex.RUnlock() + return len(fake.onStoreDataBlobArgsForCall) +} + +func (fake *FakeLocalParticipantListener) OnStoreDataBlobCalls(stub func(types.LocalParticipant, *livekit.DataBlob)) { + fake.onStoreDataBlobMutex.Lock() + defer fake.onStoreDataBlobMutex.Unlock() + fake.OnStoreDataBlobStub = stub +} + +func (fake *FakeLocalParticipantListener) OnStoreDataBlobArgsForCall(i int) (types.LocalParticipant, *livekit.DataBlob) { + fake.onStoreDataBlobMutex.RLock() + defer fake.onStoreDataBlobMutex.RUnlock() + argsForCall := fake.onStoreDataBlobArgsForCall[i] + return argsForCall.arg1, argsForCall.arg2 +} + func (fake *FakeLocalParticipantListener) OnSubscribeStatusChanged(arg1 types.LocalParticipant, arg2 livekit.ParticipantID, arg3 bool) { fake.onSubscribeStatusChangedMutex.Lock() fake.onSubscribeStatusChangedArgsForCall = append(fake.onSubscribeStatusChangedArgsForCall, struct { diff --git a/pkg/rtc/types/typesfakes/fake_participant.go b/pkg/rtc/types/typesfakes/fake_participant.go index 5b5d38919..652c8bf14 100644 --- a/pkg/rtc/types/typesfakes/fake_participant.go +++ b/pkg/rtc/types/typesfakes/fake_participant.go @@ -13,10 +13,10 @@ import ( ) type FakeParticipant struct { - AddDataTrackSchemaStub func(*livekit.DataTrackSchemaDefinition) - addDataTrackSchemaMutex sync.RWMutex - addDataTrackSchemaArgsForCall []struct { - arg1 *livekit.DataTrackSchemaDefinition + AddDataBlobStub func(*livekit.DataBlob) + addDataBlobMutex sync.RWMutex + addDataBlobArgsForCall []struct { + arg1 *livekit.DataBlob } CanSkipBroadcastStub func() bool canSkipBroadcastMutex sync.RWMutex @@ -83,16 +83,16 @@ type FakeParticipant struct { result1 float64 result2 bool } - GetDataTrackSchemaStub func(*livekit.DataTrackSchemaId) *livekit.DataTrackSchemaDefinition - getDataTrackSchemaMutex sync.RWMutex - getDataTrackSchemaArgsForCall []struct { - arg1 *livekit.DataTrackSchemaId + GetDataBlobStub func(*livekit.DataBlobKey) *livekit.DataBlob + getDataBlobMutex sync.RWMutex + getDataBlobArgsForCall []struct { + arg1 *livekit.DataBlobKey } - getDataTrackSchemaReturns struct { - result1 *livekit.DataTrackSchemaDefinition + getDataBlobReturns struct { + result1 *livekit.DataBlob } - getDataTrackSchemaReturnsOnCall map[int]struct { - result1 *livekit.DataTrackSchemaDefinition + getDataBlobReturnsOnCall map[int]struct { + result1 *livekit.DataBlob } GetLoggerStub func() logger.Logger getLoggerMutex sync.RWMutex @@ -377,35 +377,35 @@ type FakeParticipant struct { invocationsMutex sync.RWMutex } -func (fake *FakeParticipant) AddDataTrackSchema(arg1 *livekit.DataTrackSchemaDefinition) { - fake.addDataTrackSchemaMutex.Lock() - fake.addDataTrackSchemaArgsForCall = append(fake.addDataTrackSchemaArgsForCall, struct { - arg1 *livekit.DataTrackSchemaDefinition +func (fake *FakeParticipant) AddDataBlob(arg1 *livekit.DataBlob) { + fake.addDataBlobMutex.Lock() + fake.addDataBlobArgsForCall = append(fake.addDataBlobArgsForCall, struct { + arg1 *livekit.DataBlob }{arg1}) - stub := fake.AddDataTrackSchemaStub - fake.recordInvocation("AddDataTrackSchema", []interface{}{arg1}) - fake.addDataTrackSchemaMutex.Unlock() + stub := fake.AddDataBlobStub + fake.recordInvocation("AddDataBlob", []interface{}{arg1}) + fake.addDataBlobMutex.Unlock() if stub != nil { - fake.AddDataTrackSchemaStub(arg1) + fake.AddDataBlobStub(arg1) } } -func (fake *FakeParticipant) AddDataTrackSchemaCallCount() int { - fake.addDataTrackSchemaMutex.RLock() - defer fake.addDataTrackSchemaMutex.RUnlock() - return len(fake.addDataTrackSchemaArgsForCall) +func (fake *FakeParticipant) AddDataBlobCallCount() int { + fake.addDataBlobMutex.RLock() + defer fake.addDataBlobMutex.RUnlock() + return len(fake.addDataBlobArgsForCall) } -func (fake *FakeParticipant) AddDataTrackSchemaCalls(stub func(*livekit.DataTrackSchemaDefinition)) { - fake.addDataTrackSchemaMutex.Lock() - defer fake.addDataTrackSchemaMutex.Unlock() - fake.AddDataTrackSchemaStub = stub +func (fake *FakeParticipant) AddDataBlobCalls(stub func(*livekit.DataBlob)) { + fake.addDataBlobMutex.Lock() + defer fake.addDataBlobMutex.Unlock() + fake.AddDataBlobStub = stub } -func (fake *FakeParticipant) AddDataTrackSchemaArgsForCall(i int) *livekit.DataTrackSchemaDefinition { - fake.addDataTrackSchemaMutex.RLock() - defer fake.addDataTrackSchemaMutex.RUnlock() - argsForCall := fake.addDataTrackSchemaArgsForCall[i] +func (fake *FakeParticipant) AddDataBlobArgsForCall(i int) *livekit.DataBlob { + fake.addDataBlobMutex.RLock() + defer fake.addDataBlobMutex.RUnlock() + argsForCall := fake.addDataBlobArgsForCall[i] return argsForCall.arg1 } @@ -740,16 +740,16 @@ func (fake *FakeParticipant) GetAudioLevelReturnsOnCall(i int, result1 float64, }{result1, result2} } -func (fake *FakeParticipant) GetDataTrackSchema(arg1 *livekit.DataTrackSchemaId) *livekit.DataTrackSchemaDefinition { - fake.getDataTrackSchemaMutex.Lock() - ret, specificReturn := fake.getDataTrackSchemaReturnsOnCall[len(fake.getDataTrackSchemaArgsForCall)] - fake.getDataTrackSchemaArgsForCall = append(fake.getDataTrackSchemaArgsForCall, struct { - arg1 *livekit.DataTrackSchemaId +func (fake *FakeParticipant) GetDataBlob(arg1 *livekit.DataBlobKey) *livekit.DataBlob { + fake.getDataBlobMutex.Lock() + ret, specificReturn := fake.getDataBlobReturnsOnCall[len(fake.getDataBlobArgsForCall)] + fake.getDataBlobArgsForCall = append(fake.getDataBlobArgsForCall, struct { + arg1 *livekit.DataBlobKey }{arg1}) - stub := fake.GetDataTrackSchemaStub - fakeReturns := fake.getDataTrackSchemaReturns - fake.recordInvocation("GetDataTrackSchema", []interface{}{arg1}) - fake.getDataTrackSchemaMutex.Unlock() + stub := fake.GetDataBlobStub + fakeReturns := fake.getDataBlobReturns + fake.recordInvocation("GetDataBlob", []interface{}{arg1}) + fake.getDataBlobMutex.Unlock() if stub != nil { return stub(arg1) } @@ -759,45 +759,45 @@ func (fake *FakeParticipant) GetDataTrackSchema(arg1 *livekit.DataTrackSchemaId) return fakeReturns.result1 } -func (fake *FakeParticipant) GetDataTrackSchemaCallCount() int { - fake.getDataTrackSchemaMutex.RLock() - defer fake.getDataTrackSchemaMutex.RUnlock() - return len(fake.getDataTrackSchemaArgsForCall) +func (fake *FakeParticipant) GetDataBlobCallCount() int { + fake.getDataBlobMutex.RLock() + defer fake.getDataBlobMutex.RUnlock() + return len(fake.getDataBlobArgsForCall) } -func (fake *FakeParticipant) GetDataTrackSchemaCalls(stub func(*livekit.DataTrackSchemaId) *livekit.DataTrackSchemaDefinition) { - fake.getDataTrackSchemaMutex.Lock() - defer fake.getDataTrackSchemaMutex.Unlock() - fake.GetDataTrackSchemaStub = stub +func (fake *FakeParticipant) GetDataBlobCalls(stub func(*livekit.DataBlobKey) *livekit.DataBlob) { + fake.getDataBlobMutex.Lock() + defer fake.getDataBlobMutex.Unlock() + fake.GetDataBlobStub = stub } -func (fake *FakeParticipant) GetDataTrackSchemaArgsForCall(i int) *livekit.DataTrackSchemaId { - fake.getDataTrackSchemaMutex.RLock() - defer fake.getDataTrackSchemaMutex.RUnlock() - argsForCall := fake.getDataTrackSchemaArgsForCall[i] +func (fake *FakeParticipant) GetDataBlobArgsForCall(i int) *livekit.DataBlobKey { + fake.getDataBlobMutex.RLock() + defer fake.getDataBlobMutex.RUnlock() + argsForCall := fake.getDataBlobArgsForCall[i] return argsForCall.arg1 } -func (fake *FakeParticipant) GetDataTrackSchemaReturns(result1 *livekit.DataTrackSchemaDefinition) { - fake.getDataTrackSchemaMutex.Lock() - defer fake.getDataTrackSchemaMutex.Unlock() - fake.GetDataTrackSchemaStub = nil - fake.getDataTrackSchemaReturns = struct { - result1 *livekit.DataTrackSchemaDefinition +func (fake *FakeParticipant) GetDataBlobReturns(result1 *livekit.DataBlob) { + fake.getDataBlobMutex.Lock() + defer fake.getDataBlobMutex.Unlock() + fake.GetDataBlobStub = nil + fake.getDataBlobReturns = struct { + result1 *livekit.DataBlob }{result1} } -func (fake *FakeParticipant) GetDataTrackSchemaReturnsOnCall(i int, result1 *livekit.DataTrackSchemaDefinition) { - fake.getDataTrackSchemaMutex.Lock() - defer fake.getDataTrackSchemaMutex.Unlock() - fake.GetDataTrackSchemaStub = nil - if fake.getDataTrackSchemaReturnsOnCall == nil { - fake.getDataTrackSchemaReturnsOnCall = make(map[int]struct { - result1 *livekit.DataTrackSchemaDefinition +func (fake *FakeParticipant) GetDataBlobReturnsOnCall(i int, result1 *livekit.DataBlob) { + fake.getDataBlobMutex.Lock() + defer fake.getDataBlobMutex.Unlock() + fake.GetDataBlobStub = nil + if fake.getDataBlobReturnsOnCall == nil { + fake.getDataBlobReturnsOnCall = make(map[int]struct { + result1 *livekit.DataBlob }) } - fake.getDataTrackSchemaReturnsOnCall[i] = struct { - result1 *livekit.DataTrackSchemaDefinition + fake.getDataBlobReturnsOnCall[i] = struct { + result1 *livekit.DataBlob }{result1} } diff --git a/pkg/service/roommanager.go b/pkg/service/roommanager.go index 5a2208551..c498e23f8 100644 --- a/pkg/service/roommanager.go +++ b/pkg/service/roommanager.go @@ -523,9 +523,9 @@ func (r *RoomManager) StartSession( DatachannelLossyTargetLatency: r.config.RTC.DatachannelLossyTargetLatency, FireOnTrackBySdp: true, UseSinglePeerConnection: pi.UseSinglePeerConnection, - EnableDataTracks: r.config.EnableDataTracks, - EnableParticipantAsyncAttributes: r.config.EnableParticipantAsyncAttributes, - EnableRTPStreamRestartDetection: r.config.RTC.EnableRTPStreamRestartDetection, + EnableDataTracks: r.config.EnableDataTracks, + EnableParticipantDataBlob: r.config.EnableParticipantDataBlob, + EnableRTPStreamRestartDetection: r.config.RTC.EnableRTPStreamRestartDetection, }) if err != nil { return err diff --git a/test/multinode_test.go b/test/multinode_test.go index c26f9d468..fc66eed08 100644 --- a/test/multinode_test.go +++ b/test/multinode_test.go @@ -427,22 +427,22 @@ func TestCloseDisconnectedParticipantOnSignalClose(t *testing.T) { } } -func TestMultiNodeAsyncAttributes(t *testing.T) { +func TestMultiNodeDataBlob(t *testing.T) { if testing.Short() { t.SkipNow() return } - _, _, finish := setupMultiNodeTestWithConfig("TestMultiNodeAsyncAttributes", func(c *config.Config) { - c.EnableParticipantAsyncAttributes = true - c.Limit.MaxAsyncAttributesSize = 1024 + _, _, finish := setupMultiNodeTestWithConfig("TestMultiNodeDataBlob", func(c *config.Config) { + c.EnableParticipantDataBlob = true + c.Limit.MaxDataBlobSize = 1024 }) defer finish() for _, testRTCServicePath := range testRTCServicePaths { t.Run(fmt.Sprintf("testRTCServicePath=%s", testRTCServicePath.String()), func(t *testing.T) { - pubCapture := &asyncAttributesCapture{} - subCapture := &asyncAttributesCapture{} + pubCapture := &dataBlobCapture{} + subCapture := &dataBlobCapture{} // publisher on node 1, subscriber on node 2 pub := createRTCClient("pub", defaultServerPort, testRTCServicePath, &client.Options{ @@ -464,18 +464,19 @@ func TestMultiNodeAsyncAttributes(t *testing.T) { return "" }) - schemaID := &livekit.DataTrackSchemaId{ - Name: "schema-multinode", - Encoding: livekit.DataTrackSchemaEncoding_DATA_TRACK_SCHEMA_ENCODING_PROTOBUF, + key := &livekit.DataBlobKey{ + Key: &livekit.DataBlobKey_Generic{ + Generic: "blob-multinode", + }, } - definition := []byte("multinode-definition") + contents := []byte("multinode-content") require.NoError(t, pub.SendRequest(&livekit.SignalRequest{ - Message: &livekit.SignalRequest_DefineDataTrackSchema{ - DefineDataTrackSchema: &livekit.DefineDataTrackSchemaRequest{ - SchemaDefinition: &livekit.DataTrackSchemaDefinition{ - Id: schemaID, - Definition: definition, + Message: &livekit.SignalRequest_StoreDataBlobRequest{ + StoreDataBlobRequest: &livekit.StoreDataBlobRequest{ + Blob: &livekit.DataBlob{ + Key: key, + Contents: contents, }, }, }, @@ -485,40 +486,40 @@ func TestMultiNodeAsyncAttributes(t *testing.T) { time.Sleep(syncDelay) require.Equal(t, 0, pubCapture.requestResponseCount(), "publisher should not receive an error response on success") - // subscriber on a different node asks for the schema; the request routes + // subscriber on a different node asks for the blob; the request routes // across nodes to the publisher. require.NoError(t, sub.SendRequest(&livekit.SignalRequest{ - Message: &livekit.SignalRequest_GetDataTrackSchema{ - GetDataTrackSchema: &livekit.GetDataTrackSchemaRequest{ + Message: &livekit.SignalRequest_GetDataBlobRequest{ + GetDataBlobRequest: &livekit.GetDataBlobRequest{ ParticipantIdentity: "pub", - SchemaId: schemaID, + Key: key, }, }, })) testutils.WithTimeout(t, func() string { - resp := subCapture.takeSchemaResponse() + resp := subCapture.takeBlobResponse() if resp == nil { - return "subscriber did not receive schema response" + return "subscriber did not receive blob response" } - if resp.SchemaDefinition == nil { - return "schema response missing definition" + if resp.Blob == nil { + return "blob response missing blob" } - if resp.SchemaDefinition.Id.Name != schemaID.Name { - return fmt.Sprintf("expected schema name %s, got %s", schemaID.Name, resp.SchemaDefinition.Id.Name) + if resp.Blob.Key.String() != key.String() { + return fmt.Sprintf("expected data blob key %s, got %s", key.String(), resp.Blob.Key.String()) } - if string(resp.SchemaDefinition.Definition) != string(definition) { - return fmt.Sprintf("expected definition %q, got %q", definition, resp.SchemaDefinition.Definition) + if string(resp.Blob.Contents) != string(contents) { + return fmt.Sprintf("expected contents %q, got %q", contents, resp.Blob.Contents) } return "" }) // requesting an unknown publisher identity should return NOT_FOUND require.NoError(t, sub.SendRequest(&livekit.SignalRequest{ - Message: &livekit.SignalRequest_GetDataTrackSchema{ - GetDataTrackSchema: &livekit.GetDataTrackSchemaRequest{ + Message: &livekit.SignalRequest_GetDataBlobRequest{ + GetDataBlobRequest: &livekit.GetDataBlobRequest{ ParticipantIdentity: "unknown-publisher", - SchemaId: schemaID, + Key: key, }, }, })) diff --git a/test/singlenode_test.go b/test/singlenode_test.go index 8a9de9062..9a36f70a4 100644 --- a/test/singlenode_test.go +++ b/test/singlenode_test.go @@ -1527,32 +1527,32 @@ func TestTurnAuthFailure(t *testing.T) { } } -// asyncAttributesCapture buffers RequestResponse and GetDataTrackSchemaResponse messages +// dataBlobCapture buffers RequestResponse and GetDataBlobResponse messages // sent to a test client so they can be asserted on. Other messages flow through to the // default handler. -type asyncAttributesCapture struct { - mu sync.Mutex - requestResponses []*livekit.RequestResponse - schemaResponses []*livekit.GetDataTrackSchemaResponse +type dataBlobCapture struct { + mu sync.Mutex + requestResponses []*livekit.RequestResponse + blobResponses []*livekit.GetDataBlobResponse } -func (c *asyncAttributesCapture) interceptor() testclient.SignalResponseInterceptor { +func (c *dataBlobCapture) interceptor() testclient.SignalResponseInterceptor { return func(msg *livekit.SignalResponse, next testclient.SignalResponseHandler) error { switch m := msg.Message.(type) { case *livekit.SignalResponse_RequestResponse: c.mu.Lock() c.requestResponses = append(c.requestResponses, m.RequestResponse) c.mu.Unlock() - case *livekit.SignalResponse_GetDataTrackSchemaResponse: + case *livekit.SignalResponse_GetDataBlobResponse: c.mu.Lock() - c.schemaResponses = append(c.schemaResponses, m.GetDataTrackSchemaResponse) + c.blobResponses = append(c.blobResponses, m.GetDataBlobResponse) c.mu.Unlock() } return next(msg) } } -func (c *asyncAttributesCapture) takeRequestResponse() *livekit.RequestResponse { +func (c *dataBlobCapture) takeRequestResponse() *livekit.RequestResponse { c.mu.Lock() defer c.mu.Unlock() if len(c.requestResponses) == 0 { @@ -1563,28 +1563,28 @@ func (c *asyncAttributesCapture) takeRequestResponse() *livekit.RequestResponse return rr } -func (c *asyncAttributesCapture) takeSchemaResponse() *livekit.GetDataTrackSchemaResponse { +func (c *dataBlobCapture) takeBlobResponse() *livekit.GetDataBlobResponse { c.mu.Lock() defer c.mu.Unlock() - if len(c.schemaResponses) == 0 { + if len(c.blobResponses) == 0 { return nil } - sr := c.schemaResponses[0] - c.schemaResponses = c.schemaResponses[1:] + sr := c.blobResponses[0] + c.blobResponses = c.blobResponses[1:] return sr } -func (c *asyncAttributesCapture) requestResponseCount() int { +func (c *dataBlobCapture) requestResponseCount() int { c.mu.Lock() defer c.mu.Unlock() return len(c.requestResponses) } -func setupAsyncAttributesServer(t *testing.T, name string, enable bool) (*service.LivekitServer, func()) { +func setupDataBlobServer(t *testing.T, name string, enable bool) (*service.LivekitServer, func()) { logger.Infow("----------------STARTING TEST----------------", "test", name) s := createSingleNodeServer(func(c *config.Config) { - c.EnableParticipantAsyncAttributes = enable - c.Limit.MaxAsyncAttributesSize = 1024 + c.EnableParticipantDataBlob = enable + c.Limit.MaxDataBlobSize = 1024 }) go func() { if err := s.Start(); err != nil { @@ -1598,19 +1598,19 @@ func setupAsyncAttributesServer(t *testing.T, name string, enable bool) (*servic } } -func TestSingleNodeAsyncAttributes(t *testing.T) { +func TestSingleNodeDataBlob(t *testing.T) { if testing.Short() { t.SkipNow() return } - _, finish := setupAsyncAttributesServer(t, "TestSingleNodeAsyncAttributes", true) + _, finish := setupDataBlobServer(t, "TestSingleNodeDataBlob", true) defer finish() for _, testRTCServicePath := range testRTCServicePaths { t.Run(fmt.Sprintf("testRTCServicePath=%s", testRTCServicePath.String()), func(t *testing.T) { - pubCapture := &asyncAttributesCapture{} - subCapture := &asyncAttributesCapture{} + pubCapture := &dataBlobCapture{} + subCapture := &dataBlobCapture{} pub := createRTCClient("pub", defaultServerPort, testRTCServicePath, &testclient.Options{ AutoSubscribe: true, @@ -1623,19 +1623,20 @@ func TestSingleNodeAsyncAttributes(t *testing.T) { waitUntilConnected(t, pub, sub) defer stopClients(pub, sub) - schemaID := &livekit.DataTrackSchemaId{ - Name: "schema-1", - Encoding: livekit.DataTrackSchemaEncoding_DATA_TRACK_SCHEMA_ENCODING_PROTOBUF, + key := &livekit.DataBlobKey{ + Key: &livekit.DataBlobKey_Generic{ + Generic: "blob-1", + }, } - definition := []byte("definition-bytes") + contents := []byte("definition-bytes") - // publisher defines a schema + // publisher stores a blob require.NoError(t, pub.SendRequest(&livekit.SignalRequest{ - Message: &livekit.SignalRequest_DefineDataTrackSchema{ - DefineDataTrackSchema: &livekit.DefineDataTrackSchemaRequest{ - SchemaDefinition: &livekit.DataTrackSchemaDefinition{ - Id: schemaID, - Definition: definition, + Message: &livekit.SignalRequest_StoreDataBlobRequest{ + StoreDataBlobRequest: &livekit.StoreDataBlobRequest{ + Blob: &livekit.DataBlob{ + Key: key, + Contents: contents, }, }, }, @@ -1645,41 +1646,42 @@ func TestSingleNodeAsyncAttributes(t *testing.T) { time.Sleep(syncDelay) require.Equal(t, 0, pubCapture.requestResponseCount(), "publisher should not receive an error response on success") - // subscriber asks for the schema + // subscriber asks for the blob require.NoError(t, sub.SendRequest(&livekit.SignalRequest{ - Message: &livekit.SignalRequest_GetDataTrackSchema{ - GetDataTrackSchema: &livekit.GetDataTrackSchemaRequest{ + Message: &livekit.SignalRequest_GetDataBlobRequest{ + GetDataBlobRequest: &livekit.GetDataBlobRequest{ ParticipantIdentity: "pub", - SchemaId: schemaID, + Key: key, }, }, })) testutils.WithTimeout(t, func() string { - resp := subCapture.takeSchemaResponse() + resp := subCapture.takeBlobResponse() if resp == nil { - return "subscriber did not receive schema response" + return "subscriber did not receive blob response" } - if resp.SchemaDefinition == nil { - return "schema response missing definition" + if resp.Blob == nil { + return "blob response missing blob" } - if resp.SchemaDefinition.Id.Name != schemaID.Name { - return fmt.Sprintf("expected schema name %s, got %s", schemaID.Name, resp.SchemaDefinition.Id.Name) + if resp.Blob.Key.String() != key.String() { + return fmt.Sprintf("expected blob key %s, got %s", key.String(), resp.Blob.Key.String()) } - if string(resp.SchemaDefinition.Definition) != string(definition) { - return fmt.Sprintf("expected definition %q, got %q", definition, resp.SchemaDefinition.Definition) + if string(resp.Blob.Contents) != string(contents) { + return fmt.Sprintf("expected contents %q, got %q", contents, resp.Blob.Contents) } return "" }) - // subscriber asks for an unknown schema on a known publisher + // subscriber asks for an unknown blob on a known publisher require.NoError(t, sub.SendRequest(&livekit.SignalRequest{ - Message: &livekit.SignalRequest_GetDataTrackSchema{ - GetDataTrackSchema: &livekit.GetDataTrackSchemaRequest{ + Message: &livekit.SignalRequest_GetDataBlobRequest{ + GetDataBlobRequest: &livekit.GetDataBlobRequest{ ParticipantIdentity: "pub", - SchemaId: &livekit.DataTrackSchemaId{ - Name: "does-not-exist", - Encoding: livekit.DataTrackSchemaEncoding_DATA_TRACK_SCHEMA_ENCODING_PROTOBUF, + Key: &livekit.DataBlobKey{ + Key: &livekit.DataBlobKey_Generic{ + Generic: "does-not-exist", + }, }, }, }, @@ -1688,7 +1690,7 @@ func TestSingleNodeAsyncAttributes(t *testing.T) { testutils.WithTimeout(t, func() string { rr := subCapture.takeRequestResponse() if rr == nil { - return "subscriber did not receive RequestResponse for missing schema" + return "subscriber did not receive RequestResponse for missing blob" } if rr.Reason != livekit.RequestResponse_NOT_FOUND { return fmt.Sprintf("expected NOT_FOUND, got %s", rr.Reason) @@ -1696,12 +1698,12 @@ func TestSingleNodeAsyncAttributes(t *testing.T) { return "" }) - // subscriber asks for a schema on an unknown publisher identity + // subscriber asks for a blob on an unknown publisher identity require.NoError(t, sub.SendRequest(&livekit.SignalRequest{ - Message: &livekit.SignalRequest_GetDataTrackSchema{ - GetDataTrackSchema: &livekit.GetDataTrackSchemaRequest{ + Message: &livekit.SignalRequest_GetDataBlobRequest{ + GetDataBlobRequest: &livekit.GetDataBlobRequest{ ParticipantIdentity: "unknown-publisher", - SchemaId: schemaID, + Key: key, }, }, })) @@ -1717,16 +1719,12 @@ func TestSingleNodeAsyncAttributes(t *testing.T) { return "" }) - // publisher sends an invalid define (empty id name) + // publisher sends an invalid blob (empty key) require.NoError(t, pub.SendRequest(&livekit.SignalRequest{ - Message: &livekit.SignalRequest_DefineDataTrackSchema{ - DefineDataTrackSchema: &livekit.DefineDataTrackSchemaRequest{ - SchemaDefinition: &livekit.DataTrackSchemaDefinition{ - Id: &livekit.DataTrackSchemaId{ - Name: "", - Encoding: livekit.DataTrackSchemaEncoding_DATA_TRACK_SCHEMA_ENCODING_PROTOBUF, - }, - Definition: definition, + Message: &livekit.SignalRequest_StoreDataBlobRequest{ + StoreDataBlobRequest: &livekit.StoreDataBlobRequest{ + Blob: &livekit.DataBlob{ + Contents: contents, }, }, }, @@ -1746,18 +1744,18 @@ func TestSingleNodeAsyncAttributes(t *testing.T) { } } -func TestSingleNodeAsyncAttributesDisabled(t *testing.T) { +func TestSingleNodeDataBlobDisabled(t *testing.T) { if testing.Short() { t.SkipNow() return } - _, finish := setupAsyncAttributesServer(t, "TestSingleNodeAsyncAttributesDisabled", false) + _, finish := setupDataBlobServer(t, "TestSingleNodeDataBlobDisabled", false) defer finish() for _, testRTCServicePath := range testRTCServicePaths { t.Run(fmt.Sprintf("testRTCServicePath=%s", testRTCServicePath.String()), func(t *testing.T) { - pubCapture := &asyncAttributesCapture{} + pubCapture := &dataBlobCapture{} pub := createRTCClient("pub", defaultServerPort, testRTCServicePath, &testclient.Options{ AutoSubscribe: true, SignalResponseInterceptor: pubCapture.interceptor(), @@ -1766,14 +1764,15 @@ func TestSingleNodeAsyncAttributesDisabled(t *testing.T) { defer stopClients(pub) require.NoError(t, pub.SendRequest(&livekit.SignalRequest{ - Message: &livekit.SignalRequest_DefineDataTrackSchema{ - DefineDataTrackSchema: &livekit.DefineDataTrackSchemaRequest{ - SchemaDefinition: &livekit.DataTrackSchemaDefinition{ - Id: &livekit.DataTrackSchemaId{ - Name: "schema-1", - Encoding: livekit.DataTrackSchemaEncoding_DATA_TRACK_SCHEMA_ENCODING_PROTOBUF, + Message: &livekit.SignalRequest_StoreDataBlobRequest{ + StoreDataBlobRequest: &livekit.StoreDataBlobRequest{ + Blob: &livekit.DataBlob{ + Key: &livekit.DataBlobKey{ + Key: &livekit.DataBlobKey_Generic{ + Generic: "blob-1", + }, }, - Definition: []byte("definition-bytes"), + Contents: []byte("definition-bytes"), }, }, },