aboutsummaryrefslogtreecommitdiff
path: root/Emby.Server.Implementations/SyncPlay/Group.cs
diff options
context:
space:
mode:
Diffstat (limited to 'Emby.Server.Implementations/SyncPlay/Group.cs')
-rw-r--r--Emby.Server.Implementations/SyncPlay/Group.cs125
1 files changed, 122 insertions, 3 deletions
diff --git a/Emby.Server.Implementations/SyncPlay/Group.cs b/Emby.Server.Implementations/SyncPlay/Group.cs
index 923bfc67aa..6fbe46ffd6 100644
--- a/Emby.Server.Implementations/SyncPlay/Group.cs
+++ b/Emby.Server.Implementations/SyncPlay/Group.cs
@@ -11,6 +11,7 @@ using MediaBrowser.Controller.Library;
using MediaBrowser.Controller.Session;
using MediaBrowser.Controller.SyncPlay;
using MediaBrowser.Controller.SyncPlay.GroupStates;
+using MediaBrowser.Controller.SyncPlay.PlaybackRequests;
using MediaBrowser.Controller.SyncPlay.Queue;
using MediaBrowser.Controller.SyncPlay.Requests;
using MediaBrowser.Model.SyncPlay;
@@ -27,6 +28,11 @@ namespace Emby.Server.Implementations.SyncPlay
public class Group : IGroupStateContext
{
/// <summary>
+ /// The default value of <see cref="GroupWaitTimeout"/>, in milliseconds.
+ /// </summary>
+ internal const long DefaultGroupWaitTimeout = 30000;
+
+ /// <summary>
/// The logger.
/// </summary>
private readonly ILogger<Group> _logger;
@@ -54,8 +60,12 @@ namespace Emby.Server.Implementations.SyncPlay
/// <summary>
/// The participants, or members of the group.
/// </summary>
- private readonly Dictionary<string, GroupMember> _participants =
- new Dictionary<string, GroupMember>(StringComparer.OrdinalIgnoreCase);
+ private readonly Dictionary<string, GroupMember> _participants = new(StringComparer.OrdinalIgnoreCase);
+
+ /// <summary>
+ /// The sessions of the participants, which only carry identifiers.
+ /// </summary>
+ private readonly Dictionary<string, SessionInfo> _participantSessions = new(StringComparer.OrdinalIgnoreCase);
/// <summary>
/// The internal group state.
@@ -115,6 +125,19 @@ namespace Emby.Server.Implementations.SyncPlay
public long MaxPlaybackOffset { get; } = 500;
/// <summary>
+ /// Gets the maximum time, in milliseconds, the group waits for its members to report ready.
+ /// </summary>
+ /// <value>The group-wait timeout.</value>
+ internal long GroupWaitTimeout { get; init; } = DefaultGroupWaitTimeout;
+
+ /// <summary>
+ /// Gets the <see cref="Environment.TickCount64"/> value at which the group gives up waiting
+ /// for its members, or <c>null</c> when it is not waiting for anyone.
+ /// </summary>
+ /// <value>The group-wait deadline.</value>
+ internal long? GroupWaitDeadline { get; private set; }
+
+ /// <summary>
/// Gets the group identifier.
/// </summary>
/// <value>The group identifier.</value>
@@ -163,6 +186,8 @@ namespace Emby.Server.Implementations.SyncPlay
Ping = DefaultPing,
IsBuffering = false
});
+
+ _participantSessions[session.Id] = session;
}
/// <summary>
@@ -172,6 +197,8 @@ namespace Emby.Server.Implementations.SyncPlay
private void RemoveSession(SessionInfo session)
{
_participants.Remove(session.Id);
+ _participantSessions.Remove(session.Id);
+ UpdateGroupWaitDeadline(false);
}
/// <summary>
@@ -389,13 +416,20 @@ namespace Emby.Server.Implementations.SyncPlay
{
value.IgnoreGroupWait = ignoreGroupWait;
}
+
+ UpdateGroupWaitDeadline(false);
}
/// <inheritdoc />
public void SetState(IGroupState state)
{
_logger.LogInformation("Group {GroupId} switching from {FromStateType} to {ToStateType}.", GroupId.ToString(), _state.Type, state.Type);
- this._state = state;
+ _state = state;
+
+ if (state.Type != GroupStateType.Waiting)
+ {
+ GroupWaitDeadline = null;
+ }
}
/// <inheritdoc />
@@ -475,6 +509,8 @@ namespace Emby.Server.Implementations.SyncPlay
{
value.IsBuffering = isBuffering;
}
+
+ UpdateGroupWaitDeadline(false);
}
/// <inheritdoc />
@@ -484,6 +520,9 @@ namespace Emby.Server.Implementations.SyncPlay
{
session.IsBuffering = isBuffering;
}
+
+ // Resetting the status of every session starts a new waiting period.
+ UpdateGroupWaitDeadline(isBuffering);
}
/// <inheritdoc />
@@ -690,5 +729,85 @@ namespace Emby.Server.Implementations.SyncPlay
PlayQueue.ShuffleMode,
PlayQueue.RepeatMode);
}
+
+ /// <summary>
+ /// Stops waiting for the members that have not reported ready and lets the rest of the
+ /// group carry on. Does nothing until <see cref="GroupWaitDeadline"/> has passed.
+ /// </summary>
+ /// <param name="cancellationToken">The cancellation token.</param>
+ internal void HandleGroupWaitTimeout(CancellationToken cancellationToken)
+ {
+ var deadline = GroupWaitDeadline;
+ if (deadline is null || deadline > Environment.TickCount64)
+ {
+ return;
+ }
+
+ GroupWaitDeadline = null;
+
+ if (_state is not WaitingGroupState waitingState)
+ {
+ return;
+ }
+
+ var blockingSessions = _participantSessions
+ .Values
+ .Where(participant => _participants.TryGetValue(participant.Id, out var member)
+ && member.IsBuffering
+ && !member.IgnoreGroupWait)
+ .ToList();
+
+ if (blockingSessions.Count == 0)
+ {
+ return;
+ }
+
+ // The recovery below is broadcast to the whole group, so it does not matter which of
+ // the sessions that kept the group waiting is the one acting on the group's behalf.
+ var session = blockingSessions[0];
+
+ _logger.LogWarning(
+ "Group {GroupId} waited {Waited} ms for session(s) {SessionIds} to report ready, giving up.",
+ GroupId.ToString(),
+ GroupWaitTimeout + Environment.TickCount64 - deadline.Value,
+ string.Join(", ", blockingSessions.Select(participant => participant.Id)));
+
+ if (waitingState.ResumePlaying)
+ {
+ // An unpause request in the waiting state means "start now, ignoring the sessions
+ // that are not ready".
+ var unpauseRequest = new UnpauseGroupRequest();
+ waitingState.HandleRequest(unpauseRequest, this, GroupStateType.Waiting, session, cancellationToken);
+ return;
+ }
+
+ // The members have been paused for the whole waiting period, so the playback position
+ // stays where the wait started.
+ SetAllBuffering(false);
+ SetState(new PausedGroupState(_loggerFactory));
+
+ var command = NewSyncPlayCommand(SendCommandType.Pause);
+ SendCommand(session, SyncPlayBroadcastType.AllGroup, command, cancellationToken);
+
+ var stateUpdate = new GroupStateUpdate(GroupStateType.Paused, PlaybackRequestType.Pause);
+ var update = new SyncPlayStateUpdate(GroupId, stateUpdate);
+ SendGroupUpdate(session, SyncPlayBroadcastType.AllGroup, update, cancellationToken);
+ }
+
+ private void UpdateGroupWaitDeadline(bool startNewWaitingPeriod)
+ {
+ if (_state.Type != GroupStateType.Waiting || !IsBuffering())
+ {
+ GroupWaitDeadline = null;
+ return;
+ }
+
+ // A running deadline covers the waiting period as a whole, so the sessions that keep
+ // reporting buffering while they load must not push it back.
+ if (GroupWaitDeadline is null || startNewWaitingPeriod)
+ {
+ GroupWaitDeadline = Environment.TickCount64 + GroupWaitTimeout;
+ }
+ }
}
}