feat: agent and CLI control plane for sessions, worktrees, and scheduled tasks (#2408)
Add a shared OpenChamber control service with two thin adapters — a native `openchamber` tool injected into managed OpenCode, and new CLI commands — so users can manage parallel sessions, worktrees, and scheduled tasks conversationally through agents or from the terminal. Control plane: - New openchamber-control service owning a fixed action contract: projects.list, models.list, session list/create/send/fork/status/messages, and schedule list/create/run/delete/toggle. Session and worktree deletion and project registration are deliberately not exposed. - New openchamber-sessions module owning create/worktree/prompt orchestration, Goal Mode dispatch, wait semantics (initial idle never counts as completion; timeout and cancellation are failures), and explicit partial-failure results. - Scheduled-task logic extracted into a service shared by routes, CLI, and the agent tool. Agent tool: - Managed OpenCode gets a materialized plugin registering one typed tool with a loopback-only callback, per-child ephemeral bearer (timing-safe, never persisted or logged), and abort propagation into the service. - The ~1.5k-token schema applies progressive disclosure: short descriptions, server-side validation returning actionable usage errors, and intent guardrails — created sessions/tasks are user-facing work (not age self-delegation); worktree/goal/agent/variant/wait are omit-by-default; dispatches produce no completion notification, and later result r to session.messages, which now returns the authoritative sessionStatus. - session.create without a user-named model picks from favorites/re send/fork omit the selection and the service reuses the target session's last user-message model, agent, and variant before falling back t - An "Agent control tool" setting (default on, Save + Reload to apply) disables plugin injection entirely. CLI: - New `openchamber session`, `schedule`, `projects`, and `models` commands with automatic instance targeting, --wait/--timeout/--last-assist worktree flags, and Goal Mode, preserving interactive, non-TTY, --quiet, and --json contracts. The control HTTP timeout derives from the w instead of the 4-second default. UI: - New built-in "Schedule a Task" starter (/schedule-task) running a dialogue that defines a task and offers to create it via the tool after explicit confirmation; Craft a Goal and Feature Planning gain the handoff offer, and guided starters reserve the question tool for concrete option choices. Localized in all 10 locales, migrated into custom starter lists, hidden on VS Code. - Sidebar shows CLI/agent-created sessions live via the control eve - openchamber tool calls render with per-action titles and metadata.
This commit is contained in:
committed by
GitHub
parent
484fe8bc18
commit
e908db637b
@@ -0,0 +1,47 @@
|
||||
# OpenChamber Control Service
|
||||
|
||||
## Purpose
|
||||
|
||||
This module owns the typed control contract shared by the OpenChamber CLI and
|
||||
the managed OpenCode `openchamber` tool. Both adapters delegate to
|
||||
`createOpenChamberControlService()`; neither adapter may call or spawn the
|
||||
other.
|
||||
|
||||
## Boundaries
|
||||
|
||||
- `service.js` validates and executes the fixed project, model, session, and
|
||||
scheduled-task action allowlist. `actions.js` marks CLI-only actions with
|
||||
`agentExposed: false` (currently `schedule.status`); the agent tool consumes
|
||||
the filtered `OPENCHAMBER_AGENT_TOOL_*` exports. `schedule.toggle` requires
|
||||
the `disabled` boolean and replaces separate enable/disable actions;
|
||||
`schedule.list` also returns scheduler status as `scheduler`.
|
||||
- `routes.js` is the authenticated CLI HTTP adapter. It forwards one action,
|
||||
preserves service status and partial-result details, and propagates request
|
||||
cancellation.
|
||||
- `../agent-tool/runtime.js` is the managed-tool adapter. It wraps service
|
||||
results in the versioned native-tool envelope and uses a separate ephemeral
|
||||
loopback credential.
|
||||
- `../openchamber-sessions/routes.js` and `../scheduled-tasks/service.js` own
|
||||
their domain operations and are composed into this service.
|
||||
|
||||
## Invariants
|
||||
|
||||
- Session status and messages come from official directory-scoped OpenCode
|
||||
APIs. Message output includes only ordered `text` parts.
|
||||
- Wait never treats an initial idle response as completion after dispatch. It
|
||||
requires observed activity or a newly completed assistant message.
|
||||
- Timeout and cancellation are failures, never authoritative idle results.
|
||||
- Validation that protects side effects runs before session creation or
|
||||
dispatch.
|
||||
- Send and fork dispatches without an explicit model/agent/variant reuse the
|
||||
target session's last user-message selection before falling back to the
|
||||
configured defaults; only session creation resolves defaults directly.
|
||||
- Usage errors name the missing or conflicting input so CLI and agent-tool
|
||||
callers can correct an invalid request without an upfront usage manual.
|
||||
- Explicit `projectId` or `directory` scope takes precedence over the managed
|
||||
tool's current-session directory fallback; the fallback never creates a
|
||||
conflicting second scope.
|
||||
- One failed directory status lookup produces `unknown` for only that
|
||||
directory and does not erase other session results.
|
||||
- Destructive session/worktree deletion and project-path registration are not
|
||||
part of the action contract.
|
||||
@@ -0,0 +1,28 @@
|
||||
export const OPENCHAMBER_CONTROL_ACTION_DEFINITIONS = Object.freeze([
|
||||
{ action: 'projects.list', title: 'List configured projects', description: 'List configured projects; no parameters' },
|
||||
{ action: 'models.list', title: 'Show model preferences', description: 'Show default, favorite, and recent model preferences; no parameters' },
|
||||
{ action: 'session.list', title: 'List sessions', description: 'List sessions; optional directory, limit (default 10), all, or withStatus' },
|
||||
{ action: 'session.create', title: 'Create a session', description: 'Create a session in the current directory by default; prompt is optional' },
|
||||
{ action: 'session.send', title: 'Send a prompt', description: 'Send a new prompt to sessionId; scope with projectId or directory' },
|
||||
{ action: 'session.fork', title: 'Fork a session', description: 'Fork sessionId; messageId selects the boundary; prompt is optional' },
|
||||
{ action: 'session.status', title: 'Check session status', description: 'Check sessionId status; directory defaults to the current session' },
|
||||
{ action: 'session.messages', title: 'Read session messages', description: 'Read text-only messages and current sessionStatus for sessionId; directory and limit 10 are defaults' },
|
||||
{ action: 'schedule.status', title: 'Check scheduler status', description: 'Check scheduler status; no parameters', agentExposed: false },
|
||||
{ action: 'schedule.list', title: 'List scheduled tasks', description: 'List tasks and scheduler status; scope with projectId or directory' },
|
||||
{ action: 'schedule.create', title: 'Create a scheduled task', description: 'Create task; requires name, prompt, model, and one schedule selector' },
|
||||
{ action: 'schedule.run', title: 'Run a scheduled task', description: 'Run taskId; scope with projectId or directory' },
|
||||
{ action: 'schedule.delete', title: 'Delete a scheduled task', description: 'Delete taskId; scope with projectId or directory' },
|
||||
{ action: 'schedule.toggle', title: 'Enable or disable a scheduled task', description: 'Enable or disable taskId; requires the disabled boolean' },
|
||||
]);
|
||||
|
||||
export const OPENCHAMBER_CONTROL_ACTIONS = Object.freeze(
|
||||
OPENCHAMBER_CONTROL_ACTION_DEFINITIONS.map(({ action }) => action),
|
||||
);
|
||||
|
||||
export const OPENCHAMBER_AGENT_TOOL_ACTION_DEFINITIONS = Object.freeze(
|
||||
OPENCHAMBER_CONTROL_ACTION_DEFINITIONS.filter(({ agentExposed }) => agentExposed !== false),
|
||||
);
|
||||
|
||||
export const OPENCHAMBER_AGENT_TOOL_ACTIONS = Object.freeze(
|
||||
OPENCHAMBER_AGENT_TOOL_ACTION_DEFINITIONS.map(({ action }) => action),
|
||||
);
|
||||
@@ -0,0 +1,16 @@
|
||||
export class OpenChamberControlError extends Error {
|
||||
constructor(message, statusCode = 500, details = {}) {
|
||||
super(message);
|
||||
this.name = 'OpenChamberControlError';
|
||||
this.statusCode = statusCode;
|
||||
Object.assign(this, details);
|
||||
}
|
||||
}
|
||||
|
||||
export const asControlError = (error, fallbackMessage, fallbackStatus = 500) => {
|
||||
if (error instanceof OpenChamberControlError) return error;
|
||||
const message = error instanceof Error ? error.message : fallbackMessage;
|
||||
return new OpenChamberControlError(message || fallbackMessage, Number(error?.statusCode) || fallbackStatus, {
|
||||
...(error?.goalConfigured === true ? { goalConfigured: true } : {}),
|
||||
});
|
||||
};
|
||||
@@ -0,0 +1,36 @@
|
||||
import express from 'express';
|
||||
import { asControlError } from './error.js';
|
||||
|
||||
export const registerOpenChamberControlRoutes = (app, { controlService }) => {
|
||||
app.post('/api/openchamber/control', express.json({ limit: '1mb' }), async (req, res) => {
|
||||
const controller = new AbortController();
|
||||
const abortOnDisconnect = () => {
|
||||
if (!res.writableEnded) controller.abort();
|
||||
};
|
||||
req.once('aborted', abortOnDisconnect);
|
||||
res.once('close', abortOnDisconnect);
|
||||
try {
|
||||
const action = typeof req.body?.action === 'string' ? req.body.action : '';
|
||||
const requestInput = req.body?.input;
|
||||
const input = requestInput && typeof requestInput === 'object' && !Array.isArray(requestInput)
|
||||
? requestInput
|
||||
: {};
|
||||
const data = await controlService.execute(action, input, req.body?.contextDirectory, { signal: controller.signal });
|
||||
return res.json(data);
|
||||
} catch (error) {
|
||||
const controlError = asControlError(error, 'OpenChamber control action failed');
|
||||
return res.status(controlError.statusCode).json({
|
||||
error: controlError.message,
|
||||
...(controlError.partial === true ? {
|
||||
partial: true,
|
||||
partialAction: controlError.partialAction,
|
||||
sessionId: controlError.sessionId,
|
||||
directory: controlError.directory,
|
||||
} : {}),
|
||||
});
|
||||
} finally {
|
||||
req.off('aborted', abortOnDisconnect);
|
||||
res.off('close', abortOnDisconnect);
|
||||
}
|
||||
});
|
||||
};
|
||||
@@ -0,0 +1,46 @@
|
||||
import express from 'express';
|
||||
import request from 'supertest';
|
||||
import { describe, expect, it, vi } from 'vitest';
|
||||
|
||||
import { OpenChamberControlError } from './error.js';
|
||||
import { registerOpenChamberControlRoutes } from './routes.js';
|
||||
|
||||
const createApp = (execute) => {
|
||||
const app = express();
|
||||
registerOpenChamberControlRoutes(app, { controlService: { execute } });
|
||||
return app;
|
||||
};
|
||||
|
||||
describe('OpenChamber control route', () => {
|
||||
it('is a thin adapter over the control service', async () => {
|
||||
const execute = vi.fn(async () => ({ projects: [] }));
|
||||
const response = await request(createApp(execute))
|
||||
.post('/api/openchamber/control')
|
||||
.send({ action: 'projects.list', input: {}, contextDirectory: '/repo' })
|
||||
.expect(200);
|
||||
expect(response.body).toEqual({ projects: [] });
|
||||
expect(execute).toHaveBeenCalledWith('projects.list', {}, '/repo', expect.objectContaining({ signal: expect.any(AbortSignal) }));
|
||||
});
|
||||
|
||||
it('preserves service status and partial-result details', async () => {
|
||||
const execute = vi.fn(async () => {
|
||||
throw new OpenChamberControlError('dispatch failed', 500, {
|
||||
partial: true,
|
||||
partialAction: 'fork-created',
|
||||
sessionId: 'ses_fork',
|
||||
directory: '/repo',
|
||||
});
|
||||
});
|
||||
const response = await request(createApp(execute))
|
||||
.post('/api/openchamber/control')
|
||||
.send({ action: 'session.fork', input: {} })
|
||||
.expect(500);
|
||||
expect(response.body).toEqual({
|
||||
error: 'dispatch failed',
|
||||
partial: true,
|
||||
partialAction: 'fork-created',
|
||||
sessionId: 'ses_fork',
|
||||
directory: '/repo',
|
||||
});
|
||||
});
|
||||
});
|
||||
@@ -0,0 +1,398 @@
|
||||
import path from 'node:path';
|
||||
import { createOpencodeClient } from '@opencode-ai/sdk/v2';
|
||||
import { OpenChamberControlError, asControlError } from './error.js';
|
||||
import { OPENCHAMBER_CONTROL_ACTIONS } from './actions.js';
|
||||
|
||||
const DEFAULT_WAIT_TIMEOUT_SECONDS = 600;
|
||||
const MAX_WAIT_TIMEOUT_SECONDS = 86_400;
|
||||
const WAIT_POLL_INTERVAL_MS = 500;
|
||||
const CONTROL_ACTIONS = new Set(OPENCHAMBER_CONTROL_ACTIONS);
|
||||
const SCHEDULE_TASK_ID_ACTIONS = new Set([
|
||||
'schedule.run',
|
||||
'schedule.delete',
|
||||
'schedule.toggle',
|
||||
]);
|
||||
|
||||
const asNonEmptyString = (value) => {
|
||||
if (typeof value !== 'string') return null;
|
||||
const trimmed = value.trim();
|
||||
return trimmed.length > 0 ? trimmed : null;
|
||||
};
|
||||
|
||||
const positiveInteger = (value, fallback, field) => {
|
||||
if (value === undefined || value === null) return fallback;
|
||||
const number = Number(value);
|
||||
if (!Number.isSafeInteger(number) || number < 1) {
|
||||
throw new OpenChamberControlError(`${field} must be a positive integer`, 400);
|
||||
}
|
||||
return number;
|
||||
};
|
||||
|
||||
const normalizeWaitTimeoutMs = (value) => {
|
||||
const seconds = value === undefined || value === null ? DEFAULT_WAIT_TIMEOUT_SECONDS : Number(value);
|
||||
if (!Number.isSafeInteger(seconds) || seconds < 1 || seconds > MAX_WAIT_TIMEOUT_SECONDS) {
|
||||
throw new OpenChamberControlError(`timeout must be from 1 to ${MAX_WAIT_TIMEOUT_SECONDS} seconds`, 400);
|
||||
}
|
||||
return seconds * 1000;
|
||||
};
|
||||
|
||||
const extractTextMessages = (messages, role = 'all') => {
|
||||
const result = [];
|
||||
for (const record of Array.isArray(messages) ? messages : []) {
|
||||
const info = record?.info;
|
||||
const messageRole = info?.role;
|
||||
if ((messageRole !== 'user' && messageRole !== 'assistant') || (role !== 'all' && role !== messageRole)) continue;
|
||||
const text = Array.isArray(record?.parts)
|
||||
? record.parts.filter((part) => part?.type === 'text' && typeof part.text === 'string').map((part) => part.text).join('').trim()
|
||||
: '';
|
||||
if (!text) continue;
|
||||
const providerID = asNonEmptyString(info.providerID);
|
||||
const modelID = asNonEmptyString(info.modelID);
|
||||
result.push({
|
||||
id: asNonEmptyString(info.id) || '',
|
||||
role: messageRole,
|
||||
createdAt: Number.isFinite(info?.time?.created) ? info.time.created : null,
|
||||
completedAt: Number.isFinite(info?.time?.completed) ? info.time.completed : null,
|
||||
model: providerID && modelID ? `${providerID}/${modelID}` : null,
|
||||
text,
|
||||
});
|
||||
}
|
||||
return result.sort((left, right) => (left.createdAt || 0) - (right.createdAt || 0));
|
||||
};
|
||||
|
||||
const parseModel = (value) => {
|
||||
const model = asNonEmptyString(value);
|
||||
if (!model) throw new OpenChamberControlError('model is required', 400);
|
||||
const slashIndex = model.indexOf('/');
|
||||
if (slashIndex <= 0 || slashIndex === model.length - 1) {
|
||||
throw new OpenChamberControlError('model must be in provider/model format', 400);
|
||||
}
|
||||
return { providerID: model.slice(0, slashIndex), modelID: model.slice(slashIndex + 1) };
|
||||
};
|
||||
|
||||
const parseWeekdays = (value) => {
|
||||
const raw = asNonEmptyString(value);
|
||||
if (!raw) throw new OpenChamberControlError('weekly is required', 400);
|
||||
const weekdays = raw.split(',').map((entry) => Number.parseInt(entry.trim(), 10));
|
||||
if (weekdays.some((entry) => !Number.isInteger(entry) || entry < 0 || entry > 6)) {
|
||||
throw new OpenChamberControlError('weekly must contain weekdays from 0 to 6', 400);
|
||||
}
|
||||
return Array.from(new Set(weekdays)).sort((a, b) => a - b);
|
||||
};
|
||||
|
||||
const buildSchedule = (input) => {
|
||||
const daily = asNonEmptyString(input.daily);
|
||||
const weekly = asNonEmptyString(input.weekly);
|
||||
const once = asNonEmptyString(input.once);
|
||||
const cron = asNonEmptyString(input.cron);
|
||||
const selectors = [daily, weekly, once, cron].filter(Boolean);
|
||||
if (selectors.length !== 1) {
|
||||
throw new OpenChamberControlError('Provide exactly one of daily, weekly, once, or cron', 400);
|
||||
}
|
||||
const timezone = asNonEmptyString(input.timezone);
|
||||
if (daily) return { kind: 'daily', times: [daily], ...(timezone ? { timezone } : {}) };
|
||||
if (weekly) {
|
||||
const time = asNonEmptyString(input.time);
|
||||
if (!time) throw new OpenChamberControlError('time is required with weekly', 400);
|
||||
return { kind: 'weekly', weekdays: parseWeekdays(weekly), times: [time], ...(timezone ? { timezone } : {}) };
|
||||
}
|
||||
if (once) {
|
||||
const time = asNonEmptyString(input.time);
|
||||
if (!time) throw new OpenChamberControlError('time is required with once', 400);
|
||||
return { kind: 'once', date: once, time, ...(timezone ? { timezone } : {}) };
|
||||
}
|
||||
return { kind: 'cron', cron, ...(timezone ? { timezone } : {}) };
|
||||
};
|
||||
|
||||
const buildScheduledTask = (input) => {
|
||||
const name = asNonEmptyString(input.name);
|
||||
const prompt = asNonEmptyString(input.prompt);
|
||||
if (!name) throw new OpenChamberControlError('name is required', 400);
|
||||
if (!prompt) throw new OpenChamberControlError('prompt is required', 400);
|
||||
const model = parseModel(input.model);
|
||||
const goalTokenBudget = input.goalTokenBudget;
|
||||
if (goalTokenBudget !== undefined && input.goal !== true) {
|
||||
throw new OpenChamberControlError('goalTokenBudget requires goal', 400);
|
||||
}
|
||||
if (goalTokenBudget !== undefined && (!Number.isSafeInteger(goalTokenBudget) || goalTokenBudget < 1000 || goalTokenBudget > 100_000_000)) {
|
||||
throw new OpenChamberControlError('goalTokenBudget must be from 1000 to 100000000', 400);
|
||||
}
|
||||
return {
|
||||
name,
|
||||
enabled: input.disabled !== true,
|
||||
schedule: buildSchedule(input),
|
||||
execution: {
|
||||
prompt,
|
||||
...model,
|
||||
...(asNonEmptyString(input.agent) ? { agent: input.agent.trim() } : {}),
|
||||
...(asNonEmptyString(input.variant) ? { variant: input.variant.trim() } : {}),
|
||||
...(input.goal === true ? { goalEnabled: true } : {}),
|
||||
...(goalTokenBudget !== undefined ? { goalTokenBudget } : {}),
|
||||
},
|
||||
};
|
||||
};
|
||||
|
||||
export const createOpenChamberControlService = (dependencies) => {
|
||||
const {
|
||||
readSettingsFromDiskMigrated,
|
||||
sanitizeProjects,
|
||||
buildOpenCodeUrl,
|
||||
getOpenCodeAuthHeaders,
|
||||
waitForOpenCodeReady,
|
||||
sessionService,
|
||||
scheduledTaskService,
|
||||
createClient = createOpencodeClient,
|
||||
sleep = (duration) => new Promise((resolve) => setTimeout(resolve, duration)),
|
||||
now = Date.now,
|
||||
} = dependencies;
|
||||
|
||||
const wait = (duration, signal) => {
|
||||
if (!signal) return sleep(duration);
|
||||
if (signal.aborted) return Promise.reject(new OpenChamberControlError('OpenChamber action was cancelled', 499));
|
||||
return new Promise((resolve, reject) => {
|
||||
const onAbort = () => {
|
||||
signal.removeEventListener('abort', onAbort);
|
||||
reject(new OpenChamberControlError('OpenChamber action was cancelled', 499));
|
||||
};
|
||||
signal.addEventListener('abort', onAbort, { once: true });
|
||||
sleep(duration).then(() => {
|
||||
signal.removeEventListener('abort', onAbort);
|
||||
resolve();
|
||||
}, (error) => {
|
||||
signal.removeEventListener('abort', onAbort);
|
||||
reject(error);
|
||||
});
|
||||
});
|
||||
};
|
||||
|
||||
const getClient = async () => {
|
||||
if (typeof waitForOpenCodeReady === 'function') await waitForOpenCodeReady(10_000, 250);
|
||||
return createClient({
|
||||
baseUrl: buildOpenCodeUrl('/', '').replace(/\/$/, ''),
|
||||
headers: getOpenCodeAuthHeaders(),
|
||||
});
|
||||
};
|
||||
|
||||
const projects = async () => {
|
||||
const settings = await readSettingsFromDiskMigrated();
|
||||
return sanitizeProjects(settings?.projects || []).map((project) => ({
|
||||
id: project.id,
|
||||
path: path.resolve(project.path),
|
||||
label: asNonEmptyString(project.label) || path.basename(project.path) || project.path,
|
||||
}));
|
||||
};
|
||||
|
||||
const models = async () => {
|
||||
const settings = await readSettingsFromDiskMigrated();
|
||||
return {
|
||||
defaultModel: asNonEmptyString(settings?.defaultModel),
|
||||
defaultVariant: asNonEmptyString(settings?.defaultVariant),
|
||||
defaultAgent: asNonEmptyString(settings?.defaultAgent),
|
||||
favoriteModels: Array.isArray(settings?.favoriteModels) ? settings.favoriteModels : [],
|
||||
recentModels: Array.isArray(settings?.recentModels) ? settings.recentModels : [],
|
||||
};
|
||||
};
|
||||
|
||||
const sessionStatus = async (client, sessionID, directory) => {
|
||||
const response = await client.session.status({ directory });
|
||||
const statuses = response?.data;
|
||||
if (!statuses || typeof statuses !== 'object' || Array.isArray(statuses)) {
|
||||
throw new OpenChamberControlError('Invalid session status response', 500);
|
||||
}
|
||||
return statuses[sessionID] || { type: 'idle' };
|
||||
};
|
||||
|
||||
const sessionMessages = async (client, sessionID, directory, role, limit) => {
|
||||
const fetchLimit = limit === undefined ? undefined : Math.max(100, limit * 4);
|
||||
let response = await client.session.messages({ sessionID, directory, ...(fetchLimit ? { limit: fetchLimit } : {}) });
|
||||
let raw = Array.isArray(response?.data) ? response.data : [];
|
||||
let messages = extractTextMessages(raw, role);
|
||||
if (limit !== undefined && messages.length < limit && raw.length >= fetchLimit) {
|
||||
response = await client.session.messages({ sessionID, directory });
|
||||
raw = Array.isArray(response?.data) ? response.data : [];
|
||||
messages = extractTextMessages(raw, role);
|
||||
}
|
||||
return limit === undefined ? messages : messages.slice(-limit);
|
||||
};
|
||||
|
||||
const waitForIdle = async ({ client, sessionID, directory, timeoutMs, requireActivity, baselineMessageID, startedAt, signal }) => {
|
||||
const deadline = now() + timeoutMs;
|
||||
let observedActivity = false;
|
||||
while (true) {
|
||||
if (signal?.aborted) throw new OpenChamberControlError('OpenChamber action was cancelled', 499);
|
||||
const status = await sessionStatus(client, sessionID, directory);
|
||||
if (status.type === 'busy' || status.type === 'retry') {
|
||||
observedActivity = true;
|
||||
} else if (!requireActivity || observedActivity) {
|
||||
return status;
|
||||
} else {
|
||||
const messages = await sessionMessages(client, sessionID, directory, 'assistant', 1);
|
||||
const message = messages[0];
|
||||
if (message?.completedAt && (baselineMessageID ? message.id !== baselineMessageID : message.completedAt >= startedAt)) {
|
||||
return status;
|
||||
}
|
||||
}
|
||||
const remaining = deadline - now();
|
||||
if (remaining <= 0) {
|
||||
throw new OpenChamberControlError(`Session did not become idle within ${Math.ceil(timeoutMs / 1000)} seconds`, 500);
|
||||
}
|
||||
await wait(Math.min(WAIT_POLL_INTERVAL_MS, remaining), signal);
|
||||
}
|
||||
};
|
||||
|
||||
const executeSessionAction = async (action, input, contextDirectory, signal) => {
|
||||
if (input.timeout !== undefined && input.wait !== true) throw new OpenChamberControlError('timeout requires wait', 400);
|
||||
if (input.lastAssistant === true && input.wait !== true) throw new OpenChamberControlError('lastAssistant requires wait', 400);
|
||||
const directory = asNonEmptyString(input.directory) || (!input.projectId ? asNonEmptyString(contextDirectory) : null);
|
||||
const payload = {
|
||||
...(directory ? { directory } : {}),
|
||||
...(asNonEmptyString(input.projectId) ? { projectId: input.projectId.trim() } : {}),
|
||||
...(asNonEmptyString(input.title) ? { title: input.title.trim() } : {}),
|
||||
...(asNonEmptyString(input.prompt) ? { prompt: input.prompt.trim() } : {}),
|
||||
...(asNonEmptyString(input.model) ? { model: input.model.trim() } : {}),
|
||||
...(asNonEmptyString(input.agent) ? { agent: input.agent.trim() } : {}),
|
||||
...(asNonEmptyString(input.variant) ? { variant: input.variant.trim() } : {}),
|
||||
...(input.goal === true ? { goal: true } : {}),
|
||||
...(input.goalTokenBudget !== undefined ? { goalTokenBudget: input.goalTokenBudget } : {}),
|
||||
...(asNonEmptyString(input.worktree) ? { worktree: {
|
||||
name: input.worktree.trim(),
|
||||
...(asNonEmptyString(input.branch) ? { branchName: input.branch.trim() } : {}),
|
||||
...(asNonEmptyString(input.startRef) ? { startRef: input.startRef.trim() } : {}),
|
||||
} } : {}),
|
||||
...(typeof input.setUpstream === 'boolean' ? { setUpstream: input.setUpstream } : {}),
|
||||
...(asNonEmptyString(input.messageId) ? { messageId: input.messageId.trim() } : {}),
|
||||
};
|
||||
const sessionID = asNonEmptyString(input.sessionId);
|
||||
const startedAt = now();
|
||||
let result;
|
||||
if (action === 'session.create') {
|
||||
result = await sessionService.create(payload);
|
||||
} else {
|
||||
if (!sessionID) throw new OpenChamberControlError('sessionId is required', 400);
|
||||
if (action === 'session.send') {
|
||||
result = await sessionService.send(sessionID, payload);
|
||||
} else {
|
||||
result = await sessionService.fork(sessionID, payload);
|
||||
}
|
||||
}
|
||||
if (input.wait !== true) {
|
||||
const publicResult = { ...result };
|
||||
delete publicResult.baselineAssistantMessageId;
|
||||
return publicResult;
|
||||
}
|
||||
const client = await getClient();
|
||||
const status = await waitForIdle({
|
||||
client,
|
||||
sessionID: result.sessionId,
|
||||
directory: result.directory,
|
||||
timeoutMs: normalizeWaitTimeoutMs(input.timeout),
|
||||
requireActivity: result.promptDispatched === true,
|
||||
baselineMessageID: result.baselineAssistantMessageId,
|
||||
startedAt,
|
||||
signal,
|
||||
});
|
||||
const publicResult = { ...result, sessionStatus: status };
|
||||
delete publicResult.baselineAssistantMessageId;
|
||||
if (input.lastAssistant === true) {
|
||||
publicResult.lastAssistantMessage = (await sessionMessages(client, result.sessionId, result.directory, 'assistant', 1))[0] || null;
|
||||
}
|
||||
return publicResult;
|
||||
};
|
||||
|
||||
const execute = async (action, input = {}, contextDirectory, options = {}) => {
|
||||
try {
|
||||
if (!CONTROL_ACTIONS.has(action)) {
|
||||
throw new OpenChamberControlError(`Unsupported OpenChamber action: ${action || 'missing'}`, 400);
|
||||
}
|
||||
if (action === 'projects.list') return { projects: await projects() };
|
||||
if (action === 'models.list') return models();
|
||||
if (action === 'schedule.status') return scheduledTaskService.status();
|
||||
if (action.startsWith('schedule.')) {
|
||||
const taskID = asNonEmptyString(input.taskId);
|
||||
if (SCHEDULE_TASK_ID_ACTIONS.has(action) && !taskID) {
|
||||
throw new OpenChamberControlError('taskId is required', 400);
|
||||
}
|
||||
const explicitProjectID = asNonEmptyString(input.projectId);
|
||||
const explicitDirectory = asNonEmptyString(input.directory);
|
||||
const contextDirectoryFallback = explicitProjectID
|
||||
? undefined
|
||||
: asNonEmptyString(contextDirectory) || undefined;
|
||||
const projectID = await scheduledTaskService.resolveProjectID({
|
||||
projectId: explicitProjectID || undefined,
|
||||
directory: explicitDirectory || contextDirectoryFallback,
|
||||
});
|
||||
switch (action) {
|
||||
case 'schedule.list':
|
||||
return { scheduler: await scheduledTaskService.status(), tasks: await scheduledTaskService.list(projectID) };
|
||||
case 'schedule.create': {
|
||||
const result = await scheduledTaskService.upsert(projectID, buildScheduledTask(input));
|
||||
return { task: result.task, created: result.created };
|
||||
}
|
||||
case 'schedule.run':
|
||||
return scheduledTaskService.run(projectID, taskID);
|
||||
case 'schedule.delete':
|
||||
return { deleted: true, tasks: await scheduledTaskService.remove(projectID, taskID) };
|
||||
case 'schedule.toggle': {
|
||||
if (typeof input.disabled !== 'boolean') {
|
||||
throw new OpenChamberControlError('disabled is required for schedule.toggle', 400);
|
||||
}
|
||||
const enabled = input.disabled === false;
|
||||
return { task: await scheduledTaskService.setEnabled(projectID, taskID, enabled), enabled };
|
||||
}
|
||||
}
|
||||
}
|
||||
if (action === 'session.create' || action === 'session.send' || action === 'session.fork') {
|
||||
return executeSessionAction(action, input, contextDirectory, options.signal);
|
||||
}
|
||||
if (action.startsWith('session.')) {
|
||||
const directory = asNonEmptyString(input.directory) || asNonEmptyString(contextDirectory);
|
||||
const sessionID = asNonEmptyString(input.sessionId);
|
||||
const client = await getClient();
|
||||
if (action === 'session.list') {
|
||||
const limit = positiveInteger(input.limit, 10, 'limit');
|
||||
const response = await client.session.list(directory ? { directory } : {});
|
||||
let sessions = Array.isArray(response?.data) ? response.data : [];
|
||||
if (input.all !== true) sessions = sessions.filter((session) => !session?.time?.archived);
|
||||
sessions = sessions.slice(0, limit);
|
||||
if (input.withStatus === true) {
|
||||
const cache = new Map();
|
||||
sessions = await Promise.all(sessions.map(async (session) => {
|
||||
const sessionDirectory = asNonEmptyString(session?.directory);
|
||||
if (!sessionDirectory) return { ...session, status: { type: 'unknown' } };
|
||||
if (!cache.has(sessionDirectory)) {
|
||||
const statusRequest = client.session.status({ directory: sessionDirectory }).catch(() => null);
|
||||
cache.set(sessionDirectory, statusRequest);
|
||||
}
|
||||
const statusResponse = await cache.get(sessionDirectory);
|
||||
return { ...session, status: statusResponse?.data?.[session.id] || (statusResponse ? { type: 'idle' } : { type: 'unknown' }) };
|
||||
}));
|
||||
}
|
||||
return { sessions, limit, directory, archived: input.all === true ? 'included' : 'excluded' };
|
||||
}
|
||||
if (!sessionID) throw new OpenChamberControlError('sessionId is required', 400);
|
||||
if (!directory) throw new OpenChamberControlError('directory is required', 400);
|
||||
if (action === 'session.status') {
|
||||
return { sessionId: sessionID, directory, sessionStatus: await sessionStatus(client, sessionID, directory) };
|
||||
}
|
||||
if (action === 'session.messages') {
|
||||
if (input.timeout !== undefined && input.wait !== true) throw new OpenChamberControlError('timeout requires wait', 400);
|
||||
const role = input.lastAssistant === true ? 'assistant' : (asNonEmptyString(input.role) || 'all');
|
||||
if (!['all', 'user', 'assistant'].includes(role)) throw new OpenChamberControlError('role must be all, user, or assistant', 400);
|
||||
const last = input.last === true || input.lastAssistant === true;
|
||||
if (input.all === true && (last || input.limit !== undefined)) throw new OpenChamberControlError('all cannot be combined with last or limit', 400);
|
||||
if (last && input.limit !== undefined) throw new OpenChamberControlError('last cannot be combined with limit', 400);
|
||||
const currentStatus = input.wait === true
|
||||
? await waitForIdle({ client, sessionID, directory, timeoutMs: normalizeWaitTimeoutMs(input.timeout), requireActivity: false, startedAt: now(), signal: options.signal })
|
||||
: await sessionStatus(client, sessionID, directory);
|
||||
const limit = input.all === true ? undefined : (last ? 1 : positiveInteger(input.limit, 10, 'limit'));
|
||||
return { sessionId: sessionID, directory, role, sessionStatus: currentStatus, messages: await sessionMessages(client, sessionID, directory, role, limit) };
|
||||
}
|
||||
}
|
||||
throw new OpenChamberControlError(`Unsupported OpenChamber action: ${action || 'missing'}`, 400);
|
||||
} catch (error) {
|
||||
throw asControlError(error, `Failed to execute ${action || 'OpenChamber action'}`);
|
||||
}
|
||||
};
|
||||
|
||||
return { execute };
|
||||
};
|
||||
@@ -0,0 +1,226 @@
|
||||
import { describe, expect, it, vi } from 'vitest';
|
||||
|
||||
import { createOpenChamberControlService } from './service.js';
|
||||
|
||||
const createService = (overrides = {}) => {
|
||||
const client = {
|
||||
session: {
|
||||
list: vi.fn(async () => ({ data: [] })),
|
||||
status: vi.fn(async () => ({ data: {} })),
|
||||
messages: vi.fn(async () => ({ data: [] })),
|
||||
},
|
||||
};
|
||||
const sessionService = {
|
||||
create: vi.fn(async () => ({ sessionId: 'ses_1', directory: '/repo', promptDispatched: false })),
|
||||
send: vi.fn(),
|
||||
fork: vi.fn(),
|
||||
};
|
||||
const scheduledTaskService = {
|
||||
status: vi.fn(async () => ({ enabledScheduledTasksCount: 0 })),
|
||||
resolveProjectID: vi.fn(async () => 'project-1'),
|
||||
list: vi.fn(async () => []),
|
||||
upsert: vi.fn(),
|
||||
run: vi.fn(),
|
||||
remove: vi.fn(),
|
||||
setEnabled: vi.fn(),
|
||||
};
|
||||
const service = createOpenChamberControlService({
|
||||
readSettingsFromDiskMigrated: vi.fn(async () => ({
|
||||
projects: [{ id: 'project-1', path: '/repo', label: 'Repo' }],
|
||||
defaultModel: 'provider/model',
|
||||
favoriteModels: [],
|
||||
recentModels: [],
|
||||
})),
|
||||
sanitizeProjects: (projects) => projects,
|
||||
buildOpenCodeUrl: () => 'http://127.0.0.1:4096/',
|
||||
getOpenCodeAuthHeaders: () => ({ authorization: 'Basic test' }),
|
||||
waitForOpenCodeReady: vi.fn(),
|
||||
createClient: vi.fn(() => client),
|
||||
sessionService,
|
||||
scheduledTaskService,
|
||||
...overrides,
|
||||
});
|
||||
return { service, client, sessionService, scheduledTaskService };
|
||||
};
|
||||
|
||||
describe('OpenChamber control service', () => {
|
||||
it('serves project and model projections without an HTTP or CLI round trip', async () => {
|
||||
const { service } = createService();
|
||||
await expect(service.execute('projects.list')).resolves.toEqual({
|
||||
projects: [{ id: 'project-1', path: '/repo', label: 'Repo' }],
|
||||
});
|
||||
await expect(service.execute('models.list')).resolves.toEqual(expect.objectContaining({
|
||||
defaultModel: 'provider/model',
|
||||
favoriteModels: [],
|
||||
}));
|
||||
});
|
||||
|
||||
it('maps schedule creation into the shared scheduled-task service', async () => {
|
||||
const { service, scheduledTaskService } = createService();
|
||||
scheduledTaskService.upsert.mockResolvedValue({ task: { id: 'task-1' }, created: true });
|
||||
await expect(service.execute('schedule.create', {
|
||||
directory: '/repo',
|
||||
name: 'Daily',
|
||||
prompt: 'Run checks',
|
||||
model: 'provider/model',
|
||||
daily: ' 09:00 ',
|
||||
goal: true,
|
||||
goalTokenBudget: 5000,
|
||||
})).resolves.toEqual({ task: { id: 'task-1' }, created: true });
|
||||
expect(scheduledTaskService.resolveProjectID).toHaveBeenCalledWith({ projectId: undefined, directory: '/repo' });
|
||||
expect(scheduledTaskService.upsert).toHaveBeenCalledWith('project-1', expect.objectContaining({
|
||||
name: 'Daily',
|
||||
schedule: { kind: 'daily', times: ['09:00'] },
|
||||
execution: expect.objectContaining({ providerID: 'provider', modelID: 'model', goalEnabled: true, goalTokenBudget: 5000 }),
|
||||
}));
|
||||
});
|
||||
|
||||
it('does not combine an explicit schedule project with the tool context directory', async () => {
|
||||
const { service, scheduledTaskService } = createService();
|
||||
await service.execute('schedule.list', { projectId: ' project-1 ' }, '/current-session');
|
||||
expect(scheduledTaskService.resolveProjectID).toHaveBeenCalledWith({ projectId: 'project-1', directory: undefined });
|
||||
});
|
||||
|
||||
it('includes scheduler status alongside listed tasks', async () => {
|
||||
const { service, scheduledTaskService } = createService();
|
||||
scheduledTaskService.list.mockResolvedValue([{ id: 'task-1' }]);
|
||||
await expect(service.execute('schedule.list', {}, '/repo')).resolves.toEqual({
|
||||
scheduler: { enabledScheduledTasksCount: 0 },
|
||||
tasks: [{ id: 'task-1' }],
|
||||
});
|
||||
});
|
||||
|
||||
it('toggles a scheduled task through the required disabled boolean', async () => {
|
||||
const { service, scheduledTaskService } = createService();
|
||||
scheduledTaskService.setEnabled.mockResolvedValue({ id: 'task-1', enabled: false });
|
||||
await expect(service.execute('schedule.toggle', { taskId: 'task-1' }, '/repo')).rejects.toThrow('disabled is required for schedule.toggle');
|
||||
await expect(service.execute('schedule.toggle', { taskId: 'task-1', disabled: true }, '/repo')).resolves.toEqual({
|
||||
task: { id: 'task-1', enabled: false },
|
||||
enabled: false,
|
||||
});
|
||||
expect(scheduledTaskService.setEnabled).toHaveBeenCalledWith('project-1', 'task-1', false);
|
||||
});
|
||||
|
||||
it('returns an actionable taskId error before resolving schedule scope', async () => {
|
||||
const { service, scheduledTaskService } = createService();
|
||||
await expect(service.execute('schedule.run', {}, '/repo')).rejects.toThrow('taskId is required');
|
||||
expect(scheduledTaskService.resolveProjectID).not.toHaveBeenCalled();
|
||||
expect(scheduledTaskService.run).not.toHaveBeenCalled();
|
||||
});
|
||||
|
||||
it('validates wait modifiers before creating a session', async () => {
|
||||
const { service, sessionService } = createService();
|
||||
await expect(service.execute('session.create', { directory: '/repo', timeout: 30 })).rejects.toThrow('timeout requires wait');
|
||||
expect(sessionService.create).not.toHaveBeenCalled();
|
||||
});
|
||||
|
||||
it('uses the tool context directory for session actions', async () => {
|
||||
const { service, sessionService } = createService();
|
||||
await service.execute('session.create', { title: 'From tool' }, '/repo');
|
||||
expect(sessionService.create).toHaveBeenCalledWith({ directory: '/repo', title: 'From tool' });
|
||||
});
|
||||
|
||||
it.each([
|
||||
['session.send', 'send'],
|
||||
['session.fork', 'fork'],
|
||||
])('delegates %s directly to the session service', async (action, method) => {
|
||||
const { service, sessionService } = createService();
|
||||
sessionService[method].mockResolvedValue({ sessionId: 'ses_1', directory: '/repo' });
|
||||
|
||||
await service.execute(action, { sessionId: 'ses_1', directory: '/repo', prompt: 'Continue' });
|
||||
|
||||
expect(sessionService[method]).toHaveBeenCalledWith('ses_1', { directory: '/repo', prompt: 'Continue' });
|
||||
});
|
||||
|
||||
it('waits past initial idle until a completed assistant result appears', async () => {
|
||||
let timestamp = 1000;
|
||||
const { service, client, sessionService } = createService({
|
||||
now: () => timestamp,
|
||||
sleep: async (duration) => { timestamp += duration; },
|
||||
});
|
||||
sessionService.create.mockResolvedValue({
|
||||
sessionId: 'ses_1',
|
||||
directory: '/repo',
|
||||
promptDispatched: true,
|
||||
baselineAssistantMessageId: 'msg_old',
|
||||
});
|
||||
client.session.status.mockResolvedValue({ data: { ses_1: { type: 'idle' } } });
|
||||
client.session.messages
|
||||
.mockResolvedValueOnce({ data: [{ info: { id: 'msg_old', role: 'assistant', time: { completed: 900 } }, parts: [{ type: 'text', text: 'old' }] }] })
|
||||
.mockResolvedValueOnce({ data: [{ info: { id: 'msg_new', role: 'assistant', time: { completed: 1500 } }, parts: [{ type: 'text', text: 'done' }] }] })
|
||||
.mockResolvedValueOnce({ data: [{ info: { id: 'msg_new', role: 'assistant', time: { completed: 1500 } }, parts: [{ type: 'text', text: 'done' }] }] });
|
||||
|
||||
await expect(service.execute('session.create', {
|
||||
directory: '/repo',
|
||||
prompt: 'work',
|
||||
wait: true,
|
||||
lastAssistant: true,
|
||||
timeout: 2,
|
||||
})).resolves.toEqual(expect.objectContaining({
|
||||
sessionStatus: { type: 'idle' },
|
||||
lastAssistantMessage: expect.objectContaining({ id: 'msg_new', text: 'done' }),
|
||||
}));
|
||||
expect(client.session.status).toHaveBeenCalledTimes(2);
|
||||
});
|
||||
|
||||
it('filters archived sessions and adds directory-scoped statuses', async () => {
|
||||
const { service, client } = createService();
|
||||
client.session.list.mockResolvedValue({ data: [
|
||||
{ id: 'ses_active', directory: '/repo', time: {} },
|
||||
{ id: 'ses_archived', directory: '/repo', time: { archived: 100 } },
|
||||
{ id: 'ses_other', directory: '/other', time: {} },
|
||||
] });
|
||||
client.session.status
|
||||
.mockResolvedValueOnce({ data: { ses_active: { type: 'busy' } } })
|
||||
.mockRejectedValueOnce(new Error('unavailable'));
|
||||
|
||||
await expect(service.execute('session.list', { limit: 10, withStatus: true })).resolves.toEqual({
|
||||
sessions: [
|
||||
{ id: 'ses_active', directory: '/repo', time: {}, status: { type: 'busy' } },
|
||||
{ id: 'ses_other', directory: '/other', time: {}, status: { type: 'unknown' } },
|
||||
],
|
||||
limit: 10,
|
||||
directory: null,
|
||||
archived: 'excluded',
|
||||
});
|
||||
});
|
||||
|
||||
it('names limit in positive-integer validation errors', async () => {
|
||||
const { service, client } = createService();
|
||||
await expect(service.execute('session.list', { limit: 0 })).rejects.toThrow('limit must be a positive integer');
|
||||
expect(client.session.list).not.toHaveBeenCalled();
|
||||
});
|
||||
|
||||
it('projects only ordered text parts from session messages', async () => {
|
||||
const { service, client } = createService();
|
||||
client.session.messages.mockResolvedValue({ data: [
|
||||
{
|
||||
info: { id: 'msg_assistant', role: 'assistant', providerID: 'openai', modelID: 'gpt-5.4-mini', time: { created: 20, completed: 30 } },
|
||||
parts: [{ type: 'reasoning', text: 'hidden' }, { type: 'text', text: 'First ' }, { type: 'tool' }, { type: 'text', text: 'answer' }],
|
||||
},
|
||||
{ info: { id: 'msg_user', role: 'user', time: { created: 10 } }, parts: [{ type: 'text', text: 'Question' }] },
|
||||
{ info: { id: 'msg_tool', role: 'assistant', time: { created: 15 } }, parts: [{ type: 'tool' }] },
|
||||
] });
|
||||
|
||||
await expect(service.execute('session.messages', {
|
||||
sessionId: 'ses_1',
|
||||
directory: '/repo',
|
||||
role: 'all',
|
||||
all: true,
|
||||
})).resolves.toEqual({
|
||||
sessionId: 'ses_1',
|
||||
directory: '/repo',
|
||||
role: 'all',
|
||||
sessionStatus: { type: 'idle' },
|
||||
messages: [
|
||||
{ id: 'msg_user', role: 'user', createdAt: 10, completedAt: null, model: null, text: 'Question' },
|
||||
{ id: 'msg_assistant', role: 'assistant', createdAt: 20, completedAt: 30, model: 'openai/gpt-5.4-mini', text: 'First answer' },
|
||||
],
|
||||
});
|
||||
});
|
||||
|
||||
it('rejects actions outside the fixed contract', async () => {
|
||||
const { service } = createService();
|
||||
await expect(service.execute('session.delete')).rejects.toThrow('Unsupported OpenChamber action');
|
||||
});
|
||||
});
|
||||
Reference in New Issue
Block a user