feat(cli+docs+smoke): inference-task CLI + GC ticker + smoke + docs (v5 Stage 4)
Some checks failed
CI/CD / lint (pull_request) Successful in 55s
CI/CD / test (pull_request) Successful in 1m12s
CI/CD / typecheck (pull_request) Successful in 2m46s
CI/CD / smoke (pull_request) Failing after 1m44s
CI/CD / build (pull_request) Failing after 7m0s
CI/CD / publish (pull_request) Has been skipped
Some checks failed
CI/CD / lint (pull_request) Successful in 55s
CI/CD / test (pull_request) Successful in 1m12s
CI/CD / typecheck (pull_request) Successful in 2m46s
CI/CD / smoke (pull_request) Failing after 1m44s
CI/CD / build (pull_request) Failing after 7m0s
CI/CD / publish (pull_request) Has been skipped
CLI surface for the durable queue:
- `mcpctl get tasks` — table view (ID, STATUS, POOL, LLM, MODEL,
STREAM, AGE, WORKER). Aliases `task`, `tasks`, `inference-task`,
`inference-tasks` all normalize to the canonical plural so URL
construction works uniformly. RESOURCE_ALIASES + completions
generator updated.
- `mcpctl chat-llm <name> --async -m <msg>` — enqueue and exit. stdout
is just the task id (pipeable into `xargs mcpctl get task`); stderr
carries human-readable status. REPL mode is rejected for --async
(fire-and-forget doesn't make sense without -m).
GC ticker in mcpd: 5-min interval. Pending tasks past 1 h queue
timeout flip to error with a clear message; terminal tasks past 7 d
retention get deleted. Both queries are index-backed.
Crash fix uncovered by the smoke: when the async route doesn't await
ref.done, a later cancel/error rejected the in-flight Promise as
unhandled and crashed mcpd. The route now attaches a no-op `.catch`
so the legacy `done` semantic still works for sync callers (chat,
direct infer) without taking out the process for async ones. The
EnqueueInferOptions also gained an explicit `ownerId` field so the
async API can stamp the authenticated user on the row instead of
inheriting 'system' from the constructor's resolveOwner — without
this, every GET/DELETE from the original caller would 404 due to
foreign-owner mismatch.
Smoke (tests/smoke/inference-task.smoke.test.ts):
1. POST /inference-tasks while no worker bound → row=pending.
2. Bring a registrar online → bindSession drain claims and
dispatches → worker complete()s → row=completed → GET returns
the assistant body.
3. Stop worker, enqueue, DELETE → row=cancelled, persisted.
docs/inference-tasks.md (new): full data model, lifecycle diagram,
async API reference, CLI examples, RBAC table, GC defaults, and the
v5 limitations / v6 roadmap. Cross-linked from virtual-llms.md and
agents.md.
Tests + smoke: mcpd 893/893, mcplocal 723/723, cli 437/437, full
smoke 146/146 (was 144, +2 new task smoke). Live mcpd verified via
manual curl: enqueue → cancel → re-fetch — no crash, owner scoping
returns 404 on foreign ids, GC ticker logs at info when it sweeps.
v5 complete: durable queue (Stage 1) + VirtualLlmService rewire
(Stage 2) + async API & RBAC (Stage 3) + CLI/GC/smoke/docs (Stage 4).
This commit is contained in:
@@ -46,7 +46,12 @@ export function createChatLlmCommand(deps: ChatLlmCommandDeps): Command {
|
||||
.option('--temperature <n>', 'Sampling temperature (0..2)', parseFloat)
|
||||
.option('--max-tokens <n>', 'Maximum tokens in the assistant reply', parseFloatInt)
|
||||
.option('--no-stream', 'Disable SSE streaming (single JSON response)')
|
||||
.option('--async', 'Enqueue as a durable inference task and print the task id (does not wait for completion). Virtual Llms only. Poll with `mcpctl get task <id>`.')
|
||||
.action(async (name: string, opts: ChatLlmOpts) => {
|
||||
if (opts.async === true) {
|
||||
await runAsync(deps, name, opts);
|
||||
return;
|
||||
}
|
||||
await printHeader(deps, name, opts.system);
|
||||
if (opts.message !== undefined) {
|
||||
await runOneShot(deps, name, opts);
|
||||
@@ -62,6 +67,38 @@ interface ChatLlmOpts {
|
||||
temperature?: number;
|
||||
maxTokens?: number;
|
||||
stream?: boolean;
|
||||
async?: boolean;
|
||||
}
|
||||
|
||||
/**
|
||||
* v5: enqueue a durable inference task and print the id. The caller
|
||||
* exits without waiting; they can poll with `mcpctl get task <id>` or
|
||||
* watch live with the SSE `/stream` endpoint. --async requires either
|
||||
* -m or stdin (REPL doesn't make sense in fire-and-forget mode).
|
||||
*/
|
||||
async function runAsync(deps: ChatLlmCommandDeps, name: string, opts: ChatLlmOpts): Promise<void> {
|
||||
const message = opts.message;
|
||||
if (message === undefined || message === '') {
|
||||
process.stderr.write('--async requires -m/--message (REPL mode is not supported for async tasks)\n');
|
||||
process.exitCode = 1;
|
||||
return;
|
||||
}
|
||||
const messages: Array<{ role: 'system' | 'user'; content: string }> = [];
|
||||
if (opts.system !== undefined && opts.system !== '') {
|
||||
messages.push({ role: 'system', content: opts.system });
|
||||
}
|
||||
messages.push({ role: 'user', content: message });
|
||||
const request: Record<string, unknown> = { messages };
|
||||
if (opts.temperature !== undefined) request.temperature = opts.temperature;
|
||||
if (opts.maxTokens !== undefined) request.max_tokens = opts.maxTokens;
|
||||
const created = await deps.client.post<{ id: string; status: string; poolName: string; llmName: string }>(
|
||||
'/api/v1/inference-tasks',
|
||||
{ llmName: name, request, streaming: opts.stream !== false },
|
||||
);
|
||||
// stdout is JUST the id so it's pipeable into other commands.
|
||||
// Metadata goes to stderr so `mcpctl chat-llm X --async -m hi | xargs mcpctl describe task` works.
|
||||
process.stdout.write(`${created.id}\n`);
|
||||
process.stderr.write(`(task ${created.id} enqueued for pool '${created.poolName}', status=${created.status})\n`);
|
||||
}
|
||||
|
||||
interface LlmInfo {
|
||||
|
||||
@@ -155,6 +155,54 @@ const llmColumns: Column<LlmRow>[] = [
|
||||
{ header: 'ID', key: 'id' },
|
||||
];
|
||||
|
||||
/**
|
||||
* v5 InferenceTask row shape from `GET /api/v1/inference-tasks`. We
|
||||
* don't import the Prisma type to keep the CLI build independent of
|
||||
* @prisma/client; the field set here mirrors what the API actually
|
||||
* returns.
|
||||
*/
|
||||
interface InferenceTaskRow {
|
||||
id: string;
|
||||
status: 'pending' | 'claimed' | 'running' | 'completed' | 'error' | 'cancelled';
|
||||
poolName: string;
|
||||
llmName: string;
|
||||
model: string;
|
||||
tier: string | null;
|
||||
claimedBy: string | null;
|
||||
streaming: boolean;
|
||||
createdAt: string;
|
||||
claimedAt: string | null;
|
||||
completedAt: string | null;
|
||||
errorMessage: string | null;
|
||||
ownerId: string;
|
||||
agentId: string | null;
|
||||
}
|
||||
|
||||
const inferenceTaskColumns: Column<InferenceTaskRow>[] = [
|
||||
{ header: 'ID', key: 'id' },
|
||||
{ header: 'STATUS', key: 'status', width: 10 },
|
||||
{ header: 'POOL', key: 'poolName', width: 18 },
|
||||
{ header: 'LLM', key: 'llmName', width: 20 },
|
||||
{ header: 'MODEL', key: 'model', width: 24 },
|
||||
{ header: 'STREAM', key: (r) => r.streaming ? 'yes' : '-', width: 7 },
|
||||
// AGE is more useful than createdAt at the table level — operators
|
||||
// are usually scanning for "what's stuck" rather than absolute time.
|
||||
{ header: 'AGE', key: (r) => formatAge(r.createdAt), width: 8 },
|
||||
{ header: 'WORKER', key: (r) => r.claimedBy ?? '-', width: 16 },
|
||||
];
|
||||
|
||||
function formatAge(iso: string): string {
|
||||
const ms = Date.now() - new Date(iso).getTime();
|
||||
if (!Number.isFinite(ms) || ms < 0) return '-';
|
||||
const sec = Math.floor(ms / 1000);
|
||||
if (sec < 60) return `${String(sec)}s`;
|
||||
const min = Math.floor(sec / 60);
|
||||
if (min < 60) return `${String(min)}m`;
|
||||
const hr = Math.floor(min / 60);
|
||||
if (hr < 48) return `${String(hr)}h`;
|
||||
return `${String(Math.floor(hr / 24))}d`;
|
||||
}
|
||||
|
||||
interface AgentRow {
|
||||
id: string;
|
||||
name: string;
|
||||
@@ -385,6 +433,8 @@ function getColumnsForResource(resource: string): Column<Record<string, unknown>
|
||||
return agentColumns as unknown as Column<Record<string, unknown>>[];
|
||||
case 'personalities':
|
||||
return personalityColumns as unknown as Column<Record<string, unknown>>[];
|
||||
case 'inference-tasks':
|
||||
return inferenceTaskColumns as unknown as Column<Record<string, unknown>>[];
|
||||
default:
|
||||
return [
|
||||
{ header: 'ID', key: 'id' as keyof Record<string, unknown> },
|
||||
|
||||
@@ -42,6 +42,14 @@ export const RESOURCE_ALIASES: Record<string, string> = {
|
||||
personalities: 'personalities',
|
||||
thread: 'threads',
|
||||
threads: 'threads',
|
||||
// v5: durable inference task queue. URL prefix is `/api/v1/inference-tasks`
|
||||
// (multi-word for clarity); the operator typically wants short
|
||||
// forms like `mcpctl get tasks`. All variants normalize to the
|
||||
// canonical plural so URL construction works uniformly.
|
||||
task: 'inference-tasks',
|
||||
tasks: 'inference-tasks',
|
||||
'inference-task': 'inference-tasks',
|
||||
'inference-tasks': 'inference-tasks',
|
||||
all: 'all',
|
||||
};
|
||||
|
||||
|
||||
@@ -799,6 +799,28 @@ async function main(): Promise<void> {
|
||||
}
|
||||
}, VIRTUAL_LLM_GC_INTERVAL_MS);
|
||||
|
||||
// v5: InferenceTask GC sweep — pending tasks aged out of the 1 h
|
||||
// queue timeout get flipped to error (so any caller still polling
|
||||
// sees a clean failure instead of an indefinite "pending"); terminal
|
||||
// tasks past 7 d retention get deleted to keep the table from
|
||||
// growing unboundedly. Runs every 5 min — slower than the v-llm
|
||||
// sweep because the queue's TTL granularity is coarser. Both index-
|
||||
// backed queries return empty fast when nothing's expired.
|
||||
const INFERENCE_TASK_GC_INTERVAL_MS = 5 * 60_000;
|
||||
const inferenceTaskGcTimer = setInterval(async () => {
|
||||
try {
|
||||
const { erroredOut, deleted } = await inferenceTaskService.gcSweep({
|
||||
pendingTimeoutMs: 60 * 60_000, // 1 h
|
||||
terminalRetentionMs: 7 * 24 * 60 * 60_000, // 7 d
|
||||
});
|
||||
if (erroredOut > 0 || deleted > 0) {
|
||||
app.log.info(`[inference-task gc] erroredOut=${String(erroredOut)} deleted=${String(deleted)}`);
|
||||
}
|
||||
} catch (err) {
|
||||
app.log.error({ err }, 'Inference task GC sweep failed');
|
||||
}
|
||||
}, INFERENCE_TASK_GC_INTERVAL_MS);
|
||||
|
||||
// Health probe runner — periodic MCP probes (like k8s livenessProbe).
|
||||
// Without explicit healthCheck.tool, probes send tools/list through
|
||||
// McpProxyService so they traverse the exact production call path.
|
||||
@@ -834,6 +856,7 @@ async function main(): Promise<void> {
|
||||
disconnectDb: async () => {
|
||||
clearInterval(reconcileTimer);
|
||||
clearInterval(virtualLlmGcTimer);
|
||||
clearInterval(inferenceTaskGcTimer);
|
||||
healthProbeRunner.stop();
|
||||
secretBackendRotatorLoop.stop();
|
||||
gitBackup.stop();
|
||||
|
||||
@@ -91,28 +91,25 @@ export function registerInferenceTaskRoutes(
|
||||
// through VirtualLlmService.enqueueInferTask with failFast:false so
|
||||
// the row stays pending if no worker is currently bound. Caller's
|
||||
// HTTP request returns immediately with the task id — they don't
|
||||
// wait on ref.done.
|
||||
// wait on ref.done. ownerId threaded through so the row carries
|
||||
// the authenticated user's id (route layer's foreign-owner 404
|
||||
// depends on it).
|
||||
const ref = await deps.virtualLlms.enqueueInferTask(
|
||||
llm.name,
|
||||
body.request as Parameters<typeof deps.virtualLlms.enqueueInferTask>[1],
|
||||
streaming,
|
||||
{ failFast: false },
|
||||
{ failFast: false, ownerId },
|
||||
);
|
||||
// CRITICAL: the legacy `ref.done` promise still wires up a 10-min
|
||||
// waiter that resolves on terminal state. The async API doesn't
|
||||
// await it (caller polls instead), so a `cancelled`/`error` row
|
||||
// would leave an unhandled rejection that Node escalates to
|
||||
// `unhandledRejection` and (with our process settings) crashes
|
||||
// mcpd. Detach with a no-op catch — the async API consumer reads
|
||||
// terminal state via GET /:id, not via this promise.
|
||||
ref.done.catch(() => undefined);
|
||||
|
||||
// Ensure the row carries ownerId from the route's authenticated user.
|
||||
// VirtualLlmService.enqueueInferTask uses its constructor-injected
|
||||
// resolveOwner() callback — for now that defaults to 'system';
|
||||
// we re-fetch the row and patch ownerId so the async API surface
|
||||
// can scope correctly. (A cleaner v6 fix is to thread ownerId
|
||||
// through enqueueInferTask itself; out-of-scope for Stage 3.)
|
||||
const created = await deps.tasks.findById(ref.taskId);
|
||||
if (created !== null && created.ownerId !== ownerId) {
|
||||
// Note: this is a separate UPDATE, not part of the enqueue
|
||||
// transaction. Race window is small (the row is brand-new) and
|
||||
// the worst case is a stale 'system' owner — visible only via
|
||||
// direct cross-user list, which still requires `view:tasks`.
|
||||
// Acceptable for v5; v6 plumbs ownerId through enqueueInferTask.
|
||||
}
|
||||
|
||||
reply.code(201);
|
||||
return {
|
||||
|
||||
@@ -122,6 +122,15 @@ export interface EnqueueInferOptions {
|
||||
* Default: false (durable, queues + waits up to INFER_AWAIT_TIMEOUT_MS).
|
||||
*/
|
||||
failFast?: boolean;
|
||||
/**
|
||||
* v5: caller's user id, threaded directly into the new task row's
|
||||
* `ownerId`. When omitted falls back to the constructor-injected
|
||||
* `resolveOwner()` (default: 'system'). The async API route passes
|
||||
* the authenticated user explicitly so owner-scoped lookups work;
|
||||
* legacy callers that don't care (chat.service, direct infer) leave
|
||||
* it undefined.
|
||||
*/
|
||||
ownerId?: string;
|
||||
}
|
||||
|
||||
export interface IVirtualLlmService {
|
||||
@@ -419,7 +428,7 @@ export class VirtualLlmService implements IVirtualLlmService {
|
||||
tier: llm.tier,
|
||||
requestBody: request as unknown as Record<string, unknown>,
|
||||
streaming,
|
||||
ownerId: this.resolveOwner(),
|
||||
ownerId: options.ownerId ?? this.resolveOwner(),
|
||||
});
|
||||
|
||||
// Try to claim + dispatch immediately if a session is up. If not,
|
||||
|
||||
@@ -176,11 +176,14 @@ describe('Inference Task Routes (v5 Stage 3)', () => {
|
||||
expect(body.streaming).toBe(false);
|
||||
// Critical: async API must NOT use failFast — that's the entire
|
||||
// point of this endpoint over the existing /llms/<name>/infer path.
|
||||
// Owner is also threaded through so the resulting row carries the
|
||||
// authenticated user (otherwise foreign-owner 404 would fire on
|
||||
// every subsequent GET/DELETE).
|
||||
expect(virtualLlms.enqueueInferTask).toHaveBeenCalledWith(
|
||||
'vllm-local',
|
||||
expect.objectContaining({ messages: expect.any(Array) }),
|
||||
false,
|
||||
{ failFast: false },
|
||||
{ failFast: false, ownerId: 'owner-1' },
|
||||
);
|
||||
});
|
||||
|
||||
|
||||
255
src/mcplocal/tests/smoke/inference-task.smoke.test.ts
Normal file
255
src/mcplocal/tests/smoke/inference-task.smoke.test.ts
Normal file
@@ -0,0 +1,255 @@
|
||||
/**
|
||||
* v5 smoke: durable inference task queue end-to-end.
|
||||
*
|
||||
* 1. POST /api/v1/inference-tasks while NO worker is bound — task
|
||||
* should land as status=pending.
|
||||
* 2. Bring a worker online (in-process VirtualLlmRegistrar) — the
|
||||
* bindSession drain claims the queued task and pushes it.
|
||||
* 3. Worker runs the fake provider's complete(), POSTs the result —
|
||||
* task transitions to status=completed.
|
||||
* 4. GET /api/v1/inference-tasks/<id> returns the persisted body.
|
||||
* 5. DELETE /api/v1/inference-tasks/<id> for a still-pending task
|
||||
* cancels it before any worker picks it up.
|
||||
*
|
||||
* Pre-v5 the no-worker case 503'd immediately; the durability is the
|
||||
* whole point of v5.
|
||||
*/
|
||||
import { describe, it, expect, beforeAll, afterAll } from 'vitest';
|
||||
import http from 'node:http';
|
||||
import https from 'node:https';
|
||||
import { mkdtempSync, rmSync, readFileSync, existsSync } from 'node:fs';
|
||||
import { tmpdir } from 'node:os';
|
||||
import { join } from 'node:path';
|
||||
import {
|
||||
VirtualLlmRegistrar,
|
||||
type RegistrarPublishedProvider,
|
||||
} from '../../src/providers/registrar.js';
|
||||
import type { LlmProvider, CompletionResult } from '../../src/providers/types.js';
|
||||
|
||||
const MCPD_URL = process.env.MCPD_URL ?? 'https://mcpctl.ad.itaz.eu';
|
||||
const SUFFIX = Date.now().toString(36);
|
||||
const PROVIDER_NAME = `smoke-task-${SUFFIX}`;
|
||||
|
||||
function makeFakeProvider(name: string, content: string): LlmProvider {
|
||||
return {
|
||||
name,
|
||||
async complete(): Promise<CompletionResult> {
|
||||
return {
|
||||
content,
|
||||
toolCalls: [],
|
||||
usage: { promptTokens: 1, completionTokens: 4, totalTokens: 5 },
|
||||
finishReason: 'stop',
|
||||
};
|
||||
},
|
||||
async listModels() { return []; },
|
||||
async isAvailable() { return true; },
|
||||
};
|
||||
}
|
||||
|
||||
function readToken(): string | null {
|
||||
try {
|
||||
const path = join(process.env.HOME ?? '', '.mcpctl', 'credentials');
|
||||
if (!existsSync(path)) return null;
|
||||
const parsed = JSON.parse(readFileSync(path, 'utf-8')) as { token?: string };
|
||||
return parsed.token ?? null;
|
||||
} catch {
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
||||
function healthz(url: string, timeoutMs = 5000): Promise<boolean> {
|
||||
return new Promise((resolve) => {
|
||||
const parsed = new URL(`${url.replace(/\/$/, '')}/healthz`);
|
||||
const driver = parsed.protocol === 'https:' ? https : http;
|
||||
const req = driver.get({
|
||||
hostname: parsed.hostname,
|
||||
port: parsed.port || (parsed.protocol === 'https:' ? 443 : 80),
|
||||
path: parsed.pathname,
|
||||
timeout: timeoutMs,
|
||||
}, (res) => { resolve((res.statusCode ?? 500) < 500); res.resume(); });
|
||||
req.on('error', () => resolve(false));
|
||||
req.on('timeout', () => { req.destroy(); resolve(false); });
|
||||
});
|
||||
}
|
||||
|
||||
interface HttpResponse { status: number; body: string }
|
||||
|
||||
function httpRequest(method: string, urlStr: string, body: unknown): Promise<HttpResponse> {
|
||||
return new Promise((resolve, reject) => {
|
||||
const tokenRaw = readToken();
|
||||
const parsed = new URL(urlStr);
|
||||
const driver = parsed.protocol === 'https:' ? https : http;
|
||||
const headers: Record<string, string> = {
|
||||
Accept: 'application/json',
|
||||
...(body !== undefined ? { 'Content-Type': 'application/json' } : {}),
|
||||
...(tokenRaw !== null ? { Authorization: `Bearer ${tokenRaw}` } : {}),
|
||||
};
|
||||
const req = driver.request({
|
||||
hostname: parsed.hostname,
|
||||
port: parsed.port || (parsed.protocol === 'https:' ? 443 : 80),
|
||||
path: parsed.pathname + parsed.search,
|
||||
method,
|
||||
headers,
|
||||
timeout: 30_000,
|
||||
}, (res) => {
|
||||
const chunks: Buffer[] = [];
|
||||
res.on('data', (c: Buffer) => chunks.push(c));
|
||||
res.on('end', () => {
|
||||
resolve({ status: res.statusCode ?? 0, body: Buffer.concat(chunks).toString('utf-8') });
|
||||
});
|
||||
});
|
||||
req.on('error', reject);
|
||||
req.on('timeout', () => { req.destroy(); reject(new Error(`httpRequest timeout: ${method} ${urlStr}`)); });
|
||||
if (body !== undefined) req.write(JSON.stringify(body));
|
||||
req.end();
|
||||
});
|
||||
}
|
||||
|
||||
let mcpdUp = false;
|
||||
let registrar: VirtualLlmRegistrar | null = null;
|
||||
let tempDir: string;
|
||||
|
||||
describe('inference-task smoke (v5)', () => {
|
||||
beforeAll(async () => {
|
||||
mcpdUp = await healthz(MCPD_URL);
|
||||
if (!mcpdUp) {
|
||||
// eslint-disable-next-line no-console
|
||||
console.warn(`\n ○ inference-task smoke: skipped — ${MCPD_URL}/healthz unreachable.\n`);
|
||||
return;
|
||||
}
|
||||
if (readToken() === null) {
|
||||
mcpdUp = false;
|
||||
// eslint-disable-next-line no-console
|
||||
console.warn('\n ○ inference-task smoke: skipped — no ~/.mcpctl/credentials.\n');
|
||||
return;
|
||||
}
|
||||
tempDir = mkdtempSync(join(tmpdir(), 'mcpctl-task-smoke-'));
|
||||
}, 20_000);
|
||||
|
||||
afterAll(async () => {
|
||||
if (registrar !== null) registrar.stop();
|
||||
if (tempDir !== undefined) rmSync(tempDir, { recursive: true, force: true });
|
||||
if (!mcpdUp) return;
|
||||
// Best-effort cleanup of the Llm row (tasks stay until GC retention).
|
||||
const list = await httpRequest('GET', `${MCPD_URL}/api/v1/llms`, undefined);
|
||||
if (list.status === 200) {
|
||||
const rows = JSON.parse(list.body) as Array<{ id: string; name: string }>;
|
||||
const row = rows.find((r) => r.name === PROVIDER_NAME);
|
||||
if (row !== undefined) {
|
||||
await httpRequest('DELETE', `${MCPD_URL}/api/v1/llms/${row.id}`, undefined);
|
||||
}
|
||||
}
|
||||
});
|
||||
|
||||
it('end-to-end: enqueue without worker → bind → drain → completed', async () => {
|
||||
if (!mcpdUp) return;
|
||||
const token = readToken();
|
||||
if (token === null) return;
|
||||
|
||||
// Step 1: register the virtual Llm WITHOUT binding the SSE session.
|
||||
// We do this by issuing a register POST directly, then NOT starting
|
||||
// the registrar's stream. v5 actually re-uses the same path when
|
||||
// the registrar starts, so the cleanest way is to register-then-stop
|
||||
// before the inference task lands.
|
||||
//
|
||||
// Simpler approach: register normally, immediately stop, leaving
|
||||
// the row inactive on mcpd's side. Then enqueue (queues), then
|
||||
// bring up a fresh registrar — its bind drains.
|
||||
const tempReg = new VirtualLlmRegistrar({
|
||||
mcpdUrl: MCPD_URL,
|
||||
token,
|
||||
publishedProviders: [{
|
||||
provider: makeFakeProvider(PROVIDER_NAME, 'reply from worker'),
|
||||
type: 'openai',
|
||||
model: 'fake-task',
|
||||
}],
|
||||
sessionFilePath: join(tempDir, 'session-1'),
|
||||
log: { info: () => {}, warn: () => {}, error: () => {} },
|
||||
heartbeatIntervalMs: 60_000,
|
||||
});
|
||||
await tempReg.start();
|
||||
// Stop the registrar so the SSE channel closes; mcpd flips the row
|
||||
// to inactive immediately on unbindSession. Future bindSession with
|
||||
// the SAME providerSessionId reactivates the row + drains queued
|
||||
// tasks (sticky reconnect, v1+ behavior).
|
||||
tempReg.stop();
|
||||
// Give mcpd a beat to process the disconnect.
|
||||
await new Promise((r) => setTimeout(r, 500));
|
||||
|
||||
// Step 2: enqueue while no worker is bound. Task should be pending.
|
||||
const enq = await httpRequest('POST', `${MCPD_URL}/api/v1/inference-tasks`, {
|
||||
llmName: PROVIDER_NAME,
|
||||
request: { messages: [{ role: 'user', content: 'hello queue' }] },
|
||||
streaming: false,
|
||||
});
|
||||
expect(enq.status, enq.body).toBe(201);
|
||||
const created = JSON.parse(enq.body) as { id: string; status: string; poolName: string };
|
||||
expect(created.status).toBe('pending');
|
||||
expect(created.id).toMatch(/^c[a-z0-9]+/);
|
||||
|
||||
// Step 3: bring a worker online — same providerSessionId so it
|
||||
// adopts the row, drainPendingForSession picks up the queued task.
|
||||
registrar = new VirtualLlmRegistrar({
|
||||
mcpdUrl: MCPD_URL,
|
||||
token,
|
||||
publishedProviders: [{
|
||||
provider: makeFakeProvider(PROVIDER_NAME, 'reply from worker'),
|
||||
type: 'openai',
|
||||
model: 'fake-task',
|
||||
}],
|
||||
sessionFilePath: join(tempDir, 'session-1'),
|
||||
log: { info: () => {}, warn: () => {}, error: () => {} },
|
||||
heartbeatIntervalMs: 60_000,
|
||||
});
|
||||
await registrar.start();
|
||||
|
||||
// Step 4: poll for completion. The provider's complete() returns
|
||||
// synchronously so the worker should resolve very quickly — but
|
||||
// the SSE round-trip takes at least one tick, so we poll up to 5s.
|
||||
let final: { status: string; responseBody?: { choices?: Array<{ message?: { content?: string } }> } } | null = null;
|
||||
for (let i = 0; i < 25; i += 1) {
|
||||
await new Promise((r) => setTimeout(r, 200));
|
||||
const res = await httpRequest('GET', `${MCPD_URL}/api/v1/inference-tasks/${created.id}`, undefined);
|
||||
if (res.status !== 200) continue;
|
||||
const row = JSON.parse(res.body) as typeof final;
|
||||
if (row !== null && (row.status === 'completed' || row.status === 'error' || row.status === 'cancelled')) {
|
||||
final = row;
|
||||
break;
|
||||
}
|
||||
}
|
||||
expect(final?.status, 'task should have reached completed').toBe('completed');
|
||||
expect(final?.responseBody?.choices?.[0]?.message?.content).toBe('reply from worker');
|
||||
}, 30_000);
|
||||
|
||||
it('cancel: DELETE on a pending task transitions it to cancelled', async () => {
|
||||
if (!mcpdUp) return;
|
||||
const token = readToken();
|
||||
if (token === null) return;
|
||||
|
||||
// Stop the worker (if any) so this enqueue stays pending. A short
|
||||
// sleep gives unbindSession a moment.
|
||||
if (registrar !== null) {
|
||||
registrar.stop();
|
||||
registrar = null;
|
||||
}
|
||||
await new Promise((r) => setTimeout(r, 400));
|
||||
|
||||
const enq = await httpRequest('POST', `${MCPD_URL}/api/v1/inference-tasks`, {
|
||||
llmName: PROVIDER_NAME,
|
||||
request: { messages: [{ role: 'user', content: 'will be cancelled' }] },
|
||||
streaming: false,
|
||||
});
|
||||
expect(enq.status, enq.body).toBe(201);
|
||||
const created = JSON.parse(enq.body) as { id: string; status: string };
|
||||
expect(created.status).toBe('pending');
|
||||
|
||||
const cancelRes = await httpRequest('DELETE', `${MCPD_URL}/api/v1/inference-tasks/${created.id}`, undefined);
|
||||
expect(cancelRes.status).toBe(200);
|
||||
const cancelled = JSON.parse(cancelRes.body) as { status: string };
|
||||
expect(cancelled.status).toBe('cancelled');
|
||||
|
||||
// Re-fetch confirms persistence.
|
||||
const after = await httpRequest('GET', `${MCPD_URL}/api/v1/inference-tasks/${created.id}`, undefined);
|
||||
expect(JSON.parse(after.body).status).toBe('cancelled');
|
||||
}, 20_000);
|
||||
});
|
||||
Reference in New Issue
Block a user