fix(session): tie directory ownership to the live connection

Ownership is claimed with a Lua check-and-set keyed on the instance holding
the websocket, routing prefers a live controller over a local copy, the
session list deduplicates by owner, removal is ownership-checked, undelivered
routed messages surface, single-session lookups stop scanning the keyspace and
directory writes leave the request path bounded by a timeout.
This commit is contained in:
2026-09-24 23:50:01 +10:00
parent 6f362c33c9
commit 9e66708d87
17 changed files with 613 additions and 123 deletions
@@ -1,4 +1,5 @@
using System;
using System.Threading;
using System.Threading.Tasks;
namespace MediaBrowser.Controller.Session;
@@ -15,11 +16,14 @@ public interface IPodMessageBus
string PodId { get; }
/// <summary>
/// Sends a message to one instance. Delivery is best effort and never throws.
/// Sends a message to one instance and reports how many listeners took it, so that a message
/// addressed to an instance that is no longer there is not mistaken for a delivered one.
/// </summary>
/// <param name="targetPod">The instance to deliver to.</param>
/// <param name="message">The message.</param>
void Publish(string targetPod, PodMessage message);
/// <param name="cancellationToken">The cancellation token.</param>
/// <returns>The number of instances the message reached.</returns>
Task<long> PublishAsync(string targetPod, PodMessage message, CancellationToken cancellationToken = default);
/// <summary>
/// Registers a handler for the messages addressed to this instance.
@@ -11,20 +11,32 @@ namespace MediaBrowser.Controller.Session;
public interface ISessionDirectory
{
/// <summary>
/// Publishes an entry and restarts its expiry.
/// Claims a session for the publishing instance and restarts its expiry. The claim is refused when
/// another instance holds the connection, so an instance that merely served a request for the session
/// cannot take ownership of it.
/// </summary>
/// <param name="entry">The entry.</param>
/// <param name="connectedUtcTicks">When the publishing instance's connection to the session was established, or zero when it holds none.</param>
/// <param name="cancellationToken">The cancellation token.</param>
/// <returns>A task representing the operation.</returns>
Task PublishAsync(SessionDirectoryEntry entry, CancellationToken cancellationToken = default);
/// <returns><c>true</c> if the entry was written.</returns>
Task<bool> PublishAsync(SessionDirectoryEntry entry, long connectedUtcTicks, CancellationToken cancellationToken = default);
/// <summary>
/// Removes an entry.
/// Removes an entry, but only while the calling instance still owns it.
/// </summary>
/// <param name="sessionId">The session identifier.</param>
/// <param name="ownerPod">The instance requesting the removal.</param>
/// <param name="cancellationToken">The cancellation token.</param>
/// <returns>A task representing the operation.</returns>
Task RemoveAsync(string sessionId, string ownerPod, CancellationToken cancellationToken = default);
/// <summary>
/// Gets one entry by session identifier.
/// </summary>
/// <param name="sessionId">The session identifier.</param>
/// <param name="cancellationToken">The cancellation token.</param>
/// <returns>A task representing the operation.</returns>
Task RemoveAsync(string sessionId, CancellationToken cancellationToken = default);
/// <returns>The entry, or <c>null</c> when the session is in no instance's directory.</returns>
Task<SessionDirectoryEntry?> GetAsync(string sessionId, CancellationToken cancellationToken = default);
/// <summary>
/// Gets every entry that has not expired.
@@ -80,7 +80,8 @@ namespace MediaBrowser.Controller.Session
/// Used to report that a session controller has connected.
/// </summary>
/// <param name="session">The session.</param>
void OnSessionControllerConnected(SessionInfo session);
/// <returns>A task representing the operation.</returns>
Task OnSessionControllerConnected(SessionInfo session);
void UpdateDeviceName(string sessionId, string reportedDeviceName);
@@ -241,7 +242,8 @@ namespace MediaBrowser.Controller.Session
/// <param name="controllingSessionId">The controlling session identifier.</param>
/// <param name="sessionId">The session identifier.</param>
/// <param name="userId">The user identifier.</param>
void AddAdditionalUser(string controllingSessionId, string sessionId, Guid userId);
/// <returns>A task representing the operation.</returns>
Task AddAdditionalUser(string controllingSessionId, string sessionId, Guid userId);
/// <summary>
/// Removes the additional user.
@@ -249,7 +251,8 @@ namespace MediaBrowser.Controller.Session
/// <param name="controllingSessionId">The controlling session identifier.</param>
/// <param name="sessionId">The session identifier.</param>
/// <param name="userId">The user identifier.</param>
void RemoveAdditionalUser(string controllingSessionId, string sessionId, Guid userId);
/// <returns>A task representing the operation.</returns>
Task RemoveAdditionalUser(string controllingSessionId, string sessionId, Guid userId);
/// <summary>
/// Reports the now viewing item.
@@ -1,4 +1,5 @@
using System;
using System.Threading;
using System.Threading.Tasks;
namespace MediaBrowser.Controller.Session;
@@ -17,9 +18,8 @@ public sealed class NullPodMessageBus : IPodMessageBus
public string PodId => PodIdentity.Current;
/// <inheritdoc />
public void Publish(string targetPod, PodMessage message)
{
}
public Task<long> PublishAsync(string targetPod, PodMessage message, CancellationToken cancellationToken = default)
=> Task.FromResult(0L);
/// <inheritdoc />
public void Subscribe(Func<PodMessage, Task> handler)
@@ -17,12 +17,16 @@ public sealed class NullSessionDirectory : ISessionDirectory
public static NullSessionDirectory Instance { get; } = new NullSessionDirectory();
/// <inheritdoc />
public Task PublishAsync(SessionDirectoryEntry entry, CancellationToken cancellationToken = default)
public Task<bool> PublishAsync(SessionDirectoryEntry entry, long connectedUtcTicks, CancellationToken cancellationToken = default)
=> Task.FromResult(false);
/// <inheritdoc />
public Task RemoveAsync(string sessionId, string ownerPod, CancellationToken cancellationToken = default)
=> Task.CompletedTask;
/// <inheritdoc />
public Task RemoveAsync(string sessionId, CancellationToken cancellationToken = default)
=> Task.CompletedTask;
public Task<SessionDirectoryEntry?> GetAsync(string sessionId, CancellationToken cancellationToken = default)
=> Task.FromResult<SessionDirectoryEntry?>(null);
/// <inheritdoc />
public Task<IReadOnlyList<SessionDirectoryEntry>> GetAllAsync(CancellationToken cancellationToken = default)
@@ -0,0 +1,30 @@
using System;
namespace MediaBrowser.Controller.Session;
/// <summary>
/// An additional-user change for a session held by another instance, carried as a
/// <see cref="PodMessage"/>. The calling instance has already authorized it.
/// </summary>
public sealed class RoutedAdditionalUserChange
{
/// <summary>
/// The <see cref="PodMessage.Kind"/> this payload travels under.
/// </summary>
public const string Kind = "AdditionalUserChange";
/// <summary>
/// Gets or sets the session the change applies to.
/// </summary>
public string SessionId { get; set; } = string.Empty;
/// <summary>
/// Gets or sets the user to attach or detach.
/// </summary>
public Guid UserId { get; set; }
/// <summary>
/// Gets or sets a value indicating whether the user is being attached rather than detached.
/// </summary>
public bool Add { get; set; }
}
@@ -12,6 +12,12 @@ public sealed class SessionDirectoryEntry
/// </summary>
public string OwnerPod { get; set; } = string.Empty;
/// <summary>
/// Gets or sets a value indicating whether the owner holds a live connection to the session. Only an
/// owner that does can be routed a remote-control message.
/// </summary>
public bool HoldsConnection { get; set; }
/// <summary>
/// Gets or sets the session as its owner last rendered it.
/// </summary>
@@ -20,4 +20,9 @@ public sealed class SessionDirectoryOptions
/// Gets or sets how often in seconds an instance republishes the sessions it holds.
/// </summary>
public int RefreshIntervalSeconds { get; set; } = 20;
/// <summary>
/// Gets or sets how long in seconds a single directory operation may take before it is abandoned.
/// </summary>
public int OperationTimeoutSeconds { get; set; } = 5;
}