InvoiceParser/src/ParseWorker.cs
2026-09-16 13:09:19 -07:00

505 lines
21 KiB
C#

using System.Collections.Concurrent;
using System.IO;
using System.Text.Json.Nodes;
using System.Threading.Channels;
using Microsoft.Extensions.Hosting;
namespace InvoiceParser.Src;
// The single owner of the parse queue. Replaces the old triple-worker setup
// (pending_process.asp, parse_queue.ps1, browser-side parse) so there is no
// claim/steal logic left - one in-process worker drains one channel.
public class ParseWorker : BackgroundService
{
private readonly AppConfig _cfg;
private readonly PendingStore _store;
private readonly LlmClient _llm;
private readonly EventBroadcaster _events;
private readonly Channel<string> _queue = Channel.CreateUnbounded<string>();
private readonly ConcurrentDictionary<string, byte> _inQueue = new();
public ParseWorker(AppConfig cfg, PendingStore store, LlmClient llm, EventBroadcaster events)
{
_cfg = cfg;
_store = store;
_llm = llm;
_events = events;
}
public bool Enqueue(string id)
{
if (!_inQueue.TryAdd(id, 0)) return false; // already waiting
_queue.Writer.TryWrite(id);
return true;
}
// Enqueues queued + errored items (the "Process queued now" button).
// Returns how many were newly queued and how many were already waiting, so
// a caller that gets 0 can tell "nothing to do" from "already in flight".
public (int Queued, int AlreadyWaiting) EnqueueEligible(bool includeErrors)
{
int queued = 0, waiting = 0;
foreach (var item in _store.List())
{
var stage = _store.Stage(item.Id);
bool eligible = stage is "queued" or "queue" || (includeErrors && stage == "error");
if (!eligible) continue;
if (Enqueue(item.Id)) queued++;
else waiting++;
}
return (queued, waiting);
}
protected override async Task ExecuteAsync(CancellationToken ct)
{
// Recover from a previous crash: anything stranded at "processing"
// goes back to "queued", then all queued work is picked up.
int recovered = 0;
foreach (var item in _store.List())
{
if (_store.Stage(item.Id) != "processing") continue;
_store.SetStage(item.Id, "queued");
recovered++;
}
Log.Info("worker", $"started - {recovered} draft(s) reset from processing to queued");
// The drain loop must outlive any single failure. An exception escaping
// it leaves the channel with no reader: uploads still queue, _inQueue
// still dedupes them, and every invoice sits at "queued" forever with no
// error shown - the host keeps serving because nothing awaits its
// shutdown request, so the dead queue is completely silent.
while (!ct.IsCancellationRequested)
{
try
{
await DrainAsync(ct);
return; // channel completed - only happens on shutdown
}
catch (OperationCanceledException) when (ct.IsCancellationRequested)
{
return;
}
catch (Exception ex)
{
// Anything stranded in the dedupe set was never read; forget it
// so the next sweep can queue it again.
Log.Error("worker", "drain loop crashed - restarting it", ex);
_inQueue.Clear();
RequeueUnfinished();
}
}
}
private async Task DrainAsync(CancellationToken ct)
{
EnqueueEligible(includeErrors: false);
await foreach (var id in _queue.Reader.ReadAllAsync(ct))
{
_inQueue.TryRemove(id, out _);
var started = System.Diagnostics.Stopwatch.StartNew();
try
{
await ProcessOne(id, ct);
Log.Info("worker", $"{id} finished in {started.ElapsedMilliseconds}ms (stage {_store.Stage(id)})");
}
catch (OperationCanceledException) when (ct.IsCancellationRequested) { throw; }
catch (Exception ex)
{
// Includes a request timeout, which arrives as a cancellation
// that has nothing to do with our token.
Log.Error("worker", $"{id} failed after {started.ElapsedMilliseconds}ms", ex);
_store.SetError(id, ex.Message);
_events.StageChanged(id, "error", ex.Message);
}
}
}
private void RequeueUnfinished()
{
foreach (var item in _store.List())
{
if (_store.Stage(item.Id) == "processing") _store.SetStage(item.Id, "queued");
}
}
private async Task ProcessOne(string id, CancellationToken ct)
{
var state = _store.ReadState(id);
if (state == null) return;
_store.SetStage(id, "processing");
_events.StageChanged(id, "processing");
state = _store.ReadState(id)!;
Log.Info("worker", $"{id} processing \"{state["source"]?["file_name"]}\"");
// Fresh uploads arrive unprepped so the upload request returns at once;
// drafts from the old app carry no flag and are already prepped.
if (state["prepped"] is JsonValue pv && pv.TryGetValue<bool>(out var prepped) && !prepped)
{
Prep(id);
state = _store.ReadState(id)!;
}
string parseMode = state["parse_mode"]?.ToString() ?? "";
string ocrText = state["ocr"]?["text"]?.ToString() ?? "";
// One request per page, never the whole invoice at once. Asked for a
// 6-page invoice in a single call the model returned finishReason=STOP
// after 124 items and silently omitted a whole page of line items - not
// a token-limit truncation, just fidelity loss over a long output. Page
// sized requests are small enough to come back complete.
string backend;
List<string> responses;
try
{
if (parseMode == "text" && !string.IsNullOrWhiteSpace(ocrText))
{
backend = _llm.TextBackend;
var chunks = SplitTextByPage(ocrText, state);
var pageImages = _cfg.Llm.AttachPageImages ? PageImagesByNumber(id) : null;
Log.Info("worker", $"{id} text route ({backend}) - {chunks.Count} page(s), " +
$"{ocrText.Length} chars, images={(pageImages?.Count ?? 0)}");
responses = await RunPagesAsync(chunks, (c, t) =>
{
(string, string)? img = null;
if (pageImages != null && pageImages.TryGetValue(c.Page, out var found)) img = found;
return _llm.ParseTextAsync(c.Text, img, t);
}, ct);
}
else
{
backend = _llm.ImageBackend;
var images = CollectImages(id, state);
Log.Info("worker", $"{id} image route ({backend}) - {images.Count} page image(s)");
responses = await RunPagesAsync(images,
(img, t) => _llm.ParseImagesAsync(new List<(string, string)> { img }, t), ct);
}
}
catch (OperationCanceledException) when (ct.IsCancellationRequested) { throw; }
catch (Exception ex)
{
Log.Error("worker", $"{id} LLM call failed", ex);
_store.SetError(id, ex.Message);
_events.StageChanged(id, "error", ex.Message);
return;
}
// The LLM call can take minutes; if the user opened the draft and started
// editing meanwhile, do not clobber their state.
if (_store.Stage(id) != "processing")
{
Log.Warn("worker", $"{id} left processing while the LLM ran (now {_store.Stage(id)}) - result discarded");
return;
}
JsonNode api;
JsonNode parsed;
try
{
(api, parsed) = MergeResponses(responses);
}
catch (Exception ex)
{
Log.Error("worker", $"{id} response merge failed", ex);
_store.SetError(id, ex.Message);
_events.StageChanged(id, "error", ex.Message);
return;
}
var newState = BuildParsedState(state, api, parsed, backend);
_store.WriteState(id, newState);
_events.StageChanged(id, "parsed");
Log.Info("worker", $"{id} parsed: invoice={parsed["invoice_number"]} date={parsed["invoice_date"]} " +
$"total={parsed["total"]} lines={(parsed["line_items"] as JsonArray)?.Count ?? 0}");
}
// Server-side prep (was browser-side in the old app): text-layer extraction
// for digital PDFs, page rendering for scanned ones.
private void Prep(string id)
{
var src = _store.FindSource(id, out var ext)
?? throw new Exception("source file missing after upload");
string parseMode = "image";
string ocrText = "", ocrMethod = "";
JsonNode? lineMap = null;
if (ext.Equals("pdf", StringComparison.OrdinalIgnoreCase))
{
var lines = PdfPrep.ExtractLines(src);
var text = string.Join("\n", lines.Select(l => l.Text));
if (text.Trim().Length > 50)
{
parseMode = "text";
ocrText = text;
ocrMethod = "via PDF text layer";
lineMap = PdfPrep.LineMapToJson(lines);
}
else
{
PdfPrep.RenderAllPages(src, Path.Combine(_store.Dir(id), "pages"), _cfg.Llm.MaxImageDim);
}
}
var state = _store.ReadState(id) ?? throw new Exception("state.json missing");
state["parse_mode"] = parseMode;
state["ocr"] = new JsonObject
{
["text"] = ocrText,
["method"] = ocrMethod,
["line_map"] = lineMap
};
state["prepped"] = true;
state["updated_at"] = Util.NowIso();
_store.WriteState(id, state);
Log.Info("prep", $"{id} ext={ext} mode={parseMode} textChars={ocrText.Length} method=\"{ocrMethod}\"");
}
// Runs the per-page calls with bounded concurrency, preserving page order.
// Sequentially six pages took longer than the (lossy) single call, because
// each request pays its own model overhead; overlapping them removes that.
private async Task<List<string>> RunPagesAsync<T>(
IReadOnlyList<T> pages, Func<T, CancellationToken, Task<string>> call, CancellationToken ct)
{
int limit = Math.Max(1, _cfg.Llm.MaxParallelPages);
var results = new string[pages.Count];
using var gate = new SemaphoreSlim(limit);
var tasks = pages.Select(async (page, i) =>
{
await gate.WaitAsync(ct);
try { results[i] = await call(page, ct); }
finally { gate.Release(); }
}).ToList();
await Task.WhenAll(tasks);
return results.ToList();
}
// Splits the OCR text into one numbered chunk per PDF page, using the
// per-line page numbers in ocr.line_map. Line numbers stay GLOBAL across
// chunks so the source_line values the model returns still line up with
// line_map and the preview markers. Falls back to a single chunk when a
// draft has no line_map (older drafts, or image sources).
private static List<(int Page, string Text)> SplitTextByPage(string ocrText, JsonNode state)
{
var numbered = LlmClient.NumberLines(ocrText).Replace("\r\n", "\n").Split('\n');
var map = state["ocr"]?["line_map"] as JsonArray;
if (map == null || map.Count != numbered.Length)
return new List<(int, string)> { (1, string.Join("\n", numbered)) };
var chunks = new List<(int, string)>();
var current = new List<string>();
int currentPage = -1;
for (int i = 0; i < numbered.Length; i++)
{
int page = int.TryParse(map[i]?["page"]?.ToString(), out var p) ? p : 1;
if (page != currentPage && current.Count > 0)
{
chunks.Add((currentPage, string.Join("\n", current)));
current.Clear();
}
currentPage = page;
current.Add(numbered[i]);
}
if (current.Count > 0) chunks.Add((currentPage, string.Join("\n", current)));
return chunks;
}
// Rendered pages keyed by page number for hybrid text+image parsing. A
// text-layer PDF is not rendered during prep, so render it here; the
// preview endpoint reuses the same files afterwards.
private Dictionary<int, (string Mime, string B64)> PageImagesByNumber(string id)
{
var result = new Dictionary<int, (string, string)>();
var src = _store.FindSource(id, out var ext);
if (src == null || !ext.Equals("pdf", StringComparison.OrdinalIgnoreCase)) return result;
var pagesDir = Path.Combine(_store.Dir(id), "pages");
var files = Directory.Exists(pagesDir) ? Directory.GetFiles(pagesDir, "page_*.jpg") : Array.Empty<string>();
if (files.Length == 0)
{
PdfPrep.RenderAllPages(src, pagesDir, _cfg.Llm.MaxImageDim);
files = Directory.GetFiles(pagesDir, "page_*.jpg");
}
foreach (var f in files)
result[PageNumber(f)] = ("image/jpeg", Convert.ToBase64String(File.ReadAllBytes(f)));
return result;
}
// Folds the per-page responses into one invoice. Header fields are taken
// from the first page that supplies them (they live on page 1 of a real
// invoice); line items are concatenated in page order.
private static (JsonNode Api, JsonNode Parsed) MergeResponses(List<string> responses)
{
var rawArray = new JsonArray();
var merged = new JsonObject();
var items = new JsonArray();
bool grandTotalFound = false;
for (int pageNo = 0; pageNo < responses.Count; pageNo++)
{
var raw = responses[pageNo];
JsonNode api;
JsonNode parsed;
try
{
api = JsonNode.Parse(raw)!;
parsed = LlmClient.ExtractInvoiceData(api);
}
catch (Exception ex)
{
// A page that cannot be read is a failed parse, not a shorter
// invoice: skipping it would drop its line items with no trace.
throw new Exception($"Page {pageNo + 1} of {responses.Count}: {ex.Message}");
}
rawArray.Add(api.DeepClone());
// Identity fields come from the first page that has them.
foreach (var field in new[] { "vendor_name", "invoice_number", "invoice_date" })
{
if (merged[field] != null) continue;
var v = parsed[field];
if (v != null && v.GetValueKind() != System.Text.Json.JsonValueKind.Null &&
v.ToString() != "")
merged[field] = v.DeepClone();
}
// The total: a page whose label names the invoice grand total wins
// outright. Otherwise the last page with a total is kept, which is
// wrong when a recap/allowance page follows the totals page - the
// prompt asks for null there, the label check is the backstop.
var t = parsed["total"];
if (t != null && t.GetValueKind() != System.Text.Json.JsonValueKind.Null && t.ToString() != "")
{
var label = parsed["total_label"]?.ToString() ?? "";
bool grand = IsGrandTotalLabel(label);
if (!grandTotalFound || grand)
{
merged["total"] = t.DeepClone();
merged["total_label"] = label;
merged["total_page"] = pageNo + 1;
grandTotalFound = grandTotalFound || grand;
}
}
if (parsed["line_items"] is JsonArray li)
foreach (var it in li)
if (it != null) items.Add(it.DeepClone());
}
if (rawArray.Count == 0) throw new Exception("LLM returned no usable response");
merged["line_items"] = items;
return (rawArray, merged);
}
private static bool IsGrandTotalLabel(string label) =>
System.Text.RegularExpressions.Regex.IsMatch(label,
@"invoice\s*total|total\s*due|amount\s*due|balance\s*due|net\s*total|grand\s*total|net\s*invoice|total\s*invoice",
System.Text.RegularExpressions.RegexOptions.IgnoreCase);
// Vision inputs, in order of preference: rendered pages on disk,
// legacy base64 page_images inside state.json, or the raw source image.
// A scanned PDF with no rendered pages gets rendered right here.
private List<(string Mime, string B64)> CollectImages(string id, JsonNode state)
{
var pagesDir = Path.Combine(_store.Dir(id), "pages");
if (Directory.Exists(pagesDir))
{
var files = Directory.GetFiles(pagesDir, "page_*.jpg")
.OrderBy(PageNumber).ToList();
if (files.Count > 0)
return files.Select(f => ("image/jpeg", Convert.ToBase64String(File.ReadAllBytes(f)))).ToList();
}
if (state["page_images"] is JsonArray legacy && legacy.Count > 0)
return legacy.Select(n => ("image/jpeg", n?.ToString() ?? "")).ToList();
var src = _store.FindSource(id, out var ext);
if (src == null)
throw new Exception("No page images and no source file for this invoice");
if (ext.Equals("pdf", StringComparison.OrdinalIgnoreCase))
{
PdfPrep.RenderAllPages(src, pagesDir, _cfg.Llm.MaxImageDim);
var files = Directory.GetFiles(pagesDir, "page_*.jpg").OrderBy(PageNumber).ToList();
if (files.Count == 0) throw new Exception("PDF produced no page images");
return files.Select(f => ("image/jpeg", Convert.ToBase64String(File.ReadAllBytes(f)))).ToList();
}
return new List<(string, string)>
{
(Util.MimeFromExt(ext), Convert.ToBase64String(File.ReadAllBytes(src)))
};
}
private static int PageNumber(string path)
{
var name = Path.GetFileNameWithoutExtension(path);
return int.TryParse(name.Replace("page_", ""), out var n) ? n : 0;
}
// Port of jsBuildParsedState from pending_process.asp - same output shape
// the review screen expects.
private static JsonNode BuildParsedState(JsonNode state, JsonNode api, JsonNode parsed, string backend)
{
var now = Util.NowIso();
state["stage"] = "parsed";
state["error"] = "";
state["updated_at"] = now;
state["header"] = new JsonObject
{
["vendor_id"] = "",
["vendor_name"] = parsed["vendor_name"]?.ToString() ?? "",
["vendor_hint"] = parsed["vendor_name"]?.ToString() ?? "",
["invoice_number"] = parsed["invoice_number"]?.ToString() ?? "",
["invoice_date"] = parsed["invoice_date"]?.ToString() ?? "",
["po_number"] = "",
["invoice_total"] = parsed["total"]?.ToString() ?? "",
["total_label"] = parsed["total_label"]?.ToString() ?? "",
["total_page"] = parsed["total_page"]?.ToString() ?? ""
};
var items = new JsonArray();
if (parsed["line_items"] is JsonArray src)
{
foreach (var itNode in src)
{
var it = itNode as JsonObject ?? new JsonObject();
items.Add(new JsonObject
{
["upc"] = it["upc"]?.ToString() ?? "",
["desc"] = it["description"]?.ToString() ?? "",
["cert"] = it["product_code"]?.ToString() ?? "",
["qty"] = it["quantity"]?.ToString() ?? "",
["cost"] = it["unit_price"]?.ToString() ?? "",
["total"] = it["total"]?.ToString() ?? "",
["pack"] = it["pack"]?.ToString() ?? "",
["size"] = it["size"]?.ToString() ?? "",
["uom"] = string.IsNullOrEmpty(it["unit_of_measure"]?.ToString()) ? "EA" : it["unit_of_measure"]!.ToString(),
["source_line"] = it["source_line"]?.ToString() ?? "",
["status"] = "unmatched",
["parsed_desc"] = it["description"]?.ToString() ?? "",
["matched_desc"] = "",
["match_badge_text"] = "",
["match_badge_class"] = "",
["po_status"] = ""
});
}
}
state["line_items"] = items;
state["page_images"] = new JsonArray(); // drop legacy heavy images now that parsing is done
state["llm"] = new JsonObject
{
["backend"] = backend,
["at"] = now,
["raw_response"] = api.DeepClone(),
["parsed"] = parsed.DeepClone()
};
return state;
}
}