Skip to content
Merged
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
2 changes: 1 addition & 1 deletion API.IntegrationTests/GeoLocatedWebApplicationFactory.cs
Original file line number Diff line number Diff line change
Expand Up @@ -37,6 +37,6 @@ private sealed class FakeIpEnrichmentService : IIpEnrichmentService
public GeoPoint? Location { get; set; }

public IpEnrichmentData? Enrich(IPAddress ip) =>
Location is { } location ? new IpEnrichmentData(null, null, null, null, location, 20) : null;
Location is { } location ? new IpEnrichmentData(null, null, null, null, null, location, 20) : null;
}
}
6 changes: 5 additions & 1 deletion API/Controller/Admin/GetOnlineDevices.cs
Original file line number Diff line number Diff line change
Expand Up @@ -60,7 +60,9 @@ public async Task<IActionResult> GetOnlineDevices()
LatencyMs = x.LatencyMs,
Rssi = x.Rssi,
Country = x.Country,
Ip = IpAddressUtils.ParseStoredOrNull(x.Ip, _logger)
Ip = IpAddressUtils.ParseStoredOrNull(x.Ip, _logger),
Asn = x.Asn,
AsnOrg = x.AsnOrg
};
})
);
Expand All @@ -84,5 +86,7 @@ public sealed class AdminOnlineDeviceResponse
public required int? Rssi { get; init; }
public required string? Country { get; set; }
public required IPAddress? Ip { get; set; }
public required long? Asn { get; set; }
public required string? AsnOrg { get; set; }
}
}
2 changes: 1 addition & 1 deletion Common/OpenShockServiceHelper.cs
Original file line number Diff line number Diff line change
Expand Up @@ -347,7 +347,7 @@ public static IServiceCollection AddOpenShockServices(this IServiceCollection se
services.AddScoped<IAutomationTokenService, AutomationTokenService>();

// Ensure GeoOptions is always resolvable so IpEnrichmentService can activate even in hosts
// (Cron, LiveControlGateway, SeedE2E) that don't call RegisterGeoOptions(). TryAdd leaves the
// (Cron, SeedE2E) that don't call RegisterGeoOptions(). TryAdd leaves the
// API's config-bound instance untouched; other hosts get a disabled default (no DB paths).
services.TryAddSingleton(new GeoOptions());
services.AddSingleton<IIpEnrichmentService, IpEnrichmentService>();
Expand Down
2 changes: 2 additions & 0 deletions Common/Redis/DeviceOnline.cs
Original file line number Diff line number Diff line change
Expand Up @@ -23,4 +23,6 @@ public sealed class DeviceOnline
public int? Rssi { get; set; }
public string? Country { get; set; }
public string? Ip { get; set; }
public long? Asn { get; set; }
public string? AsnOrg { get; set; }
}
1 change: 1 addition & 0 deletions Common/Services/Geo/IpEnrichmentData.cs
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
namespace OpenShock.Common.Services.Geo;

public sealed record IpEnrichmentData(
long? Asn,
string? AsnOrg,
bool? IsVpn,
string? CountryCode,
Expand Down
4 changes: 3 additions & 1 deletion Common/Services/Geo/IpEnrichmentService.cs
Original file line number Diff line number Diff line change
Expand Up @@ -86,6 +86,7 @@ public static bool MatchesVpnProvider(string asnOrg)
{
if (_asnReader is null && _cityReader is null) return null;

long? asnNumber = null;
string? asnOrg = null;
// Null means "unknown" (no ASN DB, lookup miss, or failure); only a resolved ASN org yields a verdict.
bool? isVpn = null;
Expand All @@ -96,6 +97,7 @@ public static bool MatchesVpnProvider(string asnOrg)
{
if (_asnReader.TryAsn(ip, out var asn) && asn is not null)
{
asnNumber = asn.AutonomousSystemNumber;
asnOrg = asn.AutonomousSystemOrganization;
if (asnOrg is not null)
{
Expand Down Expand Up @@ -132,7 +134,7 @@ public static bool MatchesVpnProvider(string asnOrg)
}
}

return new IpEnrichmentData(asnOrg, isVpn, countryCode, city, location, accuracyRadiusKm);
return new IpEnrichmentData(asnNumber, asnOrg, isVpn, countryCode, city, location, accuracyRadiusKm);
}

public void Dispose()
Expand Down
12 changes: 11 additions & 1 deletion LiveControlGateway/Controllers/HubControllerBase.cs
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,7 @@
using System.Security.Claims;
using OpenShock.Common.Authentication;
using OpenShock.Common.Extensions;
using OpenShock.Common.Services.Geo;
using SemVersion = OpenShock.Common.Models.SemVersion;
using Timer = System.Timers.Timer;

Expand Down Expand Up @@ -65,6 +66,9 @@ protected HubLifetime HubLifetime

/// <inheritdoc cref="IHubController.Id" />
public override Guid Id => CurrentHubId;

/// <inheritdoc />
public IpEnrichmentData? IpInfo { get; private set; }

/// <summary>
/// Authentication context
Expand Down Expand Up @@ -143,6 +147,10 @@ protected override async Task<SuccessOrProblem> ConnectionPrecondition()
}

_userAgent = HttpContext.Request.Headers.UserAgent.ToString().Truncate(256);

// Resolved before the lifetime is added, so the connect attempt metric already carries it.
IpInfo = ServiceProvider.GetRequiredService<IIpEnrichmentService>().Enrich(HttpContext.GetRemoteIP());

var hubLifetimeResult = await _hubLifetimeManager.TryAddDeviceConnection(5, this, LinkedToken);

switch (hubLifetimeResult)
Expand Down Expand Up @@ -250,7 +258,9 @@ protected async Task<bool> SelfOnline(ulong uptimeMs, ushort? latency = null, in
LatencyMs = latency,
Rssi = rssi,
Country = HttpContext.GetCFIPCountry(),
Ip = HttpContext.GetRemoteIP()
Ip = HttpContext.GetRemoteIP(),
Asn = IpInfo?.Asn,
AsnOrg = IpInfo?.AsnOrg
});

return true;
Expand Down
7 changes: 7 additions & 0 deletions LiveControlGateway/Controllers/IHubController.cs
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
using OpenShock.Common.Models;
using OpenShock.Common.Services.Geo;
using OpenShock.Serialization.Gateway;

namespace OpenShock.LiveControlGateway.Controllers;
Expand All @@ -13,6 +14,12 @@ public interface IHubController : IAsyncDisposable
/// </summary>
public Guid Id { get; }

/// <summary>
/// GeoIP data for the address this connection came from, resolved once when it connected.
/// Null when no GeoIP database is configured.
/// </summary>
public IpEnrichmentData? IpInfo { get; }

/// <summary>
/// Control shockers
/// </summary>
Expand Down
15 changes: 15 additions & 0 deletions LiveControlGateway/LifetimeManager/HubLifetime.cs
Original file line number Diff line number Diff line change
Expand Up @@ -504,9 +504,12 @@ public async Task<Union2<Success, OnlineStateUpdated>> Online(Guid device, SelfO
// as we don't want to send a device online status every time, we will do it here
online.BootedAt = data.BootedAt;
online.LatencyMs = data.LatencyMs;

online.Rssi = data.Rssi;
online.Country = data.Country;
online.Ip = data.Ip?.ToString();
online.Asn = data.Asn;
online.AsnOrg = data.AsnOrg;

var sendOnlineStatusUpdate = false;

Expand Down Expand Up @@ -562,6 +565,8 @@ await deviceOnline.InsertAsync(new DeviceOnline
Rssi = data.Rssi,
Country = data.Country,
Ip = data.Ip?.ToString(),
Asn = data.Asn,
AsnOrg = data.AsnOrg,
}, Duration.DeviceKeepAliveTimeout);
}
}
Expand Down Expand Up @@ -674,4 +679,14 @@ public SelfOnlineData(
/// Remote ip address
/// </summary>
public IPAddress? Ip { get; init; } = null;

/// <summary>
/// Autonomous system number of the remote ip, if GeoIP is available
/// </summary>
public long? Asn { get; init; } = null;

/// <summary>
/// Organization owning <see cref="Asn"/>
/// </summary>
public string? AsnOrg { get; init; } = null;
}
39 changes: 30 additions & 9 deletions LiveControlGateway/LifetimeManager/HubLifetimeManager.cs
Original file line number Diff line number Diff line change
Expand Up @@ -75,11 +75,32 @@ public HubLifetimeManager(

meter.CreateObservableUpDownCounter("openshock_hub_connections", () =>
{
return new[]
var lifetimes = _lifetimes;

// Grouped on the org as well as the number: a lookup can resolve the ASN but not its org,
// and the series here must match the tags the per-hub counters were recorded with.
var perAsn = new Dictionary<(string Asn, string Org), int>();
foreach (var (_, lifetime) in lifetimes)
{
new Measurement<int>(_lifetimes.Count, gatewayFqdn)
};
}, "connections", "Current number of connected hubs");
var ipInfo = lifetime.HubController.IpInfo;
var key = ((string)GatewayMetrics.AsnTag(ipInfo).Value!, (string)GatewayMetrics.AsnOrgTag(ipInfo).Value!);
perAsn[key] = perAsn.GetValueOrDefault(key) + 1;
}

// An idle gateway still reports a zero, so a sum over gateways does not go empty.
if (perAsn.Count == 0) perAsn[(GatewayMetrics.Unknown, GatewayMetrics.Unknown)] = 0;

var measurements = new Measurement<int>[perAsn.Count];
var i = 0;
foreach (var ((asn, org), count) in perAsn)
{
measurements[i++] = new Measurement<int>(count, gatewayFqdn,
new KeyValuePair<string, object?>("asn", asn),
new KeyValuePair<string, object?>("asn_org", org));
}

return measurements;
}, "connections", "Current number of connected hubs by network (ASN)");
Comment thread
coderabbitai[bot] marked this conversation as resolved.

// Derived from the hubs themselves rather than tracked alongside them: a separately
// maintained counter drifts from the truth the moment a teardown path misses a decrement.
Expand Down Expand Up @@ -130,7 +151,7 @@ public async Task<Union3<HubLifetime, Busy, Error>> TryAddDeviceConnection(byte
// There already is a hub lifetime, lets swap!
if (!hubLifetime.TryMarkSwapping())
{
_metrics.HubConnectAttempt(GatewayMetrics.HubConnectOutcome.Busy);
_metrics.HubConnectAttempt(GatewayMetrics.HubConnectOutcome.Busy, hubController);
return new Busy();
}

Expand All @@ -149,7 +170,7 @@ public async Task<Union3<HubLifetime, Busy, Error>> TryAddDeviceConnection(byte
{
_logger.LogTrace("Swapping hub lifetime [{HubId}]", hubController.Id);
await hubLifetime.Swap(hubController);
_metrics.HubConnectAttempt(GatewayMetrics.HubConnectOutcome.Swapped);
_metrics.HubConnectAttempt(GatewayMetrics.HubConnectOutcome.Swapped, hubController);
}
else
{
Expand All @@ -159,11 +180,11 @@ public async Task<Union3<HubLifetime, Busy, Error>> TryAddDeviceConnection(byte
// If we fail to initialize, the hub must be removed
await RemoveDeviceConnection(hubController); // Here be dragons?
_logger.LogError("Failed to initialize hub lifetime [{HubId}]", hubController.Id);
_metrics.HubConnectAttempt(GatewayMetrics.HubConnectOutcome.InitFailed);
_metrics.HubConnectAttempt(GatewayMetrics.HubConnectOutcome.InitFailed, hubController);
return new Error();
}

_metrics.HubConnectAttempt(GatewayMetrics.HubConnectOutcome.Connected);
_metrics.HubConnectAttempt(GatewayMetrics.HubConnectOutcome.Connected, hubController);
}

return hubLifetime;
Expand Down Expand Up @@ -237,7 +258,7 @@ public async Task RemoveDeviceConnection(IHubController hubController)
else
{
_lifetimes = withoutHub;
_metrics.HubDisconnected();
_metrics.HubDisconnected(hubController);
}
}
}
Expand Down
32 changes: 29 additions & 3 deletions LiveControlGateway/Metrics/GatewayMetrics.cs
Original file line number Diff line number Diff line change
@@ -1,6 +1,8 @@
using System.Diagnostics.Metrics;
using OpenShock.Common.Metrics;
using OpenShock.Common.OpenShockDb;
using OpenShock.Common.Services.Geo;
using OpenShock.LiveControlGateway.Controllers;
using OpenShock.LiveControlGateway.Options;

namespace OpenShock.LiveControlGateway.Metrics;
Expand All @@ -14,6 +16,11 @@ namespace OpenShock.LiveControlGateway.Metrics;
/// Every measurement carries <c>gateway_fqdn</c> as a normal tag rather than a tag on the
/// <see cref="Meter"/> itself - a meter-level tag is a *scope* attribute, which the Prometheus
/// exporter emits prefixed as <c>otel_scope_gateway_fqdn</c>.
/// <para>
/// Hub instruments also carry the hub's network as <c>asn</c> and <c>asn_org</c>, from GeoIP. The
/// org is redundant with the number but saves every dashboard a lookup table; it adds no series
/// of its own. Both are <see cref="Unknown"/> when GeoIP is not configured or the lookup missed.
/// </para>
/// </remarks>
public sealed class GatewayMetrics
{
Expand Down Expand Up @@ -71,6 +78,9 @@ public static class FrameOutcome
public const string ShockerExclusive = "shocker_exclusive";
}

/// <summary>Tag value for a hub whose network could not be resolved.</summary>
public const string Unknown = "unknown";

private readonly KeyValuePair<string, object?> _gatewayFqdn;

private readonly Counter<long> _hubConnectAttempts;
Expand Down Expand Up @@ -117,13 +127,29 @@ public GatewayMetrics(LcgOptions lcgOptions, [FromKeyedServices("OpenShock.Gatew
/// Record the outcome of a hub connection attempt.
/// </summary>
/// <param name="outcome">One of <see cref="HubConnectOutcome"/></param>
public void HubConnectAttempt(string outcome) =>
_hubConnectAttempts.Add(1, _gatewayFqdn, new KeyValuePair<string, object?>("outcome", outcome));
/// <param name="hub">The connecting hub</param>
public void HubConnectAttempt(string outcome, IHubController hub) =>
_hubConnectAttempts.Add(1, _gatewayFqdn, new KeyValuePair<string, object?>("outcome", outcome),
AsnTag(hub.IpInfo), AsnOrgTag(hub.IpInfo));

/// <summary>
/// Record a hub lifetime being torn down.
/// </summary>
public void HubDisconnected() => _hubDisconnections.Add(1, _gatewayFqdn);
/// <param name="hub">The hub that disconnected</param>
public void HubDisconnected(IHubController hub) =>
_hubDisconnections.Add(1, _gatewayFqdn, AsnTag(hub.IpInfo), AsnOrgTag(hub.IpInfo));

/// <summary>
/// The <c>asn</c> tag for a hub's network.
/// </summary>
public static KeyValuePair<string, object?> AsnTag(IpEnrichmentData? ipInfo) =>
new("asn", ipInfo?.Asn?.ToString() ?? Unknown);

/// <summary>
/// The <c>asn_org</c> tag for a hub's network.
/// </summary>
public static KeyValuePair<string, object?> AsnOrgTag(IpEnrichmentData? ipInfo) =>
new("asn_org", ipInfo?.AsnOrg ?? Unknown);

/// <summary>
/// Record the outcome of a live control session attempt. This is churn only - the number of
Expand Down
1 change: 1 addition & 0 deletions LiveControlGateway/Program.cs
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,7 @@
var redisOptions = builder.RegisterRedisOptions();
var databaseOptions = builder.RegisterDatabaseOptions();
builder.RegisterMetricsOptions();
builder.RegisterGeoOptions();

var lcgOptions = builder.Configuration.GetRequiredSection(LcgOptions.SectionName).Get<LcgOptions>();
if (lcgOptions is null)
Expand Down
Loading