From 8fb17352554403c8264114d971df330b50202c43 Mon Sep 17 00:00:00 2001 From: Zacgoose <107489668+Zacgoose@users.noreply.github.com> Date: Fri, 2 Oct 2026 01:35:35 +0800 Subject: [PATCH 1/4] feat(hosting): serve interactive HTTP requests ahead of API clients - HTTP worker checkout waits in two FIFO tiers; API clients and background cache refreshes only get a worker when no interactive request is waiting - Queue timeout unchanged and shared by both tiers --- .../PowerShellHost/PowerShellRunnerService.cs | 15 ++-- .../PowerShellHost/PowerShellWorkerPool.cs | 16 ++-- Services/PowerShellHost/TieredHandoff.cs | 63 +++++++++++++ appsettings.example.jsonc | 3 +- tests/Craft.Tests/TieredHandoffTests.cs | 89 +++++++++++++++++++ 5 files changed, 170 insertions(+), 16 deletions(-) create mode 100644 Services/PowerShellHost/TieredHandoff.cs create mode 100644 tests/Craft.Tests/TieredHandoffTests.cs diff --git a/Services/PowerShellHost/PowerShellRunnerService.cs b/Services/PowerShellHost/PowerShellRunnerService.cs index 66d3354..a08a8c1 100644 --- a/Services/PowerShellHost/PowerShellRunnerService.cs +++ b/Services/PowerShellHost/PowerShellRunnerService.cs @@ -238,7 +238,8 @@ public async Task ExecuteHttpScript(string route, HttpContext http if (!DispatchProfiler.Enabled) { var req = await BuildRequestObject(httpContext); - return await ExecuteHttpScriptInternal(route, req, isHttp: true, clientAborted: clientAborted); + return await ExecuteHttpScriptInternal(route, req, isHttp: true, clientAborted: clientAborted, + lowPriority: CallerClassifier.IsApiClient(httpContext)); } // Profiling path: time request marshaling + the runner-side segments (checkout/invoke/extract). @@ -249,7 +250,7 @@ public async Task ExecuteHttpScript(string route, HttpContext http var marshalTicks = Stopwatch.GetTimestamp() - mStart; var timing = new DispatchTiming(); var result = await ExecuteHttpScriptInternal(route, request, isHttp: true, timing, - clientAborted: clientAborted); + clientAborted: clientAborted, lowPriority: CallerClassifier.IsApiClient(httpContext)); DispatchProfiler.Record(marshalTicks, timing.CheckoutTicks, timing.InvokeTicks, timing.ExtractTicks, Stopwatch.GetTimestamp() - totalStart); return result; @@ -262,11 +263,11 @@ public async Task ExecuteHttpScript(string route, HttpContext http /// public async Task ExecuteHttpScript(string route, Hashtable requestSnapshot) { - return await ExecuteHttpScriptInternal(route, requestSnapshot, isHttp: true); + return await ExecuteHttpScriptInternal(route, requestSnapshot, isHttp: true, lowPriority: true); } private async Task ExecuteHttpScriptInternal(string route, Hashtable request, bool isHttp, - DispatchTiming? timing = null, CancellationToken clientAborted = default) + DispatchTiming? timing = null, bool lowPriority = false, CancellationToken clientAborted = default) { var sw = Stopwatch.StartNew(); var entry = _repo.GetByRoute(route); @@ -291,12 +292,12 @@ private async Task ExecuteHttpScriptInternal(string route, Hashtab var checkoutStart = timing != null ? Stopwatch.GetTimestamp() : 0; if (isHttp) { - worker = _pool.CheckoutHttp(_httpQueueTimeout); + worker = _pool.CheckoutHttp(_httpQueueTimeout, lowPriority); if (timing != null) timing.CheckoutTicks = Stopwatch.GetTimestamp() - checkoutStart; if (worker == null) { - _logger.LogWarning("HTTP pool exhausted — no worker available within {Timeout:0}s for {Route}", - _httpQueueTimeout.TotalSeconds, route); + _logger.LogWarning("HTTP pool exhausted — no worker available within {Timeout:0}s for {Route} ({Tier})", + _httpQueueTimeout.TotalSeconds, route, lowPriority ? "low" : "priority"); return new ScriptResult { StatusCode = 503, diff --git a/Services/PowerShellHost/PowerShellWorkerPool.cs b/Services/PowerShellHost/PowerShellWorkerPool.cs index 866aeeb..61cc143 100644 --- a/Services/PowerShellHost/PowerShellWorkerPool.cs +++ b/Services/PowerShellHost/PowerShellWorkerPool.cs @@ -7,7 +7,7 @@ namespace Craft.PowerShellHost; public class PowerShellWorkerPool : IDisposable { - private readonly BlockingCollection _httpPool; + private readonly TieredHandoff _httpPool = new(); private readonly BlockingCollection _bgPool; private readonly ILogger _logger; private readonly ScriptRepository _repo; @@ -49,9 +49,6 @@ public PowerShellWorkerPool(ScriptRepository repo, ILogger _httpPoolSize = Math.Max(0, settings.Worker.HttpPoolSize); _bgPoolSize = Math.Max(1, settings.Worker.BgPoolSize); - // BlockingCollection rejects a bounded capacity of 0. The pool is never populated in that case, - // so the capacity is arbitrary — it just has to be legal. - _httpPool = new BlockingCollection(Math.Max(1, _httpPoolSize)); _bgPool = new BlockingCollection(_bgPoolSize); } @@ -456,12 +453,16 @@ private void AddToBgPool(PowerShellWorker worker) /// public bool WaitForBgReady(TimeSpan timeout) => _bgReady.Wait(timeout); - public PowerShellWorker? CheckoutHttp(TimeSpan timeout) + /// + /// Check out an HTTP worker. callers (API clients, background cache + /// refreshes) only get a worker when no interactive request is waiting for one. + /// + public PowerShellWorker? CheckoutHttp(TimeSpan timeout, bool lowPriority = false) { // Wait for HTTP pool initialization before attempting checkout if (!_httpReady.IsSet) _httpReady.Wait(timeout); - if (_httpPool.TryTake(out var w, timeout)) + if (_httpPool.TryTake(lowPriority, timeout) is { } w) { w.CheckoutTimestamp = System.Diagnostics.Stopwatch.GetTimestamp(); WorkerMetricsBridge.RecordCheckout(w.Id, isHttp: true); @@ -956,9 +957,8 @@ public void Dispose() // Nothing here owns unmanaged resources directly, but suppressing finalization keeps a // derived type that adds a finalizer from having to re-implement IDisposable to do it. GC.SuppressFinalize(this); - try { while (_httpPool.TryTake(out var w)) w.Dispose(); } catch (ObjectDisposedException) { } + while (_httpPool.TryTake(lowPriority: false, TimeSpan.Zero) is { } h) h.Dispose(); try { while (_bgPool.TryTake(out var w)) w.Dispose(); } catch (ObjectDisposedException) { } - _httpPool.Dispose(); _bgPool.Dispose(); } } diff --git a/Services/PowerShellHost/TieredHandoff.cs b/Services/PowerShellHost/TieredHandoff.cs new file mode 100644 index 0000000..f988689 --- /dev/null +++ b/Services/PowerShellHost/TieredHandoff.cs @@ -0,0 +1,63 @@ +namespace Craft.PowerShellHost; + +/// +/// A pool of items handed to waiters in two FIFO tiers: a returned item always goes to the oldest +/// priority waiter, and to a low-priority waiter only when no priority waiter is queued. This is what +/// keeps interactive HTTP requests ahead of an API-client burst on the shared HTTP worker pool. +/// +public sealed class TieredHandoff where T : class +{ + private readonly object _lock = new(); + private readonly Queue _free = new(); + private readonly LinkedList> _priority = new(); + private readonly LinkedList> _low = new(); + + public int Count { get { lock (_lock) return _free.Count; } } + + public void Add(T item) + { + TaskCompletionSource waiter; + lock (_lock) + { + var list = _priority.Count > 0 ? _priority : _low.Count > 0 ? _low : null; + if (list == null) + { + _free.Enqueue(item); + return; + } + waiter = list.First!.Value; + list.RemoveFirst(); + } + waiter.SetResult(item); + } + + /// + /// Takes an item, waiting up to behind earlier waiters of the same tier + /// (and, for , behind every priority waiter). Null on timeout. + /// + public T? TryTake(bool lowPriority, TimeSpan timeout) + { + LinkedListNode> node; + lock (_lock) + { + // Add serves waiters before refilling _free, so a free item means nobody is queued. + if (_free.TryDequeue(out var item)) return item; + if (timeout <= TimeSpan.Zero) return null; + node = (lowPriority ? _low : _priority).AddLast( + new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously)); + } + + if (node.Value.Task.Wait(timeout)) return node.Value.Task.Result; + + lock (_lock) + { + if (node.List != null) + { + node.List.Remove(node); + return null; + } + } + // Add dequeued this waiter as the timeout fired; its item is already on the way. + return node.Value.Task.GetAwaiter().GetResult(); + } +} diff --git a/appsettings.example.jsonc b/appsettings.example.jsonc index c7c28e4..76198c7 100644 --- a/appsettings.example.jsonc +++ b/appsettings.example.jsonc @@ -233,7 +233,8 @@ // shed with 503 "Server busy, please retry". A load-shedding bound, NOT a capacity knob — it // does not add throughput, it only decides how long callers wait before a spurious 503. Widen // it to absorb brief bursts; raise HttpPoolSize/MinThreads for steady saturation. Distinct from - // HttpTimeoutSeconds (which bounds execution once a worker is held). + // HttpTimeoutSeconds (which bounds execution once a worker is held). Waiters are served in two + // FIFO tiers: interactive users first, then API clients and background cache refreshes. // 0 = built-in default of 30s. Env override: CRAFT_HTTP_QUEUE_TIMEOUT (seconds). // "HttpQueueTimeoutSeconds": 30, diff --git a/tests/Craft.Tests/TieredHandoffTests.cs b/tests/Craft.Tests/TieredHandoffTests.cs new file mode 100644 index 0000000..9f29e58 --- /dev/null +++ b/tests/Craft.Tests/TieredHandoffTests.cs @@ -0,0 +1,89 @@ +using Craft.PowerShellHost; + +namespace Craft.Tests; + +/// +/// The HTTP worker pool hands workers out through this queue, so its ordering is the guarantee that an +/// API-client burst cannot push interactive requests behind it. +/// +public class TieredHandoffTests +{ + private static readonly TimeSpan Long = TimeSpan.FromSeconds(10); + + private static Task Waiter(TieredHandoff queue, bool lowPriority, TimeSpan timeout) + { + var task = Task.Factory.StartNew(() => queue.TryTake(lowPriority, timeout), TaskCreationOptions.LongRunning); + Thread.Sleep(100); // let it enqueue before the next waiter + return task; + } + + [Fact] + public void FreeItem_IsReturnedImmediately() + { + var queue = new TieredHandoff(); + queue.Add("w1"); + Assert.Equal("w1", queue.TryTake(lowPriority: true, TimeSpan.Zero)); + Assert.Null(queue.TryTake(lowPriority: false, TimeSpan.Zero)); + } + + [Fact] + public async Task PriorityWaiter_IsServedBeforeEarlierLowWaiters() + { + var queue = new TieredHandoff(); + var low1 = Waiter(queue, lowPriority: true, Long); + var low2 = Waiter(queue, lowPriority: true, Long); + var priority = Waiter(queue, lowPriority: false, Long); + + queue.Add("a"); + Assert.Equal("a", await priority); + Assert.False(low1.IsCompleted); + + queue.Add("b"); + queue.Add("c"); + Assert.Equal("b", await low1); + Assert.Equal("c", await low2); + Assert.Equal(0, queue.Count); + } + + [Fact] + public async Task PriorityWaiters_AreFifo() + { + var queue = new TieredHandoff(); + var first = Waiter(queue, lowPriority: false, Long); + var second = Waiter(queue, lowPriority: false, Long); + + queue.Add("a"); + Assert.Equal("a", await first); + queue.Add("b"); + Assert.Equal("b", await second); + } + + [Fact] + public async Task TimedOutWaiter_IsDropped_AndTheNextItemGoesToTheFreeList() + { + var queue = new TieredHandoff(); + Assert.Null(await Waiter(queue, lowPriority: false, TimeSpan.FromMilliseconds(50))); + + queue.Add("a"); + Assert.Equal(1, queue.Count); + } + + [Fact] + public async Task ConcurrentChurn_LosesNoItems() + { + var queue = new TieredHandoff(); + for (var i = 0; i < 4; i++) queue.Add($"w{i}"); + + var callers = Enumerable.Range(0, 32).Select(i => Task.Run(() => + { + for (var n = 0; n < 200; n++) + { + var item = queue.TryTake(lowPriority: i % 2 == 0, TimeSpan.FromMilliseconds(5)); + if (item != null) queue.Add(item); + } + })); + await Task.WhenAll(callers); + + Assert.Equal(4, queue.Count); + } +} From 45f75b920a25ced5280bf79a7fc8f55057217fae Mon Sep 17 00:00:00 2001 From: Zacgoose <107489668+Zacgoose@users.noreply.github.com> Date: Fri, 2 Oct 2026 02:24:34 +0800 Subject: [PATCH 2/4] feat(logging): mask UPNs and customer domains in file and console logs - Keep both ends of each label and the public suffix readable; the hidden middle is encrypted deterministically (short HMAC tag seeding an AES-CTR keystream), so a value always masks to the same short token - Instance key lives in the CraftInstanceKeys table, created on first start - LogBridge reveals masked values before filtering, so search, exclude and regex work on the originals - Platform hosts, tenant ids and GUIDs are left as is; CRAFT_LOG_REDACTION=false turns it off --- Services/Bridges/LogBridge.cs | 3 + Services/Configuration/FileLoggingSettings.cs | 9 + .../Hosting/CraftHostBuilderExtensions.cs | 13 +- Services/Hosting/FileLoggerProvider.cs | 2 + Services/Hosting/LogRedactor.Tlds.cs | 110 +++++++ Services/Hosting/LogRedactor.cs | 307 ++++++++++++++++++ .../Hosting/RedactingConsoleLoggerProvider.cs | 37 +++ Services/Program.cs | 9 + tests/Craft.Tests/LogRedactorTests.cs | 174 ++++++++++ 9 files changed, 663 insertions(+), 1 deletion(-) create mode 100644 Services/Hosting/LogRedactor.Tlds.cs create mode 100644 Services/Hosting/LogRedactor.cs create mode 100644 Services/Hosting/RedactingConsoleLoggerProvider.cs create mode 100644 tests/Craft.Tests/LogRedactorTests.cs diff --git a/Services/Bridges/LogBridge.cs b/Services/Bridges/LogBridge.cs index 7b2ec94..5bece99 100644 --- a/Services/Bridges/LogBridge.cs +++ b/Services/Bridges/LogBridge.cs @@ -99,6 +99,9 @@ public static string[] ReadLog(int tail = 0, string? level = null, string? searc return Array.Empty(); } + for (var i = 0; i < allLines.Length; i++) + allLines[i] = LogRedactor.Reveal(allLines[i]); + // Parse multi-level filter once string[]? levels = null; if (!string.IsNullOrEmpty(level)) diff --git a/Services/Configuration/FileLoggingSettings.cs b/Services/Configuration/FileLoggingSettings.cs index 2a7deaf..fb8dc40 100644 --- a/Services/Configuration/FileLoggingSettings.cs +++ b/Services/Configuration/FileLoggingSettings.cs @@ -46,6 +46,15 @@ public class FileLoggingSettings /// public string LogLevel { get; set; } = "Information"; + /// + /// Mask UPNs and customer domains in file and console output, keeping them readable at a glance + /// (see LogRedactor). LogBridge reveals them again. Env override: CRAFT_LOG_REDACTION=false. + /// + public bool Redact { get; set; } = true; + + /// Extra domains (and their subdomains) left readable, e.g. the app's own. + public List RedactAllowDomains { get; set; } = new(); + /// Resolved directory path, applying platform defaults when Directory is empty. internal string ResolvedDirectory => !string.IsNullOrEmpty(Directory) ? Directory diff --git a/Services/Hosting/CraftHostBuilderExtensions.cs b/Services/Hosting/CraftHostBuilderExtensions.cs index 2189af0..ab3f637 100644 --- a/Services/Hosting/CraftHostBuilderExtensions.cs +++ b/Services/Hosting/CraftHostBuilderExtensions.cs @@ -160,6 +160,11 @@ public static LogLevel AddCraftLogging(this WebApplicationBuilder builder) builder.Configuration.GetSection("App:FileLogging").Bind(fileLoggingSettings); var level = fileLoggingSettings.ParsedLogLevel; + var redactEnv = Environment.GetEnvironmentVariable("CRAFT_LOG_REDACTION"); + LogRedactor.Configure( + fileLoggingSettings.Redact && redactEnv is not ("0" or "false" or "False" or "FALSE"), + fileLoggingSettings.RedactAllowDomains); + var fileLoggerProvider = new FileLoggerProvider(fileLoggingSettings, level); builder.Logging.AddProvider(fileLoggerProvider); LogBridge.Initialize(fileLoggerProvider); @@ -169,10 +174,16 @@ public static LogLevel AddCraftLogging(this WebApplicationBuilder builder) options.TimestampFormat = "yyyy-MM-ddTHH:mm:ss.fffZ "; options.SingleLine = true; }); + var services = builder.Logging.Services; + services.Remove(services.Single(d => + d.ServiceType == typeof(ILoggerProvider) && d.ImplementationType == typeof(ConsoleLoggerProvider))); + services.AddSingleton(); + services.AddSingleton(sp => + new RedactingConsoleLoggerProvider(sp.GetRequiredService())); if (level > LogLevel.Debug) { - builder.Logging.AddFilter(l => l >= LogLevel.Information); + builder.Logging.AddFilter(l => l >= LogLevel.Information); // Framework logging is noise at Information and above. builder.Logging.AddFilter("Microsoft.AspNetCore", LogLevel.Warning); diff --git a/Services/Hosting/FileLoggerProvider.cs b/Services/Hosting/FileLoggerProvider.cs index ec980d1..cf16065 100644 --- a/Services/Hosting/FileLoggerProvider.cs +++ b/Services/Hosting/FileLoggerProvider.cs @@ -135,6 +135,8 @@ private void WriteToFile(string line, string? exceptionLine) try { if (_writer == null) return; + line = LogRedactor.Redact(line); + if (exceptionLine != null) exceptionLine = LogRedactor.Redact(exceptionLine); _writer.WriteLine(line); _currentFileSize += line.Length + Environment.NewLine.Length; diff --git a/Services/Hosting/LogRedactor.Tlds.cs b/Services/Hosting/LogRedactor.Tlds.cs new file mode 100644 index 0000000..e6b3d7c --- /dev/null +++ b/Services/Hosting/LogRedactor.Tlds.cs @@ -0,0 +1,110 @@ +namespace Craft.Hosting; + +public static partial class LogRedactor +{ + // IANA root zone TLDs (Version 2026093003, Last Updated Thu Oct 1 07:07:01 2026 UTC), from https://data.iana.org/TLD/tlds-alpha-by-domain.txt + private const string TldList = + "aaa aarp abb abbott abbvie abc able abogado abudhabi ac academy accenture accountant " + + "accountants aco actor ad ads adult ae aeg aero aetna af afl africa ag agakhan agency ai aig " + + "airbus airforce airtel akdn al alibaba alipay allfinanz allstate ally alsace alstom am amazon " + + "americanexpress americanfamily amex amfam amica amsterdam analytics android anquan anz ao aol " + + "apartments app apple aq aquarelle ar arab aramco archi army arpa art arte as asda asia " + + "associates at athleta attorney au auction audi audible audio auspost author auto autos aw aws " + + "ax axa az azure ba baby baidu banamex band bank bar barcelona barclaycard barclays barefoot " + + "bargains baseball basketball bauhaus bayern bb bbc bbt bbva bcg bcn bd be beats beauty beer " + + "berlin best bestbuy bet bf bg bh bharti bi bible bid bike bing bingo bio biz bj black " + + "blackfriday blockbuster blog bloomberg blue bm bms bmw bn bnpparibas bo boats boehringer bofa " + + "bom bond boo book booking bosch bostik boston bot boutique box br bradesco bridgestone broadway " + + "broker brother brussels bs bt build builders business buy buzz bv bw by bz bzh ca cab cafe cal " + + "call calvinklein cam camera camp canon capetown capital capitalone car caravan cards care " + + "career careers cars casa case cash casino cat catering catholic cba cbn cbre cc cd center ceo " + + "cern cf cfa cfd cg ch chanel channel charity chase chat cheap chintai christmas chrome church " + + "ci cipriani circle cisco citadel citi citic city ck cl claims cleaning click clinic clinique " + + "clothing cloud club clubmed cm cn co coach codes coffee college cologne com commbank community " + + "company compare computer comsec condos construction consulting contact contractors cooking cool " + + "coop corsica country coupon coupons courses cpa cr credit creditcard creditunion cricket crown " + + "crs cruise cruises cu cuisinella cv cw cx cy cymru cyou cz dad dance data date dating datsun " + + "day dclk dds de deal dealer deals degree delivery dell deloitte delta democrat dental dentist " + + "desi design dev dhl diamonds diet digital direct directory discount discover dish diy dj dk dm " + + "dnp do docs doctor dog domains dot download drive dtv dubai dupont durban dvag dvr dz earth eat " + + "ec eco edeka edu education ee eg email emerck energy engineer engineering enterprises epson " + + "equipment er ericsson erni es esq estate et eu eurovision eus events exchange expert exposed " + + "express extraspace fage fail fairwinds faith family fan fans farm farmers fashion fast fedex " + + "feedback ferrari ferrero fi fidelity fido film final finance financial fire firestone firmdale " + + "fish fishing fit fitness fj fk flickr flights flir florist flowers fly fm fo foo food football " + + "ford forex forsale forum foundation fox fr free fresenius frl frogans frontier ftr fujitsu fun " + + "fund furniture futbol fyi ga gal gallery gallo gallup game games gap garden gay gb gbiz gd gdn " + + "ge gea gent genting george gf gg ggee gh gi gift gifts gives giving gl glass gle global globo " + + "gm gmail gmbh gmo gmx gn godaddy gold goldpoint golf goodyear goog google gop got gov gp gq gr " + + "grainger graphics gratis green gripe grocery group gs gt gu gucci guge guide guitars guru gw gy " + + "hair hamburg hangout haus hbo hdfc hdfcbank health healthcare help helsinki here hermes hiphop " + + "hisamitsu hitachi hiv hk hkt hm hn hockey holdings holiday homedepot homegoods homes homesense " + + "honda horse hospital host hosting hot hotels hotmail house how hr hsbc ht hu hughes hyatt " + + "hyundai ibm icbc ice icu id ie ieee ifm ikano il im imamat imdb immo immobilien in inc " + + "industries infiniti info ing ink institute insurance insure int international intuit " + + "investments io ipiranga iq ir irish is ismaili ist istanbul it itau itv jaguar java jcb je jeep " + + "jetzt jewelry jio jll jm jmp jnj jo jobs joburg jot joy jp jpmorgan jprs juegos kaufen kddi ke " + + "kerryhotels kerryproperties kfh kg kh ki kia kids kim kindle kitchen kiwi km kn koeln komatsu " + + "kosher kp kpmg kpn kr krd kred kuokgroup kw ky kyoto kz la lacaixa lamborghini lamer land " + + "landrover lanxess lasalle lat latino latrobe law lawyer lb lc lds lease leclerc lefrak legal " + + "lego lexus lgbt li lidl life lifeinsurance lifestyle lighting like lilly limited limo lincoln " + + "link live living lk llc llp loan loans locker locus lol london lotte lotto love lpl " + + "lplfinancial lr ls lt ltd ltda lu lundbeck luxe luxury lv ly ma madrid maif maison makeup man " + + "management mango map market marketing markets marriott marshalls mattel mba mc mckinsey md me " + + "med media meet melbourne meme memorial men menu merck merckmsd mg mh miami microsoft mil mini " + + "mint mit mitsubishi mk ml mlb mls mm mma mn mo mobi mobile moda moe moi mom monash money " + + "monster mormon mortgage moscow moto motorcycles mov movie mp mq mr ms msd mt mtn mtr mu museum " + + "music mv mw mx my mz na nab nagoya name navy nba nc ne nec net netbank netflix network neustar " + + "new news next nextdirect nexus nf nfl ng ngo nhk ni nico nike nikon ninja nissan nissay nl no " + + "nokia norton now nowruz nowtv np nr nra nrw ntt nu nyc nz obi observer office okinawa olayan " + + "olayangroup ollo om omega one ong onl online ooo open oracle orange org organic origins osaka " + + "otsuka ott ovh pa page panasonic paris pars partners parts party pay pccw pe pet pf pfizer pg " + + "ph pharmacy phd philips phone photo photography photos physio pics pictet pictures pid pin ping " + + "pink pioneer pizza pk pl place play playstation plumbing plus pm pn pnc pohl poker politie porn " + + "post pr praxi press prime pro prod productions prof progressive promo properties property " + + "protection pru prudential ps pt pub pw pwc py qa qpon quebec quest racing radio re read " + + "realestate realtor realty recipes red redumbrella rehab reise reisen reit reliance ren rent " + + "rentals repair report republican rest restaurant review reviews rexroth rich richardli ricoh " + + "ril rio rip ro rocks rodeo rogers room rs rsvp ru rugby ruhr run rw rwe ryukyu sa saarland safe " + + "safety sakura sale salon samsclub samsung sandvik sandvikcoromant sanofi sap sarl sas save saxo " + + "sb sbi sbs sc scb schaeffler schmidt scholarships school schule schwarz science scot sd se " + + "search seat secure security seek select sener services seven sew sex sexy sfr sg sh shangrila " + + "sharp shell shia shiksha shoes shop shopping shouji show si silk sina singles site sj sk ski " + + "skin sky skype sl sling sm smart smile sn sncf so soccer social softbank software sohu solar " + + "solutions song sony soy spa space sport spot sr srl ss st stada staples star statebank " + + "statefarm stc stcgroup stockholm storage store stream studio study style su sucks supplies " + + "supply support surf surgery suzuki sv swatch swiss sx sy sydney systems sz tab taipei talk " + + "taobao target tatamotors tatar tattoo tax taxi tc tci td tdk team tech technology tel temasek " + + "tennis teva tf tg th thd theater theatre tiaa tickets tienda tips tires tirol tj tjmaxx tjx tk " + + "tkmaxx tl tm tmall tn to today tokyo tools top toray toshiba total tours town toyota toys tr " + + "trade trading training travel travelers travelersinsurance trust trv tt tube tui tunes tushu tv " + + "tvs tw tz ua ubank ubs ug uk unicom university uno uol ups us uy uz va vacations vana vanguard " + + "vc ve vegas ventures verisign versicherung vet vg vi viajes video vig viking villas vin vip " + + "virgin visa vision viva vivo vlaanderen vn vodka volvo vote voting voto voyage vu wales walmart " + + "walter wang wanggou watch watches weather weatherchannel web webcam weber website wed wedding " + + "weibo weir wf whoswho wien wiki williamhill win windows wine winners wme woodside work works " + + "world wow ws wtc wtf xbox xerox xihuan xin xn--11b4c3d xn--1ck2e1b xn--1qqw23a xn--2scrj9c " + + "xn--30rr7y xn--3bst00m xn--3ds443g xn--3e0b707e xn--3hcrj9c xn--3pxu8k xn--42c2d9a xn--45br5cyl " + + "xn--45brj9c xn--45q11c xn--4dbrk0ce xn--4gbrim xn--54b7fta0cc xn--55qw42g xn--55qx5d " + + "xn--5su34j936bgsg xn--5tzm5g xn--6frz82g xn--6qq986b3xl xn--80adxhks xn--80ao21a xn--80aqecdr1a " + + "xn--80asehdb xn--80aswg xn--8y0a063a xn--90a3ac xn--90ae xn--90ais xn--9dbq2a xn--9et52u " + + "xn--9krt00a xn--b4w605ferd xn--bck1b9a5dre4c xn--c1avg xn--c2br7g xn--cck2b3b xn--cckwcxetd " + + "xn--cg4bki xn--clchc0ea0b2g2a9gcd xn--czr694b xn--czrs0t xn--czru2d xn--d1acj3b xn--d1alf " + + "xn--e1a4c xn--eckvdtc9d xn--efvy88h xn--fct429k xn--fhbei xn--fiq228c5hs xn--fiq64b xn--fiqs8s " + + "xn--fiqz9s xn--fjq720a xn--flw351e xn--fpcrj9c3d xn--fzc2c9e2c xn--fzys8d69uvgm xn--g2xx48c " + + "xn--gckr3f0f xn--gecrj9c xn--gk3at1e xn--h2breg3eve xn--h2brj9c xn--h2brj9c8c xn--hxt814e " + + "xn--i1b6b1a6a2e xn--imr513n xn--io0a7i xn--j1aef xn--j1amh xn--j6w193g xn--jlq480n2rg " + + "xn--jvr189m xn--kcrx77d1x4a xn--kprw13d xn--kpry57d xn--kput3i xn--l1acc xn--lgbbat1ad8j " + + "xn--mgb9awbf xn--mgba3a3ejt xn--mgba3a4f16a xn--mgba7c0bbn0a xn--mgbaam7a8h xn--mgbab2bd " + + "xn--mgbah1a3hjkrd xn--mgbai9azgqp6j xn--mgbayh7gpa xn--mgbbh1a xn--mgbbh1a71e xn--mgbc0a9azcg " + + "xn--mgbca7dzdo xn--mgbcpq6gpa1a xn--mgberp4a5d4ar xn--mgbgu82a xn--mgbi4ecexp xn--mgbpl2fh " + + "xn--mgbt3dhd xn--mgbtx2b xn--mgbx4cd0ab xn--mix891f xn--mk1bu44c xn--mxtq1m xn--ngbc5azd " + + "xn--ngbe9e0a xn--ngbrx xn--node xn--nqv7f xn--nqv7fs00ema xn--nyqy26a xn--o3cw4h xn--ogbpf8fl " + + "xn--otu796d xn--p1acf xn--p1ai xn--pgbs0dh xn--pssy2u xn--q7ce6a xn--q9jyb4c xn--qcka1pmc " + + "xn--qxa6a xn--qxam xn--rhqv96g xn--rovu88b xn--rvc1e0am3e xn--s9brj9c xn--ses554g xn--t60b56a " + + "xn--tckwe xn--tiq49xqyj xn--unup4y xn--vermgensberater-ctb xn--vermgensberatung-pwb xn--vhquv " + + "xn--vuq861b xn--w4r85el8fhu5dnra xn--w4rs40l xn--wgbh1c xn--wgbl6a xn--xhq521b xn--xkc2al3hye2a " + + "xn--xkc2dl3a5ee0h xn--y9a3aq xn--yfro4i67o xn--ygbi2ammx xn--zfr164b xxx xyz yachts yahoo " + + "yamaxun yandex ye yodobashi yoga yokohama you youtube yt yun za zappos zara zero zip zm zone " + + "zuerich zw "; +} diff --git a/Services/Hosting/LogRedactor.cs b/Services/Hosting/LogRedactor.cs new file mode 100644 index 0000000..ed80492 --- /dev/null +++ b/Services/Hosting/LogRedactor.cs @@ -0,0 +1,307 @@ +using System.Collections.Concurrent; +using System.Security.Cryptography; +using System.Text; +using System.Text.RegularExpressions; +using Craft.Storage; + +namespace Craft.Hosting; + +/// +/// Masks UPNs and customer domains in log lines while keeping them recognisable at a glance: +/// jane.doe@contoso.onmicrosoft.com becomes ja~…~@co~…~.onmicrosoft.com. The hidden middle +/// is AES-encrypted with a fixed IV, so a value always yields the same token (raw files stay greppable +/// and correlatable) and restores it with the instance key. Tenant ids and other +/// GUIDs are left alone. +/// +public static partial class LogRedactor +{ + private const string KeyTable = "CraftInstanceKeys"; + private const string KeyPartition = "LogRedaction"; + private const string KeyRow = "v1"; + private const int MaxCacheEntries = 50_000; + + private static readonly HashSet s_tlds = + new(TldList.Split(' ', StringSplitOptions.RemoveEmptyEntries), StringComparer.Ordinal); + + // Kept readable; the label in front of them is masked. + private static readonly string[] s_suffixes = + [ + "mail.protection.outlook.com", "onmicrosoft.com", "sharepoint.com", + "co.uk", "org.uk", "ac.uk", "gov.uk", "ltd.uk", "plc.uk", "uk.net", + "com.au", "net.au", "org.au", "edu.au", "gov.au", "co.nz", "org.nz", "co.za", "com.br", + ]; + + // Platform hosts left in clear, along with their subdomains. + private static readonly string[] s_builtInAllow = + [ + "microsoft.com", "microsoftonline.com", "windows.net", "office365.com", "office.com", "office.net", + "outlook.com", "live.com", "azure.com", "azure.net", "azurewebsites.net", "azure-api.net", + "windowsazure.com", "cloud.microsoft", "microsoft", "github.com", "githubusercontent.com", "ghcr.io", + "powershellgallery.com", "nuget.org", + ]; + + // Bytes of HMAC carried in each token: it seeds the keystream and rejects tokens from another key. + private const int TagBytes = 3; + private static readonly object s_aesLock = new(); + private static byte[] s_key = RandomNumberGenerator.GetBytes(16); + private static Aes s_aes = CreateAes(s_key); + private static string[] s_allow = s_builtInAllow; + private static volatile bool s_enabled = true; + private static readonly ConcurrentDictionary s_encrypted = new(StringComparer.Ordinal); + private static readonly ConcurrentDictionary s_decrypted = new(StringComparer.Ordinal); + private static readonly ConcurrentDictionary s_domainMasks = new(StringComparer.Ordinal); + private static readonly ConcurrentDictionary s_emailHostMasks = new(StringComparer.Ordinal); + + [GeneratedRegex(@"~(?[A-Za-z0-9_-]{6,})~", RegexOptions.CultureInvariant | RegexOptions.ExplicitCapture)] + private static partial Regex TokenPattern(); + + public static void Configure(bool enabled, IEnumerable? allowDomains) + { + s_enabled = enabled; + s_allow = [.. s_builtInAllow, .. (allowDomains ?? []).Where(d => !string.IsNullOrWhiteSpace(d)) + .Select(d => d.Trim().TrimStart('.').ToLowerInvariant())]; + s_domainMasks.Clear(); + s_emailHostMasks.Clear(); + } + + public static void SetKey(byte[] key) + { + lock (s_aesLock) + { + s_key = key; + s_aes = CreateAes(key); + s_encrypted.Clear(); + s_decrypted.Clear(); + s_domainMasks.Clear(); + s_emailHostMasks.Clear(); + } + } + + /// + /// Loads the instance key from the Craft key table, creating it on first start. Every node of an + /// instance shares it, so any node can reveal any node's lines. + /// + public static async Task LoadKeyAsync(ICraftTableStore store, CancellationToken ct = default) + { + ArgumentNullException.ThrowIfNull(store); + await store.EnsureTableAsync(KeyTable, ct); + var row = await store.GetAsync(KeyTable, KeyPartition, KeyRow, ct); + if (row?.GetString("Key") is not { Length: > 0 }) + { + var created = new StoreRow(KeyPartition, KeyRow); + created["Key"] = Convert.ToBase64String(RandomNumberGenerator.GetBytes(16)); + await store.UpsertAsync(KeyTable, created, ct); + // Re-read so two nodes racing on first start settle on the same key. + row = await store.GetAsync(KeyTable, KeyPartition, KeyRow, ct); + } + SetKey(Convert.FromBase64String(row!.GetString("Key")!)); + } + + /// + /// One pass over the line, jumping between dots: each dotted run is a candidate host, and a run + /// preceded by @ or %40 takes its local part with it. + /// + public static string Redact(string line) + { + if (!s_enabled || string.IsNullOrEmpty(line)) return line; + StringBuilder? sb = null; + var copied = 0; + var dot = line.IndexOf('.'); + while (dot >= 0) + { + var start = dot; + while (start > 0 && IsHostChar(line[start - 1])) start--; + var end = dot + 1; + while (end < line.Length && IsHostChar(line[end])) end++; + var resume = end; + + start = SkipEscape(line, start, end); + while (start < end && line[start] is '.' or '-') start++; + while (end > start && line[end - 1] is '.' or '-') end--; + + if (end - start >= 3 && !IsInLocalPart(line, end)) + { + var atLength = start >= 1 && line[start - 1] == '@' ? 1 + : start >= 3 && line.AsSpan(start - 3, 3).SequenceEqual("%40") ? 3 : 0; + var host = line[start..end]; + var masked = CachedMaskHost(host, isEmail: atLength > 0); + + var from = start; + var replacement = masked; + if (atLength > 0) + { + var localEnd = start - atLength; + var localStart = localEnd; + while (localStart > copied && IsLocalChar(line[localStart - 1])) localStart--; + localStart = SkipEscape(line, localStart, localEnd); + if (localEnd > localStart) + { + from = localStart; + replacement = string.Concat(MaskLabel(line[localStart..localEnd]), + line.AsSpan(localEnd, atLength), masked ?? host); + } + } + + if (replacement != null) + { + sb ??= new StringBuilder(line.Length + 64); + sb.Append(line, copied, from - copied).Append(replacement); + copied = end; + } + } + + dot = resume < line.Length ? line.IndexOf('.', resume) : -1; + } + return sb == null ? line : sb.Append(line, copied, line.Length - copied).ToString(); + } + + private static string? CachedMaskHost(string host, bool isEmail) + { + var cache = isEmail ? s_emailHostMasks : s_domainMasks; + if (cache.TryGetValue(host, out var masked)) return masked; + masked = MaskHost(host, isEmail); + if (cache.Count >= MaxCacheEntries) cache.Clear(); + cache.TryAdd(host, masked); + return masked; + } + + private static bool IsHostChar(char c) => char.IsAsciiLetterOrDigit(c) || c is '.' or '-'; + + private static bool IsLocalChar(char c) => char.IsAsciiLetterOrDigit(c) || c is '.' or '-' or '_' or '+' or '\'' or '#'; + + // %27contoso.com%27: the two hex digits after % belong to the escape, not the name. + private static int SkipEscape(string line, int start, int end) => + start > 0 && line[start - 1] == '%' && end - start > 2 + && char.IsAsciiHexDigit(line[start]) && char.IsAsciiHexDigit(line[start + 1]) ? start + 2 : start; + + // A dotted run that carries on into an @ is the local part of an email, handled with its host. + private static bool IsInLocalPart(string line, int end) + { + var i = end; + while (i < line.Length && IsLocalChar(line[i])) i++; + return i < line.Length && (line[i] == '@' || line.AsSpan(i).StartsWith("%40")); + } + + /// Restores masked values in a line written by this instance. Unknown tokens are left as is. + public static string Reveal(string line) + { + if (string.IsNullOrEmpty(line) || line.IndexOf('~') < 0) return line; + return TokenPattern().Replace(line, static m => Decrypt(m.Groups["blob"].Value) ?? m.Value); + } + + // Null leaves the host as written. + private static string? MaskHost(string host, bool isEmail) + { + if (!host.Contains('.', StringComparison.Ordinal) || host.Contains("..", StringComparison.Ordinal) + || !char.IsAsciiLetter(host[host.LastIndexOf('.') + 1])) + return null; + var lower = host.ToLowerInvariant(); + var suffixLength = 0; + foreach (var suffix in s_suffixes) + { + if (lower == suffix) return null; + if (lower.Length > suffix.Length && lower.EndsWith(suffix, StringComparison.Ordinal) + && lower[^(suffix.Length + 1)] == '.') + { + suffixLength = suffix.Length; + break; + } + } + + if (suffixLength == 0) + { + foreach (var allowed in s_allow) + { + if (lower == allowed || (lower.EndsWith(allowed, StringComparison.Ordinal) + && lower[^(allowed.Length + 1)] == '.')) + return null; + } + suffixLength = host.Length - host.LastIndexOf('.') - 1; + if (!isEmail && !s_tlds.Contains(host[^suffixLength..])) return null; + } + + var name = host[..^(suffixLength + 1)]; + var labelStart = name.LastIndexOf('.') + 1; + return name[..labelStart] + MaskLabel(name[labelStart..]) + host[^(suffixLength + 1)..]; + } + + // Roughly half of a label stays readable, split across both ends; short ones keep only the first character. + private static string MaskLabel(string value) + { + if (value.Length < 2) return value; + var (head, tail) = value.Length switch + { + <= 3 => (1, 0), + <= 5 => (1, 1), + <= 8 => (2, 1), + _ => (2, 2), + }; + return value[..head] + "~" + Encrypt(value[head..^tail].ToLowerInvariant()) + "~" + value[^tail..]; + } + + // Deterministic and length-preserving: a short HMAC tag of the value seeds an AES-CTR keystream, + // so a token is the tag plus one byte per hidden character. + private static string Encrypt(string hidden) + { + if (s_encrypted.TryGetValue(hidden, out var cached)) return cached; + var plain = Encoding.UTF8.GetBytes(hidden); + var output = new byte[TagBytes + plain.Length]; + lock (s_aesLock) + { + HMACSHA256.HashData(s_key, plain).AsSpan(0, TagBytes).CopyTo(output); + ApplyKeystream(output.AsSpan(0, TagBytes), plain, output.AsSpan(TagBytes)); + } + var token = Convert.ToBase64String(output).TrimEnd('=').Replace('+', '-').Replace('/', '_'); + if (s_encrypted.Count >= MaxCacheEntries) s_encrypted.Clear(); + s_encrypted[hidden] = token; + return token; + } + + private static string? Decrypt(string token) + { + if (s_decrypted.TryGetValue(token, out var cached)) return cached; + string? plain = null; + var b64 = token.Replace('-', '+').Replace('_', '/'); + b64 += new string('=', (4 - b64.Length % 4) % 4); + try + { + var bytes = Convert.FromBase64String(b64); + if (bytes.Length > TagBytes) + { + var decoded = new byte[bytes.Length - TagBytes]; + bool valid; + lock (s_aesLock) + { + ApplyKeystream(bytes.AsSpan(0, TagBytes), bytes.AsSpan(TagBytes), decoded); + valid = HMACSHA256.HashData(s_key, decoded).AsSpan(0, TagBytes).SequenceEqual(bytes.AsSpan(0, TagBytes)); + } + if (valid) plain = Encoding.UTF8.GetString(decoded); + } + } + catch (FormatException) { } + if (s_decrypted.Count >= MaxCacheEntries) s_decrypted.Clear(); + s_decrypted[token] = plain; + return plain; + } + + private static void ApplyKeystream(ReadOnlySpan tag, ReadOnlySpan input, Span output) + { + Span counter = stackalloc byte[16]; + Span block = stackalloc byte[16]; + tag.CopyTo(counter); + for (var offset = 0; offset < input.Length; offset += 16) + { + counter[15] = (byte)(offset / 16); + s_aes.EncryptEcb(counter, block, PaddingMode.None); + for (var i = 0; i < 16 && offset + i < input.Length; i++) + output[offset + i] = (byte)(input[offset + i] ^ block[i]); + } + } + + private static Aes CreateAes(byte[] key) + { + var aes = Aes.Create(); + aes.Key = key; + return aes; + } +} diff --git a/Services/Hosting/RedactingConsoleLoggerProvider.cs b/Services/Hosting/RedactingConsoleLoggerProvider.cs new file mode 100644 index 0000000..39eb140 --- /dev/null +++ b/Services/Hosting/RedactingConsoleLoggerProvider.cs @@ -0,0 +1,37 @@ +using Microsoft.Extensions.Logging.Console; + +namespace Craft.Hosting; + +/// +/// Puts the console sink behind , leaving its output format unchanged. +/// +[ProviderAlias("Console")] +public sealed class RedactingConsoleLoggerProvider(ConsoleLoggerProvider inner) : ILoggerProvider, ISupportExternalScope +{ + public ILogger CreateLogger(string categoryName) => new RedactingLogger(inner.CreateLogger(categoryName)); + + public void SetScopeProvider(IExternalScopeProvider scopeProvider) => inner.SetScopeProvider(scopeProvider); + + // The container owns and disposes the inner provider. + public void Dispose() { } + + private sealed class RedactingLogger(ILogger inner) : ILogger + { + public IDisposable? BeginScope(TState state) where TState : notnull => inner.BeginScope(state); + + public bool IsEnabled(LogLevel logLevel) => inner.IsEnabled(logLevel); + + public void Log(LogLevel logLevel, EventId eventId, TState state, Exception? exception, + Func formatter) + { + if (!inner.IsEnabled(logLevel)) return; + inner.Log(logLevel, eventId, LogRedactor.Redact(formatter(state, exception)), + exception is null ? null : new RedactedException(exception), static (message, _) => message); + } + } + + private sealed class RedactedException(Exception original) : Exception(LogRedactor.Redact(original.Message)) + { + public override string ToString() => LogRedactor.Redact(original.ToString()); + } +} diff --git a/Services/Program.cs b/Services/Program.cs index 0263819..414bd3d 100644 --- a/Services/Program.cs +++ b/Services/Program.cs @@ -102,6 +102,15 @@ var repo = app.Services.GetRequiredService(); var pool = app.Services.GetRequiredService(); var logger = app.Services.GetRequiredService>(); +try +{ + using var keyTimeout = new CancellationTokenSource(TimeSpan.FromSeconds(15)); + await LogRedactor.LoadKeyAsync(app.Services.GetRequiredService(), keyTimeout.Token); +} +catch (Exception ex) +{ + logger.LogWarning("[System] Log redaction key unavailable, masked values from this run cannot be revealed: {Error}", ex.Message); +} var psRunner = app.Services.GetRequiredService(); var cache = app.Services.GetRequiredService(); var CraftSettings = app.Services.GetRequiredService(); diff --git a/tests/Craft.Tests/LogRedactorTests.cs b/tests/Craft.Tests/LogRedactorTests.cs new file mode 100644 index 0000000..38fa7e9 --- /dev/null +++ b/tests/Craft.Tests/LogRedactorTests.cs @@ -0,0 +1,174 @@ +using Craft.Configuration; +using Craft.Hosting; +using Craft.Services; +using Microsoft.AspNetCore.Builder; +using Microsoft.Extensions.DependencyInjection; +using Microsoft.Extensions.Logging; + +namespace Craft.Tests; + +/// +/// Log lines must hide UPNs and customer domains while staying greppable and reversible: the same value +/// always masks to the same token, and LogBridge reveals it with the instance key. +/// +[Collection(nameof(LogRedactorTests))] +public class LogRedactorTests +{ + private static readonly byte[] Key = Enumerable.Range(1, 16).Select(i => (byte)i).ToArray(); + + public LogRedactorTests() + { + LogRedactor.Configure(enabled: true, allowDomains: ["cipp.app"]); + LogRedactor.SetKey(Key); + } + + [Theory] + [InlineData("Processing jane.doe@contoso.com now")] + [InlineData("tenant contoso.onmicrosoft.com failed")] + [InlineData("Started: MailboxRules_contoso.onmicrosoft.com#GetMailbox")] + [InlineData("GET https://graph.microsoft.com/v1.0/users/jane.doe%40contoso.com/mailFolders")] + [InlineData("$filter=domain eq %27contoso.co.uk%27")] + [InlineData("site contoso-my.sharepoint.com and autodiscover.contoso.com")] + [InlineData("guest john_fabrikam.com#EXT#@contoso.onmicrosoft.com signed in")] + [InlineData("first user Jane.Doe@Contoso.com")] + public void MaskedLines_RevealToTheOriginal_CaseFoldedInTheHiddenPart(string line) + { + var masked = LogRedactor.Redact(line); + + Assert.NotEqual(line, masked); + Assert.DoesNotContain("contoso", masked, StringComparison.OrdinalIgnoreCase); + Assert.Equal(line, LogRedactor.Reveal(masked), StringComparer.OrdinalIgnoreCase); + } + + [Fact] + public void Mask_KeepsPrefixAndSuffixReadable() + { + var masked = LogRedactor.Redact("jane.doe@contoso.onmicrosoft.com"); + + Assert.Matches(@"^ja~[A-Za-z0-9_-]{11}~e@co~[A-Za-z0-9_-]{10}~o\.onmicrosoft\.com$", masked); + } + + [Fact] + public void SameDomain_MasksToTheSameToken_Everywhere() + { + var domain = LogRedactor.Redact("contoso.com"); + + Assert.Contains(domain, LogRedactor.Redact("jane@contoso.com"), StringComparison.Ordinal); + Assert.Contains(domain, LogRedactor.Redact("x %27contoso.com%27 y"), StringComparison.Ordinal); + Assert.Contains(LogRedactor.Redact("contoso.onmicrosoft.com")[..^".onmicrosoft.com".Length], + LogRedactor.Redact("MailboxRules_contoso.onmicrosoft.com"), StringComparison.Ordinal); + } + + [Theory] + [InlineData("https://graph.microsoft.com/v1.0/users?$top=999")] + [InlineData("https://cippstg.table.core.windows.net/CippLogs()")] + [InlineData("https://outlook.office365.com/adminapi/beta/x/InvokeCommand")] + [InlineData("https://cipp.app/docs")] + [InlineData("needs Tenant.Reports.Read and CIPP.Dashboard.Read")] + [InlineData("System.Net.Http.HttpRequestException in Microsoft.Web")] + [InlineData("rotated craft.1.log, loaded organizationmanagementroles.json, ran Push-Thing.ps1")] + [InlineData("tenant 3f9a1c2b-7d4e-4f00-9a1b-2c3d4e5f6a7b v1.0.13 10.0.0.1")] + public void InfrastructureAndNonDomains_AreLeftAlone(string line) + { + Assert.Equal(line, LogRedactor.Redact(line)); + } + + [Fact] + public void Redact_IsIdempotent() + { + var once = LogRedactor.Redact("jane.doe@contoso.com in contoso.onmicrosoft.com"); + + Assert.Equal(once, LogRedactor.Redact(once)); + } + + [Fact] + public void AllowlistedEmailDomain_StillMasksTheUser() + { + Assert.Matches(@"^a~[A-Za-z0-9_-]{8}~n@microsoft\.com$", LogRedactor.Redact("admin@microsoft.com")); + } + + [Fact] + public void Reveal_WithAnotherKey_LeavesTokensInPlace() + { + var masked = LogRedactor.Redact("jane.doe@contoso.com"); + LogRedactor.SetKey(new byte[16]); + + Assert.Equal(masked, LogRedactor.Reveal(masked)); + } + + [Fact] + public void Disabled_PassesLinesThrough() + { + LogRedactor.Configure(enabled: false, allowDomains: null); + + Assert.Equal("jane@contoso.com", LogRedactor.Redact("jane@contoso.com")); + } + + [Fact] + public void FileLog_IsMaskedOnDisk_AndRevealedThroughLogBridge() + { + var dir = Path.Combine(Path.GetTempPath(), "craft-redact-" + Guid.NewGuid().ToString("N")); + try + { + var provider = new FileLoggerProvider(new FileLoggingSettings { Directory = dir }); + provider.CreateLogger("Test").LogWarning("PS warning: mailbox jane.doe@contoso.com in contoso.onmicrosoft.com"); + provider.Dispose(); + + var raw = File.ReadAllText(Path.Combine(dir, "craft.log")); + Assert.DoesNotContain("contoso", raw, StringComparison.OrdinalIgnoreCase); + Assert.Contains("ja~", raw, StringComparison.Ordinal); + + LogBridge.Initialize(provider); + var line = Assert.Single(LogBridge.ReadLog(0, null, "jane.doe@contoso.com"), l => l.Length > 0); + Assert.Contains("mailbox jane.doe@contoso.com in contoso.onmicrosoft.com", line, StringComparison.Ordinal); + Assert.Single(LogBridge.ReadLog(0, null, null, null, null, null, null, @"contoso\.onmicrosoft"), l => l.Length > 0); + } + finally + { + Directory.Delete(dir, recursive: true); + } + } + + [Fact] + public void ConsoleSink_IsMasked() + { + var output = new StringWriter(); + var original = Console.Out; + Console.SetOut(output); + try + { + var builder = WebApplication.CreateBuilder(); + builder.Configuration["App:FileLogging:Directory"] = + Path.Combine(Path.GetTempPath(), "craft-redact-" + Guid.NewGuid().ToString("N")); + builder.AddCraftLogging(); + using (var services = builder.Services.BuildServiceProvider()) + { + services.GetRequiredService().CreateLogger("Test") + .LogError(new InvalidOperationException("denied for jane.doe@contoso.com"), "auth failed for jane.doe@contoso.com"); + } + } + finally + { + Console.SetOut(original); + } + + var text = output.ToString(); + Assert.Contains("auth failed for ja~", text, StringComparison.Ordinal); + Assert.Contains("denied for ja~", text, StringComparison.Ordinal); + Assert.DoesNotContain("contoso", text, StringComparison.OrdinalIgnoreCase); + } + + [Fact] + public async Task Key_IsCreatedOnce_AndReusedOnTheNextStart() + { + var store = new FakeTableStore(); + await LogRedactor.LoadKeyAsync(store); + var first = LogRedactor.Redact("jane@contoso.com"); + + LogRedactor.SetKey(new byte[16]); + await LogRedactor.LoadKeyAsync(store); + + Assert.Equal(first, LogRedactor.Redact("jane@contoso.com")); + Assert.Equal(1, store.Count("CraftInstanceKeys")); + } +} From 3c63d5196c2dfca168f734465ce261c563a5f0f1 Mon Sep 17 00:00:00 2001 From: Zacgoose <107489668+Zacgoose@users.noreply.github.com> Date: Fri, 2 Oct 2026 17:30:37 +0800 Subject: [PATCH 3/4] feat(egress): break API egress out per endpoint and track signed-in user traffic - record each response against its endpoint (bytes, requests, max size, cache hits, errors, shed) on daily client rows and instance 15-minute buckets - handlers can sub-label traffic with an X-Craft-Endpoint response header, which is stripped before the response is sent - signed-in users are tracked in an 'interactive' partition that never counts toward the cap - instance daily request count now includes clients idle since an earlier bucket --- .../Hosting/ApiEgressLimiterMiddleware.cs | 8 +- .../Hosting/ApiEgressWireCounterMiddleware.cs | 47 ++++- Services/Hosting/EgressLedger.cs | 170 ++++++++++++++---- Services/Hosting/EgressTableSchema.cs | 127 ++++++++++++- .../Endpoints/PowerShellDispatchEndpoint.cs | 2 + .../ApiCompressionPipelineTests.cs | 64 +++++++ .../ApiEgressWireCounterMiddlewareTests.cs | 73 ++++++++ tests/Craft.Tests/EgressLedgerTableTests.cs | 127 +++++++++++++ 8 files changed, 567 insertions(+), 51 deletions(-) diff --git a/Services/Hosting/ApiEgressLimiterMiddleware.cs b/Services/Hosting/ApiEgressLimiterMiddleware.cs index ba01351..1a2c693 100644 --- a/Services/Hosting/ApiEgressLimiterMiddleware.cs +++ b/Services/Hosting/ApiEgressLimiterMiddleware.cs @@ -41,10 +41,12 @@ public async Task InvokeAsync(HttpContext context) { ArgumentNullException.ThrowIfNull(context); - // UI and anonymous callers are never counted or capped — the feature is only about app-only - // automation. One header check and out, so interactive traffic pays essentially nothing. + // UI and anonymous callers are never capped — the cap is only about app-only automation. A + // signed-in user's traffic is still flagged so it is reported (never billed); anonymous is not. if (!CallerClassifier.IsApiClient(context)) { + if (!string.IsNullOrEmpty(context.Request.Headers["x-ms-client-principal-name"].ToString())) + context.Items[ApiEgressWireCounterMiddleware.InteractiveItemKey] = true; await _next(context); return; } @@ -59,7 +61,7 @@ public async Task InvokeAsync(HttpContext context) // concurrency cap, the overshoot is bounded to (concurrency × largest response). if (_ledger.ShouldReject()) { - _ledger.RecordShed(appId); + _ledger.RecordShed(appId, context.GetRouteValue("endpoint") as string); await RejectAsync(context); return; } diff --git a/Services/Hosting/ApiEgressWireCounterMiddleware.cs b/Services/Hosting/ApiEgressWireCounterMiddleware.cs index e38b735..fc68b6c 100644 --- a/Services/Hosting/ApiEgressWireCounterMiddleware.cs +++ b/Services/Hosting/ApiEgressWireCounterMiddleware.cs @@ -34,6 +34,22 @@ public sealed class ApiEgressWireCounterMiddleware /// public const string ChargeItemKey = "Craft.Egress.Charge"; + /// + /// Request item set by on a request from an interactive + /// (signed-in user) caller: its bytes are recorded for reporting but never billed against the cap. + /// + public const string InteractiveItemKey = "Craft.Egress.Interactive"; + + /// Request item carrying the accounting label of a dispatched endpoint — see . + public const string EndpointItemKey = "Craft.Egress.Endpoint"; + + /// + /// Optional response header a handler sets to sub-label its traffic (e.g. the Graph resource behind a + /// generic proxy endpoint, or the tool behind an MCP call). Consumed by and + /// never sent to the client. + /// + public const string EndpointHeader = "X-Craft-Endpoint"; + private readonly RequestDelegate _next; private readonly EgressLedger _ledger; @@ -62,7 +78,36 @@ public async Task InvokeAsync(HttpContext context) { context.Response.Body = original; if (context.Items.TryGetValue(ChargeItemKey, out var charge) && charge is string appId && appId.Length > 0) - _ledger.Record(counting.BytesWritten, appId); + _ledger.Record(counting.BytesWritten, appId, ResolveEndpoint(context), IsCacheHit(context), context.Response.StatusCode); + else if (context.Items.ContainsKey(InteractiveItemKey)) + _ledger.RecordInteractive(counting.BytesWritten, ResolveEndpoint(context), IsCacheHit(context), context.Response.StatusCode); } } + + /// + /// Labels a dispatched request for accounting: the endpoint name, suffixed with the handler's + /// value when it set one (ListGraphRequest:users). The header is + /// removed so it never reaches the client. Call once the endpoint is known to exist and after the + /// handler's headers are applied — both the executed and the cache-hit paths. + /// + public static void TagEndpoint(HttpContext context, string endpoint) + { + ArgumentNullException.ThrowIfNull(context); + var tag = context.Response.Headers[EndpointHeader].ToString(); + if (tag.Length > 0) context.Response.Headers.Remove(EndpointHeader); + context.Items[EndpointItemKey] = tag.Length > 0 ? $"{endpoint}:{tag}" : endpoint; + } + + // Dispatched endpoints are tagged explicitly; any other matched /api route is a literal (native or + // built-in) path, labelled by its last segment. Unmatched requests stay unlabelled. + private static string? ResolveEndpoint(HttpContext context) + { + if (context.Items.TryGetValue(EndpointItemKey, out var tagged) && tagged is string label) return label; + if (context.GetEndpoint() is RouteEndpoint { RoutePattern: { Parameters.Count: 0, RawText: { } raw } }) + return raw.TrimEnd('/').Split('/')[^1]; + return null; + } + + private static bool IsCacheHit(HttpContext context) => + context.Response.Headers["X-Cache"].ToString().StartsWith("HIT", StringComparison.Ordinal); } diff --git a/Services/Hosting/EgressLedger.cs b/Services/Hosting/EgressLedger.cs index 519627b..6a66465 100644 --- a/Services/Hosting/EgressLedger.cs +++ b/Services/Hosting/EgressLedger.cs @@ -10,6 +10,9 @@ namespace Craft.Hosting; /// daily totals are the source of truth for the cap and are persisted to a small local file so a restart /// does not reset the day; the time-bucketed history and a durable daily audit are additionally mirrored /// to a table so the product can show usage over time and whether/when/how often the cap was hit. +/// Each response is also attributed to its endpoint (daily per client, and per instance bucket), and +/// interactive user traffic is tracked the same way under its own partition without ever counting +/// toward the cap. /// /// Request path ( / ) only touches memory under a short /// lock — no IO. A background timer flushes the daily totals to disk and mirrors the current 15-minute @@ -50,13 +53,28 @@ public sealed class EgressLedger : BackgroundService private bool _tableDirty; // something changed since the last table sync private readonly Dictionary _clients = new(StringComparer.OrdinalIgnoreCase); - // Pending table buckets not yet finalised: bucketRowKey -> (appId -> accum). + // Interactive (user) traffic: reported in its own partition, never part of _bytes or the cap. + private ClientTotals _interactive = new(); + // Pending table buckets not yet finalised: bucketRowKey -> (partition -> accum). Partitions are the + // API clients' AppIds, the instance-total aggregate and the interactive partition; only the latter + // two carry an endpoint breakdown. private readonly Dictionary> _buckets = new(StringComparer.Ordinal); private bool _tableReady; - private sealed class ClientTotals { public long Bytes; public long Requests; public long Shed; public DateTime LastSeenUtc; } - private sealed class BucketAccum { public long Bytes; public long Requests; public long Shed; } + private sealed class ClientTotals + { + public long Bytes; public long Requests; public long Shed; public DateTime LastSeenUtc; + public Dictionary Endpoints = new(StringComparer.OrdinalIgnoreCase); + public ClientTotals Clone() => new() { Bytes = Bytes, Requests = Requests, Shed = Shed, LastSeenUtc = LastSeenUtc, Endpoints = EndpointStats.CloneMap(Endpoints) }; + } + + private sealed class BucketAccum + { + public long Bytes; public long Requests; public long Shed; + public Dictionary Endpoints = new(StringComparer.OrdinalIgnoreCase); + public BucketAccum Clone() => new() { Bytes = Bytes, Requests = Requests, Shed = Shed, Endpoints = EndpointStats.CloneMap(Endpoints) }; + } private static readonly JsonSerializerOptions s_json = new() { @@ -114,20 +132,51 @@ public bool ShouldReject() } } - /// Add the outbound bytes of one response to today's running totals, attributed to a client. - public void Record(long bytes, string appId) + /// Add the outbound bytes of one API-client response to today's running totals, attributed to + /// the client and, in its daily and the instance bucket breakdowns, to . + public void Record(long bytes, string appId, string? endpoint = null, bool cacheHit = false, int statusCode = 200) { if (bytes <= 0) return; appId = Normalize(appId); + var label = EgressTableSchema.EndpointLabel(endpoint); lock (_lock) { RolloverIfNeeded(); _bytes += bytes; - Client(appId).Bytes += bytes; - Client(appId).Requests += 1; - Client(appId).LastSeenUtc = _utcNow(); - Bucket(appId).Bytes += bytes; - Bucket(appId).Requests += 1; + var client = Client(appId); + client.Bytes += bytes; + client.Requests += 1; + client.LastSeenUtc = _utcNow(); + EndpointStats.In(client.Endpoints, label).Add(bytes, cacheHit, statusCode); + var bucket = Bucket(appId); + bucket.Bytes += bytes; + bucket.Requests += 1; + var system = Bucket(EgressTableSchema.SystemPartition); + system.Bytes += bytes; + system.Requests += 1; + EndpointStats.In(system.Endpoints, label).Add(bytes, cacheHit, statusCode); + _dirty = true; + _tableDirty = true; + } + } + + /// Add the outbound bytes of one interactive (user) response. Reported under + /// ; never counts toward the cap. + public void RecordInteractive(long bytes, string? endpoint = null, bool cacheHit = false, int statusCode = 200) + { + if (bytes <= 0) return; + var label = EgressTableSchema.EndpointLabel(endpoint); + lock (_lock) + { + RolloverIfNeeded(); + _interactive.Bytes += bytes; + _interactive.Requests += 1; + _interactive.LastSeenUtc = _utcNow(); + EndpointStats.In(_interactive.Endpoints, label).Add(bytes, cacheHit, statusCode); + var bucket = Bucket(EgressTableSchema.InteractivePartition); + bucket.Bytes += bytes; + bucket.Requests += 1; + EndpointStats.In(bucket.Endpoints, label).Add(bytes, cacheHit, statusCode); _dirty = true; _tableDirty = true; } @@ -135,16 +184,22 @@ public void Record(long bytes, string appId) /// Record that a request from was shed (429) once over budget. The /// shed response body is deliberately not counted as bytes; this tracks the cap-hit count. - public void RecordShed(string appId) + public void RecordShed(string appId, string? endpoint = null) { appId = Normalize(appId); + var label = EgressTableSchema.EndpointLabel(endpoint); lock (_lock) { RolloverIfNeeded(); _shedRequests += 1; _capReachedUtc ??= _utcNow(); - Client(appId).Shed += 1; + var client = Client(appId); + client.Shed += 1; + EndpointStats.In(client.Endpoints, label).Shed += 1; Bucket(appId).Shed += 1; + var system = Bucket(EgressTableSchema.SystemPartition); + system.Shed += 1; + EndpointStats.In(system.Endpoints, label).Shed += 1; _dirty = true; _tableDirty = true; } @@ -201,6 +256,7 @@ private void RolloverIfNeeded() _shedRequests = 0; _capReachedUtc = null; _clients.Clear(); + _interactive = new ClientTotals(); _buckets.Clear(); // yesterday's bucket rows stay in the table until retention purges them _dirty = true; _tableDirty = true; @@ -262,7 +318,7 @@ private async Task InitTableAsync(CancellationToken ct) } /// Load the persisted daily totals. Restores only when the file is stamped with today's UTC - /// date; a stale or unreadable file leaves today starting from zero. Tolerates the v1 format. + /// date; a stale or unreadable file leaves today starting from zero. Tolerates the v1 and v2 formats. internal void Load() { try @@ -291,8 +347,9 @@ internal void Load() if (state.Clients is not null) { foreach (var (appId, c) in state.Clients) - _clients[appId] = new ClientTotals { Bytes = c.Bytes, Requests = c.Requests, Shed = c.Shed, LastSeenUtc = c.LastSeenUtc }; + _clients[appId] = c.ToTotals(); } + _interactive = state.Interactive?.ToTotals() ?? new ClientTotals(); _dirty = false; } _logger.LogInformation("[Egress] Restored {Bytes} bytes across {Clients} client(s) already served today from {File}", @@ -316,16 +373,15 @@ internal void Flush() RolloverIfNeeded(); snapshot = new LedgerState { - Version = 2, + Version = 3, DateUtc = _dateUtc.ToString("yyyy-MM-dd", CultureInfo.InvariantCulture), Bytes = _bytes, CapBytes = _capBytes, CapReachedUtc = _capReachedUtc, ShedRequests = _shedRequests, UpdatedUtc = _utcNow(), - Clients = _clients.ToDictionary( - kv => kv.Key, - kv => new ClientState { Bytes = kv.Value.Bytes, Requests = kv.Value.Requests, Shed = kv.Value.Shed, LastSeenUtc = kv.Value.LastSeenUtc }), + Clients = _clients.ToDictionary(kv => kv.Key, kv => ClientState.From(kv.Value)), + Interactive = _interactive.Requests > 0 ? ClientState.From(_interactive) : null, }; _dirty = false; } @@ -354,41 +410,54 @@ internal async Task SyncToTableAsync(CancellationToken ct) if (!_tableReady) { try { await _store!.EnsureTableAsync(_tableName, ct).ConfigureAwait(false); _tableReady = true; } catch { return; } } // Snapshot under the lock. - DateOnly day; long dayBytes, dayShed; DateTime? capReached; bool enforcing; + DateOnly day; long dayBytes, dayRequests, dayShed; DateTime? capReached; bool enforcing; string currentBucketKey; - List<(string bucketKey, Dictionary clients)> buckets; + List<(string bucketKey, Dictionary partitions)> buckets; Dictionary clientDaily; + ClientTotals? interactiveDaily; + var instanceEndpoints = new Dictionary(StringComparer.OrdinalIgnoreCase); lock (_lock) { if (!_tableDirty && !_buckets.Keys.Any(k => k != CurrentBucketKey())) return; RolloverIfNeeded(); currentBucketKey = CurrentBucketKey(); - buckets = _buckets.Select(kv => (kv.Key, kv.Value.ToDictionary(c => c.Key, c => new BucketAccum { Bytes = c.Value.Bytes, Requests = c.Value.Requests, Shed = c.Value.Shed }, StringComparer.OrdinalIgnoreCase))).ToList(); - var active = new HashSet(buckets.SelectMany(b => b.clients.Keys), StringComparer.OrdinalIgnoreCase); - clientDaily = active.Where(a => _clients.ContainsKey(a)).ToDictionary(a => a, a => new ClientTotals { Bytes = _clients[a].Bytes, Requests = _clients[a].Requests, Shed = _clients[a].Shed, LastSeenUtc = _clients[a].LastSeenUtc }, StringComparer.OrdinalIgnoreCase); - day = _dateUtc; dayBytes = _bytes; dayShed = _shedRequests; capReached = _capReachedUtc; enforcing = _capBytes > 0; + buckets = _buckets.Select(kv => (kv.Key, kv.Value.ToDictionary(c => c.Key, c => c.Value.Clone(), StringComparer.OrdinalIgnoreCase))).ToList(); + var active = new HashSet(buckets.SelectMany(b => b.partitions.Keys), StringComparer.OrdinalIgnoreCase); + clientDaily = active.Where(a => _clients.ContainsKey(a)).ToDictionary(a => a, a => _clients[a].Clone(), StringComparer.OrdinalIgnoreCase); + interactiveDaily = active.Contains(EgressTableSchema.InteractivePartition) ? _interactive.Clone() : null; + foreach (var client in _clients.Values) + foreach (var (label, stats) in client.Endpoints) + { + if (!instanceEndpoints.TryGetValue(label, out var sum)) instanceEndpoints[label] = sum = new EndpointStats(); + sum.Merge(stats); + } + day = _dateUtc; dayBytes = _bytes; dayRequests = _clients.Values.Sum(c => c.Requests); dayShed = _shedRequests; + capReached = _capReachedUtc; enforcing = _capBytes > 0; _tableDirty = false; } try { - // Bucket rows — per client + a per-bucket instance aggregate. - foreach (var (bucketKey, clients) in buckets) + // Bucket rows — per client, the instance aggregate and interactive; the latter two carry endpoints. + foreach (var (bucketKey, partitions) in buckets) { var bucketStart = ParseBucketStart(bucketKey); - long bBytes = 0, bReq = 0, bShed = 0; - foreach (var (appId, a) in clients) + foreach (var (partition, a) in partitions) { - await _store!.UpsertAsync(_tableName, EgressTableSchema.ClientBucketRow(appId, bucketStart, a.Bytes, a.Requests, a.Shed, _capBytes), ct).ConfigureAwait(false); - bBytes += a.Bytes; bReq += a.Requests; bShed += a.Shed; + var row = partition == EgressTableSchema.SystemPartition + ? EgressTableSchema.SystemBucketRow(bucketStart, a.Bytes, a.Requests, a.Shed, _capBytes, a.Endpoints) + : EgressTableSchema.ClientBucketRow(partition, bucketStart, a.Bytes, a.Requests, a.Shed, _capBytes, a.Endpoints); + await _store!.UpsertAsync(_tableName, row, ct).ConfigureAwait(false); } - await _store!.UpsertAsync(_tableName, EgressTableSchema.SystemBucketRow(bucketStart, bBytes, bReq, bShed, _capBytes), ct).ConfigureAwait(false); } - // Daily audit rows — per active client + the instance summary (cap, when hit, how many shed). + // Daily audit rows — per active client, interactive, and the instance summary (cap, when hit, how many shed). foreach (var (appId, c) in clientDaily) - await _store!.UpsertAsync(_tableName, EgressTableSchema.ClientDailyRow(appId, day, c.Bytes, c.Requests, c.Shed, _capBytes), ct).ConfigureAwait(false); - await _store!.UpsertAsync(_tableName, EgressTableSchema.SystemDailyRow(day, dayBytes, clientDaily.Values.Sum(c => c.Requests), dayShed, _capBytes, capReached, enforcing), ct).ConfigureAwait(false); + await _store!.UpsertAsync(_tableName, EgressTableSchema.ClientDailyRow(appId, day, c.Bytes, c.Requests, c.Shed, _capBytes, c.Endpoints), ct).ConfigureAwait(false); + if (interactiveDaily is not null) + await _store!.UpsertAsync(_tableName, EgressTableSchema.ClientDailyRow(EgressTableSchema.InteractivePartition, day, + interactiveDaily.Bytes, interactiveDaily.Requests, 0, _capBytes, interactiveDaily.Endpoints), ct).ConfigureAwait(false); + await _store!.UpsertAsync(_tableName, EgressTableSchema.SystemDailyRow(day, dayBytes, dayRequests, dayShed, _capBytes, capReached, enforcing, instanceEndpoints), ct).ConfigureAwait(false); // Drop finalised buckets we just wrote (strictly older than the current one). lock (_lock) @@ -416,15 +485,20 @@ internal async Task SeedCurrentBucketAsync(CancellationToken ct) var seeded = 0; await foreach (var row in _store!.QueryTableAsync(_tableName, filter, ct).ConfigureAwait(false)) { - if (row.RowKey != key || row.PartitionKey == EgressTableSchema.SystemPartition) continue; // re-apply filter; skip aggregate + if (row.RowKey != key) continue; // re-apply the filter — the store may ignore it lock (_lock) { - if (!_buckets.TryGetValue(key, out var b)) { b = new Dictionary(StringComparer.OrdinalIgnoreCase); _buckets[key] = b; } - b[row.PartitionKey] = new BucketAccum { Bytes = AsLong(row[EgressTableSchema.PropBytes]), Requests = AsLong(row[EgressTableSchema.PropRequests]), Shed = AsLong(row[EgressTableSchema.PropShed]) }; + // Added onto, not replacing, anything this process already recorded before the seed ran. + var a = Bucket(row.PartitionKey); + a.Bytes += AsLong(row[EgressTableSchema.PropBytes]); + a.Requests += AsLong(row[EgressTableSchema.PropRequests]); + a.Shed += AsLong(row[EgressTableSchema.PropShed]); + foreach (var (label, stats) in EgressTableSchema.ParseEndpoints(row.GetString(EgressTableSchema.PropEndpoints))) + EndpointStats.In(a.Endpoints, label).Merge(stats); } seeded++; } - if (seeded > 0) _logger.LogInformation("[Egress] Reseeded current bucket from table ({Count} client rows)", seeded); + if (seeded > 0) _logger.LogInformation("[Egress] Reseeded current bucket from table ({Count} rows)", seeded); } catch (Exception ex) { @@ -498,6 +572,7 @@ private sealed class LedgerState public long ShedRequests { get; set; } public DateTime? UpdatedUtc { get; set; } public Dictionary? Clients { get; set; } + public ClientState? Interactive { get; set; } } private sealed class ClientState @@ -506,5 +581,24 @@ private sealed class ClientState public long Requests { get; set; } public long Shed { get; set; } public DateTime LastSeenUtc { get; set; } + // Same [Bytes, Requests, MaxBytes, CacheHits, Errors, Shed] layout as the table column. + public Dictionary? Endpoints { get; set; } + + public static ClientState From(ClientTotals t) => new() + { + Bytes = t.Bytes, + Requests = t.Requests, + Shed = t.Shed, + LastSeenUtc = t.LastSeenUtc, + Endpoints = t.Endpoints.ToDictionary(kv => kv.Key, kv => kv.Value.ToArray()), + }; + + public ClientTotals ToTotals() + { + var t = new ClientTotals { Bytes = Bytes, Requests = Requests, Shed = Shed, LastSeenUtc = LastSeenUtc }; + foreach (var (label, values) in Endpoints ?? []) + t.Endpoints[label] = EndpointStats.FromArray(values); + return t; + } } } diff --git a/Services/Hosting/EgressTableSchema.cs b/Services/Hosting/EgressTableSchema.cs index 460ef9e..a8b7fb3 100644 --- a/Services/Hosting/EgressTableSchema.cs +++ b/Services/Hosting/EgressTableSchema.cs @@ -1,4 +1,5 @@ using System.Globalization; +using System.Text.Json; using Craft.Storage; namespace Craft.Hosting; @@ -30,6 +31,9 @@ internal static class EgressTableSchema /// Partition holding the instance aggregate. Not a GUID, so it never collides with an AppId. public const string SystemPartition = "instance-total"; + /// Partition holding interactive (user) traffic: tracked for reporting, never billed against the cap. + public const string InteractivePartition = "interactive"; + public const string BucketPrefix = "bkt_"; public const string DailyPrefix = "day_"; // Upper bound of the bkt_ prefix range: 'u' is the next char after 't', so any bkt_* RowKey is < this. @@ -46,6 +50,15 @@ internal static class EgressTableSchema public const string PropCapReachedUtc = "CapReachedUtc"; public const string PropEnforcing = "Enforcing"; public const string PropAppId = "AppId"; + // JSON object of endpoint label -> [Bytes, Requests, MaxBytes, CacheHits, Errors, Shed]. Arrays rather + // than named fields keep a full map under the 32K-char property limit; CIPP can't reassemble Craft splits. + public const string PropEndpoints = "Endpoints"; + + /// Endpoints kept per map; the rest fold into . + public const int MaxEndpoints = 100; + public const int MaxEndpointLabelLength = 96; + public const string OtherEndpoint = "_other"; + public const string UnmatchedEndpoint = "_unmatched"; /// The UTC start of the bucket falls in, floored to /// . @@ -79,50 +92,146 @@ public static string DailyRetentionCutoffRowKey(DateTime nowUtc, int retentionDa DailyPrefix + DayStamp(DateOnly.FromDateTime(nowUtc.AddDays(-Math.Max(1, retentionDays)))); public static StoreRow ClientBucketRow(string appId, DateTime bucketStartUtc, - long bytes, long requests, long shed, long capBytes) + long bytes, long requests, long shed, long capBytes, IReadOnlyDictionary? endpoints = null) { var row = new StoreRow(appId, BucketRowKey(bucketStartUtc)); - FillCounts(row, bytes, requests, shed, capBytes); + FillCounts(row, bytes, requests, shed, capBytes, endpoints); row[PropBucketStart] = bucketStartUtc; row[PropAppId] = appId; return row; } public static StoreRow SystemBucketRow(DateTime bucketStartUtc, - long bytes, long requests, long shed, long capBytes) + long bytes, long requests, long shed, long capBytes, IReadOnlyDictionary? endpoints = null) { var row = new StoreRow(SystemPartition, BucketRowKey(bucketStartUtc)); - FillCounts(row, bytes, requests, shed, capBytes); + FillCounts(row, bytes, requests, shed, capBytes, endpoints); row[PropBucketStart] = bucketStartUtc; return row; } public static StoreRow ClientDailyRow(string appId, DateOnly dayUtc, - long bytes, long requests, long shed, long capBytes) + long bytes, long requests, long shed, long capBytes, IReadOnlyDictionary? endpoints = null) { var row = new StoreRow(appId, DailyRowKey(dayUtc)); - FillCounts(row, bytes, requests, shed, capBytes); + FillCounts(row, bytes, requests, shed, capBytes, endpoints); row[PropDateUtc] = dayUtc.ToString("yyyy-MM-dd", CultureInfo.InvariantCulture); row[PropAppId] = appId; return row; } public static StoreRow SystemDailyRow(DateOnly dayUtc, - long bytes, long requests, long shed, long capBytes, DateTime? capReachedUtc, bool enforcing) + long bytes, long requests, long shed, long capBytes, DateTime? capReachedUtc, bool enforcing, + IReadOnlyDictionary? endpoints = null) { var row = new StoreRow(SystemPartition, DailyRowKey(dayUtc)); - FillCounts(row, bytes, requests, shed, capBytes); + FillCounts(row, bytes, requests, shed, capBytes, endpoints); row[PropDateUtc] = dayUtc.ToString("yyyy-MM-dd", CultureInfo.InvariantCulture); row[PropEnforcing] = enforcing; if (capReachedUtc.HasValue) row[PropCapReachedUtc] = capReachedUtc.Value; return row; } - private static void FillCounts(StoreRow row, long bytes, long requests, long shed, long capBytes) + private static void FillCounts(StoreRow row, long bytes, long requests, long shed, long capBytes, + IReadOnlyDictionary? endpoints) { row[PropBytes] = bytes; row[PropRequests] = requests; row[PropShed] = shed; row[PropCapBytes] = capBytes; + if (endpoints is { Count: > 0 }) row[PropEndpoints] = SerializeEndpoints(endpoints); + } + + /// Clamps an endpoint label to a bounded, table-safe form; blank is . + public static string EndpointLabel(string? label) + { + var trimmed = label?.Trim(); + if (string.IsNullOrEmpty(trimmed)) return UnmatchedEndpoint; + if (trimmed.Length > MaxEndpointLabelLength) trimmed = trimmed[..MaxEndpointLabelLength]; + return string.Create(trimmed.Length, trimmed, static (dst, src) => + { + for (var i = 0; i < src.Length; i++) + { + var c = src[i]; + dst[i] = char.IsAsciiLetterOrDigit(c) || c is '_' or '.' or ':' or '/' or '{' or '}' or '(' or ')' or '$' or '-' ? c : '_'; + } + }); + } + + /// The top endpoints by bytes, the remainder folded into + /// , as the compact JSON stored in . + public static string SerializeEndpoints(IReadOnlyDictionary endpoints) + { + var ordered = endpoints.Where(kv => kv.Key != OtherEndpoint).OrderByDescending(kv => kv.Value.Bytes).ToList(); + var other = endpoints.TryGetValue(OtherEndpoint, out var o) ? o.Clone() : null; + foreach (var (_, stats) in ordered.Skip(MaxEndpoints)) + (other ??= new EndpointStats()).Merge(stats); + + var map = new Dictionary(StringComparer.Ordinal); + foreach (var (label, stats) in ordered.Take(MaxEndpoints)) map[label] = stats.ToArray(); + if (other is not null) map[OtherEndpoint] = other.ToArray(); + return JsonSerializer.Serialize(map); + } + + /// Parses ; anything unreadable is an empty map. + public static Dictionary ParseEndpoints(string? json) + { + var result = new Dictionary(StringComparer.OrdinalIgnoreCase); + if (string.IsNullOrEmpty(json)) return result; + try + { + foreach (var (label, values) in JsonSerializer.Deserialize>(json) ?? []) + result[label] = EndpointStats.FromArray(values); + } + catch (JsonException) { } + return result; + } +} + +/// Per-endpoint counters within one client-day or one bucket. +internal sealed class EndpointStats +{ + public long Bytes, Requests, MaxBytes, CacheHits, Errors, Shed; + + public void Add(long bytes, bool cacheHit, int statusCode) + { + Bytes += bytes; + Requests += 1; + if (bytes > MaxBytes) MaxBytes = bytes; + if (cacheHit) CacheHits += 1; + if (statusCode >= 400) Errors += 1; + } + + public void Merge(EndpointStats other) + { + Bytes += other.Bytes; + Requests += other.Requests; + MaxBytes = Math.Max(MaxBytes, other.MaxBytes); + CacheHits += other.CacheHits; + Errors += other.Errors; + Shed += other.Shed; + } + + public EndpointStats Clone() => (EndpointStats)MemberwiseClone(); + + /// The entry for , created on first use; once the map is full, new + /// labels share so a caller can't grow it unbounded. + public static EndpointStats In(Dictionary map, string label) + { + if (map.TryGetValue(label, out var stats)) return stats; + if (map.Count >= EgressTableSchema.MaxEndpoints) label = EgressTableSchema.OtherEndpoint; + if (!map.TryGetValue(label, out stats)) map[label] = stats = new EndpointStats(); + return stats; + } + + public static Dictionary CloneMap(Dictionary map) => + map.ToDictionary(kv => kv.Key, kv => kv.Value.Clone(), StringComparer.OrdinalIgnoreCase); + + public long[] ToArray() => [Bytes, Requests, MaxBytes, CacheHits, Errors, Shed]; + + public static EndpointStats FromArray(long[]? v) + { + long At(int i) => v is not null && i < v.Length ? v[i] : 0; + return new EndpointStats { Bytes = At(0), Requests = At(1), MaxBytes = At(2), CacheHits = At(3), Errors = At(4), Shed = At(5) }; } } diff --git a/Services/Hosting/Endpoints/PowerShellDispatchEndpoint.cs b/Services/Hosting/Endpoints/PowerShellDispatchEndpoint.cs index 9121ce4..edb49fd 100644 --- a/Services/Hosting/Endpoints/PowerShellDispatchEndpoint.cs +++ b/Services/Hosting/Endpoints/PowerShellDispatchEndpoint.cs @@ -119,6 +119,7 @@ await ServeCachedAsync(context, cache, psRunner, orchestrator, logger, context.Response.StatusCode = result.StatusCode; context.Response.ContentType = HandlerHeaders.ResolveContentType(result.ContentType); HandlerHeaders.Apply(context.Response, result.Headers); + ApiEgressWireCounterMiddleware.TagEndpoint(context, endpoint); // BYPASS vs MISS matters when someone asks why an endpoint never caches: MISS means it // was eligible and simply had no entry, BYPASS means the policy kept it out. context.Response.Headers["X-Cache"] = cacheBypassReason is null ? "MISS" : "BYPASS"; @@ -221,6 +222,7 @@ private static async Task ServeCachedAsync( context.Response.StatusCode = cached.Result.StatusCode; context.Response.ContentType = HandlerHeaders.ResolveContentType(cached.Result.ContentType); HandlerHeaders.Apply(context.Response, cached.Result.Headers); + ApiEgressWireCounterMiddleware.TagEndpoint(context, endpoint); context.Response.Headers["X-Cache"] = cached.IsStale ? "HIT-STALE" : "HIT"; context.Response.Headers["X-Cache-Age"] = $"{cached.Age.TotalSeconds:F0}s"; context.Response.Headers["X-Cache-TTL"] = $"{ttl.TotalSeconds:F0}s"; diff --git a/tests/Craft.Tests/ApiCompressionPipelineTests.cs b/tests/Craft.Tests/ApiCompressionPipelineTests.cs index 60938e5..af94e8b 100644 --- a/tests/Craft.Tests/ApiCompressionPipelineTests.cs +++ b/tests/Craft.Tests/ApiCompressionPipelineTests.cs @@ -240,4 +240,68 @@ public async Task FullRealisticChain_WithCsp_AndEgressLimiter_StillCompresses() Assert.Equal("gzip", enc); Assert.True(bytes < RawLength, $"expected compressed < {RawLength}, got {bytes}"); } + + [Fact] + public async Task RealPipeline_LabelsDispatchedAndLiteralRoutes_AndHidesTagHeader() + { + // Dispatched endpoints are tagged by the dispatcher (with the handler's X-Craft-Endpoint suffix); + // a literal route (a native endpoint) is labelled from its route pattern. Real routing, real + // compression, wire bytes billed to the routed label. + var store = new FakeTableStore(); + var ledger = new EgressLedger(NullLogger.Instance, 0, 60, + Path.Combine(_dir, Guid.NewGuid().ToString("N")[..8] + ".json"), store: store, tableName: "T"); + + var builder = WebApplication.CreateBuilder(); + builder.WebHost.ConfigureKestrel(o => o.Listen(IPAddress.Loopback, 0)); + builder.Logging.ClearProviders(); + builder.Services.AddCraftResponseCompression(); + builder.Services.AddSingleton(ledger); + var app = builder.Build(); + app.UseWhen(IsApiPath, api => + { + api.UseMiddleware(); + api.UseResponseCompression(); + }); + app.UseMiddleware(); + app.MapMethods("/API/{endpoint}", GetOnly, async (HttpContext ctx, string endpoint) => + { + ctx.Response.ContentType = "application/json"; + ctx.Response.Headers[ApiEgressWireCounterMiddleware.EndpointHeader] = "users"; + ApiEgressWireCounterMiddleware.TagEndpoint(ctx, endpoint); + await ctx.Response.WriteAsync(Payload); + }); + app.MapMethods("/API/Literal", GetOnly, async ctx => + { + ctx.Response.ContentType = "application/json"; + await ctx.Response.WriteAsync(Payload); + }); + + await app.StartAsync(); + try + { + var addr = app.Services.GetRequiredService().Features.Get()!.Addresses.First(); + using var handler = new HttpClientHandler { AutomaticDecompression = DecompressionMethods.None }; + using var client = new HttpClient(handler); + async Task Get(string path) + { + var req = new HttpRequestMessage(HttpMethod.Get, $"{addr}{path}"); + req.Headers.TryAddWithoutValidation("Accept-Encoding", "gzip"); + req.Headers.TryAddWithoutValidation("x-ms-client-principal-idp", "aad"); + req.Headers.TryAddWithoutValidation("x-ms-client-principal-name", "11111111-2222-3333-4444-555555555555"); + return await client.SendAsync(req); + } + + var dispatched = await Get("/API/ListGraphRequest"); + Assert.False(dispatched.Headers.Contains(ApiEgressWireCounterMiddleware.EndpointHeader)); + Assert.Equal("gzip", dispatched.Content.Headers.ContentEncoding.Single()); + await Get("/API/Literal"); + } + finally { await app.StopAsync(); } + + await ledger.SyncToTableAsync(CancellationToken.None); + var row = store.All("T").Single(r => r.PartitionKey == "11111111-2222-3333-4444-555555555555" && r.RowKey.StartsWith(EgressTableSchema.DailyPrefix, StringComparison.Ordinal)); + var endpoints = EgressTableSchema.ParseEndpoints(row.GetString(EgressTableSchema.PropEndpoints)); + Assert.Equal(["ListGraphRequest:users", "Literal"], endpoints.Keys.Order().ToArray()); + Assert.True(endpoints["ListGraphRequest:users"].Bytes < RawLength); // compressed wire bytes, not the raw body + } } diff --git a/tests/Craft.Tests/ApiEgressWireCounterMiddlewareTests.cs b/tests/Craft.Tests/ApiEgressWireCounterMiddlewareTests.cs index 2785c1a..6806c57 100644 --- a/tests/Craft.Tests/ApiEgressWireCounterMiddlewareTests.cs +++ b/tests/Craft.Tests/ApiEgressWireCounterMiddlewareTests.cs @@ -207,4 +207,77 @@ public async Task OverBudget_LimiterSheds_CounterDoesNotBillThe429() Assert.Equal(StatusCodes.Status429TooManyRequests, ctx.Response.StatusCode); Assert.Equal(before, ledger.CurrentBytes); // the shed 429 body is never billed } + + // ── endpoint labels + interactive callers ───────────────────────────────────────────────────────── + + private (EgressLedger Ledger, FakeTableStore Store) TableLedger() + { + var store = new FakeTableStore(); + var ledger = new EgressLedger(NullLogger.Instance, 1_000_000, 60, + Path.Combine(_dir, Guid.NewGuid().ToString("N")[..8] + ".json"), store: store, tableName: "T"); + return (ledger, store); + } + + private static async Task> DailyEndpoints(EgressLedger ledger, FakeTableStore store, string partition) + { + await ledger.SyncToTableAsync(CancellationToken.None); + var row = store.All("T").Single(r => r.PartitionKey == partition && r.RowKey.StartsWith(EgressTableSchema.DailyPrefix, StringComparison.Ordinal)); + return EgressTableSchema.ParseEndpoints(row.GetString(EgressTableSchema.PropEndpoints)); + } + + [Fact] + public async Task TagEndpoint_SuffixesHandlerHeader_StripsIt_AndCountsCacheHits() + { + var (ledger, store) = TableLedger(); + var mw = WireCounter(async ctx => + { + Flag(ctx); + ctx.Response.Headers[ApiEgressWireCounterMiddleware.EndpointHeader] = "users"; + ApiEgressWireCounterMiddleware.TagEndpoint(ctx, "ListGraphRequest"); + ctx.Response.Headers["X-Cache"] = "HIT"; + await ctx.Response.Body.WriteAsync(new byte[64]); + }, ledger); + + var ctx = Context(); + await mw.InvokeAsync(ctx); + + Assert.False(ctx.Response.Headers.ContainsKey(ApiEgressWireCounterMiddleware.EndpointHeader)); + var stats = (await DailyEndpoints(ledger, store, App))["ListGraphRequest:users"]; + Assert.Equal(64L, stats.Bytes); + Assert.Equal(1L, stats.CacheHits); + } + + [Fact] + public async Task InteractiveCaller_IsRecordedForReporting_ButNeverBilled() + { + var (ledger, store) = TableLedger(); + var limiter = new ApiEgressLimiterMiddleware(async ctx => + { + ApiEgressWireCounterMiddleware.TagEndpoint(ctx, "ListLogs"); + await ctx.Response.Body.WriteAsync(new byte[256]); + }, ledger, NullLoggerFactory.Instance); + var mw = WireCounter(ctx => limiter.InvokeAsync(ctx), ledger); + + var ctx = Context(); + ctx.Request.Headers["x-ms-client-principal-idp"] = "azureStaticWebApps"; + ctx.Request.Headers["x-ms-client-principal-name"] = "user@contoso.com"; + await mw.InvokeAsync(ctx); + + Assert.Equal(0L, ledger.CurrentBytes); + Assert.Equal(256L, (await DailyEndpoints(ledger, store, EgressTableSchema.InteractivePartition))["ListLogs"].Bytes); + } + + [Fact] + public async Task AnonymousCaller_IsNotRecorded() + { + var (ledger, store) = TableLedger(); + var limiter = new ApiEgressLimiterMiddleware( + async ctx => await ctx.Response.Body.WriteAsync(new byte[256]), ledger, NullLoggerFactory.Instance); + var mw = WireCounter(ctx => limiter.InvokeAsync(ctx), ledger); + + await mw.InvokeAsync(Context()); + await ledger.SyncToTableAsync(CancellationToken.None); + + Assert.DoesNotContain(store.All("T"), r => r.PartitionKey == EgressTableSchema.InteractivePartition); + } } diff --git a/tests/Craft.Tests/EgressLedgerTableTests.cs b/tests/Craft.Tests/EgressLedgerTableTests.cs index b705777..91e468d 100644 --- a/tests/Craft.Tests/EgressLedgerTableTests.cs +++ b/tests/Craft.Tests/EgressLedgerTableTests.cs @@ -126,4 +126,131 @@ public async Task Purge_RemovesRowsBeyondRetention_KeepsRecent() Assert.Equal(2, remaining.Count); // the two old rows are gone Assert.All(remaining, r => Assert.Contains("20260913", r.RowKey)); // only today's rows survive } + + // ── Endpoint breakdown + interactive partition ────────────────────────────────────────────────── + + private static StoreRow Row(FakeTableStore store, string partition, string prefix) => + store.All(Table).Single(r => r.PartitionKey == partition && r.RowKey.StartsWith(prefix, StringComparison.Ordinal)); + + private static Dictionary Endpoints(StoreRow row) => + EgressTableSchema.ParseEndpoints(row.GetString(EgressTableSchema.PropEndpoints)); + + [Fact] + public async Task Endpoints_BrokenOutOnDailyRows_AndInstanceBucket_NotClientBuckets() + { + var store = new FakeTableStore(); + var clock = new DateTime(2026, 9, 13, 14, 22, 0, DateTimeKind.Utc); + var ledger = New(store, cap: 1_000_000, () => clock); + + ledger.Record(100, App, "ListUsers", cacheHit: true); + ledger.Record(300, App, "ListUsers"); + ledger.Record(50, App2, "ListGraphRequest:users", statusCode: 500); + ledger.RecordShed(App2, "ListUsers"); + await ledger.SyncToTableAsync(CancellationToken.None); + + var users = Endpoints(Row(store, App, EgressTableSchema.DailyPrefix))["ListUsers"]; + Assert.Equal([400L, 2, 300, 1, 0, 0], users.ToArray()); // Bytes, Requests, MaxBytes, CacheHits, Errors, Shed + + var instanceDay = Endpoints(Row(store, EgressTableSchema.SystemPartition, EgressTableSchema.DailyPrefix)); + Assert.Equal(1L, instanceDay["ListUsers"].Shed); + Assert.Equal(1L, instanceDay["ListGraphRequest:users"].Errors); + + var instanceBucket = Endpoints(Row(store, EgressTableSchema.SystemPartition, EgressTableSchema.BucketPrefix)); + Assert.Equal(400L, instanceBucket["ListUsers"].Bytes); + Assert.Equal(50L, instanceBucket["ListGraphRequest:users"].Bytes); + + Assert.Null(Row(store, App, EgressTableSchema.BucketPrefix)[EgressTableSchema.PropEndpoints]); + } + + [Fact] + public async Task Interactive_ReportedInOwnPartition_NeverBilled() + { + var store = new FakeTableStore(); + var clock = new DateTime(2026, 9, 13, 14, 22, 0, DateTimeKind.Utc); + var ledger = New(store, cap: 1000, () => clock); + + ledger.RecordInteractive(5000, "ListLogs"); + ledger.Record(100, App, "ListUsers"); + await ledger.SyncToTableAsync(CancellationToken.None); + + Assert.Equal(100L, ledger.CurrentBytes); + Assert.False(ledger.ShouldReject()); + + Assert.Equal(5000L, Endpoints(Row(store, EgressTableSchema.InteractivePartition, EgressTableSchema.DailyPrefix))["ListLogs"].Bytes); + Assert.Equal(5000L, Endpoints(Row(store, EgressTableSchema.InteractivePartition, EgressTableSchema.BucketPrefix))["ListLogs"].Bytes); + + var instanceDay = Row(store, EgressTableSchema.SystemPartition, EgressTableSchema.DailyPrefix); + Assert.Equal(100L, Long(instanceDay, EgressTableSchema.PropBytes)); + Assert.Equal(1L, Long(instanceDay, EgressTableSchema.PropRequests)); + Assert.DoesNotContain("ListLogs", Endpoints(instanceDay).Keys); + Assert.Equal(100L, Long(Row(store, EgressTableSchema.SystemPartition, EgressTableSchema.BucketPrefix), EgressTableSchema.PropBytes)); + } + + [Fact] + public async Task InstanceDailyRequests_CountClientsIdleSinceAnEarlierBucket() + { + var store = new FakeTableStore(); + var clock = new DateTime(2026, 9, 13, 14, 0, 0, DateTimeKind.Utc); + var ledger = New(store, cap: 1_000_000, () => clock); + + ledger.Record(100, App, "ListUsers"); + await ledger.SyncToTableAsync(CancellationToken.None); + clock = clock.AddMinutes(30); + ledger.Record(100, App2, "ListUsers"); + await ledger.SyncToTableAsync(CancellationToken.None); + + var instanceDay = Row(store, EgressTableSchema.SystemPartition, EgressTableSchema.DailyPrefix); + Assert.Equal(2L, Long(instanceDay, EgressTableSchema.PropRequests)); + Assert.Equal(2L, Endpoints(instanceDay)["ListUsers"].Requests); + } + + [Fact] + public async Task Restart_KeepsDailyEndpoints_FromFile_AndBucketEndpoints_FromTable() + { + var store = new FakeTableStore(); + var clock = new DateTime(2026, 9, 13, 14, 22, 0, DateTimeKind.Utc); + var first = New(store, cap: 1_000_000, () => clock); + first.Record(100, App, "ListUsers"); + first.RecordInteractive(70, "ListLogs"); + first.Flush(); + await first.SyncToTableAsync(CancellationToken.None); + + var second = New(store, cap: 1_000_000, () => clock); + await second.SeedCurrentBucketAsync(CancellationToken.None); + second.Record(10, App, "ListUsers"); + second.RecordInteractive(5, "ListLogs"); + await second.SyncToTableAsync(CancellationToken.None); + + Assert.Equal(110L, Endpoints(Row(store, App, EgressTableSchema.DailyPrefix))["ListUsers"].Bytes); + Assert.Equal(75L, Endpoints(Row(store, EgressTableSchema.InteractivePartition, EgressTableSchema.DailyPrefix))["ListLogs"].Bytes); + Assert.Equal(110L, Endpoints(Row(store, EgressTableSchema.SystemPartition, EgressTableSchema.BucketPrefix))["ListUsers"].Bytes); + Assert.Equal(75L, Endpoints(Row(store, EgressTableSchema.InteractivePartition, EgressTableSchema.BucketPrefix))["ListLogs"].Bytes); + } + + [Fact] + public async Task EndpointMap_IsBounded_AndFitsOneTableProperty() + { + var store = new FakeTableStore(); + var clock = new DateTime(2026, 9, 13, 14, 22, 0, DateTimeKind.Utc); + var ledger = New(store, cap: 1_000_000, () => clock); + + for (var i = 0; i < 500; i++) + ledger.Record(1_000_000_000_000 + i, App, $"Endpoint{i:D3}:" + new string('x', 200)); + await ledger.SyncToTableAsync(CancellationToken.None); + + var row = Row(store, App, EgressTableSchema.DailyPrefix); + var map = Endpoints(row); + Assert.Equal(EgressTableSchema.MaxEndpoints + 1, map.Count); + Assert.Equal(500L - EgressTableSchema.MaxEndpoints, map[EgressTableSchema.OtherEndpoint].Requests); + Assert.All(map.Keys, k => Assert.True(k.Length <= EgressTableSchema.MaxEndpointLabelLength)); + Assert.True(row.GetString(EgressTableSchema.PropEndpoints)!.Length < 32_000); + } + + [Fact] + public void EndpointLabel_ReplacesUnsafeCharacters_AndDefaultsBlank() + { + Assert.Equal("ListGraphRequest:deviceManagement/managedDevices({id})/$count", EgressTableSchema.EndpointLabel("ListGraphRequest:deviceManagement/managedDevices({id})/$count")); + Assert.Equal("Exec__x_", EgressTableSchema.EndpointLabel("Exec\"")); + Assert.Equal(EgressTableSchema.UnmatchedEndpoint, EgressTableSchema.EndpointLabel(" ")); + } } From 0d7ad87f1144b9d0218dfded21129b5e7cdde150 Mon Sep 17 00:00:00 2001 From: Zacgoose <107489668+Zacgoose@users.noreply.github.com> Date: Fri, 2 Oct 2026 18:51:18 +0800 Subject: [PATCH 4/4] Update JobQueuePumpBackoffTests.cs --- tests/Craft.Tests/JobQueuePumpBackoffTests.cs | 28 ++++++++++++++++--- 1 file changed, 24 insertions(+), 4 deletions(-) diff --git a/tests/Craft.Tests/JobQueuePumpBackoffTests.cs b/tests/Craft.Tests/JobQueuePumpBackoffTests.cs index 0d29b91..2375fa8 100644 --- a/tests/Craft.Tests/JobQueuePumpBackoffTests.cs +++ b/tests/Craft.Tests/JobQueuePumpBackoffTests.cs @@ -7,6 +7,12 @@ namespace Craft.Tests; +[CollectionDefinition(Name, DisableParallelization = true)] +public class PumpTiming +{ + public const string Name = "job-queue-pump-timing"; +} + /// /// The pump polls storage for work, so an idle instance was scanning the queue table once a second /// forever. Backing off fixes that, but the interval is also a hard throughput ceiling — a refill hands @@ -18,7 +24,10 @@ namespace Craft.Tests; /// busy run is about to need its next batch. Idle means claimed nothing AND holding nothing. /// /// These tests pin both directions — that a quiet pump slows down, and that a working one does not. +/// They count scans in sub-second wall-clock windows, so they run alone: sharing a small CI runner with +/// parallel collections starved the 100ms ticks enough to read as a backoff. /// +[Collection(PumpTiming.Name)] public class JobQueuePumpBackoffTests { private static (JobQueuePump Pump, JobQueueStore Queue, JobManager Jobs) NewPump( @@ -194,14 +203,25 @@ public async Task APumpWhoseBufferIsDraining_RefillsAtTheBaseInterval() // A consumer, so the buffer actually draws down and refills are needed — without one the pump // fills once and correctly never scans again. - jobs.SetWorkResolver((_, _) => Task.FromResult?>(_ => Task.CompletedTask)); + // The window opens once the consumer has run a job, so its startup is not billed to the pump. + var draining = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + jobs.SetWorkResolver((_, _) => + { + draining.TrySetResult(); + return Task.FromResult?>(_ => Task.CompletedTask); + }); _ = Task.Run(() => jobs.StartAsync(CancellationToken.None)); - await PumpFor(pump, 900); + await pump.StartAsync(CancellationToken.None); + await draining.Task.WaitAsync(TimeSpan.FromSeconds(10)); + var before = backing.Scans; + await Task.Delay(900); + var scans = backing.Scans - before; + await Task.WhenAny(pump.StopAsync(CancellationToken.None), Task.Delay(3000)); await jobs.StopAsync(CancellationToken.None); - Assert.True(backing.Scans >= 6, - $"a draining buffer was only refilled {backing.Scans} times in 900ms — the pump backed off " + + Assert.True(scans >= 6, + $"a draining buffer was only refilled {scans} times in 900ms — the pump backed off " + "while work was flowing, which caps throughput at batchSize per idle interval"); }