docs: plan Qwen connectivity and resume recovery

This commit is contained in:
User
2026-07-14 20:33:57 +02:00
parent ad4d6963a2
commit 37cc42bebf
@@ -0,0 +1,716 @@
# 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.0` or 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_patch` for 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(): TurnState` and `SessionBridge.beginTurn(): void`.
- Consumes: real Pi `message_end.message.role`, `message_end.message.stopReason`, `agent_start`, `agent_end`, and `extension_ui_request` events.
- [ ] **Step 1: Add failing bridge lifecycle tests**
Add tests that exercise the real Pi error shape and every state transition:
```ts
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:
```ts
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:
```ts
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**
```ts
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:
```bash
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**
```bash
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(): TurnState` from Task 1.
- Consumes: `PiProcessManager.teardown(id)` and `SseHub.clear(id)`.
- Produces: Resume preserves `running`/`waiting`; Resume replaces `idle`/`failed` and executes the existing persisted-settings cold path.
- [ ] **Step 1: Add failing route tests for preserve and restart decisions**
```ts
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:
```ts
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:
```bash
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**
```bash
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 in `AppShell.doResume`.
- [ ] **Step 1: Add a failing hook test for a generation-driven reconnect**
```tsx
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:
```ts
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**
```tsx
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:
```ts
const [streamGeneration, setStreamGeneration] = useState(0);
```
Capture whether the resumed session was already active before the optimistic switch, then increment only after a successful POST:
```ts
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:
```ts
useSessionStream(activeSessionId, streamGeneration);
```
- [ ] **Step 6: Verify focused and full frontend gates**
Run:
```bash
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**
```bash
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_default` and DNS name `localllm-vllm:8000`.
- Produces: `core` membership in both `omics_portal_omics_network` and `localllm_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`:
```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:
```yaml
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:
```bash
./scripts/test-qwen-network-config.sh
docker compose config --quiet
git diff --check
```
Expected: all commands exit 0.
Commit:
```bash
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:
```bash
stat -c '%U:%G %a %n' /home/chirone/thothii-data/pi-config/agent/models.json
```
Use `apply_patch` to replace only:
```json
"baseUrl": "http://127.0.0.1:18000/v1"
```
with:
```json
"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:
```bash
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-1` and `thothii-frontend-1`, verified Qwen inference, and updated operational state.
- [ ] **Step 1: Run every repository gate from a clean worktree**
Run:
```bash
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:
```bash
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`:
```bash
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:
```bash
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);
}
}
(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:
```text
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**
```bash
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.