data blob

This commit is contained in:
boks1971
2026-06-14 01:06:25 +05:30
parent 48463c3792
commit 492422cd53
25 changed files with 1149 additions and 1286 deletions
+2 -2
View File
@@ -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
+4 -33
View File
@@ -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=
+17 -17
View File
@@ -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",
+3 -3
View File
@@ -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,
}),
}
-107
View File
@@ -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),
}
}
@@ -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()
}
@@ -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)
})
}
@@ -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)
})
}
}
+91
View File
@@ -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
}
// -------------------------------
+111
View File
@@ -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()
}
@@ -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)
})
}
+175
View File
@@ -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()
}
+3 -3
View File
@@ -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,
}))
}
+5 -5
View File
@@ -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 {
+1 -1
View File
@@ -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
}
+4 -4
View File
@@ -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
+3 -3
View File
@@ -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,
},
}
}
@@ -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
}
+10 -10
View File
@@ -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
+191 -191
View File
@@ -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
}
@@ -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 {
+68 -68
View File
@@ -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}
}
+3 -3
View File
@@ -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
+31 -30
View File
@@ -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,
},
},
}))
+72 -73
View File
@@ -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"),
},
},
},