using System;
using System.Collections.Concurrent;
using System.Collections.Generic;
using System.Net.Sockets;
using System.Text;
using System.Threading;
using Newtonsoft.Json;
using Oxide.Core;
using Oxide.Core.Plugins;
namespace Oxide.Plugins
{
///
/// The Runic Gateway bridge. Dials out to a rust-link sidecar on loopback and speaks
/// newline-delimited JSON over it: events out, commands in.
///
///
/// Threading contract, which everything later in this plugin depends on:
///
///
/// -
/// is called from the server's main thread. It formats nothing,
/// blocks on nothing, and touches no socket. It enqueues and returns. A slow, wedged,
/// or absent sidecar cannot stall the game.
///
/// -
/// 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 main thread via
/// Interface.Oxide.NextTick. The reader touches no Unity object, no BasePlayer
/// and no ConVar — every one of those is main-thread-only, and reading them from here
/// is the kind of bug that presents as a crash somewhere else entirely.
///
///
///
///
/// The outbound queue is bounded. On overflow the oldest record is dropped and counted,
/// because telemetry is worth less than the server's memory.
///
///
///
/// Loopback is the trust boundary. There is no token on this link: the plugin and
/// the sidecar share a host, and the sidecar binds 127.0.0.1. Pointing Host at
/// anything routable puts an unauthenticated command channel on the network.
///
///
[Info("Runic Gateway", "RunicGateway", "0.1.0")]
[Description("Bridges this Rust server to a Runic Gateway website through a rust-link sidecar.")]
internal class RunicGateway : RustPlugin
{
///
/// Wire-protocol version. Declared in three places that must agree: here,
/// PROTOCOL_VERSION in the sidecar, and module.json in Module-Rust. The
/// overlay's overlay.toml carries this same number so the installer can refuse
/// to pair a sidecar and an overlay that disagree.
///
private const int ProtocolVersion = 1;
///
/// 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;
/// The longest line accepted from the sidecar. Over-long lines are discarded.
private const int MaxInboundLineBytes = 1024 * 1024;
private static readonly ConcurrentQueue Outbound = new ConcurrentQueue();
private static readonly AutoResetEvent Wake = new AutoResetEvent(false);
private PluginConfig _config;
private Thread _link;
private volatile bool _running;
private volatile bool _connected;
/// Set when the peer goes away, so the writer stops trying.
private volatile bool _dead;
///
/// Whether the "cannot reach the sidecar" line has been printed for this load. It is
/// printed once, not per attempt: the retry runs every few seconds for ever.
///
private volatile bool _loggedConnectFailure;
///
/// 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 int _epoch;
private int _depth;
private long _sent, _dropped, _received, _connects, _writeErrors;
///
/// Identifies this run of the server process. It is how the website tells a game
/// restart apart from a sidecar reconnect — the same link comes back either way, and
/// only one of the two means the world it is describing started over.
///
///
/// It must therefore be stable across a plugin reload, which a fresh
/// Guid at Init is not: oxide.reload would announce a brand new
/// boot to a website that is about to reconcile its whole ledger on the strength of it,
/// while the world it describes never moved. The process's own start time is the only
/// value here that changes exactly when the thing it names changes.
///
///
private string _bootId;
// ---- lifecycle ----
private void Init()
{
_bootId = ResolveBootId();
_running = true;
_link = new Thread(LinkLoop)
{
Name = "RunicGateway Link",
IsBackground = true
};
_link.Start();
}
private void Unload()
{
_running = false;
Wake.Set();
Thread t = _link;
if (t != null && !t.Join(TimeSpan.FromSeconds(2.0)))
Puts("link thread did not stop cleanly");
_link = null;
_connected = false;
// A reload leaves the static queue holding lines addressed to a connection that no
// longer exists. Clearing it is the difference between a reload and a slow leak.
string discard;
while (Outbound.TryDequeue(out discard))
{
}
Interlocked.Exchange(ref _depth, 0);
}
// ---- configuration ----
private class PluginConfig
{
/// Where the sidecar listens. Loopback is the trust boundary; see the class docs.
public string Host = "127.0.0.1";
public int Port = 7799;
///
/// How many outbound lines may wait for a wedged sidecar before the oldest is
/// dropped. Sized in lines rather than bytes because that is what is cheap to
/// count; the inbound cap is what bounds a single line's size.
///
public int QueueCap = 5000;
///
/// This server's stable identity across wipes and restarts, as the website knows
/// it. It is not derived from the hostname: an operator renames a server for
/// a season and the site must not lose its history for it.
///
public string ServerId = "main";
}
protected override void LoadDefaultConfig()
{
Config.WriteObject(new PluginConfig(), true);
}
private void LoadConfigValues()
{
_config = Config.ReadObject() ?? new PluginConfig();
if (string.IsNullOrEmpty(_config.Host))
_config.Host = "127.0.0.1";
if (_config.Port <= 0 || _config.Port > 65535)
_config.Port = 7799;
// A cap of zero would drop every line on the way in and report a healthy link doing
// nothing, which is the worst failure available here.
if (_config.QueueCap < 16)
_config.QueueCap = 16;
if (string.IsNullOrEmpty(_config.ServerId))
_config.ServerId = "main";
}
private void Loaded()
{
LoadConfigValues();
Puts($"protocol {ProtocolVersion}, serverId '{_config.ServerId}', sidecar {_config.Host}:{_config.Port}");
}
// ---- the outbound queue ----
///
/// Main thread. Non-blocking. must already be a complete JSON
/// object with no embedded newline; the newline is appended by the writer as the frame
/// delimiter.
///
private void Emit(string line)
{
if (line == null || !_running)
return;
int cap = _config != null ? _config.QueueCap : 5000;
// 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) >= cap)
{
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 void EmitFrame(Dictionary frame)
{
try
{
Emit(JsonConvert.SerializeObject(frame, Formatting.None));
}
catch (Exception ex)
{
// A frame that cannot be serialized must not take the caller down with it — the
// caller is usually a game hook.
Puts($"could not serialize frame: {ex.Message}");
}
}
// ---- the link ----
private void LinkLoop()
{
int backoffMs = 500;
while (_running)
{
TcpClient client = null;
Thread reader = null;
int epoch = Interlocked.Increment(ref _epoch);
try
{
string host = _config != null ? _config.Host : "127.0.0.1";
int port = _config != null ? _config.Port : 7799;
client = new TcpClient();
client.NoDelay = true;
client.Connect(host, port);
NetworkStream stream = client.GetStream();
stream.WriteTimeout = 5000; // a wedged peer must surface as an error, not a hang
_dead = false;
_connected = true;
_loggedConnectFailure = false; // so a LATER outage is reported again
backoffMs = 500;
Interlocked.Increment(ref _connects);
Puts($"connected to {host}:{port}");
// Building the greeting reads the world, so it must happen on the main thread.
// It is sent on every connect rather than once at startup: the sidecar restarts
// independently of the game, and anything it needs up front has to be re-sent
// per connection.
Interface.Oxide.NextTick(SendHello);
NetworkStream localStream = stream;
reader = new Thread(() => ReadLoop(localStream, epoch))
{
Name = "RunicGateway Reader",
IsBackground = true
};
reader.Start();
WriteLoop(stream);
}
catch (Exception ex)
{
// Socket exception messages carry EMBEDDED newlines on Windows — the refusal
// text and the endpoint arrive on separate lines — so trimming the ends is not
// enough to keep the Oxide log one entry per line.
string reason = Flatten(ex.Message);
if (_connected)
{
Puts($"link error: {reason}");
}
else if (!_loggedConnectFailure)
{
// The FIRST failure to connect, and only the first: a sidecar that is not
// running yet is normal at boot and would otherwise spam the console every
// few seconds, but a wrong Host or Port is otherwise silent for ever.
_loggedConnectFailure = true;
Puts($"cannot reach the sidecar: {reason} — retrying quietly until it answers");
}
}
finally
{
// Captured before it is cleared: a connection that ended because the READER saw
// EOF leaves the writer to exit cleanly, so nothing throws and the catch above
// never runs. Without this the link going down is invisible in the game console,
// which is the one place an operator looks first.
bool wasConnected = _connected;
_connected = false;
_dead = true;
try
{
if (client != null)
client.Close();
}
catch
{
// Closing a socket that is already gone is not news.
}
if (reader != null)
reader.Join(TimeSpan.FromSeconds(1.0));
if (wasConnected)
Puts("sidecar link lost; reconnecting");
}
if (!_running)
break;
// Nothing is listening yet, or the sidecar restarted. Both are normal.
//
// `Wake`, not `Thread.Sleep`: Unload runs on the MAIN thread and joins this one, so
// an uninterruptible sleep here is a stall of up to MaxBackoffMs on every reload —
// Oxide reports it as `Calling 'Unload' took 1918ms`, and the server is frozen for
// that whole time. Unload sets `_running` false and signals `Wake`, so waiting on
// the handle turns that stall into an immediate return.
Wake.WaitOne(backoffMs);
backoffMs = Math.Min(backoffMs * 2, MaxBackoffMs);
}
_connected = false;
}
private void WriteLoop(NetworkStream stream)
{
while (_running && !_dead)
{
string line;
if (!Outbound.TryDequeue(out line))
{
Wake.WaitOne(250);
continue;
}
Interlocked.Decrement(ref _depth);
try
{
byte[] 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 void ReadLoop(NetworkStream stream, int epoch)
{
byte[] buffer = new byte[8192];
StringBuilder line = new StringBuilder(512);
bool discarding = false;
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')
{
if (discarding)
{
// The tail of a line already given up on. Swallowing it here is
// what stops the remainder being parsed as a line of its own.
discarding = false;
}
else
{
Dispatch(line.ToString());
}
line.Length = 0;
}
else if (c != '\r')
{
if (discarding)
continue;
line.Append(c);
if (line.Length > MaxInboundLineBytes)
{
Puts("inbound line over the cap; discarding to the next newline");
line.Length = 0;
discarding = true;
}
}
}
}
}
catch (System.IO.IOException)
{
// Expected when the peer vanishes mid-read.
}
catch (ObjectDisposedException)
{
// Expected when Unload closes the socket under us.
}
catch (Exception ex)
{
Puts($"reader error: {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 main thread; nothing that touches the game happens
/// here.
///
private void Dispatch(string line)
{
if (string.IsNullOrEmpty(line))
return;
Interlocked.Increment(ref _received);
Interface.Oxide.NextTick(() =>
{
try
{
HandleCommand(line);
}
catch (Exception ex)
{
// A malformed command must never escape into a game code path.
Puts($"inbound handler threw: {ex}");
}
});
}
// ---- commands ----
/// Main thread. Every world read below is safe only because of that.
private void HandleCommand(string line)
{
Dictionary cmd;
try
{
cmd = JsonConvert.DeserializeObject>(line);
}
catch (Exception ex)
{
Puts($"unparseable command: {ex.Message}");
return;
}
if (cmd == null)
return;
object verbObj;
if (!cmd.TryGetValue("cmd", out verbObj) || verbObj == null)
return;
string verb = verbObj.ToString();
object reqIdObj;
string reqId = cmd.TryGetValue("reqId", out reqIdObj) && reqIdObj != null
? reqIdObj.ToString()
: null;
switch (verb)
{
case "ping":
EmitFrame(new Dictionary
{
{ "kind", "pong" },
{ "t", NowMs() }
});
break;
case "server.status":
Dictionary status = ServerStatus();
status["kind"] = "server.status";
status["t"] = NowMs();
// The correlation id is echoed verbatim, and only when one was supplied. A
// reply that invents one would be routed to nobody; a reply that omits one the
// caller sent would leave that caller waiting out its whole timeout.
if (reqId != null)
status["reqId"] = reqId;
EmitFrame(status);
break;
default:
Puts($"unknown command '{verb}'");
break;
}
}
// ---- frames ----
/// Main thread.
private void SendHello()
{
Dictionary hello = ServerStatus();
hello["kind"] = "server.hello";
hello["t"] = NowMs();
hello["protocol"] = ProtocolVersion;
hello["bootId"] = _bootId;
hello["plugin"] = Version.ToString();
EmitFrame(hello);
}
///
/// Main thread. The shared body of server.hello and server.status — the
/// two frames differ in what wraps them, not in what they say about the server, and
/// letting them drift is how a site ends up showing two different player counts.
///
private Dictionary ServerStatus()
{
var status = new Dictionary
{
{ "serverId", _config != null ? _config.ServerId : "main" },
{ "hostname", ConVar.Server.hostname },
{ "description", ConVar.Server.description },
{ "level", ConVar.Server.level },
{ "seed", ConVar.Server.seed },
{ "worldSize", ConVar.Server.worldsize },
{ "maxPlayers", ConVar.Server.maxplayers },
{ "players", BasePlayer.activePlayerList != null ? BasePlayer.activePlayerList.Count : 0 },
{ "joining", ServerMgr.Instance != null ? ServerMgr.Instance.connectionQueue.Joining : 0 },
{ "queued", ServerMgr.Instance != null ? ServerMgr.Instance.connectionQueue.Queued : 0 },
{ "uptimeSec", (int)UnityEngine.Time.realtimeSinceStartup }
};
// The raw material a wipe id is derived from, rather than a wipe id. Deriving it is the
// website's job and it is not yet specified; emitting a guess now would bake a wrong
// one into stored rows, which is the expensive half to undo.
try
{
status["saveCreatedAt"] = SaveRestore.SaveCreatedTime.ToUniversalTime()
.ToString("yyyy-MM-ddTHH:mm:ssZ", System.Globalization.CultureInfo.InvariantCulture);
}
catch (Exception)
{
// A server that has never saved has no created time. Absent, not zero.
}
return status;
}
///
/// Collapses every run of whitespace and control characters to one space, so an
/// exception message stays one log entry.
///
///
/// The control half is not defensive padding. Mono's SocketException.Message on
/// Windows comes back from a fixed-size OS buffer and is NUL-padded — a couple
/// of hundred \0 bytes sitting in the middle of the sentence. \0 is not
/// whitespace, so Trim does not touch it and neither does a whitespace-only
/// collapse; the line looks like it contains a huge run of spaces and no amount of
/// trimming removes it.
///
///
private static string Flatten(string text)
{
if (string.IsNullOrEmpty(text))
return string.Empty;
var sb = new StringBuilder(text.Length);
bool lastWasSpace = false;
foreach (char c in text)
{
if (char.IsWhiteSpace(c) || char.IsControl(c))
{
if (!lastWasSpace && sb.Length > 0)
sb.Append(' ');
lastWasSpace = true;
}
else
{
sb.Append(c);
lastWasSpace = false;
}
}
return sb.ToString().TrimEnd();
}
///
/// The server process's start instant, as boot-yyyyMMddTHHmmssZ.
///
///
/// Process.StartTime is exact and identical on every read, which is what makes
/// the id survive a plugin reload. The fallback derives the same instant from Unity's
/// monotonic clock, and is only reached if the platform refuses the process query; it
/// is rounded to the second deliberately, because the subtraction jitters in the last
/// decimal and an id that changes on rounding is worse than one a second off.
///
///
private static string ResolveBootId()
{
DateTime startUtc;
try
{
startUtc = System.Diagnostics.Process.GetCurrentProcess().StartTime.ToUniversalTime();
}
catch (Exception)
{
double up = UnityEngine.Time.realtimeSinceStartup;
startUtc = DateTime.UtcNow.AddSeconds(-Math.Round(up));
}
startUtc = new DateTime(
startUtc.Year, startUtc.Month, startUtc.Day,
startUtc.Hour, startUtc.Minute, startUtc.Second,
DateTimeKind.Utc);
return "boot-" + startUtc.ToString(
"yyyyMMdd'T'HHmmss'Z'", System.Globalization.CultureInfo.InvariantCulture);
}
private static long NowMs()
{
return (long)(DateTime.UtcNow - new DateTime(1970, 1, 1, 0, 0, 0, DateTimeKind.Utc))
.TotalMilliseconds;
}
// ---- diagnostics ----
///
/// The link's own counters, from the server console. This is how an operator answers
/// "is the bridge working" without reading a log, and it is the first thing to ask for
/// when the website says a server is offline.
///
[ConsoleCommand("rg.link")]
private void CmdLink(ConsoleSystem.Arg arg)
{
arg.ReplyWith(
$"protocol={ProtocolVersion} serverId={(_config != null ? _config.ServerId : "?")} " +
$"connected={_connected} depth={Volatile.Read(ref _depth)} " +
$"sent={Interlocked.Read(ref _sent)} dropped={Interlocked.Read(ref _dropped)} " +
$"received={Interlocked.Read(ref _received)} connects={Interlocked.Read(ref _connects)} " +
$"writeErrors={Interlocked.Read(ref _writeErrors)} bootId={_bootId}");
}
}
}