fix: track in-flight requests by id; handle llama-swap upsert/remove events
llama-swap emits inflight events with operation 'upsert' (a single request, re-emitted as elapsed_ms updates) and id-only 'remove' events. The decoder ignored 'upsert' and dropped every id-only remove, so the tracker only ever reflected the connect-time snapshot and the count got stuck (e.g. at '1'). The tracker now keys in-flight state by request id (snapshot/add replace by id; remove deletes by id), and the decoder normalizes 'upsert' to 'add' and extracts the id from bare removes.
This commit is contained in:
@@ -9,6 +9,13 @@ export function decodeEvent(msg: SseMessage): FeedEvent | null {
|
|||||||
const inner = JSON.parse(outer.data) as Record<string, unknown>;
|
const inner = JSON.parse(outer.data) as Record<string, unknown>;
|
||||||
|
|
||||||
if (outer.type === "inflight") {
|
if (outer.type === "inflight") {
|
||||||
|
const operation = inner.operation as string;
|
||||||
|
if (operation === "remove") {
|
||||||
|
const id = typeof inner.id === "string" ? inner.id : undefined;
|
||||||
|
if (id) return { type: "inflight", operation: "remove", id };
|
||||||
|
return null;
|
||||||
|
}
|
||||||
|
if (operation !== "snapshot" && operation !== "upsert" && operation !== "add") return null;
|
||||||
const requests: InflightRequest[] = Array.isArray(inner.requests)
|
const requests: InflightRequest[] = Array.isArray(inner.requests)
|
||||||
? (inner.requests as InflightRequest[])
|
? (inner.requests as InflightRequest[])
|
||||||
: Array.isArray(inner.request)
|
: Array.isArray(inner.request)
|
||||||
@@ -18,7 +25,7 @@ export function decodeEvent(msg: SseMessage): FeedEvent | null {
|
|||||||
: [];
|
: [];
|
||||||
return {
|
return {
|
||||||
type: "inflight",
|
type: "inflight",
|
||||||
operation: inner.operation as "snapshot" | "add" | "remove",
|
operation: operation === "snapshot" ? "snapshot" : "add",
|
||||||
requests,
|
requests,
|
||||||
};
|
};
|
||||||
}
|
}
|
||||||
|
|||||||
+10
-12
@@ -13,29 +13,25 @@ export interface ModelState {
|
|||||||
}
|
}
|
||||||
|
|
||||||
export type FeedEvent =
|
export type FeedEvent =
|
||||||
| { type: "inflight"; operation: "snapshot" | "add" | "remove"; requests: InflightRequest[] }
|
| { type: "inflight"; operation: "snapshot" | "add"; requests: InflightRequest[] }
|
||||||
|
| { type: "inflight"; operation: "remove"; id: string }
|
||||||
| { type: "modelStatus"; models: ModelState[] };
|
| { type: "modelStatus"; models: ModelState[] };
|
||||||
|
|
||||||
export type ModelRuntimeState = "stopped" | "loading" | "ready";
|
export type ModelRuntimeState = "stopped" | "loading" | "ready";
|
||||||
|
|
||||||
export class InflightTracker {
|
export class InflightTracker {
|
||||||
private counts = new Map<string, number>();
|
private requests = new Map<string, string>();
|
||||||
private states = new Map<string, ModelRuntimeState>();
|
private states = new Map<string, ModelRuntimeState>();
|
||||||
|
|
||||||
apply(event: FeedEvent): void {
|
apply(event: FeedEvent): void {
|
||||||
if (event.type === "inflight") {
|
if (event.type === "inflight") {
|
||||||
if (event.operation === "snapshot") {
|
if (event.operation === "snapshot") {
|
||||||
const counts = new Map<string, number>();
|
this.requests = new Map<string, string>();
|
||||||
for (const r of event.requests) counts.set(r.model, (counts.get(r.model) ?? 0) + 1);
|
for (const r of event.requests) if (r.id) this.requests.set(r.id, r.model);
|
||||||
this.counts = counts;
|
|
||||||
} else if (event.operation === "add") {
|
} else if (event.operation === "add") {
|
||||||
for (const r of event.requests) this.counts.set(r.model, (this.counts.get(r.model) ?? 0) + 1);
|
for (const r of event.requests) if (r.id) this.requests.set(r.id, r.model);
|
||||||
} else if (event.operation === "remove") {
|
} else if (event.operation === "remove") {
|
||||||
for (const r of event.requests) {
|
this.requests.delete(event.id);
|
||||||
const c = this.counts.get(r.model) ?? 0;
|
|
||||||
if (c <= 1) this.counts.delete(r.model);
|
|
||||||
else this.counts.set(r.model, c - 1);
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
} else if (event.type === "modelStatus") {
|
} else if (event.type === "modelStatus") {
|
||||||
for (const m of event.models) this.states.set(m.id, normalizeState(m.state));
|
for (const m of event.models) this.states.set(m.id, normalizeState(m.state));
|
||||||
@@ -43,7 +39,9 @@ export class InflightTracker {
|
|||||||
}
|
}
|
||||||
|
|
||||||
count(modelId: string): number {
|
count(modelId: string): number {
|
||||||
return this.counts.get(modelId) ?? 0;
|
let n = 0;
|
||||||
|
for (const model of this.requests.values()) if (model === modelId) n++;
|
||||||
|
return n;
|
||||||
}
|
}
|
||||||
|
|
||||||
state(modelId: string): ModelRuntimeState | undefined {
|
state(modelId: string): ModelRuntimeState | undefined {
|
||||||
|
|||||||
@@ -41,6 +41,32 @@ test("decodeEvent parses an inflight add with a singular request field", () => {
|
|||||||
}
|
}
|
||||||
});
|
});
|
||||||
|
|
||||||
|
test("decodeEvent parses an inflight upsert as an add", () => {
|
||||||
|
const msg: SseMessage = {
|
||||||
|
event: "message",
|
||||||
|
data: JSON.stringify({
|
||||||
|
type: "inflight",
|
||||||
|
data: JSON.stringify({ operation: "upsert", request: { id: "820", model: "DeepSeek-V4-Flash-0731", elapsed_ms: 103 } }),
|
||||||
|
}),
|
||||||
|
};
|
||||||
|
const ev = decodeEvent(msg);
|
||||||
|
assert.ok(ev && ev.type === "inflight");
|
||||||
|
if (ev && ev.type === "inflight") {
|
||||||
|
assert.equal(ev.operation, "add");
|
||||||
|
assert.equal(ev.requests[0].id, "820");
|
||||||
|
assert.equal(ev.requests[0].elapsed_ms, 103);
|
||||||
|
}
|
||||||
|
});
|
||||||
|
|
||||||
|
test("decodeEvent parses an inflight remove with a bare id", () => {
|
||||||
|
const msg: SseMessage = {
|
||||||
|
event: "message",
|
||||||
|
data: JSON.stringify({ type: "inflight", data: JSON.stringify({ operation: "remove", id: "817" }) }),
|
||||||
|
};
|
||||||
|
const ev = decodeEvent(msg);
|
||||||
|
assert.deepEqual(ev, { type: "inflight", operation: "remove", id: "817" });
|
||||||
|
});
|
||||||
|
|
||||||
test("decodeEvent parses modelStatus", () => {
|
test("decodeEvent parses modelStatus", () => {
|
||||||
const msg: SseMessage = {
|
const msg: SseMessage = {
|
||||||
event: "message",
|
event: "message",
|
||||||
|
|||||||
@@ -18,17 +18,32 @@ test("snapshot rebuilds counts for all in-flight requests", () => {
|
|||||||
assert.equal(tracker.count("C"), 0);
|
assert.equal(tracker.count("C"), 0);
|
||||||
});
|
});
|
||||||
|
|
||||||
test("add and remove adjust per-model counts", () => {
|
test("add (upsert) is idempotent per request id", () => {
|
||||||
const tracker = new InflightTracker();
|
const tracker = new InflightTracker();
|
||||||
tracker.apply({ type: "inflight", operation: "add", requests: [{ model: "A", id: "1" }] });
|
tracker.apply({ type: "inflight", operation: "add", requests: [{ model: "A", id: "1" }] });
|
||||||
|
tracker.apply({ type: "inflight", operation: "add", requests: [{ model: "A", id: "1" }] });
|
||||||
tracker.apply({ type: "inflight", operation: "add", requests: [{ model: "A", id: "2" }] });
|
tracker.apply({ type: "inflight", operation: "add", requests: [{ model: "A", id: "2" }] });
|
||||||
assert.equal(tracker.count("A"), 2);
|
assert.equal(tracker.count("A"), 2);
|
||||||
tracker.apply({ type: "inflight", operation: "remove", requests: [{ model: "A", id: "1" }] });
|
});
|
||||||
|
|
||||||
|
test("remove is id-based and decrements the matching request", () => {
|
||||||
|
const tracker = new InflightTracker();
|
||||||
|
tracker.apply({ type: "inflight", operation: "add", requests: [{ model: "A", id: "1" }, { model: "B", id: "2" }] });
|
||||||
assert.equal(tracker.count("A"), 1);
|
assert.equal(tracker.count("A"), 1);
|
||||||
tracker.apply({ type: "inflight", operation: "remove", requests: [{ model: "A", id: "2" }] });
|
assert.equal(tracker.count("B"), 1);
|
||||||
assert.equal(tracker.count("A"), 0);
|
tracker.apply({ type: "inflight", operation: "remove", id: "1" });
|
||||||
tracker.apply({ type: "inflight", operation: "remove", requests: [{ model: "A", id: "9" }] });
|
|
||||||
assert.equal(tracker.count("A"), 0);
|
assert.equal(tracker.count("A"), 0);
|
||||||
|
assert.equal(tracker.count("B"), 1);
|
||||||
|
tracker.apply({ type: "inflight", operation: "remove", id: "999" });
|
||||||
|
assert.equal(tracker.count("B"), 1);
|
||||||
|
});
|
||||||
|
|
||||||
|
test("snapshot replaces prior in-flight state", () => {
|
||||||
|
const tracker = new InflightTracker();
|
||||||
|
tracker.apply({ type: "inflight", operation: "add", requests: [{ model: "A", id: "1" }] });
|
||||||
|
tracker.apply({ type: "inflight", operation: "snapshot", requests: [{ model: "A", id: "2" }] });
|
||||||
|
assert.equal(tracker.count("A"), 1);
|
||||||
|
assert.equal(tracker.count("B"), 0);
|
||||||
});
|
});
|
||||||
|
|
||||||
test("modelStatus normalizes states", () => {
|
test("modelStatus normalizes states", () => {
|
||||||
|
|||||||
Reference in New Issue
Block a user