Files
2026-08-17 18:18:13 +08:00

574 lines
30 KiB
C#

using System.Reflection;
using System.Security.Cryptography;
using System.Text;
using Newtonsoft.Json;
using ReplacedPerson.Core;
using ShrinkNetwork;
using ShrinkNetwork.ServerHost;
var options = ReplacedServerOptions.Parse(args);
if (options.SelfTest)
{
ReplacedServerSelfTest.Run();
return;
}
if (options.TransportSelfTest)
{
await ReplacedTransportSelfTest.RunAsync(options);
return;
}
ShrinkNetworkLogger.InfoHandler = Console.WriteLine;
ShrinkNetworkLogger.WarningHandler = value => Console.WriteLine("[Warn] " + value);
ShrinkNetworkLogger.ErrorHandler = Console.Error.WriteLine;
ShrinkNetworkLogger.ExceptionHandler = Console.Error.WriteLine;
var host = new ReplacedDedicatedHost(options);
using var shutdown = new CancellationTokenSource();
Console.CancelKeyPress += (_, eventArgs) => { eventArgs.Cancel = true; shutdown.Cancel(); };
host.Start();
Console.WriteLine($"被替代之人独立服务器已启动。TCP 0.0.0.0:{options.TcpPort} / KCP 0.0.0.0:{options.KcpPort}");
Console.WriteLine($"Mod 目录:{Path.GetFullPath(options.ModDirectory)}");
Console.WriteLine("管理命令:status | rooms | players <room> | kick <room> <player> | export <room> | reload-mods | quit");
await host.RunConsoleAsync(shutdown);
host.Stop();
internal sealed class ReplacedDedicatedHost
{
private readonly object _gate = new();
private readonly ReplacedServerOptions _options;
private readonly TcpServerTransport _tcp;
private readonly KcpServerTransport _kcp;
private readonly Dictionary<string, RoomEntry> _rooms = new(StringComparer.Ordinal);
private readonly Dictionary<ConnectionKey, PlayerBinding> _bindings = new();
private ReplacedServerContent _content;
public ReplacedDedicatedHost(ReplacedServerOptions options)
{
_options = options;
_content = ReplacedServerModLoader.Load(options.ModDirectory);
_tcp = new TcpServerTransport(System.Net.IPAddress.Any, options.TcpPort);
_kcp = new KcpServerTransport(System.Net.IPAddress.Any, options.KcpPort);
_tcp.OnEvent += value => OnTransport("tcp", _tcp, value);
_kcp.OnEvent += value => OnTransport("kcp", _kcp, value);
}
public void Start() { _tcp.Start(); _kcp.Start(); }
public void Stop() { _tcp.Stop(); _kcp.Stop(); }
public async Task RunConsoleAsync(CancellationTokenSource shutdown)
{
while (!shutdown.IsCancellationRequested)
{
string? line;
try { line = await Console.In.ReadLineAsync(shutdown.Token); }
catch (OperationCanceledException) { break; }
if (line == null) { await Task.Delay(100, shutdown.Token); continue; }
var parts = line.Split(' ', StringSplitOptions.RemoveEmptyEntries | StringSplitOptions.TrimEntries);
if (parts.Length == 0) continue;
switch (parts[0].ToLowerInvariant())
{
case "status": PrintStatus(); break;
case "rooms": PrintRooms(); break;
case "players" when parts.Length >= 2: PrintPlayers(parts[1]); break;
case "kick" when parts.Length >= 3: Kick(parts[1], parts[2]); break;
case "export" when parts.Length >= 2: Export(parts[1]); break;
case "reload-mods": ReloadMods(); break;
case "quit": shutdown.Cancel(); break;
default: Console.WriteLine("未知命令或参数不足。"); break;
}
}
}
private void OnTransport(string channel, IShrinkNetworkTransport transport, ShrinkNetworkTransportEvent value)
{
var key = new ConnectionKey(channel, value.SessionId);
if (value.Type == ShrinkNetworkTransportEventType.Disconnected)
{
lock (_gate)
{
if (_bindings.Remove(key, out var binding) && _rooms.TryGetValue(binding.RoomId, out var room))
room.Authority.Disconnect(binding.PlayerId, DateTime.UtcNow);
}
return;
}
if (value.Type != ShrinkNetworkTransportEventType.Packet) return;
ReplacedNetworkEnvelope? envelope;
try { envelope = JsonConvert.DeserializeObject<ReplacedNetworkEnvelope>(Encoding.UTF8.GetString(value.PacketData)); }
catch (Exception exception) { SendError(transport, value.SessionId, "protocol.invalid-json", exception.Message); return; }
if (envelope == null) { SendError(transport, value.SessionId, "protocol.empty", "empty envelope"); return; }
lock (_gate) HandleEnvelope(key, transport, envelope);
}
private void HandleEnvelope(ConnectionKey key, IShrinkNetworkTransport transport, ReplacedNetworkEnvelope envelope)
{
try
{
switch (envelope.Kind)
{
case ReplacedNetworkMessageKind.Create:
case ReplacedNetworkMessageKind.Join:
case ReplacedNetworkMessageKind.Reconnect:
Join(key, transport, envelope);
break;
case ReplacedNetworkMessageKind.Deck:
{
var (room, playerId) = Bound(key);
var deck = JsonConvert.DeserializeObject<ReplacedDeck>(envelope.Payload) ?? new ReplacedDeck();
var error = room.Authority.SetDeck(playerId, deck);
SendResult(transport, key.SessionId, string.IsNullOrEmpty(error), error, string.IsNullOrEmpty(error) ? "deck-accepted" : error);
break;
}
case ReplacedNetworkMessageKind.Ready:
{
var (room, playerId) = Bound(key);
var ready = JsonConvert.DeserializeObject<ReplacedReadyRequest>(envelope.Payload)?.Ready == true;
var error = room.Authority.SetReady(playerId, ready);
SendResult(transport, key.SessionId, string.IsNullOrEmpty(error), error, string.IsNullOrEmpty(error) ? "ready" : error);
Broadcast(room);
break;
}
case ReplacedNetworkMessageKind.MatchCommand:
{
var (room, playerId) = Bound(key);
var command = JsonConvert.DeserializeObject<ReplacedMatchCommand>(envelope.Payload) ?? new ReplacedMatchCommand();
var result = room.Authority.Apply(playerId, command);
SendResult(transport, key.SessionId, result.Accepted, result.ErrorCode,
result.Duplicate ? "duplicate" : result.Message, result.Duplicate);
if (result.Accepted)
{
ApplyQaTimeoutIfNeeded(room);
Broadcast(room);
}
break;
}
default:
SendError(transport, key.SessionId, "protocol.unsupported", envelope.Kind.ToString());
break;
}
}
catch (Exception exception)
{
SendError(transport, key.SessionId, "authority.exception", exception.Message);
}
}
private void Join(ConnectionKey key, IShrinkNetworkTransport transport, ReplacedNetworkEnvelope envelope)
{
var request = JsonConvert.DeserializeObject<ReplacedJoinRequest>(envelope.Payload) ?? new ReplacedJoinRequest();
var roomId = string.IsNullOrWhiteSpace(request.RoomId) ? SelectOrCreateRoomId() : request.RoomId.Trim();
if (!_rooms.TryGetValue(roomId, out var room))
{
if (envelope.Kind == ReplacedNetworkMessageKind.Join && _rooms.Count >= _options.MaxRooms)
{
SendError(transport, key.SessionId, "room.limit", "server room limit reached");
return;
}
room = new RoomEntry(new ReplacedRoomAuthority(roomId, _content.Catalog, _content.Manifests,
unchecked((ulong)DateTime.UtcNow.Ticks)));
_rooms.Add(roomId, room);
}
request.RoomId = roomId;
var compatibility = room.Authority.TryJoin(request, DateTime.UtcNow);
if (!compatibility.Compatible)
{
SendError(transport, key.SessionId, "mod.mismatch", string.Join(";", compatibility.Differences));
return;
}
_bindings[key] = new PlayerBinding(roomId, request.PlayerId, transport);
SendResult(transport, key.SessionId, true, "joined", roomId);
Broadcast(room);
}
private string SelectOrCreateRoomId()
{
foreach (var pair in _rooms)
if (pair.Value.Authority.Seats.Any(seat => seat == null)) return pair.Key;
return "RP-" + Guid.NewGuid().ToString("N")[..6].ToUpperInvariant();
}
private (RoomEntry room, string playerId) Bound(ConnectionKey key)
{
if (!_bindings.TryGetValue(key, out var binding)) throw new InvalidOperationException("session.not-joined");
if (!_rooms.TryGetValue(binding.RoomId, out var room)) throw new InvalidOperationException("room.missing");
return (room, binding.PlayerId);
}
private void Broadcast(RoomEntry room)
{
foreach (var pair in _bindings.Where(value => value.Value.RoomId == room.Authority.RoomId).ToArray())
{
var payload = new ReplacedMatchSnapshotPayload
{
Room = room.Authority.Snapshot(),
Match = room.Authority.PlayerMatchSnapshot(pair.Value.PlayerId),
LocalPlayerIndex = room.Authority.GetPlayerIndex(pair.Value.PlayerId)
};
Send(pair.Value.Transport, pair.Key.SessionId, new ReplacedNetworkEnvelope
{
Kind = ReplacedNetworkMessageKind.MatchSnapshot,
RoomId = room.Authority.RoomId,
PlayerId = pair.Value.PlayerId,
StateHash = room.Authority.Match?.ComputeStateHash() ?? string.Empty,
Payload = JsonConvert.SerializeObject(payload)
});
}
if (room.Authority.Match?.State.Phase == ReplacedMatchPhase.Completed)
Console.WriteLine($"[Match] completed room={room.Authority.RoomId} winner={room.Authority.Match.State.Winner} rounds={room.Authority.Match.State.Round} hash={room.Authority.Match.ComputeStateHash()}");
}
private void ApplyQaTimeoutIfNeeded(RoomEntry room)
{
if (_options.QaAutoTimeoutRound <= 0 || room.Authority.Match == null ||
room.Authority.Match.State.Round < _options.QaAutoTimeoutRound ||
room.Authority.Match.State.Phase == ReplacedMatchPhase.Completed || room.Authority.Seats[1] == null) return;
var playerId = room.Authority.Seats[1]!.PlayerId;
room.Authority.RecordTurnTimeout(playerId);
room.Authority.RecordTurnTimeout(playerId);
}
private void PrintStatus()
{
Console.WriteLine($"rooms={_rooms.Count}/{_options.MaxRooms} sessions={_bindings.Count} mods={_content.Manifests.Count} content={_content.Catalog.ComputeHash()}");
}
private void PrintRooms()
{
foreach (var room in _rooms.Values)
{
var snapshot = room.Authority.Snapshot();
Console.WriteLine($"{snapshot.RoomId} players={snapshot.Players.Count}/2 ready={snapshot.ReadyPlayers.Count}/2 match={snapshot.MatchStarted} hash={snapshot.MatchStateHash}");
}
}
private void PrintPlayers(string roomId)
{
if (!_rooms.TryGetValue(roomId, out var room)) { Console.WriteLine("room not found"); return; }
foreach (var seat in room.Authority.Seats)
if (seat != null) Console.WriteLine($"{seat.PlayerId} name={seat.DisplayName} connected={seat.Connected} ready={seat.Ready} timeouts={seat.ConsecutiveTimeouts}");
}
private void Kick(string roomId, string playerId)
{
if (!_rooms.TryGetValue(roomId, out var room) || !room.Authority.Kick(playerId)) { Console.WriteLine("player not found"); return; }
foreach (var pair in _bindings.Where(value => value.Value.RoomId == roomId && value.Value.PlayerId == playerId).ToArray())
{
if (pair.Value.Transport is IShrinkNetworkSessionControlTransport control)
control.DisconnectSession(pair.Key.SessionId, "kicked by local console");
_bindings.Remove(pair.Key);
}
Console.WriteLine("kicked " + playerId);
}
private void Export(string roomId)
{
if (!_rooms.TryGetValue(roomId, out var room)) { Console.WriteLine("room not found"); return; }
Directory.CreateDirectory(_options.ExportDirectory);
var path = Path.Combine(_options.ExportDirectory, roomId + "-" + DateTime.UtcNow.ToString("yyyyMMdd-HHmmss") + ".json");
File.WriteAllText(path, JsonConvert.SerializeObject(new
{
room = room.Authority.Snapshot(),
state = room.Authority.Match?.State,
replay = room.Authority.Match?.Replay
}, Formatting.Indented));
Console.WriteLine(Path.GetFullPath(path));
}
private void ReloadMods()
{
_content = ReplacedServerModLoader.Load(_options.ModDirectory);
Console.WriteLine("Mod 已重载;只影响之后创建的房间。content=" + _content.Catalog.ComputeHash());
}
private static void Send(IShrinkNetworkTransport transport, long sessionId, ReplacedNetworkEnvelope envelope) =>
transport.Send(sessionId, Encoding.UTF8.GetBytes(JsonConvert.SerializeObject(envelope)));
private static void SendResult(IShrinkNetworkTransport transport, long sessionId, bool accepted, string code, string detail, bool duplicate = false) =>
Send(transport, sessionId, new ReplacedNetworkEnvelope
{
Kind = accepted ? ReplacedNetworkMessageKind.MatchEvent : ReplacedNetworkMessageKind.Error,
Payload = JsonConvert.SerializeObject(new ReplacedNetworkResultPayload { Accepted = accepted, Duplicate = duplicate, Code = code, Detail = detail })
});
private static void SendError(IShrinkNetworkTransport transport, long sessionId, string code, string detail) =>
SendResult(transport, sessionId, false, code, detail);
private readonly record struct ConnectionKey(string Channel, long SessionId);
private sealed record PlayerBinding(string RoomId, string PlayerId, IShrinkNetworkTransport Transport);
private sealed record RoomEntry(ReplacedRoomAuthority Authority);
}
internal sealed record ReplacedServerContent(ReplacedContentCatalog Catalog, List<ReplacedPersonModManifest> Manifests,
Dictionary<string, Func<IReadOnlyList<string>, string>> Commands);
internal static class ReplacedServerModLoader
{
public static ReplacedServerContent Load(string directory)
{
var registry = new ReplacedPersonRegistry(ReplacedBuiltInContent.Create());
var manifests = new List<ReplacedPersonModManifest>();
if (!Directory.Exists(directory)) return new ReplacedServerContent(registry.CreateSnapshot(), manifests, new(registry.Commands));
foreach (var manifestPath in Directory.EnumerateFiles(directory, "manifest.json", SearchOption.AllDirectories).OrderBy(value => value, StringComparer.Ordinal))
{
var packageRoot = Path.GetDirectoryName(Path.GetFullPath(manifestPath))!;
var source = JsonConvert.DeserializeObject<PackageManifest>(File.ReadAllText(manifestPath))
?? throw new InvalidDataException("Invalid manifest: " + manifestPath);
var manifest = new ReplacedPersonModManifest { Id = source.Id, Version = source.Version, Files = source.Files };
foreach (var file in manifest.Files)
{
var fullPath = Path.GetFullPath(Path.Combine(packageRoot, file.Path));
if (!fullPath.StartsWith(packageRoot + Path.DirectorySeparatorChar, StringComparison.OrdinalIgnoreCase))
throw new InvalidDataException("Mod path escapes package: " + file.Path);
if (!File.Exists(fullPath)) throw new FileNotFoundException("Mod file missing", fullPath);
var actual = Convert.ToHexString(SHA256.HashData(File.ReadAllBytes(fullPath))).ToLowerInvariant();
if (!string.Equals(actual, file.Sha256, StringComparison.OrdinalIgnoreCase))
throw new InvalidDataException($"Mod hash mismatch: {source.Id}/{file.Path}");
}
var rulesPath = Path.GetFullPath(Path.Combine(packageRoot, source.RulesAssembly));
var assembly = Assembly.Load(File.ReadAllBytes(rulesPath));
var modType = assembly.GetTypes().Single(type => typeof(IReplacedPersonMod).IsAssignableFrom(type) && !type.IsAbstract && !type.IsInterface);
var mod = (IReplacedPersonMod)Activator.CreateInstance(modType)!;
if (mod.Id != source.Id || mod.Version != source.Version) throw new InvalidDataException("Manifest identity does not match rules assembly: " + source.Id);
mod.Register(registry);
manifests.Add(manifest);
Console.WriteLine($"[Mod] loaded {mod.Id}@{mod.Version}");
}
return new ReplacedServerContent(registry.CreateSnapshot(), manifests, new(registry.Commands));
}
private sealed class PackageManifest
{
public string Id { get; set; } = string.Empty;
public string Version { get; set; } = string.Empty;
public string RulesAssembly { get; set; } = string.Empty;
public List<ReplacedPersonModFile> Files { get; set; } = new();
}
}
internal sealed class ReplacedServerOptions
{
public int TcpPort { get; private set; } = 18140;
public int KcpPort { get; private set; } = 18141;
public int MaxRooms { get; private set; } = 64;
public string ModDirectory { get; private set; } = "Mods";
public string ExportDirectory { get; private set; } = "StateExports";
public bool SelfTest { get; private set; }
public bool TransportSelfTest { get; private set; }
public int QaAutoTimeoutRound { get; private set; }
public static ReplacedServerOptions Parse(string[] args)
{
var result = new ReplacedServerOptions();
for (var i = 0; i < args.Length; i++)
{
switch (args[i])
{
case "--tcp-port" when i + 1 < args.Length: result.TcpPort = int.Parse(args[++i]); break;
case "--kcp-port" when i + 1 < args.Length: result.KcpPort = int.Parse(args[++i]); break;
case "--max-rooms" when i + 1 < args.Length: result.MaxRooms = int.Parse(args[++i]); break;
case "--mods" when i + 1 < args.Length: result.ModDirectory = args[++i]; break;
case "--exports" when i + 1 < args.Length: result.ExportDirectory = args[++i]; break;
case "--self-test": result.SelfTest = true; break;
case "--transport-self-test": result.TransportSelfTest = true; break;
case "--qa-auto-timeout-round" when i + 1 < args.Length: result.QaAutoTimeoutRound = int.Parse(args[++i]); break;
}
}
return result;
}
}
internal static class ReplacedTransportSelfTest
{
public static async Task RunAsync(ReplacedServerOptions options)
{
ShrinkNetworkLogger.InfoHandler = _ => { };
ShrinkNetworkLogger.WarningHandler = _ => { };
ShrinkNetworkLogger.ErrorHandler = Console.Error.WriteLine;
ShrinkNetworkLogger.ExceptionHandler = Console.Error.WriteLine;
var host = new ReplacedDedicatedHost(options);
host.Start();
try
{
await RunPairAsync("TCP", () => new ShrinkTcpClientTransport("127.0.0.1", options.TcpPort));
await RunPairAsync("KCP", () => new ShrinkKcpClientTransport("127.0.0.1", options.KcpPort));
}
finally { host.Stop(); }
Console.WriteLine("TRANSPORT SELF-TEST PASS TCP+KCP create/join/deck/ready/command/duplicate/hash");
}
private static async Task RunPairAsync(string name, Func<IShrinkNetworkTransport> factory)
{
using var first = new WireClient(name.ToLowerInvariant() + "-a", factory());
using var second = new WireClient(name.ToLowerInvariant() + "-b", factory());
first.Start(); second.Start();
await Task.WhenAll(first.WaitConnected(), second.WaitConnected());
first.Send(ReplacedNetworkMessageKind.Join, Join(first.PlayerId));
second.Send(ReplacedNetworkMessageKind.Join, Join(second.PlayerId));
await Task.WhenAll(first.WaitResult(value => value.Code == "joined"), second.WaitResult(value => value.Code == "joined"));
first.Send(ReplacedNetworkMessageKind.Deck, new ReplacedDeck());
var invalid = await first.WaitResult(value => !value.Accepted);
if (!invalid.Detail.Contains("deck.normal.count", StringComparison.Ordinal)) throw new InvalidOperationException(name + " invalid deck was not rejected");
var deck = ReplacedBuiltInContent.CreateStarterDeck();
first.Send(ReplacedNetworkMessageKind.Deck, deck); second.Send(ReplacedNetworkMessageKind.Deck, deck);
await Task.WhenAll(first.WaitResult(value => value.Detail == "deck-accepted"), second.WaitResult(value => value.Detail == "deck-accepted"));
first.Send(ReplacedNetworkMessageKind.Ready, new ReplacedReadyRequest { Ready = true });
second.Send(ReplacedNetworkMessageKind.Ready, new ReplacedReadyRequest { Ready = true });
var snapshots = await Task.WhenAll(first.WaitMatch(), second.WaitMatch());
var leadClient = snapshots[0].Match!.LeadPlayer == snapshots[0].LocalPlayerIndex ? first : second;
var responseClient = ReferenceEquals(leadClient, first) ? second : first;
var leadSnapshot = ReferenceEquals(leadClient, first) ? snapshots[0] : snapshots[1];
var leadCommand = PassCommand(leadSnapshot.LocalPlayerIndex, leadSnapshot.Match!.LastAcceptedSequence[leadSnapshot.LocalPlayerIndex] + 1, name + "-lead");
leadClient.Send(ReplacedNetworkMessageKind.MatchCommand, leadCommand);
await responseClient.WaitSnapshot(value => value.Match?.Phase == ReplacedMatchPhase.AwaitingResponse);
var responseState = responseClient.LatestMatch!;
var responseCommand = PassCommand(responseClient.LocalPlayerIndex, responseState.LastAcceptedSequence[responseClient.LocalPlayerIndex] + 1, name + "-response");
responseClient.Send(ReplacedNetworkMessageKind.MatchCommand, responseCommand);
var roundTwo = await Task.WhenAll(first.WaitSnapshot(value => value.Match?.Round >= 2), second.WaitSnapshot(value => value.Match?.Round >= 2));
if (!string.Equals(first.LatestStateHash, second.LatestStateHash, StringComparison.Ordinal) || string.IsNullOrWhiteSpace(first.LatestStateHash))
throw new InvalidOperationException(name + " state hash mismatch");
leadClient.Send(ReplacedNetworkMessageKind.MatchCommand, leadCommand);
var duplicate = await leadClient.WaitResult(value => value.Duplicate);
if (!duplicate.Accepted) throw new InvalidOperationException(name + " duplicate was rejected");
Console.WriteLine(name + " PASS hash=" + first.LatestStateHash);
}
private static ReplacedJoinRequest Join(string id) => new() { PlayerId = id, DisplayName = id };
private static ReplacedMatchCommand PassCommand(int player, long sequence, string token) => new()
{
PlayerIndex = player, Sequence = sequence, IdempotencyToken = token, Submission = new ReplacedTurnSubmission { Pass = true }
};
private sealed class WireClient : IDisposable
{
private readonly IShrinkNetworkTransport _transport;
private readonly List<ReplacedNetworkEnvelope> _messages = new();
private readonly object _gate = new();
private readonly TaskCompletionSource _connected = new(TaskCreationOptions.RunContinuationsAsynchronously);
private long _sessionId;
public string PlayerId { get; }
public int LocalPlayerIndex { get; private set; } = -1;
public ReplacedMatchState? LatestMatch { get; private set; }
public string LatestStateHash { get; private set; } = string.Empty;
public WireClient(string playerId, IShrinkNetworkTransport transport)
{
PlayerId = playerId; _transport = transport; _transport.OnEvent += OnEvent;
}
public void Start() => _transport.Start();
public async Task WaitConnected() => await _connected.Task.WaitAsync(TimeSpan.FromSeconds(8));
public void Send(ReplacedNetworkMessageKind kind, object payload)
{
var envelope = new ReplacedNetworkEnvelope { Kind = kind, PlayerId = PlayerId, Payload = JsonConvert.SerializeObject(payload) };
_transport.Send(_sessionId, Encoding.UTF8.GetBytes(JsonConvert.SerializeObject(envelope)));
}
public async Task<ReplacedNetworkResultPayload> WaitResult(Func<ReplacedNetworkResultPayload, bool> predicate)
{
var envelope = await WaitEnvelope(value => value.Kind is ReplacedNetworkMessageKind.MatchEvent or ReplacedNetworkMessageKind.Error &&
PredicateResult(value, predicate));
return JsonConvert.DeserializeObject<ReplacedNetworkResultPayload>(envelope.Payload)!;
}
public Task<ReplacedMatchSnapshotPayload> WaitMatch() => WaitSnapshot(value => value.Match != null);
public async Task<ReplacedMatchSnapshotPayload> WaitSnapshot(Func<ReplacedMatchSnapshotPayload, bool> predicate)
{
var envelope = await WaitEnvelope(value => value.Kind == ReplacedNetworkMessageKind.MatchSnapshot && PredicateSnapshot(value, predicate));
return JsonConvert.DeserializeObject<ReplacedMatchSnapshotPayload>(envelope.Payload)!;
}
private async Task<ReplacedNetworkEnvelope> WaitEnvelope(Func<ReplacedNetworkEnvelope, bool> predicate)
{
var deadline = DateTime.UtcNow.AddSeconds(10);
while (DateTime.UtcNow < deadline)
{
lock (_gate)
{
var index = _messages.FindIndex(value => predicate(value));
if (index >= 0) { var result = _messages[index]; _messages.RemoveAt(index); return result; }
}
await Task.Delay(20);
}
throw new TimeoutException(PlayerId + " timed out waiting for protocol message");
}
private static bool PredicateResult(ReplacedNetworkEnvelope envelope, Func<ReplacedNetworkResultPayload, bool> predicate)
{
var result = JsonConvert.DeserializeObject<ReplacedNetworkResultPayload>(envelope.Payload);
return result != null && predicate(result);
}
private static bool PredicateSnapshot(ReplacedNetworkEnvelope envelope, Func<ReplacedMatchSnapshotPayload, bool> predicate)
{
var result = JsonConvert.DeserializeObject<ReplacedMatchSnapshotPayload>(envelope.Payload);
return result != null && predicate(result);
}
private void OnEvent(ShrinkNetworkTransportEvent value)
{
if (value.Type == ShrinkNetworkTransportEventType.Connected)
{
_sessionId = value.SessionId; _connected.TrySetResult(); return;
}
if (value.Type != ShrinkNetworkTransportEventType.Packet) return;
var envelope = JsonConvert.DeserializeObject<ReplacedNetworkEnvelope>(Encoding.UTF8.GetString(value.PacketData));
if (envelope == null) return;
if (envelope.Kind == ReplacedNetworkMessageKind.MatchSnapshot)
{
var snapshot = JsonConvert.DeserializeObject<ReplacedMatchSnapshotPayload>(envelope.Payload);
if (snapshot != null)
{
LocalPlayerIndex = snapshot.LocalPlayerIndex;
LatestMatch = snapshot.Match;
LatestStateHash = envelope.StateHash;
}
}
lock (_gate) _messages.Add(envelope);
}
public void Dispose()
{
_transport.OnEvent -= OnEvent;
// ShrinkTcpClientTransport 2022 runtime reads RemoteEndPoint in its receive-loop finally block.
// Let the authority close it during host shutdown so that endpoint remains valid for that readback.
if (_transport is not ShrinkTcpClientTransport) _transport.Stop();
}
}
}
internal static class ReplacedServerSelfTest
{
public static void Run()
{
var content = ReplacedBuiltInContent.Create();
var deck = ReplacedBuiltInContent.CreateStarterDeck();
var room = new ReplacedRoomAuthority("self-test", content, Array.Empty<ReplacedPersonModManifest>(), 1414);
room.TryJoin(Join("a"), DateTime.UtcNow); room.TryJoin(Join("b"), DateTime.UtcNow);
if (!string.IsNullOrEmpty(room.SetDeckAndReady("a", deck)) || !string.IsNullOrEmpty(room.SetDeckAndReady("b", deck)))
throw new InvalidOperationException("self-test ready failed");
var commands = new List<ReplacedMatchCommand>();
for (var i = 0; i < 2; i++)
{
var lead = room.Match!.State.LeadPlayer;
var first = Command(lead, room.Match.State.LastAcceptedSequence[lead] + 1, "self-" + i + "-a");
var second = Command(1 - lead, room.Match.State.LastAcceptedSequence[1 - lead] + 1, "self-" + i + "-b");
if (!room.Apply(lead == 0 ? "a" : "b", first).Accepted || !room.Apply(lead == 0 ? "b" : "a", second).Accepted)
throw new InvalidOperationException("self-test command failed");
commands.Add(first); commands.Add(second);
}
var replay = ReplacedReplayRunner.Replay(content, deck, deck, 1414, commands);
if (replay.ComputeStateHash() != room.Match!.ComputeStateHash()) throw new InvalidOperationException("self-test hash mismatch");
Console.WriteLine("SELF-TEST PASS stateHash=" + replay.ComputeStateHash());
}
private static ReplacedJoinRequest Join(string id) => new() { RoomId = "self-test", PlayerId = id, DisplayName = id };
private static ReplacedMatchCommand Command(int player, long sequence, string token) => new()
{
PlayerIndex = player, Sequence = sequence, IdempotencyToken = token, Submission = new ReplacedTurnSubmission { Pass = true }
};
}