Complete C# server cutover

This commit is contained in:
Nomads_Reach
2026-08-16 02:36:32 -04:00
parent ddadf87423
commit 757e614cdd
79 changed files with 923 additions and 13129 deletions
+101 -127
View File
@@ -14,18 +14,17 @@ internal static class Program
{
await Run("transport policy", TestTransportPolicy);
await Run("packet codec", TestPacketCodec);
await Run("snapshot sequence", TestSnapshotSequence);
await Run("snapshot sequencing", TestSnapshotSequencing);
await Run("snapshot envelope", TestSnapshotEnvelope);
await Run("player state validation", TestPlayerStateValidation);
await Run("npc authority epochs", TestNpcAuthority);
await Run("npc authority epochs", TestNpcAuthorityEpochs);
await Run("interest filtering", TestInterestFiltering);
await Run("config compatibility", TestConfigCompatibility);
await Run("ban persistence", TestBanStore);
await Run("authoritative server ids and interest", TestAuthoritativeServer);
await Run("durable player state relay", TestPlayerStateRelay);
await Run("npc authority validation", TestNpcAuthorityIntegration);
await Run("combat interest routing", TestCombatInterest);
await Run("ban persistence", TestBanPersistence);
await Run("server-owned ids and interest relay", TestServerOwnedIdsAndInterest);
await Run("durable player state relay", TestDurablePlayerStateRelay);
await Run("npc authority enforcement", TestNpcAuthorityEnforcement);
await Run("combat interest enforcement", TestCombatInterest);
Console.WriteLine($"C# server tests: {_passed} passed, {_failed} failed");
return _failed == 0 ? 0 : 1;
}
@@ -51,14 +50,13 @@ internal static class Program
var encoded = PacketCodec.Encode(packet);
Equal("playerState", encoded.PacketType);
Equal(Delivery.ReliableOrdered, encoded.Delivery);
var decoded = PacketCodec.Decode(encoded.Payload);
Equal("Nomad", JsonHelpers.String(decoded["characterName"]));
Throws<PacketCodecException>(() => PacketCodec.Encode(new JsonObject { ["type"] = "x", ["blob"] = new string('a', ProtocolConstants.MaxMessageBytes + 1) }));
Equal("Nomad", JsonHelpers.String(PacketCodec.Decode(encoded.Payload)["characterName"]));
Throws<PacketCodecException>(() => PacketCodec.Decode("[]"u8));
Throws<PacketCodecException>(() => PacketCodec.Encode(new JsonObject { ["type"] = "oversize", ["payload"] = new string('x', ProtocolConstants.MaxMessageBytes + 1) }));
return Task.CompletedTask;
}
private static Task TestSnapshotSequence()
private static Task TestSnapshotSequencing()
{
var window = new SequenceWindow();
True(window.Accept(1));
@@ -66,28 +64,26 @@ internal static class Program
False(window.Accept(2));
False(window.Accept(3));
True(SequenceWindow.IsNewer(1, uint.MaxValue));
var counter = new SequenceCounter();
Equal(1u, counter.Advance());
False(SequenceWindow.IsNewer(uint.MaxValue, 1));
return Task.CompletedTask;
}
private static Task TestSnapshotEnvelope()
{
var payload = Encoding.UTF8.GetBytes("{\"type\":\"transform\"}");
var wire = GnsSnapshotEnvelope.Encode("transform", payload, 7);
Equal(ProtocolConstants.MaxMessageBytes >= wire.Length, true);
True(GnsSnapshotEnvelope.TryDecode(wire, out var envelope, out var error));
var encoded = GnsSnapshotEnvelope.Encode("transform", payload, 7);
True(GnsSnapshotEnvelope.TryDecode(encoded, out var envelope, out var error));
Equal<string?>(null, error);
Equal("transform", envelope.PacketType);
Equal(7u, envelope.Sequence);
SequenceEqual(payload, envelope.Payload);
True(payload.AsSpan().SequenceEqual(envelope.Payload));
False(GnsSnapshotEnvelope.TryDecode(payload, out _, out _));
return Task.CompletedTask;
}
private static Task TestPlayerStateValidation()
{
var state = new JsonObject
var packet = new JsonObject
{
["type"] = "playerState",
["characterName"] = "Nomad",
@@ -106,28 +102,32 @@ internal static class Program
},
["actionEvents"] = new JsonArray(new JsonObject { ["sequence"] = 1, ["type"] = 3, ["eventName"] = "fireSingle" })
};
var normalized = ProtocolValidation.NormalizePlayerState(state);
var normalized = ProtocolValidation.NormalizePlayerState(packet);
NotNull(normalized);
Equal("00000123", JsonHelpers.String(((JsonObject)normalized!["appearance"]!)["hairColorFormId"]));
var bad = (JsonObject)state.DeepClone();
((JsonArray)bad["actionEvents"]!)[0]!["eventName"] = "bogus";
Equal<JsonObject?>(null, ProtocolValidation.NormalizePlayerState(bad));
var appearance = normalized!["appearance"] as JsonObject;
NotNull(appearance);
Equal("00000123", JsonHelpers.String(appearance!["hairColorFormId"]));
var invalidActions = new JsonObject
{
["type"] = "playerState",
["actionEvents"] = new JsonArray(new JsonObject { ["sequence"] = 1, ["type"] = 3, ["eventName"] = "invalid" })
};
Equal<JsonObject?>(null, ProtocolValidation.NormalizePlayerState(invalidActions));
return Task.CompletedTask;
}
private static Task TestNpcAuthority()
private static Task TestNpcAuthorityEpochs()
{
var manager = new NpcAuthorityManager();
var scope = new ScopeKey("00000001", "");
var changes = manager.Reconcile(new[] { (2u, scope), (1u, scope) });
Equal(1, changes.Count);
Equal(1u, changes[0].PlayerId);
Equal(1u, changes[0].Epoch);
var first = manager.Reconcile(new[] { (2u, scope), (1u, scope) });
Equal(1u, first.Single().PlayerId);
Equal(1u, first.Single().Epoch);
True(manager.Authorize(1, scope, 1));
changes = manager.Reconcile(Array.Empty<(uint, ScopeKey)>());
Equal(2u, changes[0].Epoch);
changes = manager.Reconcile(new[] { (3u, scope) });
Equal(3u, changes[0].Epoch);
var revoked = manager.Reconcile(Array.Empty<(uint PlayerId, ScopeKey Scope)>());
Equal(2u, revoked.Single().Epoch);
var regrant = manager.Reconcile(new[] { (3u, scope) });
Equal(3u, regrant.Single().Epoch);
False(manager.Authorize(1, scope, 1));
True(manager.Authorize(3, scope, 3));
return Task.CompletedTask;
@@ -136,14 +136,10 @@ internal static class Program
private static Task TestInterestFiltering()
{
var a = Transform("00000001", "000000AA", 0, 0);
var sameCell = Transform("00000001", "000000BB", 999999, 999999);
var near = Transform("00000002", "000000AA", 1000, 1000);
var far = Transform("00000003", "000000AA", 20000, 0);
var otherWorld = Transform("00000004", "000000BB", 0, 0);
True(ProtocolValidation.StatesShareInterest(a, sameCell));
True(ProtocolValidation.StatesShareInterest(a, near));
False(ProtocolValidation.StatesShareInterest(a, far));
False(ProtocolValidation.StatesShareInterest(a, otherWorld));
True(ProtocolValidation.StatesShareInterest(a, Transform("00000001", "000000BB", 500000, 500000)));
True(ProtocolValidation.StatesShareInterest(a, Transform("00000002", "000000AA", 1000, 1000)));
False(ProtocolValidation.StatesShareInterest(a, Transform("00000003", "000000AA", 20000, 0)));
False(ProtocolValidation.StatesShareInterest(a, Transform("00000004", "000000BB", 0, 0)));
return Task.CompletedTask;
}
@@ -163,53 +159,46 @@ internal static class Program
return Task.CompletedTask;
}
private static Task TestBanStore()
private static Task TestBanPersistence()
{
var dir = TempDir();
try
{
var path = Path.Combine(dir, "bans.json");
var store = new BanStore(path);
store.Ban("127.0.0.1", "test");
var reloaded = new BanStore(path);
Equal("test", reloaded.GetBan("127.0.0.1")?.Reason);
True(reloaded.Unban("127.0.0.1"));
new BanStore(path).Ban("127.0.0.1", "test");
var loaded = new BanStore(path);
Equal("test", loaded.GetBan("127.0.0.1")?.Reason);
True(loaded.Unban("127.0.0.1"));
Equal(0, new BanStore(path).List().Count);
}
finally { Directory.Delete(dir, true); }
return Task.CompletedTask;
}
private static async Task TestAuthoritativeServer()
private static async Task TestServerOwnedIdsAndInterest()
{
var fixture = CreateServerFixture();
await using var server = fixture.Server;
await using var fixture = new ServerFixture();
var a = new MemoryConnection("a", 31001);
var b = new MemoryConnection("b", 31002);
True(await server.AcceptConnectionAsync(a, CancellationToken.None));
True(await server.AcceptConnectionAsync(b, CancellationToken.None));
await Send(server, a, new JsonObject { ["type"] = "hello", ["protocolVersion"] = 2 });
await Send(server, b, new JsonObject { ["type"] = "hello", ["protocolVersion"] = 2 });
await Activate(fixture.Server, a, b);
var aId = ReadyId(a); var bId = ReadyId(b);
True(aId != bId);
await Send(server, a, Transform("00000001", "000000AA", 0, 0, "spawn"));
await Send(server, b, Transform("00000002", "000000BB", 0, 0, "spawn"));
await Send(fixture.Server, a, Transform("00000001", "000000AA", 0, 0, "spawn"));
await Send(fixture.Server, b, Transform("00000002", "000000BB", 0, 0, "spawn"));
a.Clear(); b.Clear();
await Send(server, a, Transform("00000001", "000000AA", 1, 0));
False(b.SentPackets.Any(p => JsonHelpers.String(p["type"]) == "transform" && JsonHelpers.TryUInt32(p["playerId"], 1, uint.MaxValue, out var id) && id == aId));
await Send(fixture.Server, a, Transform("00000001", "000000AA", 1, 0));
False(b.SentPackets.Any(p => JsonHelpers.String(p["type"]) == "transform" && PlayerId(p) == aId));
}
private static async Task TestPlayerStateRelay()
private static async Task TestDurablePlayerStateRelay()
{
var fixture = CreateServerFixture();
await using var server = fixture.Server;
await using var fixture = new ServerFixture();
var a = new MemoryConnection("a", 32001); var b = new MemoryConnection("b", 32002);
await Activate(server, a, b);
await Send(server, a, Transform("00000001", "000000AA", 0, 0, "spawn"));
await Send(server, b, Transform("00000002", "000000BB", 0, 0, "spawn"));
await Activate(fixture.Server, a, b);
await Send(fixture.Server, a, Transform("00000001", "000000AA", 0, 0, "spawn"));
await Send(fixture.Server, b, Transform("00000002", "000000BB", 0, 0, "spawn"));
a.Clear(); b.Clear();
await Send(server, a, new JsonObject
await Send(fixture.Server, a, new JsonObject
{
["type"] = "playerState",
["characterName"] = "Nomad",
@@ -220,62 +209,37 @@ internal static class Program
False(relay.ContainsKey("actionEvents"));
}
private static async Task TestNpcAuthorityIntegration()
private static async Task TestNpcAuthorityEnforcement()
{
var fixture = CreateServerFixture();
await using var server = fixture.Server;
await using var fixture = new ServerFixture();
var a = new MemoryConnection("a", 33001);
True(await server.AcceptConnectionAsync(a, CancellationToken.None));
await Send(server, a, new JsonObject { ["type"] = "hello", ["protocolVersion"] = 2 });
await Send(server, a, Transform("00000010", "", 0, 0, "spawn"));
await Activate(fixture.Server, a);
await Send(fixture.Server, a, Transform("00000010", "", 0, 0, "spawn"));
var authority = a.SentPackets.Last(p => JsonHelpers.String(p["type"]) == "npcAuthority");
JsonHelpers.TryUInt32(authority["authorityEpoch"], 1, uint.MaxValue, out var epoch);
True(JsonHelpers.TryUInt32(authority["authorityEpoch"], 1, uint.MaxValue, out var epoch));
a.Clear();
await Send(server, a, new JsonObject
await Send(fixture.Server, a, new JsonObject
{
["type"] = "npcState",
["authorityEpoch"] = epoch,
["authorityCellId"] = "00000010",
["authorityWorldspaceId"] = "",
["npcs"] = new JsonArray(new JsonObject
{
["sourceFormId"] = "00000020", ["cellId"] = "00000010", ["worldspaceId"] = "",
["x"] = 1.0, ["y"] = 2.0, ["z"] = 3.0, ["angleZ"] = 0.0
})
["type"] = "npcState", ["authorityEpoch"] = epoch, ["authorityCellId"] = "00000010", ["authorityWorldspaceId"] = "",
["npcs"] = new JsonArray(new JsonObject { ["sourceFormId"] = "00000020", ["cellId"] = "00000010", ["worldspaceId"] = "", ["x"] = 1.0, ["y"] = 2.0, ["z"] = 3.0, ["angleZ"] = 0.0 })
});
var stats = server.GetCoreStats();
True(JsonHelpers.TryUInt32(stats["npcStatePacketsReceived"], 0, uint.MaxValue, out var received) && received == 1);
await Send(server, a, new JsonObject
{
["type"] = "npcState", ["authorityEpoch"] = Math.Max(1u, epoch - 1), ["authorityCellId"] = "00000010", ["authorityWorldspaceId"] = "",
["npcs"] = new JsonArray()
});
stats = server.GetCoreStats();
True(JsonHelpers.TryUInt32(stats["npcAuthorityRejects"], 0, uint.MaxValue, out var rejects) && rejects >= 1);
Equal(1u, Stat(fixture.Server, "npcStatePacketsReceived"));
await Send(fixture.Server, a, new JsonObject { ["type"] = "npcState", ["authorityEpoch"] = epoch + 1, ["authorityCellId"] = "00000010", ["authorityWorldspaceId"] = "", ["npcs"] = new JsonArray() });
True(Stat(fixture.Server, "npcAuthorityRejects") >= 1);
}
private static async Task TestCombatInterest()
{
var fixture = CreateServerFixture();
await using var server = fixture.Server;
await using var fixture = new ServerFixture();
var a = new MemoryConnection("a", 34001); var b = new MemoryConnection("b", 34002);
await Activate(server, a, b);
await Activate(fixture.Server, a, b);
var bId = ReadyId(b);
await Send(server, a, Transform("00000100", "000000AA", 0, 0, "spawn"));
await Send(server, b, Transform("00000200", "000000BB", 0, 0, "spawn"));
await Send(fixture.Server, a, Transform("00000100", "000000AA", 0, 0, "spawn"));
await Send(fixture.Server, b, Transform("00000200", "000000BB", 0, 0, "spawn"));
a.Clear(); b.Clear();
await Send(server, a, new JsonObject { ["type"] = "combatHit", ["targetPlayerId"] = bId, ["sequence"] = 1, ["damage"] = 10.0 });
await Send(fixture.Server, a, new JsonObject { ["type"] = "combatHit", ["targetPlayerId"] = bId, ["sequence"] = 1, ["damage"] = 10.0 });
False(b.SentPackets.Any(p => JsonHelpers.String(p["type"]) == "combatHit"));
var stats = server.GetCoreStats();
True(JsonHelpers.TryUInt32(stats["packetsRejected"], 0, uint.MaxValue, out var rejected) && rejected >= 1);
}
private static (AuthoritativeServer Server, string Dir) CreateServerFixture()
{
var dir = TempDir();
var path = Path.Combine(dir, "commonwealth-server.json");
var options = new ServerOptions { ConfigPath = path, Host = "127.0.0.1", Port = 7777, AdminPort = 7779, MaxPlayers = 32 };
return (new AuthoritativeServer(options), dir);
True(Stat(fixture.Server, "packetsRejected") >= 1);
}
private static async Task Activate(AuthoritativeServer server, params MemoryConnection[] connections)
@@ -287,15 +251,7 @@ internal static class Program
}
}
private static uint ReadyId(MemoryConnection connection)
{
var ready = connection.SentPackets.Last(p => JsonHelpers.String(p["type"]) == "sessionReady");
True(JsonHelpers.TryUInt32(ready["playerId"], 1, uint.MaxValue, out var id));
return id;
}
private static Task Send(AuthoritativeServer server, MemoryConnection connection, JsonObject packet) =>
server.HandleMessageAsync(connection, PacketCodec.Encode(packet).Payload, CancellationToken.None);
private static Task Send(AuthoritativeServer server, MemoryConnection connection, JsonObject packet) => server.HandleMessageAsync(connection, PacketCodec.Encode(packet).Payload, CancellationToken.None);
private static JsonObject Transform(string cell, string world, double x, double y, string movement = "normal") => new()
{
@@ -303,6 +259,16 @@ internal static class Program
["cellId"] = cell, ["worldspaceId"] = world, ["movementType"] = movement
};
private static uint ReadyId(MemoryConnection connection)
{
var packet = connection.SentPackets.Last(p => JsonHelpers.String(p["type"]) == "sessionReady");
True(JsonHelpers.TryUInt32(packet["playerId"], 1, uint.MaxValue, out var id));
return id;
}
private static uint PlayerId(JsonObject packet) => JsonHelpers.TryUInt32(packet["playerId"], 0, uint.MaxValue, out var id) ? id : 0;
private static uint Stat(AuthoritativeServer server, string name) => JsonHelpers.TryUInt32(server.GetCoreStats()[name], 0, uint.MaxValue, out var value) ? value : 0;
private static string TempDir()
{
var path = Path.Combine(Path.GetTempPath(), "co-csharp-tests-" + Guid.NewGuid().ToString("N"));
@@ -313,35 +279,43 @@ internal static class Program
private static void True(bool value) { if (!value) throw new Exception("expected true"); }
private static void False(bool value) { if (value) throw new Exception("expected false"); }
private static void NotNull(object? value) { if (value is null) throw new Exception("expected non-null"); }
private static void Equal<T>(T expected, T actual) where T : notnull { if (!EqualityComparer<T>.Default.Equals(expected, actual)) throw new Exception($"expected {expected}, got {actual}"); }
private static void SequenceEqual(byte[] expected, byte[] actual) { if (!expected.AsSpan().SequenceEqual(actual)) throw new Exception("byte sequences differ"); }
private static void Equal<T>(T expected, T actual) { if (!EqualityComparer<T>.Default.Equals(expected, actual)) throw new Exception($"expected {expected}, got {actual}"); }
private static void Throws<T>(Action action) where T : Exception { try { action(); } catch (T) { return; } throw new Exception($"expected {typeof(T).Name}"); }
private sealed class ServerFixture : IAsyncDisposable
{
public ServerFixture()
{
Directory = TempDir();
Server = new AuthoritativeServer(new ServerOptions { ConfigPath = Path.Combine(Directory, "commonwealth-server.json"), Host = "127.0.0.1", Port = 7777, AdminPort = 7779, MaxPlayers = 32 });
}
public string Directory { get; }
public AuthoritativeServer Server { get; }
public async ValueTask DisposeAsync()
{
await Server.DisposeAsync();
try { System.IO.Directory.Delete(Directory, true); } catch { }
}
}
private sealed class MemoryConnection : IGameConnection
{
private readonly object _gate = new();
private readonly List<JsonObject> _sent = new();
private int _closed;
public MemoryConnection(string name, int port)
{
ConnectionKey = "memory:" + name;
RemoteEndpoint = new IPEndPoint(IPAddress.Parse("127.0.0.1"), port);
}
public MemoryConnection(string name, int port) { ConnectionKey = "memory:" + name; RemoteEndpoint = new IPEndPoint(IPAddress.Parse("127.0.0.1"), port); }
public string ConnectionKey { get; }
public string TransportName => "memory";
public IPEndPoint RemoteEndpoint { get; }
public bool IsClosed => Volatile.Read(ref _closed) != 0;
public IReadOnlyList<JsonObject> SentPackets { get { lock (_gate) return _sent.Select(JsonHelpers.CloneObject).ToArray(); } }
public ValueTask<SendOutcome> SendAsync(EncodedPacket packet, CancellationToken cancellationToken = default)
{
if (IsClosed) return ValueTask.FromResult(SendOutcome.NotConnected);
lock (_gate) _sent.Add(PacketCodec.Decode(packet.Payload));
return ValueTask.FromResult(SendOutcome.Sent);
}
public ValueTask DisconnectAsync(int reason, string debug) { Interlocked.Exchange(ref _closed, 1); return ValueTask.CompletedTask; }
public ValueTask DisposeAsync() => DisconnectAsync(0, "dispose");
public void Clear() { lock (_gate) _sent.Clear(); }
View File
-8
View File
@@ -1,8 +0,0 @@
from __future__ import annotations
import sys
from pathlib import Path
SERVER_DIR = Path(__file__).resolve().parents[1]
if str(SERVER_DIR) not in sys.path:
sys.path.insert(0, str(SERVER_DIR))
-109
View File
@@ -1,109 +0,0 @@
from __future__ import annotations
from pathlib import Path
import pytest
from config import (
Config,
ensure_writable_directory,
load_config,
save_config,
validate_config,
)
def test_save_and_load_utf8_lf(tmp_path: Path) -> None:
path = tmp_path / "commonwealth-server.json"
cfg = Config(
host="0.0.0.0",
port=7777,
server_name="Café Server",
server_description="résumé",
max_players=8,
log_verbosity="info",
admin_port=7779,
enable_gns_transport=True,
gns_bridge_path="/opt/commonwealth/libcommonwealth_online_gns_bridge.so",
)
save_config(cfg, str(path))
raw = path.read_bytes()
assert b"\r\n" not in raw
assert "Café Server".encode("utf-8") in raw
loaded = load_config(str(path))
assert loaded.server_name == "Café Server"
assert loaded.server_description == "résumé"
assert loaded.admin_port == 7779
assert loaded.enable_gns_transport is True
assert loaded.gns_bridge_path == "/opt/commonwealth/libcommonwealth_online_gns_bridge.so"
def test_gns_boolean_config_parsing_does_not_treat_false_string_as_true() -> None:
assert Config.from_dict({"enable_gns_transport": "false"}).enable_gns_transport is False
assert Config.from_dict({"enable_gns_transport": "off"}).enable_gns_transport is False
assert Config.from_dict({"enable_gns_transport": "0"}).enable_gns_transport is False
assert Config.from_dict({"enable_gns_transport": "true"}).enable_gns_transport is True
assert Config.from_dict({"enable_gns_transport": "yes"}).enable_gns_transport is True
assert Config.from_dict({"enable_gns_transport": 1}).enable_gns_transport is True
assert Config.from_dict({"enable_gns_transport": 0}).enable_gns_transport is False
with pytest.raises(ValueError, match="enable_gns_transport"):
Config.from_dict({"enable_gns_transport": "sometimes"})
def test_gns_bridge_path_normalizes_blank_to_none_and_rejects_nonstring() -> None:
assert Config.from_dict({"gns_bridge_path": ""}).gns_bridge_path is None
assert Config.from_dict({"gns_bridge_path": " "}).gns_bridge_path is None
assert Config.from_dict({"gns_bridge_path": None}).gns_bridge_path is None
with pytest.raises(ValueError, match="gns_bridge_path"):
Config.from_dict({"gns_bridge_path": 123})
def test_validate_rejects_bad_ports_and_max_players() -> None:
cfg = Config(port=70000, admin_port=7777, max_players=0)
ok, errors = validate_config(cfg)
assert not ok
assert any("port must be 1-65535" in error for error in errors)
assert any("max_players must be >= 1" in error for error in errors)
collide = Config(port=7777, admin_port=7777)
ok, errors = validate_config(collide)
assert not ok
assert any("admin_port must differ" in error for error in errors)
def test_validate_rejects_unresolvable_host() -> None:
cfg = Config(host="this-host-should-not-resolve.invalid")
ok, errors = validate_config(cfg)
assert not ok
assert any("not a valid IPv4" in error for error in errors)
def test_gns_requires_explicit_ipv4_bind_host() -> None:
explicit = Config(host="127.0.0.1", enable_gns_transport=True)
ok, errors = validate_config(explicit)
assert ok, errors
wildcard = Config(host="0.0.0.0", enable_gns_transport=True)
ok, errors = validate_config(wildcard)
assert ok, errors
hostname = Config(host="localhost", enable_gns_transport=True)
ok, errors = validate_config(hostname)
assert not ok
assert any("explicit IPv4 bind address" in error for error in errors)
def test_ensure_writable_directory(tmp_path: Path) -> None:
target = tmp_path / "state"
ensure_writable_directory(target)
assert target.is_dir()
assert not (target / ".commonwealth-write-probe").exists()
def test_load_invalid_json(tmp_path: Path) -> None:
path = tmp_path / "bad.json"
path.write_text("{not json", encoding="utf-8")
with pytest.raises(Exception):
load_config(str(path))
@@ -1,27 +0,0 @@
from pathlib import Path
from config import Config
from dedicated_server import build_service_config
def test_dedicated_server_preserves_gns_rollout_fields(tmp_path: Path):
config_path = tmp_path / "commonwealth-server.json"
config = Config(
host="127.0.0.1",
port=7777,
server_name="Test",
server_description="GNS",
max_players=12,
log_verbosity="debug",
admin_port=7779,
enable_gns_transport=True,
gns_bridge_path="/opt/commonwealth/libcommonwealth_online_gns_bridge.so",
)
service_config = build_service_config(config, config_path)
assert service_config.host == "127.0.0.1"
assert service_config.port == 7777
assert service_config.max_players == 12
assert service_config.enable_gns_transport is True
assert service_config.gns_bridge_path == "/opt/commonwealth/libcommonwealth_online_gns_bridge.so"
assert service_config.bans_path == str(tmp_path / "bans.json")
-281
View File
@@ -1,281 +0,0 @@
from __future__ import annotations
import json
from collections import deque
import pytest
from gns_gameplay_adapter import GnsConnectionAdapter, GnsGameplayAdapter
from gns_snapshot_envelope import decode_snapshot, encode_snapshot
from gns_transport import EventType, GnsEvent, RemoteEndpoint, SendResult
from packet_codec import encode_packet
from server_core import PROTOCOL_VERSION
from transport_policy import Delivery
from transport_server import TransportAwareFalloutTogetherServer
class FakeGnsTransport:
def __init__(self):
self.events = deque()
self.endpoints = {}
self.sent = []
self.send_results = deque()
self.disconnected = []
self.incoming_sequences = {}
self.closed = False
def poll(self):
return self.events.popleft() if self.events else None
def remote_endpoint(self, connection_id):
return self.endpoints.get(connection_id)
def send_encoded(self, connection_id, encoded):
self.sent.append((connection_id, encoded))
return self.send_results.popleft() if self.send_results else SendResult.SENT
def send_packet(self, connection_id, packet):
return self.send_encoded(connection_id, encode_packet(packet))
def disconnect(self, connection_id, *, reason=0, debug=""):
self.disconnected.append((connection_id, reason, debug))
return True
def close(self):
self.closed = True
def decoded_sent_packet(encoded):
envelope = decode_snapshot(encoded.payload)
payload = envelope.payload if envelope is not None else encoded.payload
return json.loads(payload)
def packet_types_sent(transport: FakeGnsTransport, connection_id: int) -> list[str]:
return [
decoded_sent_packet(encoded)["type"]
for target, encoded in transport.sent
if target == connection_id
]
def queue_connect(transport: FakeGnsTransport, connection_id: int, port: int = 50000):
transport.endpoints[connection_id] = RemoteEndpoint("127.0.0.1", port)
transport.events.append(GnsEvent(EventType.CONNECTED, connection_id))
def queue_packet(
transport: FakeGnsTransport,
connection_id: int,
packet: dict,
*,
sequence: int | None = None,
raw_snapshot: bool = False,
):
payload = json.dumps(packet, separators=(",", ":")).encode()
packet_type = packet.get("type")
if packet_type in {"transform", "npcState"} and not raw_snapshot:
key = (connection_id, packet_type)
if sequence is None:
sequence = transport.incoming_sequences.get(key, 0) + 1
transport.incoming_sequences[key] = sequence
payload = encode_snapshot(packet_type, payload, sequence)
transport.events.append(GnsEvent(EventType.MESSAGE, connection_id, payload))
def complete_handshake(adapter: GnsGameplayAdapter, transport: FakeGnsTransport, connection_id: int):
queue_connect(transport, connection_id, 50000 + connection_id)
assert adapter.pump_once() == 1
assert "welcome" in packet_types_sent(transport, connection_id)
queue_packet(transport, connection_id, {"type": "hello", "protocolVersion": PROTOCOL_VERSION})
assert adapter.pump_once() == 1
assert "sessionReady" in packet_types_sent(transport, connection_id)
def transform_packet(x: float) -> dict:
return {
"type": "transform",
"x": x,
"y": 0.0,
"z": 0.0,
"angleZ": 0.0,
"cellId": "00000010",
"worldspaceId": "0000003C",
"movementType": "normal",
}
def test_gns_adapter_reuses_v2_admission_and_server_owned_identity():
server = TransportAwareFalloutTogetherServer(host="127.0.0.1", port=0)
transport = FakeGnsTransport()
adapter = GnsGameplayAdapter(server, transport)
complete_handshake(adapter, transport, 101)
client = adapter._client_for_connection(101)
assert client is not None
assert client.gameplay_active
assert client.protocol_version == PROTOCOL_VERSION
packet = transform_packet(1.0)
packet["playerId"] = 999999
queue_packet(transport, 101, packet)
adapter.pump_once()
assert client.last_transform is not None
assert client.last_transform["playerId"] == client.player_id
assert client.last_transform["playerId"] != 999999
def test_gns_adapter_preserves_reliable_and_snapshot_delivery_policy():
server = TransportAwareFalloutTogetherServer(host="127.0.0.1", port=0)
transport = FakeGnsTransport()
adapter = GnsGameplayAdapter(server, transport)
complete_handshake(adapter, transport, 201)
complete_handshake(adapter, transport, 202)
transport.sent.clear()
queue_packet(transport, 202, transform_packet(10.0))
adapter.pump_once()
queue_packet(transport, 201, transform_packet(0.0))
adapter.pump_once()
transform_relays = [
encoded
for target, encoded in transport.sent
if target == 202 and decoded_sent_packet(encoded).get("type") == "transform"
]
assert transform_relays
transform_relay = transform_relays[-1]
assert transform_relay.delivery is Delivery.UNRELIABLE_SEQUENCED
transform_envelope = decode_snapshot(transform_relay.payload)
assert transform_envelope is not None
assert transform_envelope.packet_type == "transform"
assert transform_envelope.sequence > 0
queue_packet(
transport,
201,
{"type": "playerState", "characterName": "Nomad", "actionEvents": []},
)
adapter.pump_once()
player_state_relays = [
encoded
for target, encoded in transport.sent
if target == 202 and decoded_sent_packet(encoded).get("type") == "playerState"
]
assert player_state_relays
player_state_relay = player_state_relays[-1]
assert player_state_relay.delivery is Delivery.RELIABLE_ORDERED
assert decode_snapshot(player_state_relay.payload) is None
def test_gns_snapshot_reordering_and_duplicates_never_roll_state_backward():
server = TransportAwareFalloutTogetherServer(host="127.0.0.1", port=0)
transport = FakeGnsTransport()
adapter = GnsGameplayAdapter(server, transport)
complete_handshake(adapter, transport, 501)
client = adapter._client_for_connection(501)
assert client is not None
queue_packet(transport, 501, transform_packet(1.0), sequence=1)
adapter.pump_once()
assert client.last_transform is not None
assert client.last_transform["x"] == 1.0
queue_packet(transport, 501, transform_packet(3.0), sequence=3)
adapter.pump_once()
assert client.last_transform["x"] == 3.0
rejected_before = server.get_stats()["packetsRejected"]
queue_packet(transport, 501, transform_packet(2.0), sequence=2)
adapter.pump_once()
assert client.last_transform["x"] == 3.0
assert server.get_stats()["packetsRejected"] == rejected_before + 1
queue_packet(transport, 501, transform_packet(9.0), sequence=3)
adapter.pump_once()
assert client.last_transform["x"] == 3.0
assert server.get_stats()["packetsRejected"] == rejected_before + 2
def test_gns_raw_unsequenced_snapshot_is_rejected_before_state_mutation():
server = TransportAwareFalloutTogetherServer(host="127.0.0.1", port=0)
transport = FakeGnsTransport()
adapter = GnsGameplayAdapter(server, transport)
complete_handshake(adapter, transport, 601)
client = adapter._client_for_connection(601)
assert client is not None
assert client.last_transform is None
rejected_before = server.get_stats()["packetsRejected"]
queue_packet(transport, 601, transform_packet(1.0), raw_snapshot=True)
adapter.pump_once()
assert client.last_transform is None
assert server.get_stats()["packetsRejected"] == rejected_before + 1
def test_connection_adapter_may_drop_snapshots_but_never_reliable_messages_silently():
transport = FakeGnsTransport()
connection = GnsConnectionAdapter(transport, 250)
transport.send_results.append(SendResult.DROPPED)
connection.send_encoded(encode_packet({"type": "transform", "x": 1}))
transform_wire = transport.sent[-1][1]
assert decode_snapshot(transform_wire.payload) is not None
transport.send_results.append(SendResult.BACKPRESSURE)
connection.send_encoded(encode_packet({"type": "npcState", "npcs": []}))
npc_wire = transport.sent[-1][1]
assert decode_snapshot(npc_wire.payload) is not None
transport.send_results.append(SendResult.DROPPED)
with pytest.raises(OSError):
connection.send_encoded(encode_packet({"type": "playerState", "characterName": "Nomad"}))
transport.send_results.append(SendResult.BACKPRESSURE)
with pytest.raises(OSError):
connection.send_encoded(encode_packet({"type": "combatHit", "sequence": 1}))
def test_gns_adapter_enforces_existing_ban_policy(tmp_path):
server = TransportAwareFalloutTogetherServer(
host="127.0.0.1",
port=0,
bans_path=str(tmp_path / "bans.json"),
)
server.ban_ip("127.0.0.1", "test ban")
transport = FakeGnsTransport()
adapter = GnsGameplayAdapter(server, transport)
queue_connect(transport, 301)
adapter.pump_once()
assert adapter._client_for_connection(301) is None
assert transport.disconnected
session_ended = [
decoded_sent_packet(encoded)
for target, encoded in transport.sent
if target == 301 and decoded_sent_packet(encoded).get("type") == "sessionEnded"
]
assert session_ended
assert session_ended[-1]["code"] == "banned"
assert server.get_stats()["bannedConnectionsRejected"] == 1
def test_gns_adapter_ends_session_on_native_oversize_event():
server = TransportAwareFalloutTogetherServer(host="127.0.0.1", port=0)
transport = FakeGnsTransport()
adapter = GnsGameplayAdapter(server, transport)
complete_handshake(adapter, transport, 401)
transport.events.append(
GnsEvent(EventType.OVERSIZE_MESSAGE, 401, reason=0, debug="too large")
)
adapter.pump_once()
assert adapter._client_for_connection(401) is None
assert transport.disconnected
ended_packets = [
decoded_sent_packet(encoded)
for target, encoded in transport.sent
if target == 401 and decoded_sent_packet(encoded).get("type") == "sessionEnded"
]
assert ended_packets
assert ended_packets[-1]["code"] == "packet_too_large"
@@ -1,64 +0,0 @@
import pytest
from gns_snapshot_envelope import (
HEADER_SIZE,
MAGIC,
SnapshotEnvelopeError,
decode_snapshot,
encode_snapshot,
)
from packet_codec import MAX_MESSAGE_BYTES
def test_transform_and_npc_snapshots_round_trip_with_sequence():
transform = encode_snapshot("transform", b'{"type":"transform","x":1}', 7)
decoded_transform = decode_snapshot(transform)
assert decoded_transform is not None
assert decoded_transform.packet_type == "transform"
assert decoded_transform.sequence == 7
assert decoded_transform.payload == b'{"type":"transform","x":1}'
npc = encode_snapshot("npcState", b'{"type":"npcState","npcs":[]}', 9)
decoded_npc = decode_snapshot(npc)
assert decoded_npc is not None
assert decoded_npc.packet_type == "npcState"
assert decoded_npc.sequence == 9
def test_reliable_json_is_not_misidentified_as_snapshot_envelope():
assert decode_snapshot(b'{"type":"playerState"}') is None
def test_snapshot_envelope_rejects_invalid_family_sequence_and_size():
with pytest.raises(SnapshotEnvelopeError):
encode_snapshot("playerState", b"{}", 1)
with pytest.raises(SnapshotEnvelopeError):
encode_snapshot("transform", b"{}", 0)
with pytest.raises(SnapshotEnvelopeError):
encode_snapshot("transform", b"x" * MAX_MESSAGE_BYTES, 1)
def test_snapshot_decoder_rejects_truncated_or_corrupt_headers():
with pytest.raises(SnapshotEnvelopeError):
decode_snapshot(MAGIC)
valid = bytearray(encode_snapshot("transform", b"{}", 1))
valid[4] = 99
with pytest.raises(SnapshotEnvelopeError):
decode_snapshot(valid)
valid = bytearray(encode_snapshot("transform", b"{}", 1))
valid[5] = 99
with pytest.raises(SnapshotEnvelopeError):
decode_snapshot(valid)
valid = bytearray(encode_snapshot("transform", b"{}", 1))
valid[6:8] = b"\x00\x01"
with pytest.raises(SnapshotEnvelopeError):
decode_snapshot(valid)
def test_snapshot_header_leaves_payload_under_native_64k_cap():
payload = b"x" * (MAX_MESSAGE_BYTES - HEADER_SIZE)
message = encode_snapshot("transform", payload, 0xFFFFFFFF)
assert len(message) == MAX_MESSAGE_BYTES
-150
View File
@@ -1,150 +0,0 @@
from __future__ import annotations
import ctypes
from collections import deque
import pytest
from gns_transport import EventType, GnsServerTransport, GnsTransportError, RemoteEndpoint, SendResult
class FakeFunction:
def __init__(self, callback):
self.callback = callback
self.argtypes = None
self.restype = None
def __call__(self, *args):
return self.callback(*args)
class FakeNativeBridge:
def __init__(self):
self.events = deque()
self.sent = []
self.disconnects = []
self.endpoints = {42: (0x7F000001, 54321)}
self.destroyed = False
self.co_gns_server_create = FakeFunction(self._create)
self.co_gns_server_destroy = FakeFunction(self._destroy)
self.co_gns_server_local_port = FakeFunction(lambda handle: 7777)
self.co_gns_server_connection_count = FakeFunction(lambda handle: 2)
self.co_gns_server_poll = FakeFunction(self._poll)
self.co_gns_server_send = FakeFunction(self._send)
self.co_gns_server_disconnect = FakeFunction(self._disconnect)
self.co_gns_server_remote_ipv4 = FakeFunction(self._remote_ipv4)
def _create(self, host, port, out_handle, error_buffer, error_size):
assert host == b"127.0.0.1"
assert port == 0
ctypes.cast(out_handle, ctypes.POINTER(ctypes.c_void_p))[0] = ctypes.c_void_p(0x1234)
error_buffer.value = b""
return 1
def _destroy(self, handle):
self.destroyed = True
def _poll(self, handle, event_pointer, payload_buffer, payload_capacity):
if not self.events:
return 0
event_type, connection_id, payload, reason, debug = self.events.popleft()
event = event_pointer._obj
event.type = event_type
event.connection_id = connection_id
event.reason = reason
event.payload_size = len(payload)
event.debug = debug.encode()
if event_type == EventType.MESSAGE and payload:
assert len(payload) <= payload_capacity
ctypes.memmove(payload_buffer, payload, len(payload))
return 1
def _send(self, handle, connection_id, payload_buffer, payload_size, delivery):
payload = ctypes.string_at(payload_buffer, payload_size)
self.sent.append((connection_id, payload, delivery))
return SendResult.SENT
def _disconnect(self, handle, connection_id, reason, debug):
self.disconnects.append((connection_id, reason, debug))
return 1
def _remote_ipv4(self, handle, connection_id, out_ipv4, out_port):
endpoint = self.endpoints.get(connection_id)
if endpoint is None:
return 0
ipv4, port = endpoint
ctypes.cast(out_ipv4, ctypes.POINTER(ctypes.c_uint32))[0] = ipv4
ctypes.cast(out_port, ctypes.POINTER(ctypes.c_uint16))[0] = port
return 1
def test_wrapper_polls_message_events_without_line_framing():
native = FakeNativeBridge()
native.events.append((EventType.CONNECTED, 42, b"", 0, ""))
native.events.append((EventType.MESSAGE, 42, b'{"type":"transform","x":1}', 0, ""))
transport = GnsServerTransport("127.0.0.1", 0, native_library=native)
assert transport.local_port == 7777
assert transport.connection_count == 2
assert transport.remote_endpoint(42) == RemoteEndpoint("127.0.0.1", 54321)
assert transport.remote_endpoint(99) is None
connected = transport.poll()
assert connected is not None
assert connected.type is EventType.CONNECTED
assert connected.connection_id == 42
message = transport.poll()
assert message is not None
assert message.type is EventType.MESSAGE
assert message.payload == b'{"type":"transform","x":1}'
assert not message.payload.endswith(b"\n")
assert transport.poll() is None
def test_wrapper_maps_protocol_delivery_to_native_send_modes():
native = FakeNativeBridge()
transport = GnsServerTransport("127.0.0.1", 0, native_library=native)
assert transport.send_packet(42, {"type": "transform", "x": 1}) is SendResult.SENT
assert transport.send_packet(42, {"type": "playerState", "characterName": "Nomad"}) is SendResult.SENT
assert native.sent[0][0] == 42
assert native.sent[0][2] == 0
assert native.sent[0][1] == b'{"type":"transform","x":1}'
assert native.sent[1][2] == 1
assert native.sent[1][1] == b'{"type":"playerState","characterName":"Nomad"}'
def test_wrapper_preserves_disconnect_and_oversize_events():
native = FakeNativeBridge()
native.events.append((EventType.DISCONNECTED, 9, b"", 5003, "peer timeout"))
native.events.append((EventType.OVERSIZE_MESSAGE, 10, b"x" * (64 * 1024 + 1), 0, "too large"))
transport = GnsServerTransport("127.0.0.1", 0, native_library=native)
disconnected = transport.poll()
assert disconnected is not None
assert disconnected.type is EventType.DISCONNECTED
assert disconnected.connection_id == 9
assert disconnected.reason == 5003
assert disconnected.debug == "peer timeout"
oversize = transport.poll()
assert oversize is not None
assert oversize.type is EventType.OVERSIZE_MESSAGE
assert oversize.connection_id == 10
assert oversize.payload == b""
def test_wrapper_disconnect_and_close_are_idempotent():
native = FakeNativeBridge()
transport = GnsServerTransport("127.0.0.1", 0, native_library=native)
assert transport.disconnect(77, reason=1000, debug="test")
assert native.disconnects == [(77, 1000, b"test")]
transport.close()
transport.close()
assert native.destroyed
assert transport.is_closed
with pytest.raises(GnsTransportError):
transport.poll()
-404
View File
@@ -1,404 +0,0 @@
from __future__ import annotations
import json
import socket
import time
import pytest
from server_core import PROTOCOL_VERSION, FalloutTogetherServer, states_share_interest
_RECV_BUFFERS: dict[socket.socket, bytes] = {}
def recv_packet(sock: socket.socket, timeout: float = 2.0) -> dict:
sock.settimeout(timeout)
data = _RECV_BUFFERS.get(sock, b"")
while b"\n" not in data:
chunk = sock.recv(4096)
if not chunk:
raise ConnectionError("socket closed before a complete packet was received")
data += chunk
line, remainder = data.split(b"\n", 1)
_RECV_BUFFERS[sock] = remainder
return json.loads(line.decode("utf-8"))
def recv_until(sock: socket.socket, packet_type: str, timeout: float = 2.0) -> dict:
deadline = time.monotonic() + timeout
while time.monotonic() < deadline:
remaining = max(0.01, deadline - time.monotonic())
try:
packet = recv_packet(sock, timeout=remaining)
except socket.timeout:
break
if packet.get("type") == packet_type:
return packet
raise AssertionError(f"did not receive packet type {packet_type}")
def recv_authority(
sock: socket.socket,
cell: str,
world: str = "0000003C",
timeout: float = 2.0,
) -> dict:
deadline = time.monotonic() + timeout
while time.monotonic() < deadline:
remaining = max(0.01, deadline - time.monotonic())
packet = recv_until(sock, "npcAuthority", timeout=remaining)
if packet.get("authorityCellId") == cell and packet.get("authorityWorldspaceId", "") == world:
return packet
raise AssertionError(f"did not receive npcAuthority for {cell}/{world}")
def assert_no_packet_type(sock: socket.socket, packet_type: str, timeout: float = 0.2) -> None:
with pytest.raises(AssertionError):
recv_until(sock, packet_type, timeout=timeout)
def send_packet(sock: socket.socket, packet: dict) -> None:
sock.sendall(json.dumps(packet, separators=(",", ":")).encode() + b"\n")
def connect(server: FalloutTogetherServer) -> tuple[socket.socket, dict]:
sock = socket.create_connection(("127.0.0.1", server.port), timeout=2.0)
return sock, recv_packet(sock)
def hello(sock: socket.socket) -> dict:
send_packet(sock, {"type": "hello", "protocolVersion": PROTOCOL_VERSION})
return recv_until(sock, "sessionReady")
def transform(
cell: str,
x: float,
y: float,
world: str = "0000003C",
movement_type: str = "normal",
) -> dict:
return {
"type": "transform",
"x": x,
"y": y,
"z": 0.0,
"angleZ": 0.0,
"cellId": cell,
"worldspaceId": world,
"movementType": movement_type,
}
def npc_state(
cell: str,
epoch: int,
x: float,
world: str = "0000003C",
source_form_id: str = "000000AA",
) -> dict:
return {
"type": "npcState",
"authorityEpoch": epoch,
"authorityCellId": cell,
"authorityWorldspaceId": world,
"npcs": [
{
"npcId": 1,
"baseFormId": "0000000F",
"sourceFormId": source_form_id,
"x": x,
"y": 0.0,
"z": 0.0,
"angleZ": 0.0,
"cellId": cell,
"worldspaceId": world,
"isDead": False,
}
],
}
def wait_for_stat(server: FalloutTogetherServer, key: str, minimum: int, timeout: float = 2.0) -> int:
deadline = time.monotonic() + timeout
while time.monotonic() < deadline:
value = int(server.get_stats()[key])
if value >= minimum:
return value
time.sleep(0.01)
raise AssertionError(f"stat {key} did not reach {minimum}")
def start_server(max_players: int = 16) -> FalloutTogetherServer:
server = FalloutTogetherServer(host="127.0.0.1", port=0, max_players=max_players)
server.start()
deadline = time.time() + 2
while not server.is_running() and time.time() < deadline:
time.sleep(0.01)
return server
def test_interest_exact_cell_and_exterior_radius():
a = transform("00000001", 0, 0)
same = transform("00000001", 50000, 50000)
near = transform("00000002", 1000, 1000)
far = transform("00000002", 20000, 20000)
other_world = transform("00000002", 1000, 1000, world="0000003D")
assert states_share_interest(a, same)
assert states_share_interest(a, near)
assert not states_share_interest(a, far)
assert not states_share_interest(a, other_world)
def test_recv_packet_preserves_coalesced_lines():
reader, writer = socket.socketpair()
try:
writer.sendall(b'{"type":"first"}\n{"type":"second"}\n')
assert recv_packet(reader)["type"] == "first"
assert recv_packet(reader)["type"] == "second"
finally:
_RECV_BUFFERS.pop(reader, None)
reader.close()
writer.close()
def test_idle_transport_does_not_take_world_authority():
server = start_server()
idle, _ = connect(server)
active, active_welcome = connect(server)
ready = hello(active)
assert ready["playerId"] == active_welcome["playerId"]
stats = server.get_stats()
assert stats["connectedClients"] == 1
assert stats["pendingConnections"] == 1
assert server._world_state_host_player_id == active_welcome["playerId"]
idle.close()
active.close()
server.stop()
def test_max_players_applies_to_activated_sessions_not_probes():
server = start_server(max_players=1)
first, _ = connect(server)
hello(first)
second, _ = connect(server)
send_packet(second, {"type": "hello", "protocolVersion": PROTOCOL_VERSION})
ended = recv_until(second, "sessionEnded")
assert ended["code"] == "server_full"
first.close()
second.close()
server.stop()
def test_protocol_mismatch_is_rejected():
server = start_server()
sock, _ = connect(server)
send_packet(sock, {"type": "hello", "protocolVersion": PROTOCOL_VERSION + 1})
ended = recv_until(sock, "sessionEnded")
assert ended["code"] == "protocol_mismatch"
sock.close()
server.stop()
def test_nan_transform_is_rejected_and_not_cached():
server = start_server()
sock, welcome = connect(server)
hello(sock)
sock.sendall(
b'{"type":"transform","x":NaN,"y":0,"z":0,"angleZ":0,"cellId":"00000001"}\n'
)
time.sleep(0.05)
client = server._find_client_by_player_id(welcome["playerId"])
assert client is not None
assert client.last_transform is None
assert server.get_stats()["packetsRejected"] >= 1
sock.close()
server.stop()
def test_transform_interest_filters_distant_peer():
server = start_server()
a, _ = connect(server)
hello(a)
b, _ = connect(server)
hello(b)
send_packet(b, transform("00000020", 25000, 25000))
time.sleep(0.05)
send_packet(a, transform("00000010", 0, 0))
assert_no_packet_type(b, "transform")
assert server.get_stats()["transformPacketsInterestFiltered"] >= 1
a.close()
b.close()
server.stop()
def test_scoped_npc_authorities_publish_independently():
server = start_server()
a, a_welcome = connect(server)
hello(a)
b, b_welcome = connect(server)
hello(b)
send_packet(a, transform("00000010", 0.0, 0.0))
authority_a = recv_authority(a, "00000010")
assert authority_a["authorityPlayerId"] == a_welcome["playerId"]
send_packet(b, transform("00000020", 25000.0, 0.0))
authority_b = recv_authority(b, "00000020")
assert authority_b["authorityPlayerId"] == b_welcome["playerId"]
send_packet(a, npc_state("00000010", authority_a["authorityEpoch"], 0.0))
send_packet(b, npc_state("00000020", authority_b["authorityEpoch"], 25000.0, source_form_id="000000AB"))
assert wait_for_stat(server, "npcStatePacketsReceived", 2) == 2
assert server.get_stats()["npcAuthorityRejects"] == 0
a.close()
b.close()
server.stop()
def test_stale_npc_authority_epoch_cannot_publish_after_handoff():
server = start_server()
a, a_welcome = connect(server)
hello(a)
b, b_welcome = connect(server)
hello(b)
send_packet(a, transform("00000010", 0.0, 0.0))
initial_a = recv_authority(a, "00000010")
initial_b = recv_authority(b, "00000010")
assert initial_a["authorityPlayerId"] == a_welcome["playerId"]
assert initial_b["authorityEpoch"] == initial_a["authorityEpoch"]
send_packet(b, transform("00000010", 64.0, 0.0))
current_b = recv_authority(b, "00000010")
assert current_b["authorityPlayerId"] == a_welcome["playerId"]
assert current_b["authorityEpoch"] == initial_a["authorityEpoch"]
send_packet(a, npc_state("00000010", initial_a["authorityEpoch"], 0.0))
wait_for_stat(server, "npcStatePacketsReceived", 1)
send_packet(a, transform("00000020", 0.0, 0.0, movement_type="cell_change"))
handoff = recv_authority(b, "00000010")
assert handoff["authorityPlayerId"] == b_welcome["playerId"]
assert handoff["authorityEpoch"] > initial_a["authorityEpoch"]
send_packet(a, npc_state("00000010", initial_a["authorityEpoch"], 0.0))
wait_for_stat(server, "npcAuthorityRejects", 1)
assert server.get_stats()["npcStatePacketsReceived"] == 1
send_packet(b, npc_state("00000010", handoff["authorityEpoch"], 64.0, source_form_id="000000AB"))
wait_for_stat(server, "npcStatePacketsReceived", 2)
a.close()
b.close()
server.stop()
def test_scoped_npc_packet_rejects_entry_outside_authority_scope():
server = start_server()
sock, _ = connect(server)
hello(sock)
send_packet(sock, transform("00000010", 0.0, 0.0))
authority = recv_authority(sock, "00000010")
packet = npc_state("00000010", authority["authorityEpoch"], 0.0)
packet["npcs"][0]["cellId"] = "00000011"
send_packet(sock, packet)
wait_for_stat(server, "npcAuthorityRejects", 1)
assert server.get_stats()["npcStatePacketsReceived"] == 0
sock.close()
server.stop()
def test_rejected_normal_teleport_is_corrected_and_never_relayed():
server = start_server()
a, a_welcome = connect(server)
hello(a)
b, _ = connect(server)
hello(b)
send_packet(b, transform("00000010", 64.0, 0.0))
recv_until(a, "transform")
send_packet(a, transform("00000010", 0.0, 0.0))
first_relay = recv_until(b, "transform")
assert first_relay["x"] == 0.0
send_packet(a, transform("00000010", 100000.0, 0.0))
correction = recv_until(a, "positionCorrection")
assert correction["x"] == 0.0
assert correction["cellId"] == "00000010"
assert_no_packet_type(b, "transform")
client = server._find_client_by_player_id(a_welcome["playerId"])
assert client is not None
assert client.last_transform is not None
assert client.last_transform["x"] == 0.0
stats = server.get_stats()
assert stats["movementPacketsRejected"] >= 1
assert stats["movementCorrectionsSent"] >= 1
a.close()
b.close()
server.stop()
def test_scope_change_requires_explicit_transition():
server = start_server()
sock, welcome = connect(server)
hello(sock)
send_packet(sock, transform("00000010", 0.0, 0.0))
time.sleep(0.05)
send_packet(sock, transform("00000011", 10.0, 0.0))
correction = recv_until(sock, "positionCorrection")
assert correction["cellId"] == "00000010"
client = server._find_client_by_player_id(welcome["playerId"])
assert client is not None
assert client.last_transform is not None
assert client.last_transform["cellId"] == "00000010"
send_packet(sock, transform("00000011", 10.0, 0.0, movement_type="cell_change"))
deadline = time.monotonic() + 1.0
while time.monotonic() < deadline:
if client.last_transform is not None and client.last_transform.get("cellId") == "00000011":
break
time.sleep(0.01)
assert client.last_transform is not None
assert client.last_transform["cellId"] == "00000011"
sock.close()
server.stop()
def test_zero_elapsed_burst_cannot_bypass_movement_envelope():
server = start_server()
sock, welcome = connect(server)
hello(sock)
send_packet(sock, transform("00000010", 0.0, 0.0))
send_packet(sock, transform("00000010", 50000.0, 0.0))
correction = recv_until(sock, "positionCorrection")
assert correction["x"] == 0.0
client = server._find_client_by_player_id(welcome["playerId"])
assert client is not None
assert client.last_transform is not None
assert client.last_transform["x"] == 0.0
sock.close()
server.stop()
def test_oversized_unterminated_packet_closes_session():
server = start_server()
sock, _ = connect(server)
sock.sendall(b"x" * (64 * 1024 + 1))
ended = recv_until(sock, "sessionEnded")
assert ended["code"] == "packet_too_large"
sock.close()
server.stop()
-141
View File
@@ -1,141 +0,0 @@
from __future__ import annotations
import json
import socket
import time
import server_core
from server_core import PROTOCOL_VERSION, FalloutTogetherServer
_RECV_BUFFERS: dict[socket.socket, bytes] = {}
def recv_packet(sock: socket.socket, timeout: float = 2.0) -> dict:
sock.settimeout(timeout)
data = _RECV_BUFFERS.get(sock, b"")
while b"\n" not in data:
chunk = sock.recv(4096)
if not chunk:
raise ConnectionError("socket closed before a complete packet was received")
data += chunk
line, remainder = data.split(b"\n", 1)
_RECV_BUFFERS[sock] = remainder
return json.loads(line.decode("utf-8"))
def recv_until(sock: socket.socket, predicate, timeout: float = 2.0) -> dict:
deadline = time.monotonic() + timeout
while time.monotonic() < deadline:
try:
packet = recv_packet(sock, timeout=max(0.01, deadline - time.monotonic()))
except socket.timeout:
break
if predicate(packet):
return packet
raise AssertionError("did not receive matching packet")
def send_packet(sock: socket.socket, packet: dict) -> None:
sock.sendall(json.dumps(packet, separators=(",", ":")).encode() + b"\n")
def connect_v2(server: FalloutTogetherServer) -> tuple[socket.socket, int]:
sock = socket.create_connection(("127.0.0.1", server.port), timeout=2.0)
welcome = recv_until(sock, lambda packet: packet.get("type") == "welcome")
send_packet(sock, {"type": "hello", "protocolVersion": PROTOCOL_VERSION})
ready = recv_until(sock, lambda packet: packet.get("type") == "sessionReady")
assert ready["playerId"] == welcome["playerId"]
return sock, int(welcome["playerId"])
def transform(cell: str, x: float, y: float = 0.0) -> dict:
return {
"type": "transform",
"x": x,
"y": y,
"z": 0.0,
"angleZ": 0.0,
"cellId": cell,
"worldspaceId": "0000003C",
"movementType": "normal",
}
def start_server() -> FalloutTogetherServer:
server = FalloutTogetherServer(host="127.0.0.1", port=0, max_players=24)
server.start()
deadline = time.monotonic() + 2.0
while not server.is_running() and time.monotonic() < deadline:
time.sleep(0.01)
return server
def test_sixteen_clients_do_not_cross_broadcast_distant_cell_transforms(monkeypatch):
# All synthetic clients originate from localhost. Bypass only the per-IP
# connection-attempt throttle so this test measures gameplay capacity and
# interest filtering instead of the independent abuse-control policy.
monkeypatch.setattr(server_core, "MAX_CONNECT_ATTEMPTS", 64)
server = start_server()
clients: list[tuple[socket.socket, int]] = []
try:
clients = [connect_v2(server) for _ in range(16)]
assert server.get_stats()["connectedClients"] == 16
near = clients[:8]
far = clients[8:]
for index, (sock, _player_id) in enumerate(near):
send_packet(sock, transform("00000010", float(index * 32)))
for index, (sock, _player_id) in enumerate(far):
send_packet(sock, transform("00000020", 30000.0 + float(index * 32)))
deadline = time.monotonic() + 2.0
while time.monotonic() < deadline:
sessions = [server._find_client_by_player_id(player_id) for _sock, player_id in clients]
if all(session is not None and session.last_transform is not None for session in sessions):
break
time.sleep(0.01)
assert all(
session is not None and session.last_transform is not None
for session in (server._find_client_by_player_id(player_id) for _sock, player_id in clients)
)
sender_sock, sender_id = near[0]
near_peer_sock, _near_peer_id = near[1]
far_peer_sock, _far_peer_id = far[0]
marker_x = 128.0
send_packet(sender_sock, transform("00000010", marker_x))
relayed = recv_until(
near_peer_sock,
lambda packet: packet.get("type") == "transform"
and packet.get("playerId") == sender_id
and packet.get("x") == marker_x,
)
assert relayed["cellId"] == "00000010"
try:
recv_until(
far_peer_sock,
lambda packet: packet.get("type") == "transform"
and packet.get("playerId") == sender_id
and packet.get("x") == marker_x,
timeout=0.4,
)
except AssertionError:
pass
else:
raise AssertionError("distant cell received a transform that should have been interest-filtered")
stats = server.get_stats()
assert stats["connectedClients"] == 16
assert stats["transformPacketsInterestFiltered"] >= len(far)
finally:
for sock, _player_id in clients:
_RECV_BUFFERS.pop(sock, None)
try:
sock.close()
except OSError:
pass
server.stop()
-60
View File
@@ -1,60 +0,0 @@
from __future__ import annotations
from npc_authority import NpcAuthorityManager, ScopeKey, scope_from_transform
def test_first_player_in_scope_becomes_authority() -> None:
manager = NpcAuthorityManager()
scope = ScopeKey("00000010", "0000003C")
changes = manager.reconcile([(7, scope), (4, scope)])
assert len(changes) == 1
assignment = manager.get(scope)
assert assignment is not None
assert assignment.player_id == 4
assert assignment.epoch == 1
assert manager.authorize(4, scope, 1)
assert not manager.authorize(7, scope, 1)
def test_disconnect_reassigns_and_increments_epoch() -> None:
manager = NpcAuthorityManager()
scope = ScopeKey("00000010", "0000003C")
manager.reconcile([(1, scope), (2, scope)])
changes = manager.reconcile([(2, scope)])
assert len(changes) == 1
assert changes[0].previous_player_id == 1
assert changes[0].player_id == 2
assert changes[0].epoch == 2
assert not manager.authorize(1, scope, 1)
assert manager.authorize(2, scope, 2)
def test_two_scopes_have_independent_authorities() -> None:
manager = NpcAuthorityManager()
a = ScopeKey("00000010", "0000003C")
b = ScopeKey("00000020", "0000003C")
manager.reconcile([(5, a), (2, a), (7, b), (9, b)])
assert manager.get(a).player_id == 2
assert manager.get(b).player_id == 7
assert manager.get(a).epoch == 1
assert manager.get(b).epoch == 1
def test_empty_scope_revokes_and_future_reentry_uses_new_epoch() -> None:
manager = NpcAuthorityManager()
scope = ScopeKey("00000010", "")
manager.reconcile([(1, scope)])
revoked = manager.reconcile([])
assert revoked[0].player_id == 0
assert revoked[0].epoch == 2
assert manager.get(scope) is None
granted = manager.reconcile([(3, scope)])
assert granted[0].player_id == 3
assert granted[0].epoch == 3
def test_scope_from_transform_normalizes_form_ids() -> None:
scope = scope_from_transform({"cellId": "1a", "worldspaceId": "3c"})
assert scope == ScopeKey("0000001A", "0000003C")
assert scope_from_transform(None) is None
assert scope_from_transform({"cellId": ""}) is None
-48
View File
@@ -1,48 +0,0 @@
import json
import pytest
from packet_codec import MAX_MESSAGE_BYTES, PacketCodecError, decode_packet, encode_packet
from transport_policy import Delivery
def test_encoded_packet_has_no_tcp_line_framing():
encoded = encode_packet({"type": "transform", "x": 1.0})
assert encoded.packet_type == "transform"
assert encoded.delivery is Delivery.UNRELIABLE_SEQUENCED
assert encoded.payload == b'{"type":"transform","x":1.0}'
assert not encoded.payload.endswith(b"\n")
def test_reliable_packet_carries_delivery_policy_separately_from_json():
encoded = encode_packet({"type": "playerState", "characterName": "Nomad"})
assert encoded.delivery is Delivery.RELIABLE_ORDERED
assert json.loads(encoded.payload) == {"type": "playerState", "characterName": "Nomad"}
def test_decode_packet_accepts_one_complete_message_without_delimiter():
assert decode_packet(b'{"type":"sessionReady","playerId":7}') == {
"type": "sessionReady",
"playerId": 7,
}
def test_codec_rejects_nonfinite_invalid_utf8_nonobject_and_missing_type():
with pytest.raises(PacketCodecError):
encode_packet({"type": "transform", "x": float("nan")})
with pytest.raises(PacketCodecError):
decode_packet(b'{"type":"transform","x":NaN}')
with pytest.raises(PacketCodecError):
decode_packet(b"\xff")
with pytest.raises(PacketCodecError):
decode_packet(b"[]")
with pytest.raises(PacketCodecError):
decode_packet(b"{}")
def test_codec_enforces_64k_message_limit_before_transport():
payload = "x" * MAX_MESSAGE_BYTES
with pytest.raises(PacketCodecError):
encode_packet({"type": "playerState", "characterName": payload})
with pytest.raises(PacketCodecError):
decode_packet(b"x" * (MAX_MESSAGE_BYTES + 1))
-80
View File
@@ -1,80 +0,0 @@
from __future__ import annotations
from player_state import normalize_player_state_packet
from server_core import _normalize_action_events
def valid_packet() -> dict:
return {
"type": "playerState",
"playerId": 9999,
"characterName": "Sole Survivor",
"equippedItems": [
{"slot": "rightHand", "formId": "ff"},
{"slot": "body", "formId": ""},
],
"appearance": {
"version": 4,
"raceFormId": "13746",
"height": 1.0,
"morphWeight": {"thin": 0.2, "muscular": 0.3, "large": 0.5},
"bodyTintColor": {"r": 1, "g": 2, "b": 3, "a": 255},
"hairColorFormId": "",
"facialHairColorFormId": "",
"complexionFormId": "",
"isFemale": False,
"headParts": ["1a2b"],
"morphs": [{"id": "10", "value": 0.5}],
"morphRegions": [0.25],
"facialBoneMorphs": [
{
"id": "20",
"position": [0.0, 1.0, 2.0],
"rotation": [3.0, 4.0, 5.0],
"scale": [1.0, 1.0, 1.0],
}
],
"tints": [{"id": 1, "type": 2, "value": 3, "color": "ff00ff00", "swatch": 4}],
},
"actionEvents": [{"sequence": 7, "type": 3, "eventName": "fireSingle"}],
}
def test_player_state_normalizes_only_reliable_fields():
clean = normalize_player_state_packet(valid_packet(), _normalize_action_events)
assert clean is not None
assert clean["type"] == "playerState"
assert "playerId" not in clean
assert clean["equippedItems"][0]["formId"] == "000000FF"
assert clean["appearance"]["raceFormId"] == "00013746"
assert clean["appearance"]["headParts"] == ["00001A2B"]
assert clean["appearance"]["tints"][0]["color"] == "FF00FF00"
assert clean["actionEvents"][0]["sequence"] == 7
def test_empty_player_state_is_rejected():
assert normalize_player_state_packet({"type": "playerState"}, _normalize_action_events) is None
def test_invalid_equipment_is_rejected():
packet = valid_packet()
packet["equippedItems"][0]["formId"] = "not-a-form"
assert normalize_player_state_packet(packet, _normalize_action_events) is None
def test_invalid_appearance_numbers_are_rejected():
packet = valid_packet()
packet["appearance"]["height"] = float("nan")
assert normalize_player_state_packet(packet, _normalize_action_events) is None
def test_invalid_action_event_is_rejected_instead_of_silently_dropped():
packet = valid_packet()
packet["actionEvents"] = [{"sequence": 7, "type": 3, "eventName": "notAllowed"}]
assert normalize_player_state_packet(packet, _normalize_action_events) is None
def test_character_name_is_bounded():
packet = valid_packet()
packet["characterName"] = "x" * 129
assert normalize_player_state_packet(packet, _normalize_action_events) is None
@@ -1,244 +0,0 @@
from __future__ import annotations
import json
import socket
import time
import pytest
from server_core import PROTOCOL_VERSION, FalloutTogetherServer
_RECV_BUFFERS: dict[socket.socket, bytes] = {}
def recv_packet(sock: socket.socket, timeout: float = 2.0) -> dict:
sock.settimeout(timeout)
data = _RECV_BUFFERS.get(sock, b"")
while b"\n" not in data:
chunk = sock.recv(4096)
if not chunk:
raise ConnectionError("socket closed before a complete packet was received")
data += chunk
line, remainder = data.split(b"\n", 1)
_RECV_BUFFERS[sock] = remainder
return json.loads(line.decode("utf-8"))
def recv_until(sock: socket.socket, packet_type: str, timeout: float = 2.0) -> dict:
deadline = time.monotonic() + timeout
while time.monotonic() < deadline:
remaining = max(0.01, deadline - time.monotonic())
try:
packet = recv_packet(sock, timeout=remaining)
except socket.timeout:
break
if packet.get("type") == packet_type:
return packet
raise AssertionError(f"did not receive packet type {packet_type}")
def send_packet(sock: socket.socket, packet: dict) -> None:
sock.sendall(json.dumps(packet, separators=(",", ":")).encode() + b"\n")
def start_server() -> FalloutTogetherServer:
server = FalloutTogetherServer(host="127.0.0.1", port=0)
server.start()
deadline = time.monotonic() + 2.0
while not server.is_running() and time.monotonic() < deadline:
time.sleep(0.01)
return server
def connect_v2(server: FalloutTogetherServer) -> tuple[socket.socket, dict]:
sock = socket.create_connection(("127.0.0.1", server.port), timeout=2.0)
welcome = recv_until(sock, "welcome")
send_packet(sock, {"type": "hello", "protocolVersion": PROTOCOL_VERSION})
ready = recv_until(sock, "sessionReady")
assert ready["playerId"] == welcome["playerId"]
return sock, welcome
def player_state(**overrides) -> dict:
packet = {
"type": "playerState",
"playerId": 999999,
"characterName": "Sole Survivor",
"equippedItems": [{"slot": "rightHand", "formId": "ff"}],
"appearance": {
"version": 4,
"raceFormId": "13746",
"height": 1.0,
"headParts": ["1a2b"],
},
"actionEvents": [{"sequence": 3, "type": 3, "eventName": "fireSingle"}],
}
packet.update(overrides)
return packet
def transform(**overrides) -> dict:
packet = {
"type": "transform",
"x": 0.0,
"y": 0.0,
"z": 0.0,
"angleZ": 0.0,
"cellId": "00000001",
"worldspaceId": "0000003C",
"movementType": "normal",
}
packet.update(overrides)
return packet
def wait_for_stat(server: FalloutTogetherServer, key: str, minimum: int, timeout: float = 2.0) -> int:
deadline = time.monotonic() + timeout
while time.monotonic() < deadline:
value = int(server.get_stats()[key])
if value >= minimum:
return value
time.sleep(0.01)
raise AssertionError(f"stat {key} did not reach {minimum}")
def establish_transforms(
server: FalloutTogetherServer,
sender: socket.socket,
receiver: socket.socket,
*,
same_interest: bool,
) -> None:
send_packet(sender, transform())
receiver_transform = transform(
cellId="00000001" if same_interest else "00000002",
worldspaceId="0000003C" if same_interest else "00000099",
)
send_packet(receiver, receiver_transform)
wait_for_stat(server, "transformPacketsReceived", 2)
def test_player_state_relay_uses_server_owned_identity():
server = start_server()
sender, sender_welcome = connect_v2(server)
receiver, _ = connect_v2(server)
try:
establish_transforms(server, sender, receiver, same_interest=True)
send_packet(sender, player_state())
relayed = recv_until(receiver, "playerState")
assert relayed["playerId"] == sender_welcome["playerId"]
assert relayed["characterName"] == "Sole Survivor"
assert relayed["equippedItems"] == [{"slot": "rightHand", "formId": "000000FF"}]
assert relayed["appearance"]["raceFormId"] == "00013746"
assert relayed["actionEvents"][0]["sequence"] == 3
wait_for_stat(server, "playerStatePacketsReceived", 1)
wait_for_stat(server, "playerStatePacketsBroadcast", 1)
finally:
sender.close()
receiver.close()
server.stop()
def test_v2_transform_strips_reliable_player_state_fields():
server = start_server()
sender, _ = connect_v2(server)
receiver, _ = connect_v2(server)
try:
send_packet(
sender,
transform(
equippedItems=[{"slot": "rightHand", "formId": "000000FF"}],
appearance={"version": 4, "raceFormId": "00013746"},
actionEvents=[{"sequence": 3, "type": 3, "eventName": "fireSingle"}],
characterName="must-not-ride-transform",
),
)
relayed = recv_until(receiver, "transform")
assert "equippedItems" not in relayed
assert "appearance" not in relayed
assert "actionEvents" not in relayed
assert "characterName" not in relayed
finally:
sender.close()
receiver.close()
server.stop()
def test_out_of_interest_player_state_strips_action_events_but_keeps_durable_state():
server = start_server()
sender, _ = connect_v2(server)
receiver, _ = connect_v2(server)
try:
establish_transforms(server, sender, receiver, same_interest=False)
send_packet(sender, player_state())
relayed = recv_until(receiver, "playerState")
assert relayed["characterName"] == "Sole Survivor"
assert relayed["equippedItems"] == [{"slot": "rightHand", "formId": "000000FF"}]
assert "actionEvents" not in relayed
finally:
sender.close()
receiver.close()
server.stop()
def test_out_of_interest_action_only_player_state_is_not_relayed():
server = start_server()
sender, _ = connect_v2(server)
receiver, _ = connect_v2(server)
try:
establish_transforms(server, sender, receiver, same_interest=False)
send_packet(
sender,
{
"type": "playerState",
"actionEvents": [{"sequence": 4, "type": 3, "eventName": "fireSingle"}],
},
)
wait_for_stat(server, "playerStatePacketsReceived", 1)
with pytest.raises(AssertionError):
recv_until(receiver, "playerState", timeout=0.2)
finally:
sender.close()
receiver.close()
server.stop()
def test_new_v2_client_receives_cached_player_state_without_stale_action_events():
server = start_server()
sender, sender_welcome = connect_v2(server)
receiver = None
try:
send_packet(sender, transform())
wait_for_stat(server, "transformPacketsReceived", 1)
send_packet(sender, player_state())
wait_for_stat(server, "playerStatePacketsReceived", 1)
receiver, _ = connect_v2(server)
cached = recv_until(receiver, "playerState")
assert cached["playerId"] == sender_welcome["playerId"]
assert cached["characterName"] == "Sole Survivor"
assert cached["equippedItems"] == [{"slot": "rightHand", "formId": "000000FF"}]
assert "actionEvents" not in cached
finally:
sender.close()
if receiver is not None:
receiver.close()
server.stop()
def test_malformed_player_state_is_rejected_before_relay():
server = start_server()
sender, _ = connect_v2(server)
receiver, _ = connect_v2(server)
try:
before = server.get_stats()["packetsRejected"]
send_packet(sender, player_state(equippedItems=[{"slot": "rightHand", "formId": "not-a-form"}]))
wait_for_stat(server, "packetsRejected", before + 1)
assert server.get_stats()["playerStatePacketsReceived"] == 0
with pytest.raises(AssertionError):
recv_until(receiver, "playerState", timeout=0.2)
finally:
sender.close()
receiver.close()
server.stop()
-88
View File
@@ -1,88 +0,0 @@
from __future__ import annotations
import socket
import threading
import time
from protocol_v2_client import PROTOCOL_VERSION, ProtocolV2Client
from server_core import FalloutTogetherServer
def _start_server(max_players: int = 4) -> tuple[FalloutTogetherServer, threading.Thread]:
server = FalloutTogetherServer(host="127.0.0.1", port=0, max_players=max_players)
server._prepare_server_socket()
thread = threading.Thread(target=server._accept_loop, daemon=True)
thread.start()
return server, thread
def _stop_server(server: FalloutTogetherServer, thread: threading.Thread) -> None:
server.stop()
thread.join(timeout=3.0)
def _recv_until(client: ProtocolV2Client, packet_type: str, timeout: float = 2.0) -> dict:
assert client.socket is not None
previous_timeout = client.socket.gettimeout()
client.socket.settimeout(0.25)
deadline = time.monotonic() + timeout
try:
while time.monotonic() < deadline:
try:
packet = client.recv_packet()
except socket.timeout:
continue
if packet.get("type") == packet_type:
return packet
finally:
client.socket.settimeout(previous_timeout)
raise AssertionError(f"did not receive packet type {packet_type}")
def test_synthetic_client_uses_protocol_v2_without_legacy_activation() -> None:
server, thread = _start_server()
client = ProtocolV2Client("127.0.0.1", server.port, name="pytest-v2")
try:
session = client.connect()
assert session.player_id == 1
assert session.server_protocol_version == PROTOCOL_VERSION
stats = server.get_stats()
assert stats["connectedClients"] == 1
assert stats["protocolV2Connections"] == 1
assert stats["legacyConnections"] == 0
assert stats["pendingConnections"] == 0
finally:
client.close()
_stop_server(server, thread)
def test_same_cell_v2_clients_receive_each_others_transforms() -> None:
server, thread = _start_server()
first = ProtocolV2Client("127.0.0.1", server.port, name="first")
second = ProtocolV2Client("127.0.0.1", server.port, name="second")
try:
first.connect()
second.connect()
# The first client receives a host-assignment control packet after sessionReady.
_recv_until(first, "worldStateHost")
second.send_transform(
x=100.0,
y=200.0,
z=300.0,
angle_z=0.5,
cell_id=0x0000003C,
movement_speed=150.0,
animation_direction=90.0,
is_moving=True,
)
relayed = _recv_until(first, "transform")
assert relayed["playerId"] == second.session.player_id
assert relayed["cellId"] == "0000003C"
assert relayed["animationDirection"] == 90.0
finally:
first.close()
second.close()
_stop_server(server, thread)
-200
View File
@@ -1,200 +0,0 @@
from __future__ import annotations
import json
import socket
import threading
import time
from pathlib import Path
import pytest
from admin_server import DEFAULT_ADMIN_TOKEN_PATH, send_admin_command
from lan_discovery import LanDiscoveryResponder
from server_core import FalloutTogetherServer, get_lan_addresses
from server_service import ServerConfig, ServerService
def _free_port() -> int:
with socket.socket(socket.AF_INET, socket.SOCK_STREAM) as sock:
sock.bind(("127.0.0.1", 0))
return int(sock.getsockname()[1])
def test_get_lan_addresses_never_returns_wildcard() -> None:
addresses = get_lan_addresses()
assert "0.0.0.0" not in addresses
for address in addresses:
assert not address.startswith("127.")
def test_discovery_failure_does_not_stop_game_server(
tmp_path: Path, monkeypatch: pytest.MonkeyPatch
) -> None:
game_port = _free_port()
admin_port = _free_port()
def _fail_start(self: LanDiscoveryResponder) -> None:
raise OSError("Could not bind LAN discovery to 0.0.0.0:7778. simulated failure")
monkeypatch.setattr(LanDiscoveryResponder, "start", _fail_start)
service = ServerService(
ServerConfig(
host="127.0.0.1",
port=game_port,
admin_port=admin_port,
bans_path=str(tmp_path / "bans.json"),
)
)
thread = threading.Thread(target=service.serve_forever, daemon=True)
thread.start()
deadline = time.time() + 5.0
while time.time() < deadline and not service.is_running():
time.sleep(0.05)
assert service.is_running()
with socket.create_connection(("127.0.0.1", game_port), timeout=2.0) as conn:
data = conn.recv(4096)
assert b"welcome" in data
service.stop()
thread.join(timeout=3.0)
assert not service.is_running()
def test_stop_is_idempotent(tmp_path: Path) -> None:
game_port = _free_port()
admin_port = _free_port()
service = ServerService(
ServerConfig(
host="127.0.0.1",
port=game_port,
admin_port=admin_port,
bans_path=str(tmp_path / "bans.json"),
)
)
service.start()
deadline = time.time() + 5.0
while time.time() < deadline and not service.is_running():
time.sleep(0.05)
assert service.is_running()
service.stop()
service.stop()
assert not service.is_running()
def test_admin_and_game_bind_errors_include_address() -> None:
occupied = _free_port()
holder = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
try:
holder.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
holder.bind(("127.0.0.1", occupied))
holder.listen()
server = FalloutTogetherServer(host="127.0.0.1", port=occupied)
with pytest.raises(OSError, match=rf"127\.0\.0\.1:{occupied}"):
server._prepare_server_socket()
finally:
holder.close()
def test_headless_admin_roundtrip(tmp_path: Path) -> None:
game_port = _free_port()
admin_port = _free_port()
service = ServerService(
ServerConfig(
host="127.0.0.1",
port=game_port,
admin_port=admin_port,
bans_path=str(tmp_path / "bans.json"),
)
)
thread = threading.Thread(target=service.serve_forever, daemon=True)
thread.start()
deadline = time.time() + 5.0
while time.time() < deadline and not service.is_running():
time.sleep(0.05)
assert service.is_running()
response = send_admin_command({"cmd": "ping"}, port=admin_port)
assert response.get("ok") is True
assert response.get("data", {}).get("pong") is True
service.stop()
thread.join(timeout=3.0)
def test_admin_rejects_unauthenticated_requests_and_does_not_disclose_token(tmp_path: Path) -> None:
game_port = _free_port()
admin_port = _free_port()
service = ServerService(
ServerConfig(
host="127.0.0.1",
port=game_port,
admin_port=admin_port,
bans_path=str(tmp_path / "bans.json"),
)
)
thread = threading.Thread(target=service.serve_forever, daemon=True)
thread.start()
deadline = time.time() + 5.0
while time.time() < deadline and not service.is_running():
time.sleep(0.05)
assert service.is_running()
with socket.create_connection(("127.0.0.1", admin_port), timeout=2.0) as conn:
conn.sendall(b'{"cmd":"ping"}\n')
raw = conn.recv(4096)
unauthorized = json.loads(raw.split(b"\n", 1)[0].decode("utf-8"))
assert unauthorized.get("ok") is False
assert unauthorized.get("error") == "Unauthorized admin request."
authenticated = send_admin_command({"cmd": "status"}, port=admin_port)
assert authenticated.get("ok") is True
token = DEFAULT_ADMIN_TOKEN_PATH.read_text(encoding="utf-8").strip()
assert token
assert token not in json.dumps(authenticated, sort_keys=True)
service.stop()
thread.join(timeout=3.0)
def test_lan_discovery_sets_broadcast_option() -> None:
class DummyServer:
def get_stats(self):
return {
"connectedClients": 0,
"port": 7777,
"serverName": "Test",
"serverDescription": "",
"maxPlayers": 16,
}
def _log(self, message: str, *, level: str = "info") -> None:
return None
sock = socket.socket(socket.AF_INET, socket.SOCK_DGRAM)
sock.bind(("127.0.0.1", 0))
port = int(sock.getsockname()[1])
sock.close()
responder = LanDiscoveryResponder(DummyServer(), discovery_port=port)
responder.start()
try:
assert responder._socket is not None
value = responder._socket.getsockopt(socket.SOL_SOCKET, socket.SO_BROADCAST)
assert value in (0, 1)
probe = {
"type": "discover",
"protocol": "commonwealth-online",
}
with socket.socket(socket.AF_INET, socket.SOCK_DGRAM) as client:
client.settimeout(2.0)
client.sendto(json.dumps(probe).encode("utf-8"), ("127.0.0.1", port))
data, _addr = client.recvfrom(2048)
packet = json.loads(data.decode("utf-8"))
assert packet["type"] == "discoverResponse"
assert packet["port"] == 7777
finally:
responder.stop()
-114
View File
@@ -1,114 +0,0 @@
from __future__ import annotations
import json
from pathlib import Path
from typer.testing import CliRunner
import consumer_server_cli
runner = CliRunner()
def test_serve_accepts_positional_config_path(tmp_path: Path, monkeypatch) -> None:
config_path = tmp_path / "commonwealth-server.json"
config_path.write_text(
json.dumps(
{
"host": "127.0.0.1",
"port": 1,
"server_name": "Arg Test",
"max_players": 2,
"log_verbosity": "info",
"admin_port": 2,
}
)
+ "\n",
encoding="utf-8",
newline="\n",
)
# Avoid binding real sockets; validate CLI parsing only.
called: dict[str, object] = {}
class FakeService:
def __init__(self) -> None:
self.config = None
def add_log_listener(self, _callback) -> None:
return None
def serve_forever(self) -> None:
called["served"] = True
def stop(self) -> None:
return None
monkeypatch.setattr(consumer_server_cli, "get_service", FakeService)
monkeypatch.setattr(consumer_server_cli, "print_startup_banner", lambda _cfg: None)
monkeypatch.setattr(consumer_server_cli, "ensure_writable_directory", lambda _path: None)
monkeypatch.setattr(
consumer_server_cli,
"validate_config",
lambda _cfg: (True, []),
)
monkeypatch.setattr(
consumer_server_cli.signal,
"signal",
lambda *_args, **_kwargs: None,
)
result = runner.invoke(
consumer_server_cli.app,
["serve", str(config_path)],
)
assert result.exit_code == 0, result.output
assert called.get("served") is True
def test_serve_accepts_config_option(tmp_path: Path, monkeypatch) -> None:
config_path = tmp_path / "commonwealth-server.json"
config_path.write_text(
json.dumps(
{
"host": "127.0.0.1",
"port": 1,
"server_name": "Arg Test",
"max_players": 2,
"log_verbosity": "info",
"admin_port": 2,
}
)
+ "\n",
encoding="utf-8",
newline="\n",
)
called: dict[str, object] = {}
class FakeService:
def __init__(self) -> None:
self.config = None
def add_log_listener(self, _callback) -> None:
return None
def serve_forever(self) -> None:
called["served"] = True
def stop(self) -> None:
return None
monkeypatch.setattr(consumer_server_cli, "get_service", FakeService)
monkeypatch.setattr(consumer_server_cli, "print_startup_banner", lambda _cfg: None)
monkeypatch.setattr(consumer_server_cli, "ensure_writable_directory", lambda _path: None)
monkeypatch.setattr(consumer_server_cli, "validate_config", lambda _cfg: (True, []))
monkeypatch.setattr(consumer_server_cli.signal, "signal", lambda *_args, **_kwargs: None)
result = runner.invoke(
consumer_server_cli.app,
["serve", "--config", str(config_path)],
)
assert result.exit_code == 0, result.output
assert called.get("served") is True
@@ -1,120 +0,0 @@
from __future__ import annotations
import pytest
import server_service
from server_service import ServerConfig, ServerService
from transport_server import TransportAwareFalloutTogetherServer
class FakeCoreServer:
def __init__(self, port: int = 7777):
self.port = port
class FakeTransport:
instances = []
local_port_override = None
def __init__(self, host, port, *, library_path=None):
self.host = host
self.requested_port = port
self.library_path = library_path
self.local_port = self.local_port_override or port
self.closed = False
self.__class__.instances.append(self)
def close(self):
self.closed = True
class FakeAdapter:
instances = []
def __init__(self, core, transport):
self.core = core
self.transport = transport
self.started = False
self.stopped = False
self.__class__.instances.append(self)
def start(self):
self.started = True
def stop(self):
self.stopped = True
self.transport.close()
def reset_fakes():
FakeTransport.instances.clear()
FakeTransport.local_port_override = None
FakeAdapter.instances.clear()
def test_server_service_uses_transport_aware_core_with_gns_disabled_by_default(tmp_path):
config = ServerConfig(bans_path=str(tmp_path / "bans.json"))
assert not config.enable_gns_transport
service = ServerService(config)
core = service._create_server()
assert isinstance(core, TransportAwareFalloutTogetherServer)
def test_gns_feature_flag_is_noop_when_disabled(monkeypatch):
reset_fakes()
monkeypatch.setattr(server_service, "GnsServerTransport", FakeTransport)
monkeypatch.setattr(server_service, "GnsGameplayAdapter", FakeAdapter)
service = ServerService(ServerConfig(enable_gns_transport=False))
service._server = FakeCoreServer()
service._start_gns()
assert service._gns_adapter is None
assert FakeTransport.instances == []
assert FakeAdapter.instances == []
def test_gns_feature_flag_starts_udp_on_same_numeric_port_and_stops_cleanly(monkeypatch):
reset_fakes()
monkeypatch.setattr(server_service, "GnsServerTransport", FakeTransport)
monkeypatch.setattr(server_service, "GnsGameplayAdapter", FakeAdapter)
service = ServerService(
ServerConfig(
host="127.0.0.1",
port=7777,
enable_gns_transport=True,
gns_bridge_path="/opt/commonwealth/libcommonwealth_online_gns_bridge.so",
)
)
service._server = FakeCoreServer(7777)
service._start_gns()
assert len(FakeTransport.instances) == 1
transport = FakeTransport.instances[0]
assert transport.host == "127.0.0.1"
assert transport.requested_port == 7777
assert transport.library_path == "/opt/commonwealth/libcommonwealth_online_gns_bridge.so"
assert len(FakeAdapter.instances) == 1
adapter = FakeAdapter.instances[0]
assert adapter.started
assert service._gns_adapter is adapter
service._stop_gns()
assert adapter.stopped
assert transport.closed
assert service._gns_adapter is None
def test_gns_start_rejects_native_bridge_bound_to_wrong_udp_port(monkeypatch):
reset_fakes()
FakeTransport.local_port_override = 8888
monkeypatch.setattr(server_service, "GnsServerTransport", FakeTransport)
monkeypatch.setattr(server_service, "GnsGameplayAdapter", FakeAdapter)
service = ServerService(ServerConfig(port=7777, enable_gns_transport=True))
service._server = FakeCoreServer(7777)
with pytest.raises(RuntimeError, match="expected UDP 7777"):
service._start_gns()
assert len(FakeTransport.instances) == 1
assert FakeTransport.instances[0].closed
assert FakeAdapter.instances == []
assert service._gns_adapter is None
-105
View File
@@ -1,105 +0,0 @@
from __future__ import annotations
import json
import os
import signal
import socket
import subprocess
import sys
import time
from pathlib import Path
import pytest
from admin_server import send_admin_command
pytestmark = pytest.mark.skipif(os.name == "nt", reason="SIGTERM integration is POSIX-only")
SERVER_DIR = Path(__file__).resolve().parents[1]
def _free_port() -> int:
with socket.socket(socket.AF_INET, socket.SOCK_STREAM) as sock:
sock.bind(("127.0.0.1", 0))
return int(sock.getsockname()[1])
def test_sigterm_exits_cleanly(tmp_path: Path) -> None:
game_port = _free_port()
admin_port = _free_port()
config_path = tmp_path / "commonwealth-server.json"
config_path.write_text(
json.dumps(
{
"host": "127.0.0.1",
"port": game_port,
"server_name": "CI Test Server",
"server_description": "",
"max_players": 4,
"log_verbosity": "info",
"admin_port": admin_port,
},
indent=2,
)
+ "\n",
encoding="utf-8",
newline="\n",
)
env = os.environ.copy()
env["PYTHONUNBUFFERED"] = "1"
env["NO_COLOR"] = "1"
process = subprocess.Popen(
[
sys.executable,
"-u",
str(SERVER_DIR / "consumer_server_cli.py"),
"serve",
"--config",
str(config_path),
],
cwd=str(SERVER_DIR),
stdout=subprocess.PIPE,
stderr=subprocess.STDOUT,
text=True,
env=env,
)
try:
deadline = time.time() + 10.0
connected = False
while time.time() < deadline:
try:
with socket.create_connection(("127.0.0.1", game_port), timeout=0.5) as conn:
welcome = conn.recv(4096)
if b"welcome" in welcome:
connected = True
break
except OSError:
time.sleep(0.1)
assert connected, "server did not accept a test client in time"
admin_deadline = time.time() + 5.0
response: dict[str, object] | None = None
while time.time() < admin_deadline:
try:
response = send_admin_command({"cmd": "ping"}, port=admin_port, timeout_seconds=0.5)
break
except (OSError, RuntimeError, ConnectionError):
time.sleep(0.05)
assert response is not None
assert response.get("ok") is True
process.send_signal(signal.SIGTERM)
try:
exit_code = process.wait(timeout=10.0)
except subprocess.TimeoutExpired:
process.kill()
pytest.fail("server did not exit after SIGTERM")
assert exit_code == 0
finally:
if process.poll() is None:
process.kill()
process.wait(timeout=5.0)
-51
View File
@@ -1,51 +0,0 @@
from snapshot_sequence import MAX_SEQUENCE, SequenceCounter, SequenceWindow, is_newer, next_sequence
def test_next_sequence_skips_zero_after_wrap():
assert next_sequence(0) == 1
assert next_sequence(1) == 2
assert next_sequence(MAX_SEQUENCE) == 1
def test_is_newer_is_wrap_safe():
assert is_newer(1, 0)
assert is_newer(2, 1)
assert not is_newer(1, 1)
assert not is_newer(1, 2)
assert is_newer(1, MAX_SEQUENCE)
assert not is_newer(MAX_SEQUENCE, 1)
assert not is_newer(0, MAX_SEQUENCE)
def test_counter_and_window_enforce_latest_wins():
counter = SequenceCounter()
assert counter.current == 0
assert counter.advance() == 1
assert counter.advance() == 2
counter.reset()
assert counter.advance() == 1
window = SequenceWindow()
assert window.accept(1)
assert not window.accept(1)
assert window.accept(2)
assert not window.accept(1)
window.reset()
assert window.accept(MAX_SEQUENCE)
assert window.accept(1)
assert not window.accept(MAX_SEQUENCE)
def test_loss_and_reordering_do_not_stall_newer_snapshots():
window = SequenceWindow()
arrivals = (1, 3, 2, 6, 4, 5, 6, 8)
accepted = [sequence for sequence in arrivals if window.accept(sequence)]
assert accepted == [1, 3, 6, 8]
assert window.last_accepted == 8
def test_invalid_sequence_inputs_are_rejected():
assert not is_newer(-1, 0)
assert not is_newer(1, -1)
assert not is_newer(True, 0)
assert not is_newer(1, MAX_SEQUENCE + 1)
-50
View File
@@ -1,50 +0,0 @@
import socket
import pytest
from packet_codec import MAX_MESSAGE_BYTES, encode_packet
from tcp_transport import LineMessageBuffer, TcpFramingError, frame_message, send_message
def test_tcp_framing_adds_delimiter_only_at_transport_boundary():
encoded = encode_packet({"type": "playerState", "characterName": "Nomad"})
assert not encoded.payload.endswith(b"\n")
framed = frame_message(encoded.payload)
assert framed == encoded.payload + b"\n"
def test_line_buffer_preserves_fragmented_and_coalesced_messages():
buffer = LineMessageBuffer()
assert buffer.feed(b'{"type":"first"') == []
messages = buffer.feed(b'}\n{"type":"second"}\n{"type":"third"')
assert messages == [b'{"type":"first"}', b'{"type":"second"}']
assert buffer.buffered_bytes > 0
assert buffer.feed(b'}\r\n') == [b'{"type":"third"}']
assert buffer.buffered_bytes == 0
def test_line_buffer_rejects_oversize_unterminated_and_terminated_messages():
with pytest.raises(TcpFramingError):
LineMessageBuffer().feed(b"x" * (MAX_MESSAGE_BYTES + 1))
with pytest.raises(TcpFramingError):
LineMessageBuffer().feed((b"x" * (MAX_MESSAGE_BYTES + 1)) + b"\n")
def test_frame_rejects_raw_delimiters_and_oversize_payloads():
with pytest.raises(TcpFramingError):
frame_message(b"bad\nmessage")
with pytest.raises(TcpFramingError):
frame_message(b"bad\rmessage")
with pytest.raises(TcpFramingError):
frame_message(b"x" * (MAX_MESSAGE_BYTES + 1))
def test_send_message_preserves_existing_tcp_wire_format():
reader, writer = socket.socketpair()
try:
payload = encode_packet({"type": "keepAlive"}).payload
send_message(writer, payload)
assert reader.recv(4096) == payload + b"\n"
finally:
reader.close()
writer.close()
-40
View File
@@ -1,40 +0,0 @@
from transport_policy import (
Delivery,
delivery_for_packet_type,
is_snapshot_packet,
requires_application_sequence,
)
def test_snapshot_packets_use_unreliable_sequenced_delivery():
assert delivery_for_packet_type("transform") is Delivery.UNRELIABLE_SEQUENCED
assert delivery_for_packet_type("npcState") is Delivery.UNRELIABLE_SEQUENCED
assert is_snapshot_packet("transform")
assert is_snapshot_packet("npcState")
assert requires_application_sequence("transform")
assert requires_application_sequence("npcState")
def test_gameplay_and_control_packets_use_reliable_ordered_delivery():
reliable_types = (
"hello",
"sessionReady",
"worldState",
"combatHit",
"npcAuthority",
"positionCorrection",
"disconnect",
"sessionEnded",
"equipmentState",
"actionEvent",
"playerState",
)
for packet_type in reliable_types:
assert delivery_for_packet_type(packet_type) is Delivery.RELIABLE_ORDERED
assert not is_snapshot_packet(packet_type)
assert not requires_application_sequence(packet_type)
def test_unknown_packet_types_never_default_to_unreliable():
assert delivery_for_packet_type("futureControlPacket") is Delivery.RELIABLE_ORDERED
assert not requires_application_sequence("futureControlPacket")
-80
View File
@@ -1,80 +0,0 @@
from __future__ import annotations
import socket
import time
import pytest
from client_session import ClientSession
from packet_codec import EncodedPacket
from transport_policy import Delivery
from transport_server import TransportAwareFalloutTogetherServer
class MessageConnection:
def __init__(self):
self.messages: list[EncodedPacket] = []
def send_encoded(self, encoded: EncodedPacket) -> None:
self.messages.append(encoded)
def fileno(self) -> int:
return 1
def close(self) -> None:
pass
def make_client(connection) -> ClientSession:
return ClientSession(
connection=connection,
address=("127.0.0.1", 50000),
player_id=1,
connected_at=time.time(),
)
def test_transport_aware_server_preserves_tcp_wire_format():
server = TransportAwareFalloutTogetherServer(host="127.0.0.1", port=0)
reader, writer = socket.socketpair()
try:
client = make_client(writer)
server._send_packet(client, {"type": "playerState", "characterName": "Nomad"})
assert reader.recv(4096) == b'{"type":"playerState","characterName":"Nomad"}\n'
assert client.packets_sent == 1
assert server.get_stats()["packetsSent"] == 1
finally:
reader.close()
writer.close()
def test_transport_aware_server_sends_raw_message_and_delivery_metadata_to_gns_style_connection():
server = TransportAwareFalloutTogetherServer(host="127.0.0.1", port=0)
connection = MessageConnection()
client = make_client(connection)
server._send_packet(client, {"type": "transform", "x": 1.0}, broadcast=True)
server._send_packet(client, {"type": "playerState", "characterName": "Nomad"})
assert len(connection.messages) == 2
transform, player_state = connection.messages
assert transform.payload == b'{"type":"transform","x":1.0}'
assert transform.delivery is Delivery.UNRELIABLE_SEQUENCED
assert not transform.payload.endswith(b"\n")
assert player_state.delivery is Delivery.RELIABLE_ORDERED
assert player_state.payload == b'{"type":"playerState","characterName":"Nomad"}'
assert client.packets_sent == 2
assert client.packets_broadcast == 1
assert server.get_stats()["packetsSent"] == 2
assert server.get_stats()["packetsBroadcast"] == 1
def test_transport_aware_server_rejects_nonfinite_json_before_transport():
server = TransportAwareFalloutTogetherServer(host="127.0.0.1", port=0)
connection = MessageConnection()
client = make_client(connection)
with pytest.raises(ValueError):
server._send_packet(client, {"type": "transform", "x": float("nan")})
assert connection.messages == []
assert server.get_stats()["packetsSent"] == 0