Changelog, dev startup, README, AgentHub 0.4.0, model catalog, and landing site (#7)
Branch-length batch covering tooling, the model layer, the Web App and the public surfaces. Highlights: - Changelog: a per-release `changelog/<version>/` tree, grouped by the surface each change touches, with a root CHANGELOG.md holding one line per release. - Dev startup: `scripts/dev-prebuild.mjs` serializes the skills+core prebuild behind a lock and keeps `pnpm install` current; `pnpm dev` runs server+web together. - AgentHub 0.3.3 -> 0.4.0: OmniMessage complete payloads carry one opaque `fidelity` object in place of item-level `signature`/`phase`, threaded verbatim through Trace, replay and resume; malformed classification adapted to the new error types. - Model layer: a model is always referenced by an explicit `(provider, model_id)` pair. The provider is never inferred, guessed or defaulted -- both the catalog inference and the unique-match config resolution are gone, and CLI, SDK, server routes and run_subagent all require the complete pair. Catalog gains the Qwen Token Plan, Qwen Pay-As-You-Go and Fireworks AI gateways, plus an expanded OpenRouter group. - Web App: catalog preset sync and per-group speed test on the Models page, positional slash commands, a markdown renderer, skill-library update reminders, and a vertically centred draft page whose upward menus size themselves to the room available. - Public surfaces: restructured READMEs, the penguin.ooo landing site and blog, refreshed benchmark results for both suites, and the demo videos playing on the landing page. Includes the fixes from a full review of the branch: 23 confirmed findings, among them a provider-inference bug that could send one vendor's API key to another vendor's endpoint, and an Escape handler that destroyed the composer's contents unrecoverably. Verified on the branch head: pnpm test (1127 passing, 7 packages), pnpm typecheck and pnpm format:check clean, Playwright e2e 14/14.
This commit is contained in:
@@ -273,6 +273,8 @@ export interface ModelTestRequest {
|
||||
apiKey?: string;
|
||||
/** "Clear saved API key" is checked: the test does **not** fall back to the stored key (tests against the current draft). */
|
||||
clearApiKey?: boolean;
|
||||
/** Speed-test mode: raises the probe's output cap (16 -> 64 tokens) so TTFT/TPS are measurable; costs a little more quota. */
|
||||
speed?: boolean;
|
||||
/**
|
||||
* base URL (not secret; the frontend always sends the form's current value): a string
|
||||
* means use it, `null` means explicitly clear it (no fallback to the stored value),
|
||||
@@ -283,10 +285,23 @@ export interface ModelTestRequest {
|
||||
clientType?: string;
|
||||
}
|
||||
|
||||
/** Connectivity test result: carries round-trip latency when ok, and a reason on failure (truncated raw provider error). */
|
||||
/**
|
||||
* Connectivity test result: carries round-trip latency when ok, and a reason on failure
|
||||
* (truncated raw provider error). When streamed content was observed, also carries the
|
||||
* time-to-first-token and, when usage was reported (completed streams), the output rate.
|
||||
*/
|
||||
export interface ModelTestResponse {
|
||||
ok: boolean;
|
||||
latencyMs?: number;
|
||||
/** Time from request start to the first streamed content (thinking or text), ms. */
|
||||
ttftMs?: number;
|
||||
/**
|
||||
* Output tokens per second over the streaming window (first content -> stream end), 1dp.
|
||||
* Omitted unless the sample is large enough to mean anything: a reply of a few tokens is
|
||||
* dominated by the final chunk's round trip, so the rate it yields tracks network jitter
|
||||
* rather than the model. Callers render TTFT alone in that case.
|
||||
*/
|
||||
tps?: number;
|
||||
message?: string;
|
||||
}
|
||||
|
||||
@@ -342,6 +357,8 @@ export interface AgentSummary {
|
||||
vaultKeyCount: number;
|
||||
/** Schedule count (number of .toml files under agent_state/schedule/, including invalid ones). */
|
||||
scheduleCount: number;
|
||||
/** Installed Skill count (number of agent_state/skills/<name>/ directories with a SKILL.md). */
|
||||
skillCount: number;
|
||||
}
|
||||
|
||||
export interface AgentsResponse {
|
||||
@@ -460,12 +477,12 @@ export interface DirListResponse {
|
||||
}
|
||||
|
||||
export interface SessionCreateRequest {
|
||||
/** Upstream id of the session's model (paired with provider); defaults to the Project's default Model. */
|
||||
/** Upstream id of the session's model; always sent together with provider. Omit both for the Project's default Model. */
|
||||
modelId?: string;
|
||||
/**
|
||||
* Provider group for `modelId`; when omitted, resolved via resolveModelRef semantics —
|
||||
* modelId can only be resolved if it's globally unique by exact match in the config;
|
||||
* 0 or multiple matches return 400.
|
||||
* Provider group for `modelId`. A model reference is always a complete
|
||||
* (provider, modelId) pair — the provider is never inferred, so sending one field
|
||||
* without the other returns 400 instead of being resolved.
|
||||
*/
|
||||
provider?: string;
|
||||
/** Any existing directory on the server; defaults to auto-creating a temporary Workspace. */
|
||||
@@ -930,9 +947,9 @@ export interface ScheduleItem {
|
||||
/** Bound target Session; defaults to creating a new Session each time. */
|
||||
sessionId?: string;
|
||||
workspace?: string;
|
||||
/** Model for new-Session mode (upstream id, paired with provider); defaults to the Project's default reference. */
|
||||
/** Model for new-Session mode (upstream id, always paired with provider); absent means the Project's default reference. */
|
||||
modelId?: string;
|
||||
/** Provider group for `modelId`; when omitted, resolved via resolveModelRef semantics (resolvable only on a unique match). */
|
||||
/** Provider group for `modelId`; present exactly when `modelId` is — a model reference is always a pair. */
|
||||
provider?: string;
|
||||
status: ScheduleStatus;
|
||||
invalidReason?: string;
|
||||
@@ -959,9 +976,12 @@ export interface ScheduleUpsertRequest {
|
||||
endAt?: string;
|
||||
sessionId?: string;
|
||||
workspace?: string;
|
||||
/** Model for new-Session mode (upstream id); defaults to the Project's default reference. */
|
||||
/** Model for new-Session mode (upstream id); always sent together with provider, omit both for the Project's default reference. */
|
||||
modelId?: string;
|
||||
/** Provider group for `modelId`; when omitted, validated as uniquely resolvable via resolveModelRef semantics at save/reconciliation time. */
|
||||
/**
|
||||
* Provider group for `modelId`. Both fields are sent as a pair (400 otherwise); the
|
||||
* pair is checked against the Project config at save/reconciliation time.
|
||||
*/
|
||||
provider?: string;
|
||||
}
|
||||
|
||||
|
||||
@@ -3,7 +3,7 @@
|
||||
* login / logout / password change / session validation.
|
||||
*
|
||||
* - No open registration: on startup, if there are no users at all, the built-in
|
||||
* admin `admin` is seeded (initial password admin123), and it adopts
|
||||
* admin `admin` is seeded (initial password penguin-2026), and it adopts
|
||||
* `default_project`; all other users are created by an admin via the user
|
||||
* backend (admin-service).
|
||||
* - An initial password (whether seeded or set by an admin) is flagged with
|
||||
@@ -22,7 +22,7 @@ export const MIN_PASSWORD_LENGTH = 8;
|
||||
|
||||
/** Built-in admin: user_id and initial password (matches the README and login-page hint). */
|
||||
export const ADMIN_USER_ID = "admin";
|
||||
export const ADMIN_INITIAL_PASSWORD = "admin123";
|
||||
export const ADMIN_INITIAL_PASSWORD = "penguin-2026";
|
||||
|
||||
function sha256Hex(value: string): string {
|
||||
return createHash("sha256").update(value).digest("hex");
|
||||
|
||||
@@ -160,6 +160,10 @@ export function modelsRoutes(deps: AppDeps): Hono<AppEnv> {
|
||||
if (typeof body.clearApiKey !== "boolean") throw badRequest("clearApiKey must be a boolean.");
|
||||
req.clearApiKey = body.clearApiKey;
|
||||
}
|
||||
if (body.speed !== undefined) {
|
||||
if (typeof body.speed !== "boolean") throw badRequest("speed must be a boolean.");
|
||||
req.speed = body.speed;
|
||||
}
|
||||
// null = explicit clear (test against the draft, don't fall back to the stored value); empty string is treated as null.
|
||||
if (body.baseUrl !== undefined) {
|
||||
if (body.baseUrl !== null && typeof body.baseUrl !== "string") {
|
||||
|
||||
@@ -213,7 +213,8 @@ async function upsert(
|
||||
const raw = serializeSchedule(fields);
|
||||
const parsed = parseScheduleFile(name, raw);
|
||||
if (!parsed.ok) throw badRequest(`Invalid schedule configuration: ${parsed.error}`);
|
||||
// At save time, verify the model reference resolves (resolveModelRef semantics; same rules as reconciliation) so we never persist a broken file.
|
||||
// At save time, verify the (provider, modelId) pair names a configured model (same rules as reconciliation) so we never persist a broken file.
|
||||
// The pairing rule itself is enforced by parseScheduleFile above, which rejects half a reference.
|
||||
const refError = await validateScheduleModelRef(deps.config.root, projectId, parsed.def);
|
||||
if (refError !== null) throw badRequest(`Invalid schedule configuration: ${refError}`);
|
||||
await writeScheduleFile(deps.config.root, projectId, agentId, name, raw);
|
||||
|
||||
@@ -108,10 +108,13 @@ export function agentSessionsRoutes(deps: AppDeps): Hono<AppEnv> {
|
||||
const body = await readJson(c);
|
||||
const modelId = optionalString(body, "modelId", { minLen: 1, label: "modelId" });
|
||||
const provider = optionalString(body, "provider", { minLen: 1, label: "provider" });
|
||||
// Model reference is submitted as a pair: provider can't appear without modelId (core does the same validation; this catches it early).
|
||||
if (provider !== undefined && modelId === undefined) {
|
||||
// Model reference is submitted as a pair — both or neither. Neither half is ever
|
||||
// inferred from the other, so half a reference is rejected here instead of being
|
||||
// resolved (core does the same validation; this catches it early). Omitting both
|
||||
// falls back to the Project's default model.
|
||||
if ((modelId === undefined) !== (provider === undefined)) {
|
||||
throw badRequest(
|
||||
"provider is specified but modelId is not: a model reference must be given as a pair.",
|
||||
"modelId and provider must be given together as a model reference pair: specify both, or neither to use the Project's default model.",
|
||||
);
|
||||
}
|
||||
const approvalMode = optionalEnum(body, "approvalMode", APPROVAL_MODES);
|
||||
@@ -361,6 +364,11 @@ export function sessionsRoutes(deps: AppDeps): Hono<AppEnv> {
|
||||
const row = resolveSession(c);
|
||||
const rel = c.req.query("path") ?? "";
|
||||
const download = c.req.query("download") === "1";
|
||||
// Sandboxed top-level preview ("open in a new tab" for html): the document keeps its REAL
|
||||
// content type but carries a CSP sandbox WITHOUT allow-same-origin — it renders and runs
|
||||
// fully in an opaque origin, so agent-generated markup cannot reach this origin's cookies
|
||||
// or API. The request itself still authenticates (top-level GET sends the Lax cookie).
|
||||
const preview = !download && c.req.query("preview") === "1";
|
||||
const { data, fileName, contentType, scriptable } = await deps.workspaceFiles.read(
|
||||
row.workspace,
|
||||
rel,
|
||||
@@ -368,15 +376,22 @@ export function sessionsRoutes(deps: AppDeps): Hono<AppEnv> {
|
||||
const disposition = download ? "attachment" : "inline";
|
||||
// Same-origin XSS defense: html/svg inline previews are always returned as plain
|
||||
// text (Workspace files may be Agent-generated and untrusted); downloads
|
||||
// (attachment) keep the real content type. Paired with nosniff to prevent MIME
|
||||
// sniffing from undoing this.
|
||||
const effectiveType = !download && scriptable ? "text/plain; charset=utf-8" : contentType;
|
||||
// (attachment) keep the real content type, and sandboxed previews keep it under the
|
||||
// CSP above. Paired with nosniff to prevent MIME sniffing from undoing this.
|
||||
const effectiveType =
|
||||
!download && scriptable && !preview ? "text/plain; charset=utf-8" : contentType;
|
||||
return new Response(new Uint8Array(data), {
|
||||
status: 200,
|
||||
headers: {
|
||||
"Content-Type": effectiveType,
|
||||
"Content-Disposition": `${disposition}; filename*=UTF-8''${encodeURIComponent(fileName)}`,
|
||||
"X-Content-Type-Options": "nosniff",
|
||||
...(preview && scriptable
|
||||
? {
|
||||
"Content-Security-Policy":
|
||||
"sandbox allow-scripts allow-popups allow-modals allow-forms",
|
||||
}
|
||||
: {}),
|
||||
},
|
||||
});
|
||||
});
|
||||
|
||||
@@ -48,9 +48,10 @@ export function vaultRoutes(deps: AppDeps): Hono<AppEnv> {
|
||||
deps.projectService.requireProjectOwner(c.var.user.userId, projectId);
|
||||
const req = parseVaultUpdate(await readJson(c));
|
||||
const res = await deps.agentConfigService.updateVault(projectId, agentId, req);
|
||||
// Effective-value semantics: no hot update — an already-built runtime is neither
|
||||
// evicted nor reloaded; the new value only applies to Sessions created or resumed
|
||||
// afterward.
|
||||
// Effective-value semantics: no hot swap into a Task already in flight, but every
|
||||
// runtime built before this update is invalidated — the next Task on any Session
|
||||
// of this Agent re-resumes and picks up the new vault values.
|
||||
deps.manager.invalidateAgentRuntimes(projectId, agentId);
|
||||
return c.json(res);
|
||||
});
|
||||
|
||||
|
||||
@@ -20,7 +20,7 @@ const config = resolveServerConfig();
|
||||
const deps = buildAppDeps(config);
|
||||
const app = createApp(deps);
|
||||
|
||||
// Built-in admin seed (idempotent): creates admin (initial password admin123) and adopts default_project when the users table is empty.
|
||||
// Built-in admin seed (idempotent): creates admin (initial password penguin-2026) and adopts default_project when the users table is empty.
|
||||
await deps.authService.seedAdmin();
|
||||
|
||||
// Schedule scheduler: startup reconciliation (missed, don't backfill) + periodic scan; only active while the server is running.
|
||||
|
||||
@@ -8,7 +8,7 @@
|
||||
* recycled/recreated or the process restarts, so a stale Last-Event-ID always misses
|
||||
* and falls through to resync — this prevents a silent false-hit event loss when the
|
||||
* new epoch's event count happens to exceed the old id;
|
||||
* - A bounded ring buffer (most recent 1000 entries or 2MB, whichever comes first,
|
||||
* - A bounded ring buffer (most recent 10,000 entries or 8MB, whichever comes first,
|
||||
* evicting the oldest on overflow) serves replay-on-reconnect via `Last-Event-ID`;
|
||||
* an evicted/unknown id is handled by the caller sending `resync_required`;
|
||||
* - Unicast (sendTo) is used for one-off replay at subscribe time (pending approvals /
|
||||
@@ -37,8 +37,12 @@ export interface ChannelOptions {
|
||||
maxBufferBytes?: number;
|
||||
}
|
||||
|
||||
const DEFAULT_MAX_COUNT = 1000;
|
||||
const DEFAULT_MAX_BYTES = 2 * 1024 * 1024;
|
||||
// Sized so a mid-stream reconnect during a fast large-code reply still hits replay instead of
|
||||
// resync_required: the server publishes one event per provider delta (a 240KB reply ≈ 5k events
|
||||
// / 1.4MB), and 1000 events covered only ~6s of such a stream — every longer blip forced a full
|
||||
// client-side history rebuild.
|
||||
const DEFAULT_MAX_COUNT = 10_000;
|
||||
const DEFAULT_MAX_BYTES = 8 * 1024 * 1024;
|
||||
|
||||
/** Buffered entry: seq is stored separately so hit checks never need to parse the string id. */
|
||||
interface BufferedEvent {
|
||||
|
||||
@@ -35,13 +35,13 @@ export interface ScheduleDefinition {
|
||||
sessionId?: string;
|
||||
/** Workspace for new-Session mode (same semantics as manually starting a session; auto-creates a temp directory if unspecified). */
|
||||
workspace?: string;
|
||||
/** Model for new-Session mode (upstream id, paired with provider; defaults to the Project's default reference). */
|
||||
/** Model for new-Session mode (upstream id, always paired with provider; omit both for the Project's default reference). */
|
||||
modelId?: string;
|
||||
/**
|
||||
* Vendor grouping for `model_id` (paired reference); when omitted, resolved per
|
||||
* resolveModelRef semantics — whether the reference is resolvable is validated by the
|
||||
* caller against config at reconciliation/save time (this module does pure parsing
|
||||
* and never touches config).
|
||||
* Vendor grouping for `model_id`; always present when `model_id` is (the pairing rule
|
||||
* is enforced here). Whether that pair names a configured model is validated by the
|
||||
* caller against config at reconciliation/save time — this module does pure parsing
|
||||
* and never touches config.
|
||||
*/
|
||||
provider?: string;
|
||||
}
|
||||
@@ -148,13 +148,8 @@ export function parseScheduleFile(name: string, raw: string): ScheduleParseResul
|
||||
}
|
||||
provider = t["provider"];
|
||||
}
|
||||
if (provider !== undefined && modelId === undefined) {
|
||||
return {
|
||||
ok: false,
|
||||
error:
|
||||
"provider is only used together with model_id (a model reference must be given as a pair)",
|
||||
};
|
||||
}
|
||||
// Target conflicts are reported first: when session_id is set the whole model reference
|
||||
// is out of place, so complaining about its shape would point at the wrong fix.
|
||||
if (
|
||||
sessionId !== undefined &&
|
||||
(workspace !== undefined || modelId !== undefined || provider !== undefined)
|
||||
@@ -164,6 +159,18 @@ export function parseScheduleFile(name: string, raw: string): ScheduleParseResul
|
||||
error: "Pick one target: workspace and provider / model_id are only for new-Session mode",
|
||||
};
|
||||
}
|
||||
// A model reference is always a complete (provider, model_id) pair: neither half is
|
||||
// ever inferred from the other, so a file carrying only one of them is invalid rather
|
||||
// than something to resolve later. Omit both to run on the Project's default model.
|
||||
// A file written before this rule existed (model_id alone) therefore parses as
|
||||
// invalid: it is listed in invalidFiles and skipped, never scheduled.
|
||||
if ((modelId === undefined) !== (provider === undefined)) {
|
||||
return {
|
||||
ok: false,
|
||||
error:
|
||||
"provider and model_id must be given together (a model reference is always a pair); omit both to use the Project's default model",
|
||||
};
|
||||
}
|
||||
|
||||
return {
|
||||
ok: true,
|
||||
|
||||
@@ -88,19 +88,23 @@ export function serializeSchedule(fields: {
|
||||
}
|
||||
|
||||
/**
|
||||
* Resolvability check for a schedule's model reference (shared by save and
|
||||
* reconciliation): when the definition has `model_id`, it's
|
||||
* resolved against Project config per resolveModelRef semantics — omitting provider is
|
||||
* only resolvable when model_id matches exactly one entry globally; zero hits or
|
||||
* ambiguity means unresolvable. Returns an error message (unresolvable / config read
|
||||
* failure), or null if resolvable (or no model reference at all).
|
||||
* Config check for a schedule's model reference (shared by save and reconciliation):
|
||||
* the reference is a complete `(provider, model_id)` pair and must name an entry in the
|
||||
* Project config — nothing is inferred, so a pair that matches no entry is simply an
|
||||
* error. Returns an error message (no such model / half a reference / config read
|
||||
* failure), or null when the pair is configured (or there is no model reference at all).
|
||||
*/
|
||||
export async function validateScheduleModelRef(
|
||||
root: string,
|
||||
projectId: string,
|
||||
def: Pick<ScheduleDefinition, "modelId" | "provider">,
|
||||
): Promise<string | null> {
|
||||
if (def.modelId === undefined) return null;
|
||||
if (def.modelId === undefined && def.provider === undefined) return null;
|
||||
// parseScheduleFile already rejects half a reference, so this only guards callers that
|
||||
// build a definition by hand — the missing half is reported, never guessed.
|
||||
if (def.modelId === undefined || def.provider === undefined) {
|
||||
return "provider and model_id must be given together (a model reference is always a pair); omit both to use the Project's default model";
|
||||
}
|
||||
try {
|
||||
const cfg = await loadProjectConfig(root, projectId);
|
||||
resolveModelRef(cfg, def.modelId, def.provider);
|
||||
|
||||
Binary file not shown.
@@ -8,6 +8,9 @@
|
||||
* createSession using the index row's workspace/modelId, yielding a new session_id
|
||||
* and updating the index's primary key; the Task response body always returns the
|
||||
* current actual id;
|
||||
* - Vault effectiveness: a vault update bumps the Agent's config generation
|
||||
* (invalidateAgentRuntimes); runtimes built earlier are discarded on their next
|
||||
* idle access and re-resumed, so the next Task always runs with current values;
|
||||
* - Per-Session mutual exclusion: only one Task/compaction may be in progress at a
|
||||
* time;
|
||||
* - run/compact drive: consumes the output stream in the background, publishing each
|
||||
@@ -187,6 +190,12 @@ interface RuntimeEntry {
|
||||
abort: AbortController | null;
|
||||
/** The in-flight drive Promise (awaited during graceful shutdown). */
|
||||
running: Promise<void> | null;
|
||||
/**
|
||||
* Agent config generation this runtime was built under (see
|
||||
* invalidateAgentRuntimes): once it falls behind the Agent's current generation,
|
||||
* the entry is discarded on its next idle access and re-resumed via the loader.
|
||||
*/
|
||||
generation: number;
|
||||
/** Timestamp of last activity (refreshed on load / status flip / drive completion), used for idle-eviction checks. */
|
||||
lastActivityMs: number;
|
||||
}
|
||||
@@ -197,6 +206,13 @@ const ENTRY_SWEEP_INTERVAL_MS = 60 * 1000;
|
||||
|
||||
/** Cap on collected model text for title material (accumulation stops beyond this; the generator side also truncates further). */
|
||||
const TITLE_EXCERPT_LIMIT = 4000;
|
||||
/**
|
||||
* Early title trigger: once this many characters of main-session body text have streamed,
|
||||
* title generation starts right away instead of waiting for the Task to finish — the core
|
||||
* Session has captured its (capped) material by then, and a long answer would only overrun
|
||||
* it. Short conversations are still covered by the completion trigger in drive's finally.
|
||||
*/
|
||||
const EARLY_TITLE_BODY_CHARS = 1000;
|
||||
|
||||
/** Composite Agent key (used as a Set key, avoiding projectId/agentId concatenation ambiguity). */
|
||||
function agentKey(projectId: string, agentId: string): string {
|
||||
@@ -280,6 +296,8 @@ export class SessionManager {
|
||||
private readonly deletingAgents = new Set<string>();
|
||||
/** Sessions currently being deleted (guards against the entry/Trace file being rebuilt and reviving it inside the deletion race window). */
|
||||
private readonly deletingSessions = new Set<string>();
|
||||
/** Per-Agent config generation (key = agentKey), bumped by invalidateAgentRuntimes on vault updates. */
|
||||
private readonly agentGenerations = new Map<string, number>();
|
||||
private readonly sweepTimer: NodeJS.Timeout;
|
||||
|
||||
constructor(private readonly deps: SessionManagerDeps) {
|
||||
@@ -324,10 +342,25 @@ export class SessionManager {
|
||||
approvals: new ApprovalRegistry(),
|
||||
abort: null,
|
||||
running: null,
|
||||
generation: this.generationOf(row.projectId, row.agentId),
|
||||
lastActivityMs: Date.now(),
|
||||
});
|
||||
}
|
||||
|
||||
/**
|
||||
* After an Agent's vault is updated: bump the Agent's config generation so every
|
||||
* runtime built before the update is discarded on its next idle access and
|
||||
* re-resumed via the loader — resume re-reads agent_state/.vault.toml, so the next
|
||||
* Task on any of this Agent's Sessions runs with the new values (history is
|
||||
* preserved through the Trace). A Task already in flight is neither aborted nor
|
||||
* hot-swapped: it keeps the values it started with, and its entry is rebuilt on
|
||||
* the first access after it returns to idle (see ensureEntry).
|
||||
*/
|
||||
invalidateAgentRuntimes(projectId: string, agentId: string): void {
|
||||
const key = agentKey(projectId, agentId);
|
||||
this.agentGenerations.set(key, this.generationOf(projectId, agentId) + 1);
|
||||
}
|
||||
|
||||
// —— Task / compaction drive ——
|
||||
|
||||
/**
|
||||
@@ -578,10 +611,30 @@ export class SessionManager {
|
||||
}
|
||||
}
|
||||
|
||||
private generationOf(projectId: string, agentId: string): number {
|
||||
return this.agentGenerations.get(agentKey(projectId, agentId)) ?? 0;
|
||||
}
|
||||
|
||||
/** get-or-resume-or-heal: use directly on an active-table hit; otherwise load via the loader, updating the index's primary key on self-heal. */
|
||||
private async ensureEntry(sessionId: string): Promise<RuntimeEntry> {
|
||||
const existing = this.entries.get(sessionId);
|
||||
if (existing) return existing;
|
||||
if (existing) {
|
||||
if (existing.generation === this.generationOf(existing.projectId, existing.agentId)) {
|
||||
return existing;
|
||||
}
|
||||
// Built before the last vault update: discard once idle and fall through to a
|
||||
// fresh load (resume re-reads the vault). A busy entry is returned as-is — the
|
||||
// in-flight run keeps its values and assertIdle rejects the new Task anyway;
|
||||
// it is rebuilt on the first access after it finishes.
|
||||
if (
|
||||
existing.status !== "idle" ||
|
||||
existing.running !== null ||
|
||||
existing.approvals.size !== 0
|
||||
) {
|
||||
return existing;
|
||||
}
|
||||
this.entries.delete(sessionId);
|
||||
}
|
||||
const row = this.deps.sessions.findById(sessionId);
|
||||
if (!row) {
|
||||
throw new HttpError(
|
||||
@@ -590,6 +643,9 @@ export class SessionManager {
|
||||
"Session does not exist or you do not have access.",
|
||||
);
|
||||
}
|
||||
// Captured before the (awaited) load: a vault update racing with the load leaves
|
||||
// this entry stale, so the access after next rebuilds it with the new values.
|
||||
const generation = this.generationOf(row.projectId, row.agentId);
|
||||
const session = await this.deps.loader.load(row);
|
||||
// The Session/Agent was marked for deletion while loading: discard the load result,
|
||||
// don't rebuild the entry (avoids reviving an orphaned Trace).
|
||||
@@ -613,6 +669,7 @@ export class SessionManager {
|
||||
approvals: new ApprovalRegistry(),
|
||||
abort: null,
|
||||
running: null,
|
||||
generation,
|
||||
lastActivityMs: Date.now(),
|
||||
};
|
||||
this.entries.set(currentId, entry);
|
||||
@@ -630,6 +687,8 @@ export class SessionManager {
|
||||
gen: AsyncGenerator<OmniMessage>,
|
||||
titleSource?: { userExcerpt: string },
|
||||
): Promise<void> {
|
||||
let earlyTitleFired = false;
|
||||
let mainBodyChars = 0;
|
||||
const ctx: UsageContext = {
|
||||
projectId: entry.projectId,
|
||||
agentId: entry.agentId,
|
||||
@@ -715,6 +774,27 @@ export class SessionManager {
|
||||
child.assistantExcerpt += (child.assistantExcerpt ? "\n" : "") + nested.text;
|
||||
}
|
||||
}
|
||||
// Early title trigger (see EARLY_TITLE_BODY_CHARS): fire as soon as enough main-
|
||||
// session body text has streamed; maybeGenerate self-guards (NULL title, single
|
||||
// flight), so the completion trigger in finally stays as the short-answer fallback.
|
||||
if (!earlyTitleFired && titleSource?.userExcerpt.trim()) {
|
||||
const p = msg.payload as { type?: string; role?: string; text?: string };
|
||||
if (
|
||||
(!msg.origin || msg.origin.length === 0) &&
|
||||
msg.type === "model_msg" &&
|
||||
p.type === "text" &&
|
||||
p.role === "assistant" &&
|
||||
p.text
|
||||
) {
|
||||
mainBodyChars += p.text.length;
|
||||
if (mainBodyChars >= EARLY_TITLE_BODY_CHARS) {
|
||||
earlyTitleFired = true;
|
||||
this.deps.titles?.maybeGenerate(ctx, entry.session, {
|
||||
fallbackText: titleSource.userExcerpt,
|
||||
});
|
||||
}
|
||||
}
|
||||
}
|
||||
// Re-fetch the channel before every publish (matches publishEvent): the channel
|
||||
// may have been recycled and recreated during a long wait on approval, and
|
||||
// holding a stale reference would send output to an orphaned, detached channel.
|
||||
|
||||
@@ -14,7 +14,12 @@
|
||||
* - Silent failure (logged): the title stays NULL and naturally retries after the next
|
||||
* Task completes.
|
||||
*/
|
||||
import { emptyTokenCounts, sanitizeTitle, tokenUsage } from "@prismshadow/penguin-core";
|
||||
import {
|
||||
emptyTokenCounts,
|
||||
sanitizeTitle,
|
||||
stripConversationMarkers,
|
||||
tokenUsage,
|
||||
} from "@prismshadow/penguin-core";
|
||||
import type { SessionsRepo } from "../db/repos/sessions.js";
|
||||
import type { ChannelHub } from "./channel.js";
|
||||
import type { ErrorSink } from "./error-recorder.js";
|
||||
@@ -132,7 +137,11 @@ export class TitleGenerator implements TitleNotifier {
|
||||
|
||||
/** Fallback title: take the material's first non-empty line, sanitize and truncate; if sanitizing empties it out (pure punctuation, etc.) fall back to the truncated original text; returns null if all-whitespace. */
|
||||
function fallbackTitle(text: string): string | null {
|
||||
const firstLine = text.split("\n").find((l) => l.trim().length > 0);
|
||||
// Strip machine markers first: a skill invocation prepends a `<use_skills>` block, so the
|
||||
// raw first non-empty line would otherwise be that marker rather than the user's request.
|
||||
const firstLine = stripConversationMarkers(text)
|
||||
.split("\n")
|
||||
.find((l) => l.trim().length > 0);
|
||||
if (!firstLine) return null;
|
||||
// sanitizeTitle strips a pure-punctuation line down to empty — in that case keep the truncated original text, guaranteeing "a title is always obtained".
|
||||
return sanitizeTitle(firstLine) ?? firstLine.trim().slice(0, 30);
|
||||
|
||||
@@ -10,6 +10,8 @@
|
||||
* template's comments).
|
||||
*/
|
||||
import fs from "node:fs/promises";
|
||||
import type { Dirent } from "node:fs";
|
||||
import path from "node:path";
|
||||
import { HttpError } from "../http/errors.js";
|
||||
import {
|
||||
agentDir,
|
||||
@@ -20,6 +22,7 @@ import {
|
||||
isValidId,
|
||||
loadAgentVault,
|
||||
scheduleDir,
|
||||
skillsDir,
|
||||
systemConfigPath,
|
||||
} from "@prismshadow/penguin-core";
|
||||
import type { AgentsRepo } from "../db/repos/agents.js";
|
||||
@@ -41,6 +44,8 @@ export interface AgentListItem {
|
||||
vaultKeyCount: number;
|
||||
/** Number of scheduled tasks (count of .toml files under schedule/, including invalid ones). */
|
||||
scheduleCount: number;
|
||||
/** Number of installed Skills (count of skills/<name>/ directories that contain a SKILL.md). */
|
||||
skillCount: number;
|
||||
}
|
||||
|
||||
export class AgentService {
|
||||
@@ -86,11 +91,12 @@ export class AgentService {
|
||||
);
|
||||
return Promise.all(
|
||||
sorted.map(async (row) => {
|
||||
const [meta, updatedAt, vaultKeyCount, scheduleCount] = await Promise.all([
|
||||
const [meta, updatedAt, vaultKeyCount, scheduleCount, skillCount] = await Promise.all([
|
||||
this.agentConfig.readCardMeta(projectId, row.agentId),
|
||||
this.configUpdatedAt(projectId, row.agentId),
|
||||
this.vaultKeyCount(projectId, row.agentId),
|
||||
this.scheduleCount(projectId, row.agentId),
|
||||
this.skillCount(projectId, row.agentId),
|
||||
]);
|
||||
return {
|
||||
agentId: row.agentId,
|
||||
@@ -99,6 +105,7 @@ export class AgentService {
|
||||
...(updatedAt !== undefined ? { updatedAt } : {}),
|
||||
vaultKeyCount,
|
||||
scheduleCount,
|
||||
skillCount,
|
||||
};
|
||||
}),
|
||||
);
|
||||
@@ -123,6 +130,30 @@ export class AgentService {
|
||||
}
|
||||
}
|
||||
|
||||
/** Number of installed Skills: count of skills/<name>/ directories containing a SKILL.md (0 if the directory doesn't exist). */
|
||||
private async skillCount(projectId: string, agentId: string): Promise<number> {
|
||||
const base = skillsDir(this.root, projectId, agentId);
|
||||
let dirents: Dirent[];
|
||||
try {
|
||||
dirents = await fs.readdir(base, { withFileTypes: true });
|
||||
} catch {
|
||||
return 0;
|
||||
}
|
||||
const present = await Promise.all(
|
||||
dirents
|
||||
.filter((d) => d.isDirectory())
|
||||
.map(async (d) => {
|
||||
try {
|
||||
await fs.access(path.join(base, d.name, "SKILL.md"));
|
||||
return true;
|
||||
} catch {
|
||||
return false;
|
||||
}
|
||||
}),
|
||||
);
|
||||
return present.filter(Boolean).length;
|
||||
}
|
||||
|
||||
/** Last config modification time: the later of system_config.yaml and AGENTS.md mtime; omitted if neither is readable. */
|
||||
private async configUpdatedAt(projectId: string, agentId: string): Promise<string | undefined> {
|
||||
const paths = [
|
||||
@@ -218,6 +249,8 @@ export class AgentService {
|
||||
version: meta.version,
|
||||
vaultKeyCount: 0,
|
||||
scheduleCount: 0,
|
||||
// Read the real count: coreCreateAgent seeds the default skill set for default_agent.
|
||||
skillCount: await this.skillCount(projectId, agentId),
|
||||
};
|
||||
}
|
||||
}
|
||||
|
||||
@@ -27,7 +27,7 @@ import {
|
||||
resolveModelEnv,
|
||||
userText,
|
||||
} from "@prismshadow/penguin-core";
|
||||
import type { ModelRef } from "@prismshadow/penguin-core";
|
||||
import type { LLMOutcome, ModelRef, OmniMessage } from "@prismshadow/penguin-core";
|
||||
import type {
|
||||
ModelInfo,
|
||||
ModelPricingDto,
|
||||
@@ -94,6 +94,26 @@ function showRef(provider: string, modelId: string): string {
|
||||
return `(provider=${provider}, model_id=${modelId})`;
|
||||
}
|
||||
|
||||
/**
|
||||
* Connectivity probe prompt: asks for one word, so the whole exchange fits in a
|
||||
* single-digit output budget. The wording discourages reasoning and the trailing empty
|
||||
* <think></think> makes many reasoning models treat their thinking phase as already
|
||||
* closed - keeping the probe's tiny budget on actual output instead of burning it on
|
||||
* thinking.
|
||||
*/
|
||||
const PROBE_PROMPT =
|
||||
'ping - reply with the single word "pong" and nothing else. Do not think or explain.\n<think></think>';
|
||||
|
||||
/**
|
||||
* Speed probe prompt: a one-word answer can't be timed (a compliant model emits 1-3
|
||||
* tokens, and the window is then dominated by the final usage chunk's round trip), so
|
||||
* speed mode asks for something long enough to run into its raised cap. Counting to 50 is
|
||||
* deterministic and needs no knowledge, so every model produces the same token stream at
|
||||
* its own decoding rate. Carries the same anti-thinking hint as the connectivity prompt.
|
||||
*/
|
||||
const SPEED_PROBE_PROMPT =
|
||||
"Count from 1 to 50 as a comma-separated list, and nothing else. Do not think or explain.\n<think></think>";
|
||||
|
||||
export class ProjectConfigService {
|
||||
constructor(private readonly root: string) {}
|
||||
|
||||
@@ -203,12 +223,19 @@ export class ProjectConfigService {
|
||||
* as a pair in the request body; sends one minimal request using that model's
|
||||
* config (optionally overridden with an unsaved apiKey / baseUrl) — no tools, no
|
||||
* system prompt, thinking disabled, a tiny output cap, 20s timeout — just to see
|
||||
* whether it completes normally. The model id sent to AgentHub is `modelId`
|
||||
* whether the endpoint answers. The model id sent to AgentHub is `modelId`
|
||||
* itself (the upstream id verbatim; client_type inference follows it).
|
||||
*
|
||||
* A reasoning-heavy model may ignore the disabled thinking level and burn the
|
||||
* whole tiny output cap on thinking (finish_reason=length with no text — AgentHub
|
||||
* raises EmptyResponseError, collapsed to a malformed outcome): the endpoint
|
||||
* demonstrably streamed model output, which is everything a connectivity test
|
||||
* proves, so that case counts as ok too (see probeVerdict).
|
||||
*
|
||||
* Never throws: the LLM layer collapses auth/parameter/network errors into an
|
||||
* `LLMOutcome`, which is translated here into ok / message. Consumes very few
|
||||
* Tokens (single-digit output), and writes no Trace and records no usage.
|
||||
* Tokens (single-digit output; speed mode spends up to its 64-token cap to have a
|
||||
* window worth timing), and writes no Trace and records no usage.
|
||||
*/
|
||||
async testModel(projectId: string, req: ModelTestRequest): Promise<ModelTestResponse> {
|
||||
const raw = await this.readRaw(projectId);
|
||||
@@ -236,18 +263,43 @@ export class ProjectConfigService {
|
||||
...(clientType ? { clientType } : {}),
|
||||
tools: [],
|
||||
thinkingLevel: "none",
|
||||
maxTokens: 16,
|
||||
// Speed mode pairs a raised cap with a prompt that keeps generating (see
|
||||
// SPEED_PROBE_PROMPT), so the stream lasts long enough for TTFT/TPS to describe
|
||||
// decoding rather than one round trip; the plain connectivity test keeps the
|
||||
// single-digit-token budget.
|
||||
maxTokens: req.speed ? 64 : 16,
|
||||
requestTimeoutMs: 20_000,
|
||||
});
|
||||
const gen = llm.streamGenerate({ newMessages: [userText("ping")] });
|
||||
const gen = llm.streamGenerate({
|
||||
newMessages: [userText(req.speed ? SPEED_PROBE_PROMPT : PROBE_PROMPT)],
|
||||
});
|
||||
let sawContent = false;
|
||||
let firstContentAt: number | null = null;
|
||||
let outputTokens = 0;
|
||||
for (;;) {
|
||||
const step = await gen.next();
|
||||
if (step.done) {
|
||||
const outcome = step.value;
|
||||
if (outcome.status === "completed")
|
||||
return { ok: true, latencyMs: Date.now() - startedAt };
|
||||
const detail = "message" in outcome && outcome.message ? outcome.message : outcome.status;
|
||||
return { ok: false, message: String(detail).slice(0, 300) };
|
||||
const verdict = probeVerdict(step.value, sawContent);
|
||||
if (!verdict.ok) return verdict;
|
||||
const res: ModelTestResponse = { ok: true, latencyMs: Date.now() - startedAt };
|
||||
if (firstContentAt !== null) {
|
||||
res.ttftMs = firstContentAt - startedAt;
|
||||
// Output rate over the streaming window (first content -> stream end), dropped
|
||||
// when the sample is too small to mean anything (see probeTps): usage is only
|
||||
// reported on completed streams, so thinking-only malformed endings carry TTFT
|
||||
// but no rate, and so does a model that answers in a couple of tokens.
|
||||
const tps = probeTps(outputTokens, Date.now() - firstContentAt);
|
||||
if (tps !== undefined) res.tps = tps;
|
||||
}
|
||||
return res;
|
||||
}
|
||||
if (isProbeContent(step.value)) {
|
||||
sawContent = true;
|
||||
if (firstContentAt === null) firstContentAt = Date.now();
|
||||
}
|
||||
const p = step.value.payload as { type?: string; request?: { output?: number } };
|
||||
if (p.type === "token_usage" && typeof p.request?.output === "number") {
|
||||
outputTokens = p.request.output;
|
||||
}
|
||||
}
|
||||
} catch (err) {
|
||||
@@ -496,3 +548,56 @@ export class ProjectConfigService {
|
||||
return this.getModels(projectId);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Whether a streamed message carries genuine model content (thinking or text, partial delta
|
||||
* or complete backfill) — the probe's "the endpoint really answered" signal. Tool calls and
|
||||
* event messages don't count (the probe declares no tools).
|
||||
*/
|
||||
export function isProbeContent(msg: OmniMessage): boolean {
|
||||
const p = msg.payload as { type?: string; thinking?: string; text?: string };
|
||||
if (p.type === "partial_thinking" || p.type === "thinking") return Boolean(p.thinking);
|
||||
if (p.type === "partial_text" || p.type === "text") return Boolean(p.text);
|
||||
return false;
|
||||
}
|
||||
|
||||
/**
|
||||
* Probe verdict from the terminal LLM outcome. `completed` always passes. A `malformed`
|
||||
* ending after genuine streamed content also passes: the typical case is a reasoning-heavy
|
||||
* model that ignores the disabled thinking level and burns the probe's tiny max_tokens
|
||||
* entirely on thinking (finish_reason=length -> AgentHub's EmptyResponseError) — the
|
||||
* endpoint, credential, and model id all demonstrably work, which is what a connectivity
|
||||
* test measures. Everything else (auth/parameter failures, timeouts, malformed with nothing
|
||||
* received) fails with the outcome's message.
|
||||
*/
|
||||
export function probeVerdict(
|
||||
outcome: LLMOutcome,
|
||||
sawContent: boolean,
|
||||
): { ok: true } | { ok: false; message: string } {
|
||||
if (outcome.status === "completed") return { ok: true };
|
||||
if (outcome.status === "malformed" && sawContent) return { ok: true };
|
||||
const detail = "message" in outcome && outcome.message ? outcome.message : outcome.status;
|
||||
return { ok: false, message: String(detail).slice(0, 300) };
|
||||
}
|
||||
|
||||
/**
|
||||
* Sample floors below which the streaming window says nothing about decoding rate. The
|
||||
* speed probe's 64-token cap clears both by a wide margin, so hitting a floor means the
|
||||
* model didn't really stream (a one-word answer, or usage that never arrived).
|
||||
*/
|
||||
const PROBE_TPS_MIN_TOKENS = 16;
|
||||
const PROBE_TPS_MIN_WINDOW_MS = 100;
|
||||
|
||||
/**
|
||||
* Output rate (tokens/s) over the probe's streaming window (first content -> stream end),
|
||||
* rounded to 1dp; undefined when the sample is too small to be meaningful. A stream's
|
||||
* closing usage chunk costs a round trip on its own, so a two-token answer measures network
|
||||
* jitter and nothing else — 2 tokens in 30ms reads as 66.7 tok/s and the same model 30ms
|
||||
* later reads as 33.3, which the card badges would paint green vs yellow. Callers report
|
||||
* TTFT alone rather than a fabricated rate; a malformed ending carries no usage at all and
|
||||
* lands here as 0 tokens.
|
||||
*/
|
||||
export function probeTps(outputTokens: number, windowMs: number): number | undefined {
|
||||
if (outputTokens < PROBE_TPS_MIN_TOKENS || windowMs <= PROBE_TPS_MIN_WINDOW_MS) return undefined;
|
||||
return Math.round((outputTokens / (windowMs / 1000)) * 10) / 10;
|
||||
}
|
||||
|
||||
@@ -7,9 +7,9 @@
|
||||
* (provider, model_id) / workspace, which is backfilled into a DB row
|
||||
* (approval_mode defaults, createdAt is taken from the timestamp embedded in
|
||||
* session_id).
|
||||
* Create: via core's `agent.createSession` (model reference as a provider + modelId
|
||||
* pair; defaults to the Project's default reference, 400 if there is none; omitting
|
||||
* provider goes through resolveModelRef for unique resolution); the new Session is
|
||||
* Create: via core's `agent.createSession` (the model reference is always a complete
|
||||
* (provider, modelId) pair — both or neither; omitting both falls back to the
|
||||
* Project's default reference, 400 if there is none); the new Session is
|
||||
* added to session-manager's active table (state idle).
|
||||
*/
|
||||
import path from "node:path";
|
||||
@@ -22,6 +22,7 @@ import {
|
||||
} from "@prismshadow/penguin-core";
|
||||
import type { ApprovalMode, SessionInfo } from "../api/types.js";
|
||||
import { HttpError, isMissingCredential, modelCredentialMissing } from "../http/errors.js";
|
||||
import { badRequest } from "../http/validate.js";
|
||||
import type { SessionRow, SessionsRepo } from "../db/repos/sessions.js";
|
||||
import type { SessionManager } from "../runtime/session-manager.js";
|
||||
import type { ProjectConfigService } from "./project-config-service.js";
|
||||
@@ -141,28 +142,38 @@ export class SessionService {
|
||||
}
|
||||
|
||||
/**
|
||||
* Create a Session: model reference `(provider, modelId)` as a pair; defaults to
|
||||
* the Project's default reference (400 prompting to configure a model first if
|
||||
* there is none); when provider is omitted, core's resolveModelRef performs
|
||||
* unique resolution (400 on zero matches / ambiguity). `workspace` is already
|
||||
* Create a Session: the model reference is a complete `(provider, modelId)` pair.
|
||||
* Half a reference is a client error, never something to resolve — the missing half
|
||||
* is never guessed, since a guessed provider would send an entry's credential to a
|
||||
* vendor nobody named. Omitting both falls back to the Project's default reference
|
||||
* (400 prompting to configure a model first if there is none). `workspace` is already
|
||||
* validated by the route guard. The new Session is added to the active table
|
||||
* (idle).
|
||||
*/
|
||||
async createSession(args: {
|
||||
projectId: string;
|
||||
agentId: string;
|
||||
/** Upstream id of the session's model (paired with provider); defaults to the Project's default reference. */
|
||||
/** Upstream id of the session's model; always paired with provider. Omit both for the Project's default reference. */
|
||||
modelId?: string;
|
||||
/** The provider group for `modelId`; if omitted, resolveModelRef performs unique resolution. */
|
||||
/** The provider group for `modelId`; always paired with modelId, never inferred. */
|
||||
provider?: string;
|
||||
workspace?: string;
|
||||
approvalMode?: ApprovalMode;
|
||||
/** Session source marker (schedule when triggered by a scheduled task; defaults to user-created). */
|
||||
source?: "schedule";
|
||||
}): Promise<SessionInfo> {
|
||||
let modelId = args.modelId;
|
||||
let provider = args.provider;
|
||||
if (modelId === undefined) {
|
||||
if ((args.modelId === undefined) !== (args.provider === undefined)) {
|
||||
throw badRequest(
|
||||
"modelId and provider must be given together as a (provider, modelId) pair: specify both, or neither to use the Project's default model.",
|
||||
);
|
||||
}
|
||||
let modelId: string;
|
||||
let provider: string;
|
||||
if (args.modelId !== undefined && args.provider !== undefined) {
|
||||
modelId = args.modelId;
|
||||
provider = args.provider;
|
||||
} else {
|
||||
// The guard above leaves only "both omitted" here: fall back to the Project default.
|
||||
const def = await this.deps.projectConfig.getDefaultModelRef(args.projectId);
|
||||
if (def === undefined) {
|
||||
throw new HttpError(
|
||||
@@ -183,12 +194,12 @@ export class SessionService {
|
||||
try {
|
||||
session = await agent.createSession({
|
||||
modelId,
|
||||
...(provider !== undefined ? { provider } : {}),
|
||||
provider,
|
||||
...(args.workspace !== undefined ? { workspaceDir: args.workspace } : {}),
|
||||
});
|
||||
} catch (err) {
|
||||
// A missing credential is its own category (the frontend shows localized text
|
||||
// by code); other core errors (zero matches / ambiguous reference, Workspace
|
||||
// by code); other core errors (the pair naming no configured entry, Workspace
|
||||
// not existing, etc.) are collapsed to 400 — the guard already blocks most cases.
|
||||
if (isMissingCredential(err)) throw modelCredentialMissing(modelId);
|
||||
throw new HttpError(
|
||||
|
||||
Reference in New Issue
Block a user