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); } }); } } }