refactor(runtime): move session lifetime and ordered transport
Move the canonical WorldSession generation, connect/enter/tick/stop transaction, inbound subscription owner, and retryable teardown acknowledgements into AcDream.Runtime. Keep App as a borrowing graphical host with a single inertable command projection and no mirrored session state. Validated by 79 Runtime tests, 3,776 App tests with three existing skips, the Release solution build, and 8,428 complete Release tests with five existing skips. Co-authored-by: Codex <codex@openai.com>
This commit is contained in:
parent
ecc4816c5a
commit
7593078774
37 changed files with 884 additions and 355 deletions
|
|
@ -1,425 +0,0 @@
|
|||
using AcDream.Core.Chat;
|
||||
using AcDream.Core.Combat;
|
||||
using AcDream.Core.Items;
|
||||
using AcDream.Core.Net;
|
||||
using AcDream.Core.Net.Messages;
|
||||
using AcDream.Core.Player;
|
||||
using AcDream.Core.Social;
|
||||
using AcDream.Core.Spells;
|
||||
|
||||
namespace AcDream.App.Net;
|
||||
|
||||
internal sealed record LiveEntitySessionSink(
|
||||
Action<WorldSession.EntitySpawn> Spawned,
|
||||
Action<DeleteObject.Parsed> Deleted,
|
||||
Action<PickupEvent.Parsed> PickedUp,
|
||||
Action<WorldSession.EntityMotionUpdate> MotionUpdated,
|
||||
Action<WorldSession.EntityPositionUpdate> PositionUpdated,
|
||||
Action<VectorUpdate.Parsed> VectorUpdated,
|
||||
Action<SetState.Parsed> StateUpdated,
|
||||
Action<ParentEvent.Parsed> ParentUpdated,
|
||||
Action<uint> TeleportStarted,
|
||||
Action<ObjDescEvent.Parsed> AppearanceUpdated,
|
||||
Action<PlayPhysicsScript> PlayPhysicsScript,
|
||||
Action<PlayPhysicsScriptType> PlayPhysicsScriptType);
|
||||
|
||||
internal sealed record LiveEnvironmentSessionSink(
|
||||
Action<uint> EnvironChanged,
|
||||
Action<double> ServerTimeUpdated);
|
||||
|
||||
internal sealed record LiveInventorySessionBindings(
|
||||
ClientObjectTable Objects,
|
||||
LocalPlayerState LocalPlayer,
|
||||
Func<uint> PlayerGuid,
|
||||
Action<IReadOnlyList<ShortcutEntry>>? OnShortcuts,
|
||||
Action<uint>? OnUseDone,
|
||||
ItemManaState? ItemMana,
|
||||
Action<IReadOnlyList<(uint Id, uint Amount)>>? OnDesiredComponents,
|
||||
ExternalContainerState? ExternalContainers,
|
||||
Action<AppraiseInfoParser.Parsed>? OnAppraisal = null);
|
||||
|
||||
internal sealed record LiveCharacterSessionBindings(
|
||||
CombatState Combat,
|
||||
Spellbook Spellbook,
|
||||
Func<uint, IReadOnlyDictionary<uint, uint>, uint>? ResolveSkillFormulaBonus,
|
||||
Action<int, int>? OnSkillsUpdated,
|
||||
Action<GameEvents.CharacterConfirmationRequest>? OnConfirmationRequest,
|
||||
Action<GameEvents.CharacterConfirmationDone>? OnConfirmationDone,
|
||||
Action<uint, uint>? OnCharacterOptions,
|
||||
Func<double>? ClientTime);
|
||||
|
||||
internal sealed record LiveSocialSessionBindings(
|
||||
ChatLog Chat,
|
||||
TurbineChatState TurbineChat,
|
||||
FriendsState? Friends,
|
||||
SquelchState? Squelch);
|
||||
|
||||
/// <summary>
|
||||
/// Owns every inbound subscription for one exact live session. Domain state
|
||||
/// remains in the supplied sinks; this class owns only routing and teardown.
|
||||
/// </summary>
|
||||
internal sealed class LiveSessionEventRouter : ILiveSessionEventRouting
|
||||
{
|
||||
private readonly LiveSessionSubscriptionSet _subscriptions = new();
|
||||
private readonly Action<int>? _constructionCheckpoint;
|
||||
private readonly WorldSession _session;
|
||||
private readonly LiveEntitySessionSink _entities;
|
||||
private readonly LiveEnvironmentSessionSink _environment;
|
||||
private readonly LiveInventorySessionBindings _inventory;
|
||||
private readonly LiveCharacterSessionBindings _character;
|
||||
private readonly LiveSocialSessionBindings _social;
|
||||
private int _constructionStep;
|
||||
private int _accepting;
|
||||
private int _lifecycleState; // 0 created, 1 attaching, 2 attached, 3 disposed
|
||||
|
||||
public LiveSessionEventRouter(
|
||||
WorldSession session,
|
||||
LiveEntitySessionSink entities,
|
||||
LiveEnvironmentSessionSink environment,
|
||||
LiveInventorySessionBindings inventory,
|
||||
LiveCharacterSessionBindings character,
|
||||
LiveSocialSessionBindings social,
|
||||
Action<int>? constructionCheckpoint = null)
|
||||
{
|
||||
ArgumentNullException.ThrowIfNull(session);
|
||||
Validate(entities, environment, inventory, character, social);
|
||||
_session = session;
|
||||
_entities = entities;
|
||||
_environment = environment;
|
||||
_inventory = inventory;
|
||||
_character = character;
|
||||
_social = social;
|
||||
_constructionCheckpoint = constructionCheckpoint;
|
||||
}
|
||||
|
||||
public void Attach()
|
||||
{
|
||||
if (Interlocked.CompareExchange(ref _lifecycleState, 1, 0) != 0)
|
||||
throw new InvalidOperationException(
|
||||
"Live-session event routing can only attach once.");
|
||||
|
||||
Interlocked.Exchange(ref _accepting, 1);
|
||||
WorldSession session = _session;
|
||||
LiveEntitySessionSink entities = _entities;
|
||||
LiveEnvironmentSessionSink environment = _environment;
|
||||
LiveInventorySessionBindings inventory = _inventory;
|
||||
LiveCharacterSessionBindings character = _character;
|
||||
LiveSocialSessionBindings social = _social;
|
||||
|
||||
try
|
||||
{
|
||||
// Preserve the shipped pre-Connect registration order. Property
|
||||
// state is installed before lifecycle packets can be dispatched.
|
||||
_subscriptions.Add(ObjectTableWiring.Wire(
|
||||
session,
|
||||
inventory.Objects,
|
||||
inventory.PlayerGuid,
|
||||
inventory.LocalPlayer,
|
||||
IsAccepting));
|
||||
ConstructionCheckpoint();
|
||||
_subscriptions.Add(CombatStateWiring.Wire(
|
||||
session,
|
||||
character.Combat,
|
||||
IsAccepting));
|
||||
ConstructionCheckpoint();
|
||||
|
||||
Subscribe(h => session.EntitySpawned += h, h => session.EntitySpawned -= h, entities.Spawned);
|
||||
Subscribe(h => session.EntityDeleted += h, h => session.EntityDeleted -= h, entities.Deleted);
|
||||
Subscribe(h => session.EntityPickedUp += h, h => session.EntityPickedUp -= h, entities.PickedUp);
|
||||
Subscribe(h => session.MotionUpdated += h, h => session.MotionUpdated -= h, entities.MotionUpdated);
|
||||
Subscribe(h => session.PositionUpdated += h, h => session.PositionUpdated -= h, entities.PositionUpdated);
|
||||
Subscribe(h => session.VectorUpdated += h, h => session.VectorUpdated -= h, entities.VectorUpdated);
|
||||
Subscribe(h => session.StateUpdated += h, h => session.StateUpdated -= h, entities.StateUpdated);
|
||||
Subscribe(h => session.ParentUpdated += h, h => session.ParentUpdated -= h, entities.ParentUpdated);
|
||||
Subscribe(h => session.TeleportStarted += h, h => session.TeleportStarted -= h, entities.TeleportStarted);
|
||||
Subscribe(h => session.AppearanceUpdated += h, h => session.AppearanceUpdated -= h, entities.AppearanceUpdated);
|
||||
Subscribe(
|
||||
h => session.PlayPhysicsScriptReceived += h,
|
||||
h => session.PlayPhysicsScriptReceived -= h,
|
||||
entities.PlayPhysicsScript);
|
||||
Subscribe(
|
||||
h => session.PlayPhysicsScriptTypeReceived += h,
|
||||
h => session.PlayPhysicsScriptTypeReceived -= h,
|
||||
entities.PlayPhysicsScriptType);
|
||||
|
||||
Subscribe(h => session.EnvironChanged += h, h => session.EnvironChanged -= h, environment.EnvironChanged);
|
||||
Subscribe(h => session.ServerTimeUpdated += h, h => session.ServerTimeUpdated -= h, environment.ServerTimeUpdated);
|
||||
|
||||
_subscriptions.Add(GameEventWiring.WireAll(
|
||||
session.GameEvents,
|
||||
inventory.Objects,
|
||||
character.Combat,
|
||||
character.Spellbook,
|
||||
social.Chat,
|
||||
inventory.LocalPlayer,
|
||||
social.TurbineChat,
|
||||
onSkillsUpdated: character.OnSkillsUpdated,
|
||||
resolveSkillFormulaBonus: character.ResolveSkillFormulaBonus,
|
||||
onShortcuts: inventory.OnShortcuts,
|
||||
playerGuid: inventory.PlayerGuid,
|
||||
onUseDone: inventory.OnUseDone,
|
||||
onAppraisal: inventory.OnAppraisal,
|
||||
itemMana: inventory.ItemMana,
|
||||
onConfirmationRequest: character.OnConfirmationRequest,
|
||||
onConfirmationDone: character.OnConfirmationDone,
|
||||
friends: social.Friends,
|
||||
squelch: social.Squelch,
|
||||
onDesiredComponents: inventory.OnDesiredComponents,
|
||||
onCharacterOptions: character.OnCharacterOptions,
|
||||
clientTime: character.ClientTime,
|
||||
externalContainers: inventory.ExternalContainers,
|
||||
accepting: IsAccepting));
|
||||
ConstructionCheckpoint();
|
||||
_subscriptions.Add(new CombatChatTranslator(
|
||||
character.Combat,
|
||||
social.Chat,
|
||||
IsAccepting));
|
||||
ConstructionCheckpoint();
|
||||
|
||||
Subscribe<HearSpeech.Parsed>(h => session.SpeechHeard += h, h => session.SpeechHeard -= h, speech =>
|
||||
social.Chat.OnLocalSpeech(
|
||||
speech.SenderName,
|
||||
speech.Text,
|
||||
speech.SenderGuid,
|
||||
speech.IsRanged));
|
||||
Subscribe<ServerMessage.Parsed>(
|
||||
h => session.ServerMessageReceived += h,
|
||||
h => session.ServerMessageReceived -= h,
|
||||
message => social.Chat.OnSystemMessage(message.Message, message.ChatType));
|
||||
Subscribe<EmoteText.Parsed>(h => session.EmoteHeard += h, h => session.EmoteHeard -= h, emote =>
|
||||
social.Chat.OnEmote(emote.SenderName, emote.Text, emote.SenderGuid));
|
||||
Subscribe<SoulEmote.Parsed>(h => session.SoulEmoteHeard += h, h => session.SoulEmoteHeard -= h, emote =>
|
||||
social.Chat.OnSoulEmote(emote.SenderName, emote.Text, emote.SenderGuid));
|
||||
Subscribe<PlayerKilled.Parsed>(
|
||||
h => session.PlayerKilledReceived += h,
|
||||
h => session.PlayerKilledReceived -= h,
|
||||
killed => social.Chat.OnPlayerKilled(
|
||||
killed.DeathMessage,
|
||||
killed.VictimGuid,
|
||||
killed.KillerGuid));
|
||||
Subscribe<TurbineChat.Parsed>(
|
||||
h => session.TurbineChatReceived += h,
|
||||
h => session.TurbineChatReceived -= h,
|
||||
parsed => RouteTurbineChat(social.Chat, parsed));
|
||||
Subscribe<PrivateUpdateVital.ParsedFull>(h => session.VitalUpdated += h, h => session.VitalUpdated -= h, vital =>
|
||||
inventory.LocalPlayer.OnVitalUpdate(
|
||||
vital.VitalId,
|
||||
vital.Ranks,
|
||||
vital.Start,
|
||||
vital.Xp,
|
||||
vital.Current));
|
||||
Subscribe<PrivateUpdateVital.ParsedCurrent>(
|
||||
h => session.VitalCurrentUpdated += h,
|
||||
h => session.VitalCurrentUpdated -= h,
|
||||
vital => inventory.LocalPlayer.OnVitalCurrent(vital.VitalId, vital.Current));
|
||||
|
||||
if (Interlocked.CompareExchange(ref _lifecycleState, 2, 1) != 1)
|
||||
throw new ObjectDisposedException(nameof(LiveSessionEventRouter));
|
||||
}
|
||||
catch
|
||||
{
|
||||
Interlocked.Exchange(ref _accepting, 0);
|
||||
Interlocked.Exchange(ref _lifecycleState, 3);
|
||||
throw;
|
||||
}
|
||||
}
|
||||
|
||||
public bool Accepting => IsAccepting();
|
||||
|
||||
public void Dispose()
|
||||
{
|
||||
Interlocked.Exchange(ref _accepting, 0);
|
||||
Interlocked.Exchange(ref _lifecycleState, 3);
|
||||
_subscriptions.Dispose();
|
||||
}
|
||||
|
||||
private void Subscribe<T>(
|
||||
Action<Action<T>> attach,
|
||||
Action<Action<T>> detach,
|
||||
Action<T> sink)
|
||||
{
|
||||
Action<T> handler = value =>
|
||||
{
|
||||
if (Volatile.Read(ref _accepting) != 0)
|
||||
sink(value);
|
||||
};
|
||||
|
||||
attach(handler);
|
||||
// Add assumes cleanup ownership before it can call an external detach.
|
||||
// If the set is already closing, it retains/retries that exact edge;
|
||||
// calling detach again here would replay a successful removal.
|
||||
_subscriptions.Add(() => detach(handler));
|
||||
|
||||
ConstructionCheckpoint();
|
||||
}
|
||||
|
||||
private void ConstructionCheckpoint() =>
|
||||
_constructionCheckpoint?.Invoke(++_constructionStep);
|
||||
|
||||
private bool IsAccepting() => Volatile.Read(ref _accepting) != 0;
|
||||
|
||||
private static void RouteTurbineChat(ChatLog chat, TurbineChat.Parsed parsed)
|
||||
{
|
||||
if (parsed.Body is not TurbineChat.Payload.EventSendToRoom message)
|
||||
return;
|
||||
|
||||
chat.OnChannelBroadcast(
|
||||
message.RoomId,
|
||||
message.SenderName,
|
||||
message.Message,
|
||||
TurbineChatRouting.DisplayName(message.RoomId, message.ChatType));
|
||||
}
|
||||
|
||||
private static void Validate(
|
||||
LiveEntitySessionSink entities,
|
||||
LiveEnvironmentSessionSink environment,
|
||||
LiveInventorySessionBindings inventory,
|
||||
LiveCharacterSessionBindings character,
|
||||
LiveSocialSessionBindings social)
|
||||
{
|
||||
ArgumentNullException.ThrowIfNull(entities);
|
||||
ArgumentNullException.ThrowIfNull(environment);
|
||||
ArgumentNullException.ThrowIfNull(inventory);
|
||||
ArgumentNullException.ThrowIfNull(character);
|
||||
ArgumentNullException.ThrowIfNull(social);
|
||||
ArgumentNullException.ThrowIfNull(entities.Spawned);
|
||||
ArgumentNullException.ThrowIfNull(entities.Deleted);
|
||||
ArgumentNullException.ThrowIfNull(entities.PickedUp);
|
||||
ArgumentNullException.ThrowIfNull(entities.MotionUpdated);
|
||||
ArgumentNullException.ThrowIfNull(entities.PositionUpdated);
|
||||
ArgumentNullException.ThrowIfNull(entities.VectorUpdated);
|
||||
ArgumentNullException.ThrowIfNull(entities.StateUpdated);
|
||||
ArgumentNullException.ThrowIfNull(entities.ParentUpdated);
|
||||
ArgumentNullException.ThrowIfNull(entities.TeleportStarted);
|
||||
ArgumentNullException.ThrowIfNull(entities.AppearanceUpdated);
|
||||
ArgumentNullException.ThrowIfNull(entities.PlayPhysicsScript);
|
||||
ArgumentNullException.ThrowIfNull(entities.PlayPhysicsScriptType);
|
||||
ArgumentNullException.ThrowIfNull(environment.EnvironChanged);
|
||||
ArgumentNullException.ThrowIfNull(environment.ServerTimeUpdated);
|
||||
ArgumentNullException.ThrowIfNull(inventory.Objects);
|
||||
ArgumentNullException.ThrowIfNull(inventory.LocalPlayer);
|
||||
ArgumentNullException.ThrowIfNull(inventory.PlayerGuid);
|
||||
ArgumentNullException.ThrowIfNull(character.Combat);
|
||||
ArgumentNullException.ThrowIfNull(character.Spellbook);
|
||||
ArgumentNullException.ThrowIfNull(social.Chat);
|
||||
ArgumentNullException.ThrowIfNull(social.TurbineChat);
|
||||
}
|
||||
}
|
||||
|
||||
internal sealed class LiveSessionSubscriptionSet : IDisposable
|
||||
{
|
||||
private readonly object _gate = new();
|
||||
private readonly List<RetryableSubscription> _subscriptions = [];
|
||||
private bool _disposeRequested;
|
||||
|
||||
public void Add(IDisposable subscription)
|
||||
{
|
||||
ArgumentNullException.ThrowIfNull(subscription);
|
||||
AddRetained(new RetryableSubscription(subscription.Dispose));
|
||||
}
|
||||
|
||||
public void Add(Action unsubscribe)
|
||||
{
|
||||
ArgumentNullException.ThrowIfNull(unsubscribe);
|
||||
AddRetained(new RetryableSubscription(unsubscribe));
|
||||
}
|
||||
|
||||
private void AddRetained(RetryableSubscription retained)
|
||||
{
|
||||
bool disposeNow;
|
||||
lock (_gate)
|
||||
{
|
||||
disposeNow = _disposeRequested;
|
||||
_subscriptions.Add(retained);
|
||||
}
|
||||
|
||||
if (!disposeNow)
|
||||
return;
|
||||
retained.Dispose();
|
||||
throw new ObjectDisposedException(nameof(LiveSessionSubscriptionSet));
|
||||
}
|
||||
|
||||
public void Dispose()
|
||||
{
|
||||
RetryableSubscription[] subscriptions;
|
||||
lock (_gate)
|
||||
{
|
||||
_disposeRequested = true;
|
||||
subscriptions = _subscriptions.ToArray();
|
||||
}
|
||||
|
||||
List<Exception>? errors = null;
|
||||
for (int index = subscriptions.Length - 1; index >= 0; index--)
|
||||
{
|
||||
try
|
||||
{
|
||||
subscriptions[index].Dispose();
|
||||
}
|
||||
catch (Exception error)
|
||||
{
|
||||
(errors ??= []).Add(error);
|
||||
}
|
||||
}
|
||||
|
||||
if (errors is not null)
|
||||
throw new AggregateException(
|
||||
"one or more live-session subscriptions failed to detach",
|
||||
errors);
|
||||
}
|
||||
|
||||
private sealed class RetryableSubscription(Action dispose) : IDisposable
|
||||
{
|
||||
private readonly object _gate = new();
|
||||
private Action? _dispose = dispose;
|
||||
private bool _executing;
|
||||
private int _executingThreadId;
|
||||
|
||||
public void Dispose()
|
||||
{
|
||||
Action? operation;
|
||||
int threadId = Environment.CurrentManagedThreadId;
|
||||
lock (_gate)
|
||||
{
|
||||
while (_executing)
|
||||
{
|
||||
if (_executingThreadId == threadId)
|
||||
{
|
||||
throw new InvalidOperationException(
|
||||
"Live-session subscription cleanup cannot complete reentrantly.");
|
||||
}
|
||||
Monitor.Wait(_gate);
|
||||
}
|
||||
|
||||
operation = _dispose;
|
||||
if (operation is null)
|
||||
return;
|
||||
_executing = true;
|
||||
_executingThreadId = threadId;
|
||||
}
|
||||
|
||||
try
|
||||
{
|
||||
operation();
|
||||
}
|
||||
catch
|
||||
{
|
||||
CompleteAttempt(succeeded: false);
|
||||
throw;
|
||||
}
|
||||
|
||||
CompleteAttempt(succeeded: true);
|
||||
}
|
||||
|
||||
private void CompleteAttempt(bool succeeded)
|
||||
{
|
||||
lock (_gate)
|
||||
{
|
||||
if (succeeded)
|
||||
_dispose = null;
|
||||
_executing = false;
|
||||
_executingThreadId = 0;
|
||||
Monitor.PulseAll(_gate);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
Loading…
Add table
Add a link
Reference in a new issue