diff --git a/overlay/Scripts/Custom/Bridge/BridgeBoot.cs b/overlay/Scripts/Custom/Bridge/BridgeBoot.cs new file mode 100644 index 0000000..9fcab73 --- /dev/null +++ b/overlay/Scripts/Custom/Bridge/BridgeBoot.cs @@ -0,0 +1,180 @@ +using System; +using System.Collections.Generic; + +using Server.Commands; + +namespace Server.Custom.Bridge +{ + /// + /// Lifecycle wiring. Boot order (Server/Main.cs:544-562, all on the Core thread): + /// + /// Configure() -> World.Load() -> Initialize() -> EventSink.ServerStarted + /// + /// Config is read in Configure. Handlers are attached in Initialize. The socket opens on + /// ServerStarted, once the world is actually there to describe. + /// + /// EventSink.Shutdown does NOT fire on a crash (Server/Main.cs:198,313), so the sidecar + /// must treat socket EOF as normal and re-handshake rather than waiting for a goodbye. + /// + public static class BridgeBoot + { + private static readonly Dictionary>> _handlers = + new Dictionary>>(StringComparer.Ordinal); + + /// + /// Identifies this run of the shard. It is stable across sidecar reconnects and changes + /// on every shard restart, which is how the sidecar tells "I reconnected" (keep my + /// cached state) from "the shard restarted" (discard it). + /// + private static string _bootId; + + public static void Configure() + { + BridgeConfig.Configure(); + } + + public static void Initialize() + { + if (!BridgeConfig.Enabled) + { + Console.WriteLine("[Bridge] disabled by config"); + return; + } + + CommandSystem.Register("bridge", AccessLevel.Administrator, Bridge_OnCommand); + + RegisterHandler("ping", OnPing); + + BridgeLink.InboundLine += OnInboundLine; + BridgeLink.Connected_Core += EmitHello; + + EventSink.ServerStarted += OnServerStarted; + EventSink.Shutdown += OnShutdown; + EventSink.Crashed += OnCrashed; + + Console.WriteLine("[Bridge] {0}", BridgeConfig.Describe()); + } + + /// Handlers run on the Core thread. They may touch the world freely. + public static void RegisterHandler(string kind, Action> handler) + { + _handlers[kind] = handler; + } + + private static void OnServerStarted() + { + _bootId = Guid.NewGuid().ToString("N"); + BridgeLink.Start(); + } + + /// + /// Core thread, once per connection. The sidecar restarts independently of the shard, + /// so this is sent on every connect rather than once at boot — otherwise a sidecar that + /// came up second would never learn which shard it is talking to. + /// + private static void EmitHello() + { + BridgeLink.Emit(BridgeJson.Begin("server.hello") + .Str("shard", Server.Misc.ServerList.ServerName) + .Str("bootId", _bootId) + .Num("connects", BridgeLink.Connects) + .Num("items", World.Items.Count) + .Num("mobiles", World.Mobiles.Count) + .Num("accounts", Accounting.Accounts.Count) + .End()); + } + + private static void OnShutdown(ShutdownEventArgs e) + { + BridgeLink.Emit(BridgeJson.Begin("server.shutdown").End()); + + // Stop() joins the link thread for up to 2s, which gives the writer a chance to drain + // the goodbye. Best effort: the sidecar must not depend on receiving it. + BridgeLink.Stop(); + } + + private static void OnCrashed(CrashedEventArgs e) + { + try + { + BridgeLink.Emit(BridgeJson.Begin("server.crashed") + .Str("error", e.Exception == null ? null : e.Exception.Message) + .End()); + + BridgeLink.Stop(); + } + catch + { + // The process is already going down. Never make a crash worse. + } + } + + /// Core thread, one call per inbound line. + private static void OnInboundLine(string line) + { + var obj = BridgeJson.Parse(line); + + if (obj == null) + { + Console.WriteLine("[Bridge] malformed inbound line, ignoring"); + return; + } + + var kind = BridgeJson.GetString(obj, "kind"); + + if (kind == null) + return; + + Action> handler; + + if (!_handlers.TryGetValue(kind, out handler)) + { + Console.WriteLine("[Bridge] no handler for inbound kind '{0}'", kind); + return; + } + + handler(obj); + } + + private static void OnPing(Dictionary o) + { + var sb = BridgeJson.Begin("pong"); + + var id = BridgeJson.GetString(o, "id"); + + if (id != null) + sb.Str("id", id); + + BridgeLink.Emit(sb.End()); + } + + [Usage("bridge [status | reload | ping]")] + [Description("Inspects and controls the sidecar link.")] + private static void Bridge_OnCommand(CommandEventArgs e) + { + var arg = e.Length > 0 ? e.GetString(0).ToLowerInvariant() : "status"; + + switch (arg) + { + case "reload": + BridgeConfig.Load(); + e.Mobile.SendMessage("Bridge: {0}", BridgeConfig.Describe()); + e.Mobile.SendMessage("Bridge: endpoint changes take effect on reconnect."); + break; + + case "ping": + BridgeLink.Emit(BridgeJson.Begin("ping").End()); + e.Mobile.SendMessage("Bridge: ping queued."); + break; + + default: + e.Mobile.SendMessage("Bridge: {0}", BridgeConfig.Describe()); + e.Mobile.SendMessage( + "Bridge: connected={0} depth={1} sent={2} dropped={3} received={4} connects={5} writeErrors={6}", + BridgeLink.Connected, BridgeLink.Depth, BridgeLink.Sent, BridgeLink.Dropped, + BridgeLink.Received, BridgeLink.Connects, BridgeLink.WriteErrors); + break; + } + } + } +} diff --git a/overlay/Scripts/Custom/Bridge/BridgeConfig.cs b/overlay/Scripts/Custom/Bridge/BridgeConfig.cs new file mode 100644 index 0000000..c5e0207 --- /dev/null +++ b/overlay/Scripts/Custom/Bridge/BridgeConfig.cs @@ -0,0 +1,52 @@ +using System; + +namespace Server.Custom.Bridge +{ + /// + /// Tunables from Config/Bridge.cfg. Key scope is the filename, so `Port=7788` there + /// reads as "Bridge.Port" here. + /// + /// Loaded in Configure(), which ScriptCompiler invokes before World.Load. + /// + public static class BridgeConfig + { + public static string Host { get; private set; } + public static int Port { get; private set; } + public static int QueueCap { get; private set; } + + public static int StatSweepSeconds { get; private set; } + public static int DecaySweepSeconds { get; private set; } + public static int EconomySweepSeconds { get; private set; } + + public static bool Enabled { get; private set; } + + public static void Configure() + { + Load(); + } + + /// Re-readable at runtime via `[bridge reload`. + public static void Load() + { + Enabled = Config.Get("Bridge.Enabled", true); + + Host = Config.Get("Bridge.Host", "127.0.0.1"); + Port = Config.Get("Bridge.Port", 7788); + QueueCap = Config.Get("Bridge.QueueCap", 10000); + + StatSweepSeconds = Config.Get("Bridge.StatSweepSeconds", 30); + DecaySweepSeconds = Config.Get("Bridge.DecaySweepSeconds", 60); + EconomySweepSeconds = Config.Get("Bridge.EconomySweepSeconds", 300); + + if (QueueCap < 16) + QueueCap = 16; + } + + public static string Describe() + { + return String.Format( + "enabled={0} endpoint={1}:{2} queueCap={3} sweeps(stat={4}s decay={5}s econ={6}s)", + Enabled, Host, Port, QueueCap, StatSweepSeconds, DecaySweepSeconds, EconomySweepSeconds); + } + } +} diff --git a/overlay/Scripts/Custom/Bridge/BridgeJson.cs b/overlay/Scripts/Custom/Bridge/BridgeJson.cs new file mode 100644 index 0000000..8c63b06 --- /dev/null +++ b/overlay/Scripts/Custom/Bridge/BridgeJson.cs @@ -0,0 +1,169 @@ +using System; +using System.Collections.Generic; +using System.Globalization; +using System.Text; +using System.Web.Script.Serialization; + +namespace Server.Custom.Bridge +{ + /// + /// Outbound JSON is written by hand into a StringBuilder. It runs on the Core thread for + /// every emitted event, and the measured budget in docs/PLAN.md assumes this cost, not a + /// reflection serializer's. + /// + /// Inbound JSON is parsed with JavaScriptSerializer. Commands arrive at human rates, so + /// correctness beats speed there, and parsing happens on the reader thread anyway. + /// + public static class BridgeJson + { + private static readonly DateTime Epoch = new DateTime(1970, 1, 1, 0, 0, 0, DateTimeKind.Utc); + + [ThreadStatic] + private static JavaScriptSerializer _parser; + + public static long NowMs() + { + return (long)(DateTime.UtcNow - Epoch).TotalMilliseconds; + } + + // ---- outbound ---- + + /// Opens an object and writes the `t` and `kind` fields. + public static StringBuilder Begin(string kind) + { + var sb = new StringBuilder(256); + sb.Append("{\"t\":").Append(NowMs()); + sb.Append(",\"kind\":\"").Append(kind).Append('"'); + return sb; + } + + public static StringBuilder Str(this StringBuilder sb, string name, string value) + { + sb.Append(",\"").Append(name).Append("\":"); + + if (value == null) + sb.Append("null"); + else + Escape(sb, value); + + return sb; + } + + public static StringBuilder Num(this StringBuilder sb, string name, long value) + { + sb.Append(",\"").Append(name).Append("\":").Append(value); + return sb; + } + + public static StringBuilder Num(this StringBuilder sb, string name, double value) + { + sb.Append(",\"").Append(name).Append("\":") + .Append(value.ToString("R", CultureInfo.InvariantCulture)); + return sb; + } + + public static StringBuilder Bool(this StringBuilder sb, string name, bool value) + { + sb.Append(",\"").Append(name).Append("\":").Append(value ? "true" : "false"); + return sb; + } + + /// Serial as the canonical "0x1A2B" string the sidecar keys on. + public static StringBuilder Ser(this StringBuilder sb, string name, Serial serial) + { + sb.Append(",\"").Append(name).Append("\":\"0x") + .Append(serial.Value.ToString("X")).Append('"'); + return sb; + } + + /// Closes the object. The trailing newline is the frame delimiter. + public static string End(this StringBuilder sb) + { + sb.Append('}'); + return sb.ToString(); + } + + public static void Escape(StringBuilder sb, string value) + { + sb.Append('"'); + + for (int i = 0; i < value.Length; i++) + { + char c = value[i]; + + switch (c) + { + case '"': sb.Append("\\\""); break; + case '\\': sb.Append("\\\\"); break; + case '\n': sb.Append("\\n"); break; + case '\r': sb.Append("\\r"); break; + case '\t': sb.Append("\\t"); break; + case '\b': sb.Append("\\b"); break; + case '\f': sb.Append("\\f"); break; + default: + if (c < ' ') + sb.Append("\\u").Append(((int)c).ToString("x4")); + else + sb.Append(c); + break; + } + } + + sb.Append('"'); + } + + // ---- inbound ---- + + /// + /// Parses one line into a dictionary. Returns null on malformed input rather than + /// throwing: a bad line from the sidecar must never reach a game code path. + /// + public static Dictionary Parse(string line) + { + if (String.IsNullOrEmpty(line)) + return null; + + try + { + if (_parser == null) + { + _parser = new JavaScriptSerializer(); + _parser.MaxJsonLength = 1 << 20; + } + + return _parser.Deserialize>(line); + } + catch + { + return null; + } + } + + public static string GetString(Dictionary o, string key) + { + object v; + + if (o == null || !o.TryGetValue(key, out v) || v == null) + return null; + + return v as string ?? Convert.ToString(v, CultureInfo.InvariantCulture); + } + + public static int GetInt(Dictionary o, string key, int fallback) + { + object v; + + if (o == null || !o.TryGetValue(key, out v) || v == null) + return fallback; + + try + { + return Convert.ToInt32(v, CultureInfo.InvariantCulture); + } + catch + { + return fallback; + } + } + } +} diff --git a/overlay/Scripts/Custom/Bridge/BridgeLink.cs b/overlay/Scripts/Custom/Bridge/BridgeLink.cs new file mode 100644 index 0000000..fe1f193 --- /dev/null +++ b/overlay/Scripts/Custom/Bridge/BridgeLink.cs @@ -0,0 +1,331 @@ +using System; +using System.Collections.Concurrent; +using System.IO; +using System.Net.Sockets; +using System.Text; +using System.Threading; + +namespace Server.Custom.Bridge +{ + /// + /// The loopback link to the Rust sidecar. Newline-delimited JSON, bidirectional. + /// + /// Threading contract, which the whole bridge depends on: + /// + /// * is called from the Core thread. It formats nothing, blocks on + /// nothing, and touches no socket. It enqueues and returns. A slow, wedged, or absent + /// sidecar cannot stall the shard. + /// * One link thread owns the socket. It connects, drains the queue, and reconnects with + /// backoff. A single writer keeps event ordering intact. + /// * A reader thread parses inbound lines and hands each to the Core thread via + /// Timer.DelayCall. The reader never touches World, Mobile, Item, or Account. + /// + /// The outbound queue is bounded. On overflow the oldest record is dropped and counted, + /// because telemetry is worth less than the shard's memory. + /// + public static class BridgeLink + { + private static readonly ConcurrentQueue _outbound = new ConcurrentQueue(); + private static readonly AutoResetEvent _wake = new AutoResetEvent(false); + + private static Thread _link; + private static volatile bool _running; + private static volatile bool _connected; + + // Set when the peer goes away, so the writer stops trying. + private static volatile bool _dead; + + // Incremented per connection attempt. A reader from a previous connection must not be + // able to mark a newer one dead — reader.Join can time out, and the stale thread's + // finally block would otherwise tear down the connection that replaced it. + private static int _epoch; + + /// + /// Loopback reconnects are cheap, so the ceiling is low. A sidecar restart should cost + /// a few seconds of buffering, not half a minute. + /// + private const int MaxBackoffMs = 5000; + + private static int _depth; + private static long _sent, _dropped, _received, _connects, _writeErrors; + + /// Raised on the Core thread, one call per inbound line. + public static event Action InboundLine; + + /// + /// Raised on the Core thread after each successful connect. The sidecar may + /// restart independently of the shard, so anything it needs to know up front has to be + /// re-sent per connection, not once at ServerStarted. + /// + public static event Action Connected_Core; + + public static bool Connected { get { return _connected; } } + public static int Depth { get { return Volatile.Read(ref _depth); } } + public static long Sent { get { return Interlocked.Read(ref _sent); } } + public static long Dropped { get { return Interlocked.Read(ref _dropped); } } + public static long Received { get { return Interlocked.Read(ref _received); } } + /// Total successful connections, so the first connect counts as 1. + public static long Connects { get { return Interlocked.Read(ref _connects); } } + public static long WriteErrors { get { return Interlocked.Read(ref _writeErrors); } } + + public static void Start() + { + if (_running) + return; + + _running = true; + + _link = new Thread(LinkLoop) + { + Name = "Bridge Link", + IsBackground = true + }; + + _link.Start(); + } + + public static void Stop() + { + if (!_running) + return; + + _running = false; + _wake.Set(); + + var t = _link; + + if (t != null && !t.Join(TimeSpan.FromSeconds(2.0))) + Console.WriteLine("[Bridge] link thread did not stop cleanly"); + + _link = null; + _connected = false; + } + + /// + /// Core thread. Non-blocking. `line` must already be a complete JSON object with no + /// embedded newline; the newline is appended by the writer as the frame delimiter. + /// + public static void Emit(string line) + { + if (!_running || line == null) + return; + + // Drop-oldest. Bound first, then enqueue, so the queue can transiently sit one over + // the cap but never grows without limit. + while (Volatile.Read(ref _depth) >= BridgeConfig.QueueCap) + { + string discard; + + if (!_outbound.TryDequeue(out discard)) + break; + + Interlocked.Decrement(ref _depth); + Interlocked.Increment(ref _dropped); + } + + _outbound.Enqueue(line); + Interlocked.Increment(ref _depth); + + _wake.Set(); + } + + private static void LinkLoop() + { + int backoffMs = 500; + + while (_running) + { + TcpClient client = null; + Thread reader = null; + + int epoch = Interlocked.Increment(ref _epoch); + + try + { + client = new TcpClient(); + client.NoDelay = true; + client.Connect(BridgeConfig.Host, BridgeConfig.Port); + + var stream = client.GetStream(); + stream.WriteTimeout = 5000; // a wedged peer must surface as an error, not a hang + + _dead = false; + _connected = true; + backoffMs = 500; + + Interlocked.Increment(ref _connects); + Console.WriteLine("[Bridge] connected to {0}:{1}", BridgeConfig.Host, BridgeConfig.Port); + + // Building the greeting reads the world, so it must happen on the Core thread. + Timer.DelayCall(TimeSpan.Zero, () => + { + try + { + var handler = Connected_Core; + + if (handler != null) + handler(); + } + catch (Exception ex) + { + Console.WriteLine("[Bridge] connect handler threw: {0}", ex); + } + }); + + var localStream = stream; + reader = new Thread(() => ReadLoop(localStream, epoch)) + { + Name = "Bridge Reader", + IsBackground = true + }; + reader.Start(); + + WriteLoop(stream); + } + catch (Exception ex) + { + if (_connected) + Console.WriteLine("[Bridge] link error: {0}", ex.Message); + } + finally + { + _connected = false; + _dead = true; + + try { if (client != null) client.Close(); } + catch { } + + if (reader != null) + reader.Join(TimeSpan.FromSeconds(1.0)); + } + + if (!_running) + break; + + // Nothing is listening yet, or the sidecar restarted. Both are normal. + Thread.Sleep(backoffMs); + backoffMs = Math.Min(backoffMs * 2, MaxBackoffMs); + } + + _connected = false; + } + + private static void WriteLoop(NetworkStream stream) + { + while (_running && !_dead) + { + string line; + + if (!_outbound.TryDequeue(out line)) + { + _wake.WaitOne(250); + continue; + } + + Interlocked.Decrement(ref _depth); + + try + { + var bytes = Encoding.UTF8.GetBytes(line + "\n"); + stream.Write(bytes, 0, bytes.Length); + Interlocked.Increment(ref _sent); + } + catch (Exception) + { + // The record is already off the queue. Count it and let the outer loop + // reconnect; re-queueing risks an unbounded retry storm against a dead peer. + Interlocked.Increment(ref _writeErrors); + _dead = true; + throw; + } + } + } + + private static void ReadLoop(NetworkStream stream, int epoch) + { + var buffer = new byte[8192]; + var line = new StringBuilder(512); + + try + { + while (_running && !_dead) + { + int read = stream.Read(buffer, 0, buffer.Length); + + if (read <= 0) + break; // clean EOF: the sidecar closed. Normal. + + for (int i = 0; i < read; i++) + { + char c = (char)buffer[i]; + + if (c == '\n') + { + Dispatch(line.ToString()); + line.Clear(); + } + else if (c != '\r') + { + line.Append(c); + + if (line.Length > (1 << 20)) + { + Console.WriteLine("[Bridge] inbound line too long, dropping"); + line.Clear(); + } + } + } + } + } + catch (IOException) + { + // Expected when the peer vanishes mid-read. + } + catch (ObjectDisposedException) + { + // Expected when Stop() closes the socket under us. + } + catch (Exception ex) + { + Console.WriteLine("[Bridge] reader error: {0}", ex.Message); + } + finally + { + // Only tear down the connection this reader actually owned. + if (Volatile.Read(ref _epoch) == epoch) + { + _dead = true; + _wake.Set(); // let the writer notice and fall through to reconnect + } + } + } + + /// + /// Reader thread. Marshals to the Core thread. Timer.DelayCall's scheduling path is + /// lock-protected and safe to call from any thread; the callback runs on the main loop. + /// + private static void Dispatch(string line) + { + if (line.Length == 0) + return; + + Interlocked.Increment(ref _received); + + Timer.DelayCall(TimeSpan.Zero, () => + { + try + { + var handler = InboundLine; + + if (handler != null) + handler(line); + } + catch (Exception ex) + { + // A malformed command must never escape into a game code path. + Console.WriteLine("[Bridge] inbound handler threw: {0}", ex); + } + }); + } + } +} diff --git a/overlay/Scripts/Scripts.csproj b/overlay/Scripts/Scripts.csproj index 3a5a655..105a3d9 100644 --- a/overlay/Scripts/Scripts.csproj +++ b/overlay/Scripts/Scripts.csproj @@ -31,6 +31,8 @@ + + diff --git a/tools/stub_sidecar.ps1 b/tools/stub_sidecar.ps1 new file mode 100644 index 0000000..2c65604 --- /dev/null +++ b/tools/stub_sidecar.ps1 @@ -0,0 +1,31 @@ +param( + [int] $Port = 7788, + [string] $Log = "$PSScriptRoot\sidecar_loop.log" +) + +$ErrorActionPreference = 'Stop' +"[sidecar] listening on 127.0.0.1:$Port" | Out-File $Log -Encoding utf8 + +$listener = New-Object System.Net.Sockets.TcpListener([System.Net.IPAddress]::Loopback, $Port) +$listener.Start() + +while ($true) { + try { + $client = $listener.AcceptTcpClient() + "[sidecar] === shard connected ===" | Add-Content $Log + + $stream = $client.GetStream() + $reader = New-Object System.IO.StreamReader($stream) + + while ($null -ne ($line = $reader.ReadLine())) { + "[sidecar] <- $line" | Add-Content $Log + } + + "[sidecar] === shard disconnected ===" | Add-Content $Log + $client.Close() + } + catch { + "[sidecar] error: $_" | Add-Content $Log + Start-Sleep -Milliseconds 200 + } +}