Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 5 additions & 1 deletion demos/02-lumacare/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -23,8 +23,10 @@ The checked-in `config/agents.json` maps the `care_companion` feature to its pur
## What is included

- Calm, responsive family-facing care guidance
- A stable guide viewport across waiting, live streaming, and completed states
- Structured plain-language guides and care-team questions
- Data-only streaming with immediate recovery-ID persistence
- Data-only streaming with isolated per-run buffers and immediate recovery-ID persistence
- Explicit per-run cancellation that remains safe during fast stop/retry and caregiver switching
- Session recall and continue/revise flows
- Ownership checks for session details, cancellation, previews, and downloads
- Sandboxed HTML previews and authenticated artifact/archive proxying
Expand All @@ -41,3 +43,5 @@ npm run test --workspace=@harnessrouter/lumacare
npm run build --workspace=@harnessrouter/lumacare
npm run preview:lumacare
```

The LumaCare tests include stream parsing, cross-user ownership, family-facing guide formatting, checklist-state preservation, and responsive guide-layout contracts.
55 changes: 51 additions & 4 deletions demos/02-lumacare/server/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,9 @@ dotenv.config({ path: path.resolve(root, '../..', '.env') });
dotenv.config({ path: path.join(root, '.env'), override: true });
const agentMap = JSON.parse(fs.readFileSync(path.join(root, 'config/agents.json'), 'utf8')) as Record<string, string>;
const store = new OwnershipStore(path.join(root, 'data/ownership.json'));
type ActiveRun = { userId: string; controller: AbortController; reader: ReadableStreamDefaultReader<Uint8Array> | null; sessionId: string | null };
const activeRuns = new Map<string, ActiveRun>();
const cancelledRuns = new Map<string, string>();
const apiBase = 'https://api.harnessrouter.ai';
const apiKey = process.env.HR_API_KEY;

Expand Down Expand Up @@ -78,17 +81,30 @@ app.post('/api/runs', async (req, res) => {
const input = String(req.body.input || '').trim();
const previousResponseId = req.body.previousResponseId ? String(req.body.previousResponseId) : null;
const sessionId = req.body.sessionId ? String(req.body.sessionId) : null;
const clientRunId = String(req.body.clientRunId || '');
const harnessId = agentMap[featureKey];

if (!harnessId) return res.status(400).json({ error: 'Unknown feature' });
if (!input || input.length > 20_000) return res.status(400).json({ error: 'Enter a task under 20,000 characters.' });
if (!/^[0-9a-f-]{36}$/i.test(clientRunId)) return res.status(400).json({ error: 'A valid client run ID is required.' });
if (cancelledRuns.get(clientRunId) === req.productUser.id) {
cancelledRuns.delete(clientRunId);
return res.status(409).json({ error: 'This run was cancelled before it started.' });
}
if (activeRuns.has(clientRunId)) return res.status(409).json({ error: 'This run is already active.' });
if ((previousResponseId || sessionId) && !(previousResponseId && sessionId)) {
return res.status(400).json({ error: 'A continuation requires both recovery identifiers.' });
}
if (sessionId && !store.getAuthorized(sessionId, req.productUser.id)) {
return res.status(404).json({ error: 'Session not found' });
}

const upstreamController = new AbortController();
const activeRun: ActiveRun = { userId: req.productUser.id, controller: upstreamController, reader: null, sessionId };
activeRuns.set(clientRunId, activeRun);
let recoveredSessionId = sessionId;
let recoveredResponseId = previousResponseId;

const requestBody: Record<string, unknown> = { input, stream: true };
if (previousResponseId && sessionId) {
requestBody.previous_response_id = previousResponseId;
Expand All @@ -104,13 +120,18 @@ app.post('/api/runs', async (req, res) => {
'Idempotency-Key': crypto.randomUUID(),
}),
body: JSON.stringify(requestBody),
signal: upstreamController.signal,
});
} catch {
activeRuns.delete(clientRunId);
if (upstreamController.signal.aborted) return res.end();
if (res.closed) return;
return res.status(502).json({ error: 'HarnessRouter could not be reached.' });
}

if (!upstream.ok || !upstream.body) {
const detail = await upstream.json().catch(() => ({ detail: 'Agent request failed.' })) as { detail?: string };
activeRuns.delete(clientRunId);
return res.status(upstream.status || 502).json({ error: detail.detail || 'Agent request failed.' });
}

Expand All @@ -120,11 +141,10 @@ app.post('/api/runs', async (req, res) => {
res.setHeader('Connection', 'keep-alive');
res.flushHeaders();

const reader = upstream.body.getReader();
const streamReader = upstream.body.getReader();
activeRun.reader = streamReader;
const decoder = new TextDecoder();
let buffer = '';
let recoveredSessionId = sessionId;
let recoveredResponseId = previousResponseId;

const handleFrame = (frame: string) => {
const dataLines = frame.split(/\r?\n/).filter((line) => line.startsWith('data:'));
Expand All @@ -136,6 +156,7 @@ app.post('/api/runs', async (req, res) => {
if (event.type === 'response.created') {
recoveredResponseId = event.response?.id || recoveredResponseId;
recoveredSessionId = event.response?.metadata?.session_id || recoveredSessionId;
activeRun.sessionId = recoveredSessionId;
if (recoveredResponseId && recoveredSessionId) {
const now = new Date().toISOString();
const existing = store.getAuthorized(recoveredSessionId, req.productUser.id);
Expand Down Expand Up @@ -170,7 +191,7 @@ app.post('/api/runs', async (req, res) => {

try {
while (true) {
const { done, value } = await reader.read();
const { done, value } = await streamReader.read();
if (done) break;
const text = decoder.decode(value, { stream: true });
if (!res.closed) res.write(text);
Expand All @@ -191,6 +212,32 @@ app.post('/api/runs', async (req, res) => {
}
} finally {
if (!res.closed) res.end();
if (activeRuns.get(clientRunId) === activeRun) activeRuns.delete(clientRunId);
}
});

app.post('/api/runs/:clientRunId/cancel', async (req, res) => {
const clientRunId = String(req.params.clientRunId);
const run = activeRuns.get(clientRunId);
if (!run) {
cancelledRuns.set(clientRunId, req.productUser.id);
setTimeout(() => { if (cancelledRuns.get(clientRunId) === req.productUser.id) cancelledRuns.delete(clientRunId); }, 30_000).unref();
return res.json({ cancelled: true });
}
if (run.userId !== req.productUser.id) return res.status(404).json({ error: 'Active run not found' });
run.controller.abort();
await run.reader?.cancel().catch(() => {});
try {
if (run.sessionId) {
const result = await upstreamJson(`/v1/sessions/${encodeURIComponent(run.sessionId)}/cancel`, { method: 'POST' });
store.update(run.sessionId, { status: 'cancelled' });
return res.json(result);
}
return res.json({ cancelled: true });
} catch (error) {
errorResponse(res, error);
} finally {
activeRuns.delete(clientRunId);
}
});

Expand Down
Loading
Loading