From 0d752d79588d4fa147d2830ce657ccc89cedfcdc Mon Sep 17 00:00:00 2001 From: zzstoatzz Date: Fri, 10 Apr 2026 11:06:34 -0500 Subject: [PATCH] =?UTF-8?q?add=20FIFO=20job=20queue=20for=20indexing=20?= =?UTF-8?q?=E2=80=94=20sequential=20processing=20with=20position=20feedbac?= =?UTF-8?q?k?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit replaces the previous model (one thread per job, all fighting over the embedder mutex) with a single worker thread draining a FIFO queue. - jobs run one at a time: CAR download + embed, then next job - status response includes queue_position and queue_depth - frontend shows "N jobs ahead of you" while waiting - predictable throughput, honest wait times Co-Authored-By: Claude Opus 4.6 (1M context) --- backend/src/assets/main.js | 5 ++ backend/src/main.zig | 10 ++++ backend/src/server.zig | 94 +++++++++++++++++++++++++++++++++++--- 3 files changed, 102 insertions(+), 7 deletions(-) diff --git a/backend/src/assets/main.js b/backend/src/assets/main.js index 3d1a4ed..b372746 100644 --- a/backend/src/assets/main.js +++ b/backend/src/assets/main.js @@ -408,6 +408,11 @@ async function pollStatus() { } else { updateProgress({ fetched, embedded: 0, reused: 0, walking: true }); } + // show queue position if waiting behind other jobs + const pos = j.queue_position; + if (pos != null && pos > 0) { + progressRate.textContent = `${pos} ${pos === 1 ? "job" : "jobs"} ahead of you`; + } } } catch (e) { stopProgress(); diff --git a/backend/src/main.zig b/backend/src/main.zig index 7119d58..f9e2c40 100644 --- a/backend/src/main.zig +++ b/backend/src/main.zig @@ -73,15 +73,25 @@ pub fn main(init: std.process.Init) !void { std.log.warn("EMBED_DEV_NOAUTH=1 — consent gate bypassed. DO NOT DEPLOY WITH THIS.", .{}); } + var job_queue: server.JobQueue = .{}; + var app = server.App{ .io = io, .allocator = allocator, .embedder = &embedder, .embedder_mutex = &embedder_mutex, .cache = &cache, + .job_queue = &job_queue, .dev_noauth = dev_noauth, }; + // start the single job queue worker thread + const worker = std.Thread.spawn(.{}, server.jobQueueWorker, .{&app}) catch |err| { + std.log.err("failed to spawn job queue worker: {t}", .{err}); + return err; + }; + worker.detach(); + var addr = try Io.net.IpAddress.parse("::", port); var listener = addr.listen(io, .{ .reuse_address = true }) catch |err| { std.log.err("failed to listen on port {d}: {t}", .{ port, err }); diff --git a/backend/src/server.zig b/backend/src/server.zig index 86d8bac..87a01a9 100644 --- a/backend/src/server.zig +++ b/backend/src/server.zig @@ -52,11 +52,65 @@ pub const App = struct { embedder: *Llama.Embedder, embedder_mutex: *Io.Mutex, cache: *indexer.Cache, + job_queue: *JobQueue, /// dev-only: skip the consent gate in handleIndex. set from /// EMBED_DEV_NOAUTH=1 at startup. never true in prod. dev_noauth: bool = false, }; +/// FIFO job queue processed by a single worker thread. jobs run one at a +/// time so embedding throughput is predictable and the frontend can show +/// an accurate queue position ("3 jobs ahead of you"). +pub const JobQueue = struct { + mutex: Io.Mutex = .init, + items: std.ArrayListUnmanaged(*indexer.Job) = .empty, + /// DID of the job currently being processed (null if idle) + current_did: ?[]const u8 = null, + + fn push(self: *JobQueue, io: Io, allocator: Allocator, job: *indexer.Job) !void { + self.mutex.lockUncancelable(io); + defer self.mutex.unlock(io); + try self.items.append(allocator, job); + } + + fn pop(self: *JobQueue, io: Io) ?*indexer.Job { + self.mutex.lockUncancelable(io); + defer self.mutex.unlock(io); + if (self.items.items.len == 0) return null; + return self.items.orderedRemove(0); + } + + /// how many jobs are ahead of the given DID? returns null if the DID + /// isn't in the queue or currently running. + fn position(self: *JobQueue, io: Io, did: []const u8) ?u32 { + self.mutex.lockUncancelable(io); + defer self.mutex.unlock(io); + + // currently being processed = position 0 + if (self.current_did) |cur| { + if (mem.eql(u8, cur, did)) return 0; + } + + // search the queue + for (self.items.items, 0..) |job, i| { + if (mem.eql(u8, job.pack.did, did)) { + // +1 because the currently-running job is ahead of everything in the queue + return @intCast(if (self.current_did != null) i + 1 else i); + } + } + return null; + } + + /// total number of jobs waiting + the one currently running + fn depth(self: *JobQueue, io: Io) u32 { + self.mutex.lockUncancelable(io); + defer self.mutex.unlock(io); + var d: u32 = @intCast(self.items.items.len); + if (self.current_did != null) d += 1; + return d; + } +}; + // ---------- connection handling ---------- pub fn handleConnection(stream: Io.net.Stream, app: *App) void { @@ -582,13 +636,12 @@ fn kickoffIndexing(request: *http.Server.Request, app: *App, handle: []const u8) .max_per_collection = 0, }; - const t = Thread.spawn(.{}, indexerWorker, .{job}) catch |err| { - std.log.err("failed to spawn indexer thread: {t}", .{err}); + app.job_queue.push(app.io, app.allocator, job) catch |err| { + std.log.err("failed to enqueue indexer job: {t}", .{err}); app.allocator.destroy(job); - try sendJsonStatus(request, .internal_server_error, "{\"error\":\"failed to spawn worker\"}"); + try sendJsonStatus(request, .internal_server_error, "{\"error\":\"failed to enqueue job\"}"); return; }; - t.detach(); try writeStatusResponse(request, app, pack); } @@ -646,9 +699,31 @@ fn handleShareLoad(request: *http.Server.Request, app: *App, handle: []const u8) try kickoffIndexing(request, app, handle); } -fn indexerWorker(job: *indexer.Job) void { - defer job.allocator.destroy(job); - indexer.runJob(job); +/// single worker thread that drains the job queue sequentially. started +/// once at init and runs for the lifetime of the process. polls every +/// 100ms when idle. +pub fn jobQueueWorker(app: *App) void { + while (true) { + const job = app.job_queue.pop(app.io) orelse { + // idle — sleep briefly then check again + app.io.sleep(.{ .nanoseconds = 100_000_000 }, .real) catch {}; + continue; + }; + { + app.job_queue.mutex.lockUncancelable(app.io); + app.job_queue.current_did = job.pack.did; + app.job_queue.mutex.unlock(app.io); + } + + indexer.runJob(job); + + { + app.job_queue.mutex.lockUncancelable(app.io); + app.job_queue.current_did = null; + app.job_queue.mutex.unlock(app.io); + } + job.allocator.destroy(job); + } } // ---------- route: /api/status/:handle ---------- @@ -703,6 +778,11 @@ fn writeStatusResponse(request: *http.Server.Request, app: *App, pack: *indexer. try buf.print(alloc, "\"build_ms\":{d},", .{pack.build_ms}); try buf.print(alloc, "\"prior_build_ms\":{d},", .{pack.prior_build_ms}); try buf.print(alloc, "\"prior_count\":{d},", .{pack.prior_count}); + // queue position for this job (null if not queued / already done) + if (app.job_queue.position(app.io, pack.did)) |pos| { + try buf.print(alloc, "\"queue_position\":{d},", .{pos}); + } + try buf.print(alloc, "\"queue_depth\":{d},", .{app.job_queue.depth(app.io)}); if (pack.persisted_uri) |uri| { try buf.appendSlice(alloc, "\"persisted\":true,\"persisted_uri\":"); try writeJsonString(&buf, alloc, uri); -- 2.51.2