Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
7 changes: 5 additions & 2 deletions src/RustServerMetrics/Config/ConfigData.cs
Original file line number Diff line number Diff line change
Expand Up @@ -19,7 +19,7 @@ class ConfigData
#endregion

[JsonProperty(PropertyName = "Enabled")]
public bool Enabled;
public bool Enabled = false;

[JsonProperty(PropertyName = "Influx Database Url")]
public string DatabaseUrl = DefaultInfluxDbUrl;
Expand All @@ -37,11 +37,14 @@ class ConfigData
public string ServerTag = DefaultServerTag;

[JsonProperty(PropertyName = "Debug Logging")]
public bool DebugLogging;
public bool DebugLogging = false;

[JsonProperty(PropertyName = "Amount of metrics to submit in each request")]
public ushort BatchSize = 1000;

[JsonProperty(PropertyName = "Gather Player Averages (Client FPS, Client Latency, Player FPS, Player Memory, Player Latency, Player Packet Loss)")]
public bool GatherPlayerMetrics = true;

[JsonProperty(PropertyName = "Compress submitted metrics with gzip")]
public bool CompressRequests = true;
}
23 changes: 14 additions & 9 deletions src/RustServerMetrics/HarmonyPatches/Utility/MetricsTimeStorage.cs
Original file line number Diff line number Diff line change
Expand Up @@ -6,22 +6,27 @@ namespace RustServerMetrics.HarmonyPatches.Utility;

public class MetricsTimeStorage<TKey>(string metricKey, Action<StringBuilder, TKey> stringBuilderSerializer)
{
private readonly Dictionary<TKey, double> _dict = new ();

private sealed class Accumulator
{
public double Duration;
}

private readonly Dictionary<TKey, Accumulator> _dict = new ();

private readonly StringBuilder _sb = new();

public void LogTime(TKey key, double milliseconds)
{
if (!MetricsLogger.IsReady)
return;
if (!_dict.TryGetValue(key, out var currentDuration))

if (_dict.TryGetValue(key, out var accumulator))
{
_dict.Add(key, milliseconds);
accumulator.Duration += milliseconds;
return;
}
_dict[key] = currentDuration + milliseconds;

_dict.Add(key, new Accumulator { Duration = milliseconds });
}

public void SerializeToStringBuilder()
Expand All @@ -44,10 +49,10 @@ public void SerializeToStringBuilder()
stringBuilderSerializer.Invoke(_sb, item.Key);

_sb.Append("\" duration=");
_sb.Append((float)item.Value);
_sb.Append((float)item.Value.Duration);
_sb.Append(" ");
_sb.Append(epochNow);
instance.AddToSendBuffer(_sb.ToString());
instance.AddToSendBuffer(_sb);
}

_dict.Clear();
Expand Down
159 changes: 97 additions & 62 deletions src/RustServerMetrics/MetricsLogger.cs
Original file line number Diff line number Diff line change
Expand Up @@ -17,9 +17,12 @@ public class MetricsLogger : SingletonComponent<MetricsLogger>
{
private const string ConfigurationPath = "HarmonyMods_Data/ServerMetrics/Configuration.json";
private readonly StringBuilder _stringBuilder = new();
private readonly Dictionary<ulong, Action> _playerStatsActions = new();
private readonly Dictionary<ulong, uint> _perfReportDelayCounter = new();

private const int PlayerStatsBucketCount = 10;
private const float PlayerStatsBucketInterval = 1f / PlayerStatsBucketCount;
private int _playerStatsBucket;

private class NetworkUpdateData
{
public int Count;
Expand All @@ -33,17 +36,39 @@ public NetworkUpdateData(int count, long bytes)
}
}

private readonly Dictionary<Message.Type, NetworkUpdateData> _networkUpdates = Enum.GetValues(typeof(Message.Type))
.Cast<Message.Type>()
.Distinct()
.ToDictionary(x => x,
_ => new NetworkUpdateData(0, 0));
private static readonly Message.Type[] MessageTypes = Enum.GetValues(typeof(Message.Type))
.Cast<Message.Type>()
.Distinct()
.ToArray();

private static readonly IReadOnlyDictionary<Message.Type, string> MessageTypeNames = Enum.GetValues(typeof(Message.Type))
.Cast<Message.Type>()
.Distinct()
.ToDictionary(x => x,
x => x.ToString());
private static readonly int MessageTypeSlotOffset = MessageTypes.Min(x => (int)x);
private static readonly int MessageTypeSlotCount = MessageTypes.Max(x => (int)x) - MessageTypeSlotOffset + 1;

private static readonly string[] MessageTypeNames = BuildMessageTypeNames();

private readonly NetworkUpdateData[] _networkUpdates = BuildNetworkUpdates();

private static string[] BuildMessageTypeNames()
{
var names = new string[MessageTypeSlotCount];
foreach (var messageType in MessageTypes)
{
names[(int)messageType - MessageTypeSlotOffset] = messageType.ToString();
}

return names;
}

private static NetworkUpdateData[] BuildNetworkUpdates()
{
var networkUpdates = new NetworkUpdateData[MessageTypeSlotCount];
foreach (var messageType in MessageTypes)
{
networkUpdates[(int)messageType - MessageTypeSlotOffset] = new NetworkUpdateData(0, 0);
}

return networkUpdates;
}

public readonly MetricsTimeStorage<MethodInfo> ServerInvokes = new("invoke_execution", LogMethodInfo);
public readonly MetricsTimeStorage<string> ServerRpcCalls = new("rpc_calls", LogMethodName);
Expand Down Expand Up @@ -139,6 +164,7 @@ public override void Awake()
public void StartLoggingMetrics()
{
InvokeRepeating(LogNetworkUpdates, UnityEngine.Random.Range(0.25f, 0.75f), 0.5f);
InvokeRepeating(GatherPlayerStatsBucket, UnityEngine.Random.Range(0.5f, 1.5f), PlayerStatsBucketInterval);

InvokeRepeating(ServerInvokes.SerializeToStringBuilder, UnityEngine.Random.Range(0f, 1f), 1f);
InvokeRepeating(ServerRpcCalls.SerializeToStringBuilder, UnityEngine.Random.Range(0f, 1f), 1f);
Expand All @@ -155,20 +181,13 @@ internal void OnPlayerInit(BasePlayer player)
{
if (!Ready) return;
if (!Configuration.GatherPlayerMetrics) return;
var action = new Action(() => GatherPlayerSecondStats(player));
if (_playerStatsActions.TryGetValue(player.userID, out var existingAction))
player.CancelInvoke(existingAction);
_playerStatsActions[player.userID] = action;
player.InvokeRepeating(action, UnityEngine.Random.Range(0.5f, 1.5f), 1f);

_perfReportDelayCounter[player.userID] = (uint)UnityEngine.Random.Range(0, 5);
}

internal void OnPlayerDisconnected(BasePlayer player)
{
if (!Ready) return;
if (!Configuration.GatherPlayerMetrics) return;
if (_playerStatsActions.TryGetValue(player.userID, out var action))
player.CancelInvoke(action);
_playerStatsActions.Remove(player.userID);
_perfReportDelayCounter.Remove(player.userID);
}

Expand All @@ -189,7 +208,18 @@ internal void OnNetWriteSend(NetWrite write, SendInfo sendInfo)
return;
}

var data = _networkUpdates[_lastMessageType];
var slot = (int)_lastMessageType - MessageTypeSlotOffset;
if ((uint)slot >= (uint)_networkUpdates.Length)
{
return;
}

var data = _networkUpdates[slot];
if (data == null)
{
return;
}

if (sendInfo.connection != null)
{
data.Count++;
Expand All @@ -199,7 +229,7 @@ internal void OnNetWriteSend(NetWrite write, SendInfo sendInfo)
{
var count = sendInfo.connections.Count;
data.Count += count;
data.Bytes += write.Length * count;
data.Bytes += (long)write.Length * count;
}
}

Expand Down Expand Up @@ -240,8 +270,28 @@ internal bool OnClientPerformanceReport(ProtoBuf.PerformanceReport clientPerform
return true;
}

private void GatherPlayerStatsBucket()
{
if (!Ready) return;
if (!Configuration.GatherPlayerMetrics) return;

var players = BasePlayer.activePlayerList;
var bucket = _playerStatsBucket;
_playerStatsBucket = bucket + 1 < PlayerStatsBucketCount ? bucket + 1 : 0;

for (var i = bucket; i < players.Count; i += PlayerStatsBucketCount)
{
var player = players[i];
if (player == null) continue;

GatherPlayerSecondStats(player);
}
}

private void GatherPlayerSecondStats(BasePlayer player)
{
if (player.net?.connection == null) return;

if (!player.IsReceivingSnapshot)
{
_perfReportDelayCounter.TryGetValue(player.userID, out var perfReportCounter);
Expand All @@ -259,11 +309,12 @@ private void GatherPlayerSecondStats(BasePlayer player)
UploadPacket("connection_latency", player, (builder, basePlayer) =>
{
var ip = basePlayer.net.connection.ipaddress;
var portSeparator = ip.LastIndexOf(':');

builder.Append(",steamid=");
builder.Append(basePlayer.UserIDString);
builder.Append(",ip=");
builder.Append(ip[..ip.LastIndexOf(':')]);
builder.Append(ip, 0, portSeparator < 0 ? ip.Length : portSeparator);
builder.Append(" ping=");
builder.Append(Net.sv.GetAveragePing(basePlayer.net.connection));
builder.Append("i,packet_loss=");
Expand All @@ -274,20 +325,29 @@ private void GatherPlayerSecondStats(BasePlayer player)

private void LogNetworkUpdates()
{
if (_networkUpdates.Count < 1) return;
if (_networkUpdates.Length < 1) return;
var serverTag = Configuration.ServerTag;
var epochNow = DateTimeOffset.UtcNow.ToUnixTimeMilliseconds();
_stringBuilder.Clear();
_stringBuilder.Append("network_updates,server=");
_stringBuilder.Append(serverTag);
_stringBuilder.Append(" ");

var enumerator = _networkUpdates.GetEnumerator();
if (enumerator.MoveNext())
var isFirstField = true;
for (var slot = 0; slot < _networkUpdates.Length; slot++)
{
var networkUpdate = enumerator.Current;
var key = MessageTypeNames[networkUpdate.Key];
var value = networkUpdate.Value;
var value = _networkUpdates[slot];
if (value == null) continue;

var key = MessageTypeNames[slot];

if (!isFirstField)
{
_stringBuilder.Append(",");
}

isFirstField = false;

// Count first named {type}
_stringBuilder.Append(key);
_stringBuilder.Append("=");
Expand All @@ -303,35 +363,11 @@ private void LogNetworkUpdates()
_stringBuilder.Append(value.Bytes);
_stringBuilder.Append("i");
value.Bytes = 0;

while (enumerator.MoveNext())
{
networkUpdate = enumerator.Current;
key = MessageTypeNames[networkUpdate.Key];
value = networkUpdate.Value;

// Count first named {type}
_stringBuilder.Append(",");
_stringBuilder.Append(key);
_stringBuilder.Append("=");
_stringBuilder.Append(value.Count);
_stringBuilder.Append("i");
value.Count = 0;

// Bytes second named as "{type}_bytes"
_stringBuilder.Append(",");
_stringBuilder.Append(key);
_stringBuilder.Append("_bytes");
_stringBuilder.Append("=");
_stringBuilder.Append(value.Bytes);
_stringBuilder.Append("i");
value.Bytes = 0;
}
}

_stringBuilder.Append(" ");
_stringBuilder.Append(epochNow);
_reportUploader.AddToSendBuffer(_stringBuilder.ToString());
_reportUploader.AddToSendBuffer(_stringBuilder);
}

internal void OnPerformanceReportGenerated()
Expand Down Expand Up @@ -438,7 +474,7 @@ private void LogPerformanceReport(Performance.Tick current, string epochNow, str
_stringBuilder.Append("i ");
_stringBuilder.Append(epochNow);

_reportUploader.AddToSendBuffer(_stringBuilder.ToString());
_reportUploader.AddToSendBuffer(_stringBuilder);
}


Expand All @@ -458,11 +494,13 @@ public void UploadPacket<T>(string id, T data, Action<StringBuilder, T> serializ
_stringBuilder.Append(" ");
_stringBuilder.Append(epochNow);

AddToSendBuffer(_stringBuilder.ToString());
AddToSendBuffer(_stringBuilder);
}

public void AddToSendBuffer(string toString) => _reportUploader.AddToSendBuffer(toString);

public void AddToSendBuffer(StringBuilder payload) => _reportUploader.AddToSendBuffer(payload);

private long GetMemoryUsage(Performance.Tick performanceTick)
{
if (performanceTick.memoryUsageSystem > 0)
Expand Down Expand Up @@ -550,6 +588,8 @@ private void StatusCommand(ConsoleSystem.Arg arg)
_stringBuilder.AppendLine("Report Uploader:");
_stringBuilder.Append("\tRunning: "); _stringBuilder.Append(_reportUploader.IsRunning); _stringBuilder.AppendLine();
_stringBuilder.Append("\tIn Buffer: "); _stringBuilder.Append(_reportUploader.BufferSize); _stringBuilder.AppendLine();
_stringBuilder.Append("\tIn Flight: "); _stringBuilder.Append(_reportUploader.PendingBatches); _stringBuilder.AppendLine();
_stringBuilder.Append("\tDropped (total): "); _stringBuilder.Append(_reportUploader.TotalDroppedReports); _stringBuilder.AppendLine();
arg.ReplyWith(_stringBuilder.ToString());
}

Expand All @@ -568,12 +608,7 @@ private void ReloadCfgCommand(ConsoleSystem.Arg arg)
CancelInvoke(invoke.action);
}

foreach (var player in _playerStatsActions)
{
var basePlayer = BasePlayer.FindByID(player.Key);
if (basePlayer == null) continue;
basePlayer.CancelInvoke(player.Value);
}
_perfReportDelayCounter.Clear();
_reportUploader.Stop();

if (!Configuration.Enabled)
Expand Down
Loading