Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
27 changes: 27 additions & 0 deletions docs/protocol.md
Original file line number Diff line number Diff line change
Expand Up @@ -335,8 +335,35 @@ Um client precisa enviar apenas `type`, `to`, `payload` e, opcionalmente, `messa
- `webrtc.answer`
- `webrtc.ice_candidate`
- `webrtc.renegotiate`
- `video.receiver_stats`
- `pong`

`video.receiver_stats` is sent by a viewer to the publisher participant every 2 seconds. Version
1 has this payload shape (all counters are non-negative integers):

```json
{
"version": 1,
"intervalMilliseconds": 2000,
"rtpPacketsReceived": 1200,
"rtpPacketsLost": 8,
"accessUnitsReceived": 60,
"incompleteAccessUnits": 1,
"decodedFrames": 58,
"targetFramesPerSecond": 30.0
}
```

The client publisher validates version 1, a 1000–5000 ms interval, counter bounds (each at most
10,000,000; received plus lost packets at most 10,000,000), `incompleteAccessUnits` no greater
than `accessUnitsReceived`, and finite target FPS from 1 through 60. The backend treats the
payload as opaque JSON: it enforces the normal WebSocket message-size limit and authenticated
same-session recipient routing, but does not parse or interpret these fields. The routed envelope
sets `from` from the authenticated socket; clients must not trust a sender identity supplied in
the payload. The publisher further accepts feedback only from a viewer with an active peer in
that sharing session. This is signaling metadata, not media; RTP/RTCP media remains between the
peers (or traverses TURN).

These types describe the sender's own state instead of addressing one peer, so they take **no** `to`. The server validates them, persists the result and broadcasts the authoritative version to the whole session, sender included:

- `participant.capabilities`
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -15,7 +15,7 @@ public static class SignalingWebSocketEndpoint
private static readonly HashSet<string> RoutedMessageTypes =
[
"publisher.ready", "viewer.ready", "webrtc.offer", "webrtc.answer",
"webrtc.ice_candidate", "webrtc.renegotiate", "pong"
"webrtc.ice_candidate", "webrtc.renegotiate", "video.receiver_stats", "pong"
];

/// <summary>
Expand Down
71 changes: 71 additions & 0 deletions tests/SonicRelay.Api.IntegrationTests/SignalingWebSocketTests.cs
Original file line number Diff line number Diff line change
Expand Up @@ -110,6 +110,77 @@ public async Task Signaling_normalizes_the_envelope_and_overwrites_client_metada
Assert.NotEqual(DateTimeOffset.UnixEpoch, routed.GetProperty("timestamp").GetDateTimeOffset());
}

[Fact]
public async Task Receiver_stats_are_forwarded_with_authenticated_envelope_metadata()
{
var sender = await CreateParticipantAsync("receiver-stats-sender");
var receiver = await CreateViewerAsync(sender, "receiver-stats-receiver");
using var senderSocket = await ConnectAsync(sender);
using var receiverSocket = await ConnectAsync(receiver);
await ReceiveAsync(senderSocket);
await ReceiveAsync(receiverSocket);
await ReceiveAsync(senderSocket); // receiver joined
await ReceiveAsync(senderSocket); // receiver capability state
await ReceiveAsync(receiverSocket); // existing sender capability roster entry

var messageId = Guid.NewGuid();
await SendAsync(senderSocket, new
{
type = "video.receiver_stats",
messageId,
sessionId = Guid.NewGuid(),
from = Guid.NewGuid(),
to = receiver.ParticipantId,
timestamp = DateTimeOffset.UnixEpoch,
payload = new
{
version = 1,
intervalMs = 2000,
rtpPacketsReceived = 90,
rtpPacketsLost = 10,
accessUnitsReceived = 30,
incompleteAccessUnits = 2,
decodedFrames = 24,
targetFramesPerSecond = 30
}
});

using var responseTimeout = new CancellationTokenSource(TimeSpan.FromSeconds(5));
var senderResponse = ReceiveAsync(senderSocket, responseTimeout.Token);
var receiverResponse = ReceiveReceiverStatsAsync(receiverSocket, responseTimeout.Token);
var firstResponse = await Task.WhenAny(senderResponse, receiverResponse);
responseTimeout.Cancel();
var routed = await firstResponse;
if (routed.GetProperty("type").GetString() == "error")
Assert.Equal("unsupported_message_type", routed.GetProperty("payload").GetProperty("code").GetString());
AssertEnvelope(routed, "video.receiver_stats", sender.SessionId);
Assert.Equal(messageId, routed.GetProperty("messageId").GetGuid());
Assert.Equal(sender.ParticipantId, routed.GetProperty("from").GetGuid());
Assert.Equal(receiver.ParticipantId, routed.GetProperty("to").GetGuid());
AssertReceiverStatsPayload(routed);
}

private static void AssertReceiverStatsPayload(JsonElement routed)
{
var payload = routed.GetProperty("payload");
Assert.Equal(1, payload.GetProperty("version").GetInt32());
Assert.Equal(2000, payload.GetProperty("intervalMs").GetInt32());
Assert.Equal(90, payload.GetProperty("rtpPacketsReceived").GetInt32());
Assert.Equal(10, payload.GetProperty("rtpPacketsLost").GetInt32());
Assert.Equal(30, payload.GetProperty("accessUnitsReceived").GetInt32());
Assert.Equal(2, payload.GetProperty("incompleteAccessUnits").GetInt32());
Assert.Equal(24, payload.GetProperty("decodedFrames").GetInt32());
Assert.Equal(30, payload.GetProperty("targetFramesPerSecond").GetInt32());
}

private static async Task<JsonElement> ReceiveReceiverStatsAsync(WebSocket socket, CancellationToken ct)
{
var message = await ReceiveAsync(socket, ct);
while (message.GetProperty("type").GetString() == "participant.capabilities")
message = await ReceiveAsync(socket, ct);
return message;
}

[Fact]
public async Task Signaling_announces_a_new_participant_to_existing_session_peers()
{
Expand Down
Loading