fix(SEC-31,SEC-32): re-partition the API-key failure limiter on (peer, key id) with probe admission

The gRPC auth failure limiter partitioned on the key id parsed out of the
*unauthenticated* token and rejected with ResourceExhausted before VerifyAsync
ran. Key ids are not secret — they ride in every token and are listed on the
dashboard — so any network peer could send 10 garbage-secret requests per minute
and deny that key indefinitely: the legitimate holder's correct secret was
refused before it was ever checked, and the success-path Reset that would clear
the block sat behind the verification the block prevented (SEC-31). The tracked
map was also flushable — any `a_b_c`-shaped junk minted a fresh partition (the
`mxgw` literal was never compared), so ~4096 throwaway tokens evicted a blocked
entry and reset the window (SEC-32).

ApiKeyFailureLimiter moves from IsBlocked/RecordFailure/Reset(string peer) to a
partition-pair API: Check/RecordFailure/Reset(ApiKeyThrottlePartition) with an
ApiKeyThrottleDecision result. Two layers share one sliding window — a composite
(transport peer, key id) partition at ApiKeyFailureLimit, and a per-key-id
aggregate across all peers at the new ApiKeyFailureAggregateLimit (default 30)
that bounds a source-rotating sprayer. An over-limit state is now a valve rather
than a wall: one request per the new ApiKeyFailureProbeIntervalSeconds (default
5) is admitted through to the real verifier, so the correct secret always reaches
the constant-time compare and resets both layers. Guarantees preserved: guessing
stays bounded per window, and the failure path still spends no store read per
attempt.

SEC-32 rides the same change set: the interceptor validates token shape (literal
`mxgw` prefix, >= 3 non-empty `_` segments, key id <= 64 chars) before minting a
key-id partition, each transport peer may mint at most 32 of them before the
overflow collapses onto its fallback partition, and eviction prefers fully
expired windows and never drops an over-limit partition below a 2x transient
overshoot ceiling. Throttled attempts increment mxgateway.auth.throttled, tagged
stage=peer|aggregate only — /metrics is unauthenticated (open SEC-14), so no key
material may appear there.

Docs in the same commit: GatewayConfiguration limiter rows plus the two new keys,
the Authentication hot-path paragraph, the Authorization SEC-11 section, and the
limiter / SecurityOptions XML remarks (the old NAT rationale described the
defective keying). Tracking rows flipped to Done with a change-log entry.

Tests: new ApiKeyFailureLimiterTests (11) covering window pruning, composite vs
aggregate trip points, probe cadence, absolute-block mode, reset across both
layers, junk-spray eviction resistance, the per-peer cap, and expired-window
eviction preference; GatewayGrpcAuthorizationInterceptorTests gains the four
SEC-31 contract tests plus NonMxgwToken_FallsBackToTransportPeerPartition (20
total); GatewayOptionsValidatorTests covers both new keys including 0 as a
supported disable value (66 total).
This commit is contained in:
Joseph Doherty
2026-08-07 05:39:10 -04:00
parent ead921cace
commit df710e18a9
16 changed files with 1148 additions and 137 deletions
@@ -87,6 +87,18 @@ public sealed class GatewayOptionsValidator : OptionsValidatorBase<GatewayOption
options.ApiKeyFailureTrackedPeers,
"MxGateway:Security:ApiKeyFailureTrackedPeers must be greater than zero.",
builder);
// The two-layer limiter knobs (SEC-31) accept 0 as "disable this layer": a zero aggregate
// limit turns off cross-peer counting, and a zero probe interval restores absolute blocking.
// Negatives express no intent.
AddIfNegative(
options.ApiKeyFailureAggregateLimit,
"MxGateway:Security:ApiKeyFailureAggregateLimit must be greater than or equal to zero (0 disables the per-key aggregate layer).",
builder);
AddIfNegative(
options.ApiKeyFailureProbeIntervalSeconds,
"MxGateway:Security:ApiKeyFailureProbeIntervalSeconds must be greater than or equal to zero (0 blocks absolutely instead of admitting probes).",
builder);
}
private static void ValidateAuthentication(AuthenticationOptions options, ValidationBuilder builder)
@@ -43,23 +43,49 @@ public sealed class SecurityOptions
public int LoginRateLimitWindowSeconds { get; init; } = 60;
/// <summary>
/// Gets the number of consecutive failed API-key verifications, per peer, within
/// <see cref="ApiKeyFailureWindowSeconds"/> that trips the in-process short-circuit. Once tripped,
/// the gRPC auth path rejects further attempts before the store read; a successful verification
/// resets the peer's counter. Default is 10.
/// Gets the number of failed API-key verifications, per <c>(transport peer, key id)</c>
/// partition, within <see cref="ApiKeyFailureWindowSeconds"/> that trips the in-process
/// short-circuit. Once tripped, the gRPC auth path rejects further attempts from that partition
/// before the store read, except for the probe admitted every
/// <see cref="ApiKeyFailureProbeIntervalSeconds"/>; a successful verification resets the
/// partition. Default is 10.
/// </summary>
public int ApiKeyFailureLimit { get; init; } = 10;
/// <summary>
/// Gets the sliding-window length, in seconds, over which API-key verification failures are
/// counted per peer. Default is 60 seconds.
/// counted (for both the per-partition and the per-key-id aggregate layer). Default is 60
/// seconds.
/// </summary>
public int ApiKeyFailureWindowSeconds { get; init; } = 60;
/// <summary>
/// Gets the maximum number of distinct peers tracked by the API-key failure counter. The counter
/// is a bounded LRU so a spray of unique peer keys cannot grow memory without limit. Default is
/// 4096.
/// Gets the number of failed API-key verifications for one key id, counted across <em>all</em>
/// transport peers within <see cref="ApiKeyFailureWindowSeconds"/>, that puts the key id into
/// probe mode. This second layer bounds a distributed or source-rotating sprayer that never
/// trips any single <c>(peer, key id)</c> partition. Set to <c>0</c> to disable the aggregate
/// layer. Default is 30.
/// </summary>
public int ApiKeyFailureAggregateLimit { get; init; } = 30;
/// <summary>
/// Gets the minimum interval, in seconds, between probe admissions for an over-limit partition
/// or key-id aggregate. An over-limit state is a valve rather than a wall: at most one request
/// per interval is admitted through to the real verifier, so the holder of the correct secret
/// can always reach the constant-time compare (and reset the state) while an attacker is
/// spraying. Set to <c>0</c> to block absolutely instead — not recommended, because an
/// unauthenticated peer can then deny the key to its legitimate holder for the whole window.
/// Default is 5 seconds.
/// </summary>
public int ApiKeyFailureProbeIntervalSeconds { get; init; } = 5;
/// <summary>
/// Gets the maximum number of distinct partitions tracked by the API-key failure counter. The
/// counter is a bounded LRU so a spray of unique tokens cannot grow memory without limit, and —
/// since only a validly shaped token mints a key-id partition, capped per transport peer —
/// cannot be flushed to clear an active block either: eviction prefers fully expired windows and
/// never removes a partition that is currently over its limit, up to a transient overshoot
/// ceiling of twice this value. Default is 4096.
/// </summary>
public int ApiKeyFailureTrackedPeers { get; init; } = 4096;
}
@@ -24,6 +24,7 @@ public sealed class GatewayMetrics : IDisposable
private readonly Counter<long> _streamDisconnectsCounter;
private readonly Counter<long> _retryAttemptsCounter;
private readonly Counter<long> _alarmProviderSwitchesCounter;
private readonly Counter<long> _authThrottledCounter;
private readonly Histogram<double> _workerStartupLatencyHistogram;
private readonly Histogram<double> _commandLatencyHistogram;
private readonly Histogram<double> _eventStreamSendLatencyHistogram;
@@ -80,6 +81,7 @@ public sealed class GatewayMetrics : IDisposable
_streamDisconnectsCounter = _meter.CreateCounter<long>("mxgateway.grpc.streams.disconnected");
_retryAttemptsCounter = _meter.CreateCounter<long>("mxgateway.retries.attempted");
_alarmProviderSwitchesCounter = _meter.CreateCounter<long>("mxgateway.alarms.provider_switches");
_authThrottledCounter = _meter.CreateCounter<long>("mxgateway.auth.throttled");
_workerStartupLatencyHistogram = _meter.CreateHistogram<double>("mxgateway.workers.startup.duration", "s");
_commandLatencyHistogram = _meter.CreateHistogram<double>("mxgateway.commands.duration", "s");
_eventStreamSendLatencyHistogram = _meter.CreateHistogram<double>("mxgateway.events.stream_send.duration", "s");
@@ -342,6 +344,19 @@ public sealed class GatewayMetrics : IDisposable
_queueOverflowsCounter.Add(1, new KeyValuePair<string, object?>("queue", queueName));
}
/// <summary>
/// Records that an API-key authentication attempt was short-circuited by the failure limiter
/// before verification.
/// </summary>
/// <param name="stage">Limiter layer that refused the attempt: <c>peer</c> or <c>aggregate</c>.</param>
public void RecordAuthThrottled(string stage)
{
// Tagged by stage only. Neither the key id nor the peer address may become a tag: key ids
// are credential material and peer addresses are unbounded cardinality, and /metrics is not
// authenticated (open SEC-14), so anything tagged here is world-readable.
_authThrottledCounter.Add(1, new KeyValuePair<string, object?>("stage", stage));
}
/// <summary>
/// Records that a fault occurred in the given category.
/// </summary>
@@ -4,34 +4,62 @@ using ZB.MOM.WW.MxGateway.Server.Configuration;
namespace ZB.MOM.WW.MxGateway.Server.Security.Authorization;
/// <summary>
/// Cheap, in-process per-peer sliding-window failure counter for the gRPC auth path. It is
/// checked BEFORE the API-key verification store read and short-circuits a peer that has exceeded
/// <see cref="SecurityOptions.ApiKeyFailureLimit"/> failed attempts within
/// <see cref="SecurityOptions.ApiKeyFailureWindowSeconds"/>; a successful verification resets the
/// peer's counter.
/// Cheap, in-process sliding-window failure counter for the gRPC auth path. It is checked BEFORE the
/// API-key verification store read, so online secret guessing cannot spend a SQLite read per attempt,
/// and it admits a periodic probe so a throttle can never deny a key to the holder of the correct
/// secret.
/// </summary>
/// <remarks>
/// <para>
/// The peer key is the API key id when the presented token parses, falling back to the transport
/// peer address otherwise. Keying on key id (per the NAT caveat in <c>glauth.md</c>) means a single
/// abusive credential behind a shared NAT is throttled without locking out unrelated clients on the
/// same address.
/// Two layers share one window mechanism. <b>Layer 1</b> partitions on the composite
/// <c>(transport peer, key id)</c>: an attacker's failures bind to the address that produced them,
/// so a peer spraying guesses at a key id it does not own throttles only itself, while co-located
/// clients using other key ids are untouched. <b>Layer 2</b> counts failures for a key id across all
/// peers (<see cref="SecurityOptions.ApiKeyFailureAggregateLimit"/>), which bounds a distributed or
/// source-rotating sprayer that never trips any single partition.
/// </para>
/// <para>
/// The tracked-peer set is a bounded LRU (<see cref="SecurityOptions.ApiKeyFailureTrackedPeers"/>) so
/// a spray of unique peer keys cannot grow memory without limit. Successful peers are removed on
/// reset, so in steady state the map holds only peers with recent failures — the common success path
/// is a lock-free dictionary miss.
/// An over-limit state is a valve, not a wall (SEC-31): at most one request per
/// <see cref="SecurityOptions.ApiKeyFailureProbeIntervalSeconds"/> is admitted through to the real
/// verifier, and a successful verification resets both layers. The reset path therefore stays
/// reachable while throttled — the earlier design keyed solely on the attacker-supplied key id and
/// rejected before verification, so the legitimate holder could never clear the block.
/// </para>
/// <para>
/// The tracked partition set is a bounded LRU (<see cref="SecurityOptions.ApiKeyFailureTrackedPeers"/>).
/// Because only a validly shaped token mints a key-id partition and each transport peer may mint at
/// most <see cref="MaxKeyIdPartitionsPerPeer"/> of them (the overflow collapses onto that address's
/// fallback partition), the map bounds memory <em>and</em> resists a junk-token flush (SEC-32):
/// eviction prefers fully expired windows, never removes an over-limit partition below a transient
/// overshoot ceiling of twice the cap, and only then falls back to least-recently-active. Successful
/// peers are removed on reset, so in steady state the map holds only partitions with recent failures
/// — the common success path is a lock-free dictionary miss.
/// </para>
/// </remarks>
public sealed class ApiKeyFailureLimiter
{
/// <summary>
/// Maximum distinct key-id partitions one transport peer may mint. Not configurable: no
/// legitimate address fails against dozens of distinct keys inside one window, and the cap is
/// what bounds an attacker's total partitions to this value plus its fallback partition.
/// </summary>
internal const int MaxKeyIdPartitionsPerPeer = 32;
// NUL appears in neither a gRPC peer string ("ipv4:10.0.0.1:5000") nor an ASCII metadata
// value, so a composite key can never collide with a transport-peer-only fallback key.
// A collision would merge two partitions, never widen an admission.
private const char PartitionSeparator = '\0';
private readonly int _limit;
private readonly int _aggregateLimit;
private readonly long _windowTicks;
private readonly int _maxPeers;
private readonly long _probeIntervalTicks;
private readonly int _maxPartitions;
private readonly TimeProvider _clock;
private readonly ConcurrentDictionary<string, PeerState> _peers = new(StringComparer.Ordinal);
private readonly ConcurrentDictionary<string, WindowState> _partitions = new(StringComparer.Ordinal);
private readonly ConcurrentDictionary<string, WindowState> _aggregates = new(StringComparer.Ordinal);
private readonly ConcurrentDictionary<string, PeerKeyIds> _peerKeyIds = new(StringComparer.Ordinal);
/// <summary>Initializes a new instance of the <see cref="ApiKeyFailureLimiter"/> class.</summary>
/// <param name="security">Security options carrying the failure-limit knobs.</param>
@@ -41,79 +69,219 @@ public sealed class ApiKeyFailureLimiter
(security ?? throw new ArgumentNullException(nameof(security))).ApiKeyFailureLimit,
TimeSpan.FromSeconds(security.ApiKeyFailureWindowSeconds),
security.ApiKeyFailureTrackedPeers,
security.ApiKeyFailureAggregateLimit,
TimeSpan.FromSeconds(security.ApiKeyFailureProbeIntervalSeconds),
clock)
{
}
/// <summary>Initializes a new instance of the <see cref="ApiKeyFailureLimiter"/> class. Test/explicit seam.</summary>
/// <param name="limit">The maximum number of failures allowed within <paramref name="window"/>.</param>
/// <param name="limit">The maximum number of failures allowed per partition within <paramref name="window"/>.</param>
/// <param name="window">The sliding window over which failures are counted.</param>
/// <param name="maxPeers">The maximum number of tracked peers before least-recently-active eviction kicks in.</param>
/// <param name="maxPartitions">The maximum number of tracked partitions before eviction kicks in.</param>
/// <param name="aggregateLimit">The per-key-id cross-peer failure limit; <c>0</c> disables the aggregate layer.</param>
/// <param name="probeInterval">The minimum interval between probe admissions; <see cref="TimeSpan.Zero"/> blocks absolutely.</param>
/// <param name="clock">The time provider.</param>
internal ApiKeyFailureLimiter(int limit, TimeSpan window, int maxPeers, TimeProvider clock)
internal ApiKeyFailureLimiter(
int limit,
TimeSpan window,
int maxPartitions,
int aggregateLimit,
TimeSpan probeInterval,
TimeProvider clock)
{
ArgumentNullException.ThrowIfNull(clock);
_limit = limit;
_windowTicks = window.Ticks;
_maxPeers = maxPeers;
_maxPartitions = maxPartitions;
_aggregateLimit = aggregateLimit;
_probeIntervalTicks = probeInterval.Ticks;
_clock = clock;
}
/// <summary>Returns whether the peer has reached the failure limit within the current window.</summary>
/// <param name="peer">The peer key (key id or peer address).</param>
/// <returns><see langword="true"/> when the peer should be short-circuited.</returns>
public bool IsBlocked(string peer)
/// <summary>Gets the number of tracked <c>(peer, key id)</c> partitions. Test seam.</summary>
internal int TrackedPartitionCount => _partitions.Count;
/// <summary>Gets the number of tracked per-key-id aggregates. Test seam.</summary>
internal int TrackedAggregateCount => _aggregates.Count;
/// <summary>Decides whether an authentication attempt may reach the verifier.</summary>
/// <param name="partition">The throttle partition derived from the request.</param>
/// <returns>The admission decision for this attempt.</returns>
public ApiKeyThrottleDecision Check(ApiKeyThrottlePartition partition)
{
ArgumentNullException.ThrowIfNull(peer);
string peer = RequirePeer(partition);
if (_limit <= 0)
{
return false;
}
if (!_peers.TryGetValue(peer, out PeerState? state))
{
return false;
return ApiKeyThrottleDecision.Allowed;
}
long now = _clock.GetUtcNow().UtcTicks;
lock (state)
(string partitionKey, string? effectiveKeyId) = ResolvePartitionKey(peer, partition.KeyId, mint: false);
WindowState? peerState = _partitions.TryGetValue(partitionKey, out WindowState? tracked) ? tracked : null;
WindowState? aggregateState = null;
if (effectiveKeyId is not null
&& _aggregateLimit > 0
&& _aggregates.TryGetValue(effectiveKeyId, out WindowState? aggregate))
{
Prune(state, now);
return state.FailureTicks.Count >= _limit;
aggregateState = aggregate;
}
bool peerOver = peerState is not null && IsOverLimit(peerState, now, _limit);
bool aggregateOver = aggregateState is not null && IsOverLimit(aggregateState, now, _aggregateLimit);
if (!peerOver && !aggregateOver)
{
return ApiKeyThrottleDecision.Allowed;
}
if (_probeIntervalTicks <= 0)
{
return peerOver ? ApiKeyThrottleDecision.ThrottledByPeer : ApiKeyThrottleDecision.ThrottledByAggregate;
}
// Both over-limit layers must have a slot before either is consumed, so a request cannot burn
// the peer's probe and then be refused by the aggregate.
if (peerOver && !IsProbeDue(peerState!, now))
{
return ApiKeyThrottleDecision.ThrottledByPeer;
}
if (aggregateOver && !IsProbeDue(aggregateState!, now))
{
return ApiKeyThrottleDecision.ThrottledByAggregate;
}
if (peerOver)
{
ConsumeProbe(peerState!, now);
}
if (aggregateOver)
{
ConsumeProbe(aggregateState!, now);
}
return ApiKeyThrottleDecision.ProbeAdmitted;
}
/// <summary>Records a failed verification attempt for the peer.</summary>
/// <param name="peer">The peer key (key id or peer address).</param>
public void RecordFailure(string peer)
/// <summary>Records a failed verification attempt against both limiter layers.</summary>
/// <param name="partition">The throttle partition derived from the request.</param>
public void RecordFailure(ApiKeyThrottlePartition partition)
{
ArgumentNullException.ThrowIfNull(peer);
string peer = RequirePeer(partition);
if (_limit <= 0)
{
return;
}
long now = _clock.GetUtcNow().UtcTicks;
PeerState state = _peers.GetOrAdd(peer, static _ => new PeerState());
(string partitionKey, string? effectiveKeyId) = ResolvePartitionKey(peer, partition.KeyId, mint: true);
WindowState state = _partitions.GetOrAdd(partitionKey, _ => new WindowState(peer, effectiveKeyId));
RecordInto(state, now, _limit);
// Only a key id that earned its own partition feeds the aggregate: an id squeezed out by the
// per-peer cap is spray, and letting it through would make aggregate cardinality unbounded.
if (effectiveKeyId is not null && _aggregateLimit > 0)
{
WindowState aggregate = _aggregates.GetOrAdd(effectiveKeyId, static _ => new WindowState(null, null));
RecordInto(aggregate, now, _aggregateLimit);
}
EvictIfOverCapacity(_partitions, now, _limit, releaseKeyIds: true);
EvictIfOverCapacity(_aggregates, now, _aggregateLimit, releaseKeyIds: false);
}
/// <summary>Clears both limiter layers for the partition after a successful verification.</summary>
/// <param name="partition">The throttle partition derived from the request.</param>
public void Reset(ApiKeyThrottlePartition partition)
{
string peer = RequirePeer(partition);
(string partitionKey, _) = ResolvePartitionKey(peer, partition.KeyId, mint: false);
RemovePartition(partitionKey);
// The aggregate is cleared on the presented key id even when the composite partition
// collapsed to the fallback: a verified secret is proof the key is not under successful
// attack, and clearing is what makes the block recoverable.
if (partition.KeyId is not null)
{
_aggregates.TryRemove(partition.KeyId, out _);
}
}
/// <summary>Returns whether the partition currently has tracked failure state. Test seam.</summary>
/// <param name="partition">The throttle partition to look up.</param>
/// <returns><see langword="true"/> when the partition has a tracked window.</returns>
internal bool IsTracked(ApiKeyThrottlePartition partition)
{
string peer = RequirePeer(partition);
(string partitionKey, _) = ResolvePartitionKey(peer, partition.KeyId, mint: false);
return _partitions.ContainsKey(partitionKey);
}
private static string RequirePeer(ApiKeyThrottlePartition partition)
{
if (string.IsNullOrEmpty(partition.TransportPeer))
{
throw new ArgumentException("A throttle partition requires a transport peer.", nameof(partition));
}
return partition.TransportPeer;
}
private static string Composite(string peer, string keyId) => peer + PartitionSeparator + keyId;
private void RecordInto(WindowState state, long now, int limit)
{
lock (state)
{
Prune(state, now);
state.FailureTicks.Enqueue(now);
state.LastActivityTicks = now;
// Arm (or push out) the probe slot whenever the state is at or over its limit, so the
// attempt that trips the limit is not itself followed by an immediate free probe.
if (limit > 0 && state.FailureTicks.Count >= limit)
{
state.NextProbeAtTicks = now + _probeIntervalTicks;
}
}
}
private bool IsOverLimit(WindowState state, long now, int limit)
{
if (limit <= 0)
{
return false;
}
EvictIfOverCapacity();
lock (state)
{
Prune(state, now);
return state.FailureTicks.Count >= limit;
}
}
/// <summary>Clears the peer's failure count after a successful verification.</summary>
/// <param name="peer">The peer key (key id or peer address).</param>
public void Reset(string peer)
private bool IsProbeDue(WindowState state, long now)
{
ArgumentNullException.ThrowIfNull(peer);
_peers.TryRemove(peer, out _);
lock (state)
{
return now >= state.NextProbeAtTicks;
}
}
private void Prune(PeerState state, long now)
private void ConsumeProbe(WindowState state, long now)
{
lock (state)
{
state.NextProbeAtTicks = now + _probeIntervalTicks;
state.LastActivityTicks = now;
}
}
private void Prune(WindowState state, long now)
{
while (state.FailureTicks.Count > 0 && now - state.FailureTicks.Peek() >= _windowTicks)
{
@@ -121,37 +289,186 @@ public sealed class ApiKeyFailureLimiter
}
}
private void EvictIfOverCapacity()
/// <summary>
/// Maps a partition onto its storage key, enforcing the per-peer key-id cap. Returns the
/// transport-peer fallback key (and a null effective key id) when the token carried no key id or
/// when this address has already minted <see cref="MaxKeyIdPartitionsPerPeer"/> of them.
/// </summary>
private (string PartitionKey, string? EffectiveKeyId) ResolvePartitionKey(string peer, string? keyId, bool mint)
{
// Best-effort eviction: only runs when the map exceeds the cap (rare, since only peers with
// recent failures are tracked). Removes the least-recently-active peer. Racy under
// concurrency, which is acceptable for a bound rather than an exact policy.
while (_peers.Count > _maxPeers)
if (keyId is null)
{
string? oldest = null;
long oldestTicks = long.MaxValue;
foreach (KeyValuePair<string, PeerState> entry in _peers)
return (peer, null);
}
while (true)
{
if (!_peerKeyIds.TryGetValue(peer, out PeerKeyIds? keyIds))
{
long activity = Volatile.Read(ref entry.Value.LastActivityTicks);
if (activity < oldestTicks)
if (!mint)
{
oldestTicks = activity;
oldest = entry.Key;
// Nothing tracked for this address yet, so the composite lookup simply misses.
return (Composite(peer, keyId), keyId);
}
keyIds = _peerKeyIds.GetOrAdd(peer, static _ => new PeerKeyIds());
}
if (oldest is null || !_peers.TryRemove(oldest, out _))
lock (keyIds)
{
break;
if (keyIds.Removed)
{
// Lost a race with cleanup; re-read the dictionary and try again.
continue;
}
if (keyIds.KeyIds.Contains(keyId))
{
return (Composite(peer, keyId), keyId);
}
if (keyIds.KeyIds.Count >= MaxKeyIdPartitionsPerPeer)
{
return (peer, null);
}
if (mint)
{
keyIds.KeyIds.Add(keyId);
}
return (Composite(peer, keyId), keyId);
}
}
}
private sealed class PeerState
private bool RemovePartition(string partitionKey)
{
if (!_partitions.TryRemove(partitionKey, out WindowState? state))
{
return false;
}
if (state.TransportPeer is not null && state.KeyId is not null)
{
ReleaseKeyId(state.TransportPeer, state.KeyId);
}
return true;
}
private void ReleaseKeyId(string peer, string keyId)
{
if (!_peerKeyIds.TryGetValue(peer, out PeerKeyIds? keyIds))
{
return;
}
lock (keyIds)
{
keyIds.KeyIds.Remove(keyId);
if (keyIds.KeyIds.Count == 0)
{
keyIds.Removed = true;
_peerKeyIds.TryRemove(peer, out _);
}
}
}
/// <summary>
/// Best-effort eviction, run only when a map exceeds the cap. Preference order: fully expired
/// windows, then least-recently-active entries that are still under their limit, and only above
/// the 2x transient-overshoot ceiling the oldest over-limit entry. Racy under concurrency, which
/// is acceptable for a bound rather than an exact policy.
/// </summary>
private void EvictIfOverCapacity(
ConcurrentDictionary<string, WindowState> map,
long now,
int limit,
bool releaseKeyIds)
{
while (map.Count > _maxPartitions)
{
bool overHardCeiling = map.Count > (long)_maxPartitions * 2;
string? expired = null;
string? leastRecent = null;
string? oldestOverLimit = null;
long leastRecentTicks = long.MaxValue;
long oldestOverLimitTicks = long.MaxValue;
foreach (KeyValuePair<string, WindowState> entry in map)
{
int count;
long activity;
lock (entry.Value)
{
Prune(entry.Value, now);
count = entry.Value.FailureTicks.Count;
activity = entry.Value.LastActivityTicks;
}
if (count == 0)
{
expired = entry.Key;
break;
}
if (limit > 0 && count >= limit)
{
if (activity < oldestOverLimitTicks)
{
oldestOverLimitTicks = activity;
oldestOverLimit = entry.Key;
}
continue;
}
if (activity < leastRecentTicks)
{
leastRecentTicks = activity;
leastRecent = entry.Key;
}
}
string? victim = expired ?? leastRecent ?? (overHardCeiling ? oldestOverLimit : null);
if (victim is null)
{
// Everything left is load-bearing and the ceiling is not breached: accept the
// documented transient overshoot rather than clearing an active block.
return;
}
bool removed = releaseKeyIds ? RemovePartition(victim) : map.TryRemove(victim, out _);
if (!removed)
{
return;
}
}
}
private sealed class WindowState(string? transportPeer, string? keyId)
{
/// <summary>Gets the transport peer this partition belongs to; null for key-id aggregates.</summary>
public string? TransportPeer { get; } = transportPeer;
/// <summary>Gets the key id this partition belongs to; null for fallback partitions and aggregates.</summary>
public string? KeyId { get; } = keyId;
/// <summary>Timestamps (in ticks) of failures still within the sliding window.</summary>
public Queue<long> FailureTicks { get; } = new();
public long LastActivityTicks;
public long NextProbeAtTicks;
}
private sealed class PeerKeyIds
{
/// <summary>Key ids this transport peer has minted a partition for.</summary>
public HashSet<string> KeyIds { get; } = new(StringComparer.Ordinal);
/// <summary>Set once the entry has been dropped from the peer map, so racing minters retry.</summary>
public bool Removed { get; set; }
}
}
@@ -0,0 +1,28 @@
namespace ZB.MOM.WW.MxGateway.Server.Security.Authorization;
/// <summary>
/// Outcome of an <see cref="ApiKeyFailureLimiter"/> admission check for one authentication attempt.
/// </summary>
public enum ApiKeyThrottleDecision
{
/// <summary>No layer is over its limit; the request proceeds to verification normally.</summary>
Allowed = 0,
/// <summary>
/// At least one layer is over its limit, but this request took the probe slot for the current
/// interval and proceeds to verification. A success resets every layer it passed.
/// </summary>
ProbeAdmitted = 1,
/// <summary>
/// The <c>(transport peer, key id)</c> partition is over its limit and no probe slot is
/// available; the request is rejected before the verification store read.
/// </summary>
ThrottledByPeer = 2,
/// <summary>
/// The cross-peer aggregate for the presented key id is over its limit and no probe slot is
/// available; the request is rejected before the verification store read.
/// </summary>
ThrottledByAggregate = 3,
}
@@ -0,0 +1,16 @@
namespace ZB.MOM.WW.MxGateway.Server.Security.Authorization;
/// <summary>
/// Partition identity for the API-key failure limiter: the transport peer that sent the request and,
/// when the presented token is validly shaped, the key id it claims.
/// </summary>
/// <remarks>
/// Both halves come from unauthenticated input, which is why neither may stand alone. The transport
/// peer is always present, so a throttle can never outlive the address that earned it; the key id is
/// present only after the shape check in
/// <see cref="GatewayGrpcAuthorizationInterceptor"/> (literal <c>mxgw</c> prefix, at least three
/// non-empty underscore segments, bounded key-id length) so junk tokens cannot mint partitions.
/// </remarks>
/// <param name="TransportPeer">The gRPC transport peer address (<c>ServerCallContext.Peer</c>).</param>
/// <param name="KeyId">The key id claimed by a validly shaped token, or <see langword="null"/>.</param>
public readonly record struct ApiKeyThrottlePartition(string TransportPeer, string? KeyId);
@@ -3,6 +3,7 @@ using Grpc.Core.Interceptors;
using Microsoft.Extensions.Options;
using ZB.MOM.WW.Auth.Abstractions.ApiKeys;
using ZB.MOM.WW.MxGateway.Server.Configuration;
using ZB.MOM.WW.MxGateway.Server.Metrics;
using ZB.MOM.WW.MxGateway.Server.Security.Authentication;
// The handler pushes the gateway's constraint-bearing identity; alias away the shared library's
@@ -16,8 +17,13 @@ public sealed class GatewayGrpcAuthorizationInterceptor(
GatewayGrpcScopeResolver scopeResolver,
IGatewayRequestIdentityAccessor identityAccessor,
IOptions<GatewayOptions> options,
ApiKeyFailureLimiter failureLimiter) : Interceptor
ApiKeyFailureLimiter failureLimiter,
GatewayMetrics metrics) : Interceptor
{
// Generated key ids are far shorter; the cap only exists to stop an invented "key id" of
// arbitrary length from becoming a limiter partition.
private const int MaxKeyIdLength = 64;
/// <inheritdoc />
public override async Task<TResponse> UnaryServerHandler<TRequest, TResponse>(
TRequest request,
@@ -64,14 +70,19 @@ public sealed class GatewayGrpcAuthorizationInterceptor(
string? authorizationHeader = context.RequestHeaders.GetValue("authorization");
// Short-circuit a peer that has already failed too many times inside the sliding
// window BEFORE the verification store read, so online guessing cannot spend a SQLite read
// per attempt. The peer key prefers the presented key id over the transport address (NAT
// caveat). ResourceExhausted signals throttling without revealing whether any particular
// secret was valid.
string peerKey = ResolvePeerKey(authorizationHeader, context);
if (failureLimiter.IsBlocked(peerKey))
// Short-circuit a partition that has already failed too many times inside the sliding window
// BEFORE the verification store read, so online guessing cannot spend a SQLite read per
// attempt. The partition is the composite (transport peer, key id) — never the key id alone,
// which is public and would let any peer deny a key it does not own — plus a cross-peer
// aggregate for the key id. An over-limit state still admits one probe per interval, so the
// holder of the correct secret always reaches the verifier and resets the state.
// ResourceExhausted signals throttling without revealing whether any secret was valid.
ApiKeyThrottlePartition throttlePartition = ResolveThrottlePartition(authorizationHeader, context);
ApiKeyThrottleDecision decision = failureLimiter.Check(throttlePartition);
if (decision is ApiKeyThrottleDecision.ThrottledByPeer or ApiKeyThrottleDecision.ThrottledByAggregate)
{
metrics.RecordAuthThrottled(
decision == ApiKeyThrottleDecision.ThrottledByPeer ? "peer" : "aggregate");
throw new RpcException(new Status(
StatusCode.ResourceExhausted,
"Too many authentication attempts. Try again later."));
@@ -86,15 +97,17 @@ public sealed class GatewayGrpcAuthorizationInterceptor(
if (!verification.Succeeded || verification.Identity is null)
{
failureLimiter.RecordFailure(peerKey);
failureLimiter.RecordFailure(throttlePartition);
throw new RpcException(new Status(
StatusCode.Unauthenticated,
"Missing or invalid API key."));
}
// Successful authentication clears the peer's failure count so a legitimate client that
// fat-fingered a few attempts is not penalised once it recovers.
failureLimiter.Reset(peerKey);
// Successful authentication clears both limiter layers, so a legitimate client that
// fat-fingered a few attempts is not penalised once it recovers — and, because the check
// above admits a probe rather than blocking absolutely, this reset stays reachable while the
// key is under an active spray.
failureLimiter.Reset(throttlePartition);
ApiKeyIdentity identity = GatewayApiKeyIdentityMapper.ToGatewayIdentity(verification.Identity);
@@ -109,29 +122,50 @@ public sealed class GatewayGrpcAuthorizationInterceptor(
return identity;
}
// Resolves the failure-limiter partition key: the presented key id (token shape
// mxgw_<keyId>_<secret>) when the header parses, otherwise the transport peer address. Keying on
// key id throttles a single abusive credential without locking out co-located clients behind NAT.
private static string ResolvePeerKey(string? authorizationHeader, ServerCallContext context)
// Resolves the failure-limiter partition: always the transport peer, plus the presented key id
// when — and only when the token is validly shaped. Both halves are unauthenticated input, so
// the peer half is what keeps a throttle bound to the address that earned it.
private static ApiKeyThrottlePartition ResolveThrottlePartition(
string? authorizationHeader,
ServerCallContext context)
{
if (!string.IsNullOrWhiteSpace(authorizationHeader))
{
ReadOnlySpan<char> header = authorizationHeader.AsSpan().Trim();
const string bearer = "Bearer ";
ReadOnlySpan<char> token = header.StartsWith(bearer, StringComparison.OrdinalIgnoreCase)
? header[bearer.Length..].Trim()
: header;
return new ApiKeyThrottlePartition(context.Peer, TryResolveKeyId(authorizationHeader));
}
// mxgw_<keyId>_<secret>: the key id is the second underscore-delimited segment. The
// secret may itself contain underscores, but the key id is unaffected.
string tokenText = token.ToString();
string[] parts = tokenText.Split('_');
if (parts.Length >= 3 && parts[0].Length > 0 && parts[1].Length > 0)
{
return "key:" + parts[1];
}
// Token shape mxgw_<keyId>_<secret>: the key id is the second underscore-delimited segment (the
// secret may itself contain underscores, which does not affect the key id). The shape is checked
// before a key-id partition is minted so a spray of invented tokens cannot mint one tracked
// partition each and flush the limiter's bounded map (SEC-32). Anything that fails the check
// falls back to the sender's transport-peer partition.
private static string? TryResolveKeyId(string? authorizationHeader)
{
if (string.IsNullOrWhiteSpace(authorizationHeader))
{
return null;
}
return "peer:" + context.Peer;
ReadOnlySpan<char> header = authorizationHeader.AsSpan().Trim();
const string bearer = "Bearer ";
ReadOnlySpan<char> token = header.StartsWith(bearer, StringComparison.OrdinalIgnoreCase)
? header[bearer.Length..].Trim()
: header;
string[] parts = token.ToString().Split('_');
if (parts.Length < 3)
{
return null;
}
if (!string.Equals(parts[0], AuthStoreServiceCollectionExtensions.TokenPrefix, StringComparison.Ordinal))
{
return null;
}
if (parts[1].Length == 0 || parts[1].Length > MaxKeyIdLength || parts[2].Length == 0)
{
return null;
}
return parts[1];
}
}
@@ -794,6 +794,45 @@ public sealed class GatewayOptionsValidatorTests
Assert.Contains(result.Failures!, f => f.Contains("ApiKeyFailureTrackedPeers"));
}
/// <summary>Verifies a negative per-key aggregate failure limit fails validation.</summary>
[Fact]
public void Validate_Fails_WhenApiKeyFailureAggregateLimitNegative()
{
ValidateOptionsResult result = new GatewayOptionsValidator().Validate(
null,
WithSecurity(new SecurityOptions { ApiKeyFailureAggregateLimit = -1 }));
Assert.True(result.Failed);
Assert.Contains(result.Failures!, f => f.Contains("ApiKeyFailureAggregateLimit"));
}
/// <summary>Verifies a negative probe interval fails validation.</summary>
[Fact]
public void Validate_Fails_WhenApiKeyFailureProbeIntervalSecondsNegative()
{
ValidateOptionsResult result = new GatewayOptionsValidator().Validate(
null,
WithSecurity(new SecurityOptions { ApiKeyFailureProbeIntervalSeconds = -1 }));
Assert.True(result.Failed);
Assert.Contains(result.Failures!, f => f.Contains("ApiKeyFailureProbeIntervalSeconds"));
}
/// <summary>
/// Zero is a supported (documented) value for both new limiter knobs: it disables the aggregate
/// layer and probe admission respectively, so validation must accept it.
/// </summary>
[Fact]
public void Validate_Succeeds_WhenAggregateLimitAndProbeIntervalAreZero()
{
ValidateOptionsResult result = new GatewayOptionsValidator().Validate(
null,
WithSecurity(new SecurityOptions
{
ApiKeyFailureAggregateLimit = 0,
ApiKeyFailureProbeIntervalSeconds = 0,
}));
Assert.True(result.Succeeded);
}
private static GatewayOptions WithWorkerAndProtocol(WorkerOptions worker, ProtocolOptions protocol)
{
GatewayOptions source = ValidOptions();
@@ -0,0 +1,265 @@
using ZB.MOM.WW.MxGateway.Server.Security.Authorization;
using ZB.MOM.WW.MxGateway.Tests.TestSupport;
namespace ZB.MOM.WW.MxGateway.Tests.Security.Authorization;
/// <summary>
/// Unit tests for the two-layer API-key failure limiter (SEC-31 / SEC-32): composite
/// <c>(transport peer, key id)</c> partitions, the cross-peer per-key-id aggregate, probe
/// admission, the per-peer key-id partition cap, and the eviction preference order.
/// </summary>
public sealed class ApiKeyFailureLimiterTests
{
private static readonly TimeSpan Window = TimeSpan.FromSeconds(60);
private static readonly TimeSpan ProbeInterval = TimeSpan.FromSeconds(5);
/// <summary>A partition below the failure limit is admitted without consulting a probe slot.</summary>
[Fact]
public void Check_BelowLimit_Allows()
{
ManualTimeProvider clock = new(DateTimeOffset.UnixEpoch);
ApiKeyFailureLimiter limiter = CreateLimiter(clock, limit: 3);
ApiKeyThrottlePartition partition = new("ipv4:10.0.0.1:1", "victim");
limiter.RecordFailure(partition);
limiter.RecordFailure(partition);
Assert.Equal(ApiKeyThrottleDecision.Allowed, limiter.Check(partition));
}
/// <summary>Reaching the limit inside the window throttles the composite partition.</summary>
[Fact]
public void Check_AtLimit_ThrottlesCompositePartition()
{
ManualTimeProvider clock = new(DateTimeOffset.UnixEpoch);
ApiKeyFailureLimiter limiter = CreateLimiter(clock, limit: 3);
ApiKeyThrottlePartition partition = new("ipv4:10.0.0.1:1", "victim");
RecordFailures(limiter, partition, 3);
Assert.Equal(ApiKeyThrottleDecision.ThrottledByPeer, limiter.Check(partition));
}
/// <summary>Failures older than the sliding window are pruned, releasing the throttle.</summary>
[Fact]
public void Window_PrunesExpiredFailures_ReleasesThrottle()
{
ManualTimeProvider clock = new(DateTimeOffset.UnixEpoch);
ApiKeyFailureLimiter limiter = CreateLimiter(clock, limit: 3);
ApiKeyThrottlePartition partition = new("ipv4:10.0.0.1:1", "victim");
RecordFailures(limiter, partition, 3);
Assert.Equal(ApiKeyThrottleDecision.ThrottledByPeer, limiter.Check(partition));
clock.Advance(Window + TimeSpan.FromSeconds(1));
Assert.Equal(ApiKeyThrottleDecision.Allowed, limiter.Check(partition));
}
/// <summary>
/// The composite partition binds a throttle to the failing address: the same key id presented
/// from a different transport peer is unaffected. This is the structural half of the SEC-31 fix.
/// </summary>
[Fact]
public void CompositePartition_ThrottleDoesNotFollowKeyIdToAnotherPeer()
{
ManualTimeProvider clock = new(DateTimeOffset.UnixEpoch);
ApiKeyFailureLimiter limiter = CreateLimiter(clock, limit: 3, aggregateLimit: 1000);
RecordFailures(limiter, new ApiKeyThrottlePartition("ipv4:10.0.0.1:1", "victim"), 3);
Assert.Equal(
ApiKeyThrottleDecision.ThrottledByPeer,
limiter.Check(new ApiKeyThrottlePartition("ipv4:10.0.0.1:1", "victim")));
Assert.Equal(
ApiKeyThrottleDecision.Allowed,
limiter.Check(new ApiKeyThrottlePartition("ipv4:10.0.0.2:1", "victim")));
}
/// <summary>
/// Failures for one key id spread across many peers trip the per-key aggregate layer, so a
/// rotating-source sprayer is still bounded even though no single composite partition trips.
/// </summary>
[Fact]
public void AggregateLayer_TripsAcrossDistinctPeers()
{
ManualTimeProvider clock = new(DateTimeOffset.UnixEpoch);
ApiKeyFailureLimiter limiter = CreateLimiter(clock, limit: 100, aggregateLimit: 5);
for (int peer = 0; peer < 5; peer++)
{
limiter.RecordFailure(new ApiKeyThrottlePartition($"ipv4:10.0.0.{peer}:1", "victim"));
}
Assert.Equal(
ApiKeyThrottleDecision.ThrottledByAggregate,
limiter.Check(new ApiKeyThrottlePartition("ipv4:10.0.9.9:1", "victim")));
// The aggregate is per key id: an unrelated key from the same fresh peer is untouched.
Assert.Equal(
ApiKeyThrottleDecision.Allowed,
limiter.Check(new ApiKeyThrottlePartition("ipv4:10.0.9.9:1", "other-key")));
}
/// <summary>
/// A throttled partition is a valve, not a wall: exactly one request per probe interval is
/// admitted to the real verifier, so the holder of the correct secret can always get through.
/// </summary>
[Fact]
public void ProbeAdmission_AdmitsOneRequestPerInterval()
{
ManualTimeProvider clock = new(DateTimeOffset.UnixEpoch);
ApiKeyFailureLimiter limiter = CreateLimiter(clock, limit: 3);
ApiKeyThrottlePartition partition = new("ipv4:10.0.0.1:1", "victim");
RecordFailures(limiter, partition, 3);
Assert.Equal(ApiKeyThrottleDecision.ThrottledByPeer, limiter.Check(partition));
clock.Advance(ProbeInterval);
Assert.Equal(ApiKeyThrottleDecision.ProbeAdmitted, limiter.Check(partition));
// The granted slot is consumed: the next request inside the same interval is throttled.
Assert.Equal(ApiKeyThrottleDecision.ThrottledByPeer, limiter.Check(partition));
clock.Advance(ProbeInterval);
Assert.Equal(ApiKeyThrottleDecision.ProbeAdmitted, limiter.Check(partition));
}
/// <summary>A zero probe interval restores absolute blocking (documented as not recommended).</summary>
[Fact]
public void ProbeIntervalZero_BlocksAbsolutely()
{
ManualTimeProvider clock = new(DateTimeOffset.UnixEpoch);
ApiKeyFailureLimiter limiter = CreateLimiter(clock, limit: 3, probeInterval: TimeSpan.Zero);
ApiKeyThrottlePartition partition = new("ipv4:10.0.0.1:1", "victim");
RecordFailures(limiter, partition, 3);
// Well past several probe intervals but still inside the failure window: no slot opens.
clock.Advance(ProbeInterval + ProbeInterval + ProbeInterval);
Assert.Equal(ApiKeyThrottleDecision.ThrottledByPeer, limiter.Check(partition));
}
/// <summary>A successful verification clears both the composite partition and the key aggregate.</summary>
[Fact]
public void Reset_ClearsCompositeAndAggregateLayers()
{
ManualTimeProvider clock = new(DateTimeOffset.UnixEpoch);
ApiKeyFailureLimiter limiter = CreateLimiter(clock, limit: 3, aggregateLimit: 5);
ApiKeyThrottlePartition partition = new("ipv4:10.0.0.1:1", "victim");
RecordFailures(limiter, partition, 3);
for (int peer = 1; peer < 3; peer++)
{
limiter.RecordFailure(new ApiKeyThrottlePartition($"ipv4:10.0.1.{peer}:1", "victim"));
}
Assert.Equal(ApiKeyThrottleDecision.ThrottledByPeer, limiter.Check(partition));
limiter.Reset(partition);
Assert.Equal(ApiKeyThrottleDecision.Allowed, limiter.Check(partition));
Assert.Equal(0, limiter.TrackedAggregateCount);
Assert.Equal(
ApiKeyThrottleDecision.Allowed,
limiter.Check(new ApiKeyThrottlePartition("ipv4:10.0.9.9:1", "victim")));
}
/// <summary>
/// SEC-32: a spray of unique junk partitions (junk tokens resolve to the sender's fallback
/// partition, one per address) must not evict a partition that is currently throttled — the
/// LRU cap bounds memory, it must not be a reset button for the block.
/// </summary>
[Fact]
public void JunkTokenSpray_DoesNotEvictBlockedEntry()
{
ManualTimeProvider clock = new(DateTimeOffset.UnixEpoch);
ApiKeyFailureLimiter limiter = CreateLimiter(clock, limit: 3, maxPartitions: 8);
ApiKeyThrottlePartition blocked = new("ipv4:10.0.0.1:1", "victim");
RecordFailures(limiter, blocked, 3);
Assert.Equal(ApiKeyThrottleDecision.ThrottledByPeer, limiter.Check(blocked));
for (int i = 0; i < 16; i++)
{
clock.Advance(TimeSpan.FromMilliseconds(1));
limiter.RecordFailure(new ApiKeyThrottlePartition($"ipv4:10.9.{i / 256}.{i % 256}:1", KeyId: null));
}
Assert.True(limiter.IsTracked(blocked));
Assert.Equal(ApiKeyThrottleDecision.ThrottledByPeer, limiter.Check(blocked));
}
/// <summary>
/// SEC-32: one address may mint at most <see cref="ApiKeyFailureLimiter.MaxKeyIdPartitionsPerPeer"/>
/// key-id partitions; the overflow collapses into that address's fallback partition, which then
/// throttles the address wholesale instead of letting the spray mint unbounded state.
/// </summary>
[Fact]
public void UniqueMxgwKeyIdSpray_FromOnePeer_CollapsesAtPerPeerCap()
{
ManualTimeProvider clock = new(DateTimeOffset.UnixEpoch);
ApiKeyFailureLimiter limiter = CreateLimiter(clock, limit: 3, maxPartitions: 4096);
const string peer = "ipv4:10.0.0.1:1";
for (int i = 0; i < 1000; i++)
{
limiter.RecordFailure(new ApiKeyThrottlePartition(peer, $"key{i}"));
}
Assert.True(
limiter.TrackedPartitionCount <= ApiKeyFailureLimiter.MaxKeyIdPartitionsPerPeer + 1,
$"expected at most {ApiKeyFailureLimiter.MaxKeyIdPartitionsPerPeer + 1} partitions, saw {limiter.TrackedPartitionCount}");
Assert.Equal(
ApiKeyThrottleDecision.ThrottledByPeer,
limiter.Check(new ApiKeyThrottlePartition(peer, "key999")));
}
/// <summary>
/// SEC-32 eviction preference: with the map at capacity a new failure evicts an entry whose
/// window has fully expired rather than an entry that is still counting.
/// </summary>
[Fact]
public void Eviction_PrefersExpiredWindows()
{
ManualTimeProvider clock = new(DateTimeOffset.UnixEpoch);
ApiKeyFailureLimiter limiter = CreateLimiter(clock, limit: 3, maxPartitions: 2);
ApiKeyThrottlePartition expired = new("ipv4:10.0.0.1:1", KeyId: null);
ApiKeyThrottlePartition active = new("ipv4:10.0.0.2:1", KeyId: null);
ApiKeyThrottlePartition arriving = new("ipv4:10.0.0.3:1", KeyId: null);
limiter.RecordFailure(expired);
clock.Advance(Window + TimeSpan.FromSeconds(1));
limiter.RecordFailure(active);
limiter.RecordFailure(arriving);
Assert.Equal(2, limiter.TrackedPartitionCount);
Assert.False(limiter.IsTracked(expired));
Assert.True(limiter.IsTracked(active));
Assert.True(limiter.IsTracked(arriving));
}
private static void RecordFailures(ApiKeyFailureLimiter limiter, ApiKeyThrottlePartition partition, int count)
{
for (int i = 0; i < count; i++)
{
limiter.RecordFailure(partition);
}
}
private static ApiKeyFailureLimiter CreateLimiter(
TimeProvider clock,
int limit,
int aggregateLimit = 0,
int maxPartitions = 1024,
TimeSpan? probeInterval = null)
{
return new ApiKeyFailureLimiter(
limit,
Window,
maxPartitions,
aggregateLimit,
probeInterval ?? ProbeInterval,
clock);
}
}
@@ -22,6 +22,12 @@ namespace ZB.MOM.WW.MxGateway.Tests.Security.Authorization;
public sealed class GatewayGrpcAuthorizationInterceptorTests
{
private const string AttackerPeer = "ipv4:203.0.113.7:5000";
private const string HolderPeer = "ipv4:198.51.100.4:5000";
private static readonly TimeSpan FailureWindow = TimeSpan.FromMinutes(1);
private static readonly TimeSpan ProbeInterval = TimeSpan.FromSeconds(5);
/// <summary>Verifies that missing API key returns unauthenticated status.</summary>
/// <returns>A task that represents the asynchronous operation.</returns>
[Fact]
@@ -359,21 +365,18 @@ public sealed class GatewayGrpcAuthorizationInterceptorTests
}
/// <summary>
/// Once a peer has exceeded the failure limit, the interceptor short-circuits with
/// <see cref="StatusCode.ResourceExhausted"/> BEFORE calling the verifier, so an online guessing
/// loop stops spending a store read per attempt. A verifier that always fails is used; after the
/// limit is reached the verifier is no longer invoked.
/// SEC-31: once an attacking peer has exceeded the failure limit for a key id, the interceptor
/// short-circuits with <see cref="StatusCode.ResourceExhausted"/> BEFORE calling the verifier, so
/// an online guessing loop stops spending a store read per attempt. The composite
/// <c>(peer, key id)</c> partition keeps that bound per attacking address.
/// </summary>
/// <returns>A task that represents the asynchronous operation.</returns>
[Fact]
public async Task UnaryServerHandler_ExceedsFailureLimit_ShortCircuitsBeforeVerify()
public async Task BruteForceBound_StillEnforcedPerAttackingPeer()
{
CountingFailureVerifier verifier = new(Failure(ApiKeyFailure.SecretMismatch));
ApiKeyFailureLimiter limiter = new(
limit: 3,
window: TimeSpan.FromMinutes(1),
maxPeers: 16,
clock: TimeProvider.System);
ManualTimeProvider clock = new(DateTimeOffset.UnixEpoch);
ApiKeyFailureLimiter limiter = CreateLimiter(clock, limit: 3);
GatewayGrpcAuthorizationInterceptor interceptor = CreateInterceptor(
verifier,
new GatewayRequestIdentityAccessor(),
@@ -385,7 +388,7 @@ public sealed class GatewayGrpcAuthorizationInterceptorTests
RpcException failure = await Assert.ThrowsAsync<RpcException>(
() => interceptor.UnaryServerHandler(
new OpenSessionRequest(),
ContextWithAuthorization("Bearer mxgw_operator01_bad-secret"),
ContextWithAuthorization("Bearer mxgw_operator01_bad-secret", AttackerPeer),
(_, _) => Task.FromResult(new OpenSessionReply())));
Assert.Equal(StatusCode.Unauthenticated, failure.StatusCode);
}
@@ -396,13 +399,219 @@ public sealed class GatewayGrpcAuthorizationInterceptorTests
RpcException throttled = await Assert.ThrowsAsync<RpcException>(
() => interceptor.UnaryServerHandler(
new OpenSessionRequest(),
ContextWithAuthorization("Bearer mxgw_operator01_bad-secret"),
ContextWithAuthorization("Bearer mxgw_operator01_bad-secret", AttackerPeer),
(_, _) => Task.FromResult(new OpenSessionReply())));
Assert.Equal(StatusCode.ResourceExhausted, throttled.StatusCode);
Assert.Equal(3, verifier.CallCount);
}
/// <summary>
/// SEC-31 (the lockout inversion): an attacker who floods failures for a victim's key id from its
/// own address must not deny that key to the legitimate holder. The holder presents the correct
/// secret from a different transport peer and authenticates on the first attempt — the verifier
/// is reached and the RPC succeeds.
/// </summary>
/// <returns>A task that represents the asynchronous operation.</returns>
[Fact]
public async Task AttackerSpamOnVictimKeyId_FromDifferentPeer_DoesNotBlockLegitimateHolderPresentingCorrectSecret()
{
ManualTimeProvider clock = new(DateTimeOffset.UnixEpoch);
ApiKeyFailureLimiter limiter = CreateLimiter(clock, limit: 3, aggregateLimit: 30);
GatewayGrpcAuthorizationInterceptor attacked = CreateInterceptor(
new CountingFailureVerifier(Failure(ApiKeyFailure.SecretMismatch)),
new GatewayRequestIdentityAccessor(),
failureLimiter: limiter);
// Flood the victim's key id from the attacker's address until that partition is throttled.
StatusCode lastAttackerStatus = StatusCode.OK;
for (int attempt = 0; attempt < 6; attempt++)
{
RpcException failure = await Assert.ThrowsAsync<RpcException>(
() => attacked.UnaryServerHandler(
new OpenSessionRequest(),
ContextWithAuthorization("Bearer mxgw_operator01_guess", AttackerPeer),
(_, _) => Task.FromResult(new OpenSessionReply())));
lastAttackerStatus = failure.StatusCode;
}
Assert.Equal(StatusCode.ResourceExhausted, lastAttackerStatus);
// The legitimate holder, on a different address, is verified and admitted immediately.
FakeApiKeyVerifier holderVerifier = new(SuccessWithScopes(GatewayScopes.SessionOpen));
GatewayGrpcAuthorizationInterceptor holder = CreateInterceptor(
holderVerifier,
new GatewayRequestIdentityAccessor(),
failureLimiter: limiter);
OpenSessionReply reply = await holder.UnaryServerHandler(
new OpenSessionRequest(),
ContextWithAuthorization("Bearer mxgw_operator01_correct", HolderPeer),
(_, _) => Task.FromResult(new OpenSessionReply { SessionId = "session-1" }));
Assert.True(holderVerifier.WasCalled);
Assert.Equal("session-1", reply.SessionId);
}
/// <summary>
/// SEC-31 layer 2: failures for one key id sprayed across more distinct peers than
/// <c>ApiKeyFailureAggregateLimit</c> put that key id into probe mode globally, so a
/// rotating-source attacker gets at most one verifier call per probe interval.
/// </summary>
/// <returns>A task that represents the asynchronous operation.</returns>
[Fact]
public async Task AggregateSpray_AcrossManyPeers_TripsPerKeyProbeMode()
{
CountingFailureVerifier verifier = new(Failure(ApiKeyFailure.SecretMismatch));
ManualTimeProvider clock = new(DateTimeOffset.UnixEpoch);
ApiKeyFailureLimiter limiter = CreateLimiter(clock, limit: 100, aggregateLimit: 5);
GatewayGrpcAuthorizationInterceptor interceptor = CreateInterceptor(
verifier,
new GatewayRequestIdentityAccessor(),
failureLimiter: limiter);
for (int peer = 0; peer < 5; peer++)
{
await Assert.ThrowsAsync<RpcException>(
() => interceptor.UnaryServerHandler(
new OpenSessionRequest(),
ContextWithAuthorization("Bearer mxgw_operator01_guess", $"ipv4:10.9.0.{peer}:5000"),
(_, _) => Task.FromResult(new OpenSessionReply())));
}
Assert.Equal(5, verifier.CallCount);
// A never-seen peer is now probe-limited: no composite failures of its own, but the key id's
// aggregate is tripped, so the request never reaches the verifier.
RpcException throttled = await Assert.ThrowsAsync<RpcException>(
() => interceptor.UnaryServerHandler(
new OpenSessionRequest(),
ContextWithAuthorization("Bearer mxgw_operator01_guess", "ipv4:10.9.1.1:5000"),
(_, _) => Task.FromResult(new OpenSessionReply())));
Assert.Equal(StatusCode.ResourceExhausted, throttled.StatusCode);
Assert.Equal(5, verifier.CallCount);
// One probe slot opens per interval, and it is consumed by the first arrival.
clock.Advance(ProbeInterval);
RpcException probed = await Assert.ThrowsAsync<RpcException>(
() => interceptor.UnaryServerHandler(
new OpenSessionRequest(),
ContextWithAuthorization("Bearer mxgw_operator01_guess", "ipv4:10.9.1.2:5000"),
(_, _) => Task.FromResult(new OpenSessionReply())));
Assert.Equal(StatusCode.Unauthenticated, probed.StatusCode);
Assert.Equal(6, verifier.CallCount);
RpcException throttledAgain = await Assert.ThrowsAsync<RpcException>(
() => interceptor.UnaryServerHandler(
new OpenSessionRequest(),
ContextWithAuthorization("Bearer mxgw_operator01_guess", "ipv4:10.9.1.3:5000"),
(_, _) => Task.FromResult(new OpenSessionReply())));
Assert.Equal(StatusCode.ResourceExhausted, throttledAgain.StatusCode);
Assert.Equal(6, verifier.CallCount);
}
/// <summary>
/// SEC-31: the success-reset path stays reachable while throttled. A throttled partition admits
/// one probe per interval; the correct secret rides that slot, authenticates, and fully clears
/// both limiter layers, so the next wrong attempt is <see cref="StatusCode.Unauthenticated"/>
/// rather than <see cref="StatusCode.ResourceExhausted"/>.
/// </summary>
/// <returns>A task that represents the asynchronous operation.</returns>
[Fact]
public async Task CorrectSecret_DuringProbeMode_AuthenticatesViaProbeSlotAndResets()
{
ManualTimeProvider clock = new(DateTimeOffset.UnixEpoch);
ApiKeyFailureLimiter limiter = CreateLimiter(clock, limit: 3, aggregateLimit: 30);
GatewayGrpcAuthorizationInterceptor failing = CreateInterceptor(
new FakeApiKeyVerifier(Failure(ApiKeyFailure.SecretMismatch)),
new GatewayRequestIdentityAccessor(),
failureLimiter: limiter);
FakeApiKeyVerifier holderVerifier = new(SuccessWithScopes(GatewayScopes.SessionOpen));
GatewayGrpcAuthorizationInterceptor succeeding = CreateInterceptor(
holderVerifier,
new GatewayRequestIdentityAccessor(),
failureLimiter: limiter);
for (int attempt = 0; attempt < 3; attempt++)
{
await Assert.ThrowsAsync<RpcException>(
() => failing.UnaryServerHandler(
new OpenSessionRequest(),
ContextWithAuthorization("Bearer mxgw_operator01_bad", HolderPeer),
(_, _) => Task.FromResult(new OpenSessionReply())));
}
RpcException throttled = await Assert.ThrowsAsync<RpcException>(
() => failing.UnaryServerHandler(
new OpenSessionRequest(),
ContextWithAuthorization("Bearer mxgw_operator01_bad", HolderPeer),
(_, _) => Task.FromResult(new OpenSessionReply())));
Assert.Equal(StatusCode.ResourceExhausted, throttled.StatusCode);
clock.Advance(ProbeInterval);
OpenSessionReply reply = await succeeding.UnaryServerHandler(
new OpenSessionRequest(),
ContextWithAuthorization("Bearer mxgw_operator01_correct", HolderPeer),
(_, _) => Task.FromResult(new OpenSessionReply { SessionId = "session-1" }));
Assert.True(holderVerifier.WasCalled);
Assert.Equal("session-1", reply.SessionId);
RpcException afterReset = await Assert.ThrowsAsync<RpcException>(
() => failing.UnaryServerHandler(
new OpenSessionRequest(),
ContextWithAuthorization("Bearer mxgw_operator01_bad", HolderPeer),
(_, _) => Task.FromResult(new OpenSessionReply())));
Assert.Equal(StatusCode.Unauthenticated, afterReset.StatusCode);
}
/// <summary>
/// SEC-32: only a validly shaped <c>mxgw_&lt;keyId&gt;_&lt;secret&gt;</c> token mints a key-id
/// partition. Junk tokens of varied shapes all collapse onto the sender's transport-peer fallback
/// partition, so a spray cannot mint one tracked entry per invented token.
/// </summary>
/// <returns>A task that represents the asynchronous operation.</returns>
[Fact]
public async Task NonMxgwToken_FallsBackToTransportPeerPartition()
{
ManualTimeProvider clock = new(DateTimeOffset.UnixEpoch);
ApiKeyFailureLimiter limiter = CreateLimiter(clock, limit: 1000);
GatewayGrpcAuthorizationInterceptor interceptor = CreateInterceptor(
new FakeApiKeyVerifier(Failure(ApiKeyFailure.MissingOrMalformed)),
new GatewayRequestIdentityAccessor(),
failureLimiter: limiter);
string[] junkTokens =
[
"Bearer garbage",
"Bearer a_b_c",
"Bearer notmxgw_operator01_secret",
"Bearer MXGW_operator01_secret",
"Bearer mxgw__secret",
"Bearer mxgw_operator01_",
"Bearer mxgw_" + new string('k', 65) + "_secret",
"Bearer mxgw_onlytwo",
];
foreach (string token in junkTokens)
{
await Assert.ThrowsAsync<RpcException>(
() => interceptor.UnaryServerHandler(
new OpenSessionRequest(),
ContextWithAuthorization(token, AttackerPeer),
(_, _) => Task.FromResult(new OpenSessionReply())));
}
Assert.Equal(1, limiter.TrackedPartitionCount);
Assert.True(limiter.IsTracked(new ApiKeyThrottlePartition(AttackerPeer, KeyId: null)));
}
/// <summary>
/// A successful verification resets the peer's failure counter, so accumulated failures
/// from a fat-fingered secret do not lock out a client that subsequently authenticates.
@@ -411,11 +620,7 @@ public sealed class GatewayGrpcAuthorizationInterceptorTests
[Fact]
public async Task UnaryServerHandler_SuccessResetsFailureCounter()
{
ApiKeyFailureLimiter limiter = new(
limit: 3,
window: TimeSpan.FromMinutes(1),
maxPeers: 16,
clock: TimeProvider.System);
ApiKeyFailureLimiter limiter = CreateLimiter(new ManualTimeProvider(DateTimeOffset.UnixEpoch), limit: 3);
// Two failures against the same key id, then a success (which resets), then two more
// failures — without the reset the fifth attempt would be blocked at the limit of 3.
@@ -487,11 +692,23 @@ public sealed class GatewayGrpcAuthorizationInterceptorTests
Mode = authenticationMode
}
}),
failureLimiter ?? new ApiKeyFailureLimiter(
limit: 1000,
window: TimeSpan.FromMinutes(1),
maxPeers: 1000,
clock: TimeProvider.System));
failureLimiter ?? CreateLimiter(TimeProvider.System, limit: 1000),
new GatewayMetrics());
}
private static ApiKeyFailureLimiter CreateLimiter(
TimeProvider clock,
int limit,
int aggregateLimit = 0,
int maxPartitions = 1024)
{
return new ApiKeyFailureLimiter(
limit,
FailureWindow,
maxPartitions,
aggregateLimit,
ProbeInterval,
clock);
}
private static ApiKeyVerification SuccessWithScopes(params string[] scopes)
@@ -511,9 +728,11 @@ public sealed class GatewayGrpcAuthorizationInterceptorTests
return new ApiKeyVerification(Succeeded: false, Identity: null, Failure: failure);
}
private static TestServerCallContext ContextWithAuthorization(string authorizationHeader)
private static TestServerCallContext ContextWithAuthorization(string authorizationHeader, string? peer = null)
{
return new TestServerCallContext([new Metadata.Entry("authorization", authorizationHeader)]);
return new TestServerCallContext(
[new Metadata.Entry("authorization", authorizationHeader)],
peer: peer);
}
/// <summary>Records whether the gateway service ran past the interceptor for composition tests.</summary>
@@ -8,20 +8,32 @@ namespace ZB.MOM.WW.MxGateway.Tests.TestSupport;
/// </summary>
public sealed class TestServerCallContext : ServerCallContext
{
private const string DefaultPeer = "ipv4:127.0.0.1:5000";
private readonly Metadata _requestHeaders;
private readonly Metadata _responseTrailers = [];
private readonly Dictionary<object, object> _userState = [];
private readonly CancellationToken _cancellationToken;
private readonly string _peer;
private Status _status;
private WriteOptions? _writeOptions;
/// <summary>Initializes the context with the supplied request headers and cancellation token.</summary>
/// <param name="requestHeaders">Request headers visible to the service; defaults to empty.</param>
/// <param name="cancellationToken">Cancellation token surfaced to the service.</param>
public TestServerCallContext(Metadata? requestHeaders = null, CancellationToken cancellationToken = default)
/// <param name="peer">
/// Transport peer address surfaced as <see cref="ServerCallContext.Peer"/>; defaults to a
/// loopback address. Tests that exercise per-peer partitioning (for example the API-key failure
/// limiter) pass distinct values to model separate network sources.
/// </param>
public TestServerCallContext(
Metadata? requestHeaders = null,
CancellationToken cancellationToken = default,
string? peer = null)
{
_requestHeaders = requestHeaders ?? [];
_cancellationToken = cancellationToken;
_peer = peer ?? DefaultPeer;
}
/// <inheritdoc />
@@ -31,7 +43,7 @@ public sealed class TestServerCallContext : ServerCallContext
protected override string HostCore => "localhost";
/// <inheritdoc />
protected override string PeerCore => "ipv4:127.0.0.1:5000";
protected override string PeerCore => _peer;
/// <inheritdoc />
protected override DateTime DeadlineCore => DateTime.UtcNow.AddMinutes(1);