23 KiB
Qwen Connectivity and Resume Recovery Implementation Plan
For agentic workers: REQUIRED SUB-SKILL: Use superpowers:subagent-driven-development (recommended) or superpowers:executing-plans to implement this plan task-by-task. Steps use checkbox (
- [ ]) syntax for tracking.
Goal: Make local Qwen reachable from production ThothII and make Resume restart only idle or failed Pi runtimes while preserving running turns and pending reviewer gates.
Architecture: thothii-core joins vLLM's existing private Docker network and Pi addresses Qwen through Docker DNS. SessionBridge owns an explicit four-state turn lifecycle; the resume route uses that state to preserve or replace a runtime, while the frontend uses a stream generation key to reconnect EventSource for a same-ID resume.
Tech Stack: Docker Compose, Fastify, TypeScript, Pi RPC, React 18, Zustand, Vitest, Testing Library/MSW.
Global Constraints
- Existing failed Qwen sessions are not migrated, recovered, replayed, or deleted.
- vLLM remains private; do not publish it on
0.0.0.0or proxy it through the public Omics network. - Running Pi turns and pending reviewer gates must never be replaced by Resume.
- Provider failures shown in the browser must be sanitized and contain no stack, credential, request body, or secret endpoint data.
- Persisted provider, model, and thinking settings remain authoritative on Resume.
- UI strings remain English.
- Use
apply_patchfor repository and mounted-profile edits.
Task 1: Pi turn lifecycle and sanitized provider errors
Files:
- Modify:
backend/src/bridge/session-bridge.ts - Modify:
backend/src/pi/pi-process-manager.ts - Test:
backend/test/session-bridge.test.ts - Test:
backend/test/pi-process-manager.test.ts
Interfaces:
-
Produces:
TurnState = "idle" | "running" | "waiting" | "failed". -
Produces:
SessionBridge.turnState(): TurnStateandSessionBridge.beginTurn(): void. -
Consumes: real Pi
message_end.message.role,message_end.message.stopReason,agent_start,agent_end, andextension_ui_requestevents. -
Step 1: Add failing bridge lifecycle tests
Add tests that exercise the real Pi error shape and every state transition:
test("assistant provider errors are sanitized and leave the turn failed", () => {
const { rpc, fire } = fakeRpc();
const bridge = new SessionBridge(rpc);
const seen: any[] = [];
bridge.onClientEvent((event) => seen.push(event));
bridge.beginTurn();
fire({
type: "message_end",
message: {
role: "assistant",
stopReason: "error",
errorMessage: "Connection failed for https://secret.invalid/?api_key=DO_NOT_LEAK",
},
});
fire({ type: "agent_end", messages: [] });
expect(bridge.turnState()).toBe("failed");
expect(seen).toContainEqual({
type: "info",
level: "error",
text: "Model request failed. Check provider connectivity, then Resume the session.",
});
expect(JSON.stringify(seen)).not.toContain("DO_NOT_LEAK");
});
test("reviewer wait and response transition waiting back to running", () => {
const { rpc, fire } = fakeRpc();
const bridge = new SessionBridge(rpc);
bridge.beginTurn();
fire({
type: "extension_ui_request",
id: "pi-1",
method: "input",
title: JSON.stringify({ id: "gate-1", widget: "select" }),
});
expect(bridge.turnState()).toBe("waiting");
bridge.respond({ id: "gate-1", choices: ["approve"] });
expect(bridge.turnState()).toBe("running");
fire({ type: "agent_end", messages: [] });
expect(bridge.turnState()).toBe("idle");
});
- Step 2: Run the focused tests and verify RED
Run: cd backend && npx vitest run test/session-bridge.test.ts
Expected: FAIL because beginTurn() and turnState() do not exist and provider errors are ignored.
- Step 3: Implement the minimal lifecycle in
SessionBridge
Add the state type and methods, and update only lifecycle-bearing branches:
export type TurnState = "idle" | "running" | "waiting" | "failed";
export class SessionBridge {
private state: TurnState = "idle";
turnState(): TurnState { return this.state; }
beginTurn(): void { this.state = "running"; }
}
Within the existing RPC event handler:
if (m.type === "extension_ui_request" && m.method === "input") {
// retain descriptor parsing and pending correlation
this.state = "waiting";
} else if (
m.type === "message_end" &&
m.message?.role === "assistant" &&
m.message?.stopReason === "error"
) {
this.state = "failed";
this.fan({
type: "info",
level: "error",
text: "Model request failed. Check provider connectivity, then Resume the session.",
});
} else if (m.type === "agent_start") {
this.state = "running";
this.fan({ type: "system_event", event: "agent_start" });
} else if (m.type === "agent_end") {
if (this.state !== "failed" && !this.pending) this.state = "idle";
this.fan({ type: "system_event", event: "agent_end" });
}
In respond() and steer(), set this.state = "running" immediately before sending the RPC command. In PiProcessManager.createFor() call bridge.beginTurn() before registering the runtime so the configure/bootstrap window is treated as active; in start() call it again before sending the prompt.
- Step 4: Add manager coverage for the initial running state
test("a created runtime is active during configure/bootstrap", () => {
const child = recordingChild();
const mgr = new PiProcessManager(loadConfig({}), { spawnFn: () => child as any });
const runtime = mgr.createFor("starting-id", {});
expect(runtime.bridge.turnState()).toBe("running");
mgr.teardown("starting-id");
});
- Step 5: Run focused and full backend tests
Run:
cd backend
npx vitest run test/session-bridge.test.ts test/pi-process-manager.test.ts
npx vitest run
npx tsc --noEmit -p .
Expected: all tests pass and TypeScript exits 0.
- Step 6: Commit the lifecycle change
git add backend/src/bridge/session-bridge.ts backend/src/pi/pi-process-manager.ts \
backend/test/session-bridge.test.ts backend/test/pi-process-manager.test.ts
git commit -m "fix(backend): track Pi turn lifecycle"
Task 2: State-aware Resume without stale SSE replay
Files:
- Modify:
backend/src/routes/sessions.ts - Test:
backend/test/routes-sessions.test.ts
Interfaces:
-
Consumes:
SessionRuntime.bridge.turnState(): TurnStatefrom Task 1. -
Consumes:
PiProcessManager.teardown(id)andSseHub.clear(id). -
Produces: Resume preserves
running/waiting; Resume replacesidle/failedand executes the existing persisted-settings cold path. -
Step 1: Add failing route tests for preserve and restart decisions
test.each(["running", "waiting"])(
"POST resume preserves a %s runtime",
async (state) => {
let tornDown = false;
const existing = { bridge: { turnState: () => state } } as any;
const app = buildApp(loadConfig({ THT_HARNESS_DIR: "../harness" }), {
mgr: {
get: () => existing,
teardown: () => { tornDown = true; },
} as any,
thtRunner: {} as any,
getSettings: () => ({ workspace: "psd" }) as any,
});
const response = await app.inject({ method: "POST", url: "/sessions/s1/resume" });
expect(response.json()).toEqual({ id: "s1", alreadyActive: true });
expect(tornDown).toBe(false);
},
);
test.each(["idle", "failed"])(
"POST resume replaces a %s runtime and starts the persisted session",
async (state) => {
const order: string[] = [];
const oldRuntime = { bridge: { turnState: () => state } } as any;
const newRuntime = { bridge: { onClientEvent: () => {}, emitClientEvent: () => {} } } as any;
const app = buildApp(loadConfig({ THT_HARNESS_DIR: "../harness" }), {
mgr: {
get: () => oldRuntime,
teardown: (id: string) => order.push(`teardown:${id}`),
createFor: () => { order.push("create"); return newRuntime; },
configure: async () => {},
start: () => order.push("start"),
} as any,
thtRunner: {
sessionShow: async () => ({
status: "open", archived: false,
provider: "local-qwen", model: "qwen3.6-35b-a3b", thinking: "low",
}),
reopenSession: async () => order.push("reopen"),
} as any,
readiness: { ensure: async () => ({ ok: true }) } as any,
getSettings: () => ({ workspace: "local", thinking: "medium" }) as any,
});
const response = await app.inject({ method: "POST", url: "/sessions/s1/resume" });
await new Promise((resolve) => setImmediate(resolve));
expect(response.json()).toEqual({ id: "s1" });
expect(order).toEqual(["teardown:s1", "reopen", "create", "start"]);
},
);
- Step 2: Run route tests and verify RED
Run: cd backend && npx vitest run test/routes-sessions.test.ts -t "resume"
Expected: idle/failed cases return alreadyActive instead of tearing down and starting.
- Step 3: Implement state-aware handling at the top of the Resume route
Replace the unconditional existing-runtime return with:
const existing = d.mgr.get(id);
if (existing) {
const state = existing.bridge.turnState();
if (state === "running" || state === "waiting") {
return reply.code(200).send({ id, alreadyActive: true });
}
d.mgr.teardown(id);
d.hub.clear(id);
}
Leave finalized/archived checks, readiness, persisted provider/model/thinking selection, reopenSession, runtime binding, and /riprendi-sessione bootstrap unchanged.
- Step 4: Verify focused and full backend gates
Run:
cd backend
npx vitest run test/routes-sessions.test.ts -t "resume"
npx vitest run
npx tsc --noEmit -p .
Expected: all tests pass and TypeScript exits 0.
- Step 5: Commit the Resume change
git add backend/src/routes/sessions.ts backend/test/routes-sessions.test.ts
git commit -m "fix(backend): restart idle Pi sessions on resume"
Task 3: Reconnect SSE when resuming the same session
Files:
- Modify:
frontend/src/stream/useSessionStream.ts - Modify:
frontend/src/shell/AppShell.tsx - Test:
frontend/src/stream/useSessionStream.test.tsx - Test:
frontend/src/shell/AppShell.session-mgmt.test.tsx
Interfaces:
-
Produces:
useSessionStream(sessionId: string | null, generation?: number). -
Consumes: successful
resumeSession(id)completion inAppShell.doResume. -
Step 1: Add a failing hook test for a generation-driven reconnect
test("reconnects the same session when the generation changes", () => {
const { rerender } = renderHook(
({ generation }) => useSessionStream("s1", generation),
{ initialProps: { generation: 0 } },
);
const first = FakeEventSource.instances[0];
rerender({ generation: 1 });
expect(first.closed).toBe(true);
expect(FakeEventSource.instances).toHaveLength(2);
expect(FakeEventSource.instances[1].url).toBe(first.url);
});
- Step 2: Run the hook test and verify RED
Run: cd frontend && npx vitest run src/stream/useSessionStream.test.tsx
Expected: FAIL because changing the unused generation does not recreate EventSource.
- Step 3: Add the generation dependency
Change the hook signature and effect dependency:
export function useSessionStream(sessionId: string | null, generation = 0) {
// existing implementation
useEffect(() => {
// existing EventSource setup and cleanup
}, [sessionId, generation, applyEvent]);
}
- Step 4: Add an AppShell integration test for a second same-ID Resume
test("resuming the active session reconnects its EventSource", async () => {
server.use(
http.post("http://localhost:8787/sessions/:id/resume", () =>
new HttpResponse(null, { status: 204 })),
http.get("http://localhost:8787/sessions/:id", () =>
HttpResponse.json({ id: "s1", status: "open", phase: 1 })),
);
wrap();
await userEvent.click(await screen.findByText("Attiva uno"));
await userEvent.click(await screen.findByRole("button", { name: /resume/i }));
await waitFor(() => expect(FakeEventSource.instances).toHaveLength(1));
const first = FakeEventSource.instances[0];
await userEvent.click(screen.getByText("Attiva uno"));
await userEvent.click(await screen.findByRole("button", { name: /resume/i }));
await waitFor(() => expect(FakeEventSource.instances).toHaveLength(2));
expect(first.closed).toBe(true);
});
- Step 5: Wire successful same-ID resume to the generation
In AppShell add:
const [streamGeneration, setStreamGeneration] = useState(0);
Capture whether the resumed session was already active before the optimistic switch, then increment only after a successful POST:
async function doResume(id: string) {
const reconnectSameSession = activeSessionId === id;
// existing optimistic state updates
try {
await resumeSession(id);
if (reconnectSameSession) setStreamGeneration((value) => value + 1);
// existing phase refresh
} catch {
// existing rollback
}
}
Pass the generation to the hook:
useSessionStream(activeSessionId, streamGeneration);
- Step 6: Verify focused and full frontend gates
Run:
cd frontend
npx vitest run src/stream/useSessionStream.test.tsx src/shell/AppShell.session-mgmt.test.tsx
npx vitest run
npx tsc -b
npm run build
Expected: all tests pass, TypeScript exits 0, and Vite builds successfully.
- Step 7: Commit the frontend reconnect
git add frontend/src/stream/useSessionStream.ts frontend/src/shell/AppShell.tsx \
frontend/src/stream/useSessionStream.test.tsx frontend/src/shell/AppShell.session-mgmt.test.tsx
git commit -m "fix(frontend): reconnect stream on same-session resume"
Task 4: Private Docker network for Qwen
Files:
- Modify:
compose.yaml - Create:
scripts/test-qwen-network-config.sh - Modify outside Git:
/home/chirone/thothii-data/pi-config/agent/models.json
Interfaces:
-
Consumes: external Docker network
localllm_defaultand DNS namelocalllm-vllm:8000. -
Produces:
coremembership in bothomics_portal_omics_networkandlocalllm_default. -
Produces: effective Pi base URL
http://localllm-vllm:8000/v1. -
Step 1: Add a failing Compose contract test
Create scripts/test-qwen-network-config.sh:
#!/bin/sh
set -eu
cd "$(dirname "$0")/.."
tmp=$(mktemp -d)
trap 'rm -rf "$tmp"' EXIT HUP INT TERM
mkdir -p "$tmp/deploy"
cp compose.yaml "$tmp/compose.yaml"
: >"$tmp/deploy/thothii.env"
docker compose --project-directory "$tmp" config --format json >"$tmp/config.json"
node - "$tmp/config.json" <<'NODE'
const fs = require("fs");
const config = JSON.parse(fs.readFileSync(process.argv[2], "utf8"));
if (!config.networks?.localllm_default?.external) {
throw new Error("localllm_default must be an external network");
}
if (!config.services?.core?.networks?.localllm_default) {
throw new Error("core must join localllm_default");
}
if (config.services?.frontend?.networks?.localllm_default) {
throw new Error("frontend must not join the model network");
}
NODE
Make it executable with chmod +x scripts/test-qwen-network-config.sh.
- Step 2: Run the contract test and verify RED
Run: ./scripts/test-qwen-network-config.sh
Expected: FAIL with localllm_default must be an external network.
- Step 3: Attach only core to the private model network
Change the core network mapping and declare the external network:
services:
core:
networks:
omics_portal_omics_network:
aliases: ["thothii-core"]
localllm_default: {}
networks:
omics_portal_omics_network:
external: true
localllm_default:
external: true
- Step 4: Verify Compose and commit the repository change
Run:
./scripts/test-qwen-network-config.sh
docker compose config --quiet
git diff --check
Expected: all commands exit 0.
Commit:
git add compose.yaml scripts/test-qwen-network-config.sh
git commit -m "fix(deploy): connect core to local Qwen"
- Step 5: Update the mounted Pi profile without exposing credentials
Record only metadata before editing:
stat -c '%U:%G %a %n' /home/chirone/thothii-data/pi-config/agent/models.json
Use apply_patch to replace only:
"baseUrl": "http://127.0.0.1:18000/v1"
with:
"baseUrl": "http://localllm-vllm:8000/v1"
Re-run the same stat command and verify owner/group/mode are unchanged. Parse only the non-secret provider summary:
node - <<'NODE'
const fs = require("fs");
const file = "/home/chirone/thothii-data/pi-config/agent/models.json";
const config = JSON.parse(fs.readFileSync(file, "utf8"));
const qwen = config.providers["local-qwen"];
console.log(JSON.stringify({ baseUrl: qwen.baseUrl, models: qwen.models.map((m) => m.id) }));
NODE
Expected: base URL is http://localllm-vllm:8000/v1 and the only local model ID remains qwen3.6-35b-a3b.
Task 5: Full verification, deployment, and live smoke
Files:
- Modify:
PROJECT_STATE.md - Modify:
brain/codebase/pi-model-selection.md - Modify:
brain/codebase/workflow-ui-contracts.md
Interfaces:
-
Consumes: Tasks 1–4 and live Docker services.
-
Produces: rebuilt/recreated
thothii-core-1andthothii-frontend-1, verified Qwen inference, and updated operational state. -
Step 1: Run every repository gate from a clean worktree
Run:
cd backend && npx vitest run && npx tsc --noEmit -p . && npm run build
cd ../frontend && npx vitest run && npx tsc -b && npm run build
cd .. && ./scripts/test-qwen-network-config.sh && git diff --check
Expected: every command exits 0 with zero failing tests.
- Step 2: Build and force-recreate the affected services
Before recreation, list active Pi session IDs without printing any other environment values. Then run:
docker compose build core frontend
docker compose up -d --force-recreate core frontend
docker compose ps core frontend
Expected: both services are running; core becomes healthy. Restart Omics Portal web workers only if their cached Vite manifest still references the previous frontend assets.
- Step 3: Verify private Qwen reachability and real inference from core
Run a sanitized Node probe inside core:
docker exec -i thothii-core-1 node - <<'NODE'
(async () => {
const models = await fetch("http://localllm-vllm:8000/v1/models");
if (!models.ok) throw new Error(`models status ${models.status}`);
const catalog = await models.json();
if (!catalog.data?.some((model) => model.id === "qwen3.6-35b-a3b")) {
throw new Error("Qwen model missing");
}
const completion = await fetch("http://localllm-vllm:8000/v1/chat/completions", {
method: "POST",
headers: { "content-type": "application/json" },
body: JSON.stringify({
model: "qwen3.6-35b-a3b",
messages: [{ role: "user", content: "Reply with exactly QWEN_OK" }],
temperature: 0,
max_tokens: 16,
}),
});
if (!completion.ok) throw new Error(`completion status ${completion.status}`);
const body = await completion.json();
const text = body.choices?.[0]?.message?.content?.trim();
console.log(JSON.stringify({ modelReachable: true, responseReceived: Boolean(text) }));
if (!text) throw new Error("empty Qwen response");
})();
NODE
Expected: {"modelReachable":true,"responseReceived":true} without logging prompts, credentials, or raw provider responses.
- Step 4: Run an application-level Qwen session smoke
Run this complete one-off script from the host. It saves /settings, selects Qwen, creates a uniquely named smoke question, waits up to 180 seconds for the first ui_request, validates the persisted provider/model, deletes only the smoke session, and restores the exact original settings object:
docker exec -i thothii-core-1 node - <<'NODE'
const base = "http://127.0.0.1:8787";
let smokeId = null;
let original = null;
let gateSeen = false;
let restored = false;
async function json(path, init) {
const response = await fetch(base + path, init);
if (!response.ok) throw new Error(`${path} returned ${response.status}`);
return response.status === 204 ? null : response.json();
}
async function waitForGate(id) {
const controller = new AbortController();
const timer = setTimeout(() => controller.abort(), 180_000);
try {
const response = await fetch(`${base}/sessions/${id}/events`, { signal: controller.signal });
if (!response.ok || !response.body) throw new Error(`SSE returned ${response.status}`);
const reader = response.body.getReader();
const decoder = new TextDecoder();
let buffer = "";
while (true) {
const { value, done } = await reader.read();
if (done) throw new Error("SSE ended before ui_request");
buffer += decoder.decode(value, { stream: true });
let boundary;
while ((boundary = buffer.indexOf("\n\n")) >= 0) {
const frame = buffer.slice(0, boundary);
buffer = buffer.slice(boundary + 2);
const event = frame.split("\n").find((line) => line.startsWith("event: "))?.slice(7);
if (event === "ui_request") return;
}
}
} finally {
clearTimeout(timer);
controller.abort();
}
}
(async () => {
try {
original = await json("/settings");
await json("/settings", {
method: "PUT",
headers: { "content-type": "application/json" },
body: JSON.stringify({
...original,
provider: "local-qwen",
model: "qwen3.6-35b-a3b",
thinking: "low",
}),
});
const created = await json("/sessions", {
method: "POST",
headers: { "content-type": "application/json" },
body: JSON.stringify({ question: `Qwen smoke ${new Date().toISOString()}: count patients by sex` }),
});
smokeId = created.id;
await waitForGate(smokeId);
const manifest = await json(`/sessions/${smokeId}`);
if (manifest.provider !== "local-qwen" || manifest.model !== "qwen3.6-35b-a3b") {
throw new Error("session did not persist Qwen");
}
gateSeen = true;
} finally {
if (smokeId) {
await fetch(`${base}/sessions/${smokeId}/close`, { method: "POST" }).catch(() => undefined);
await fetch(`${base}/sessions/${smokeId}`, { method: "DELETE" }).catch(() => undefined);
}
if (original) {
await json("/settings", {
method: "PUT",
headers: { "content-type": "application/json" },
body: JSON.stringify(original),
});
const current = await json("/settings");
restored = JSON.stringify(current) === JSON.stringify(original);
}
}
if (!gateSeen || !restored) throw new Error("Qwen smoke did not complete cleanly");
console.log("qwen_session_gate=ok");
console.log("settings_restored=ok");
})().catch((error) => {
console.error(`qwen_session_smoke=failed:${error.message}`);
process.exitCode = 1;
});
NODE
Expected output:
qwen_session_gate=ok
settings_restored=ok
It must fail nonzero if no gate arrives within 180 seconds or restoration fails.
- Step 5: Record verified deployment state
Update PROJECT_STATE.md with test counts, deployment timestamp, running image IDs, network membership, mounted Qwen base URL, and smoke result. Replace the diagnostic caveats in the two existing brain notes with the verified final behavior; do not add duplicate notes.
- Step 6: Commit documentation and run final audit
git add PROJECT_STATE.md brain/codebase/pi-model-selection.md brain/codebase/workflow-ui-contracts.md
git commit -m "docs: record Qwen connectivity deployment"
git status --short
git log -6 --oneline
docker compose ps core frontend
Expected: only the user's pre-existing .vite/ artifact remains untracked, recent commits match Tasks 1–5, and both services are healthy/running.