Initialize repository with harness code and assets
Initial import of all source code, config, and README assets: the packages workspace (cli, core, server, web, docs, landing, skills), build scripts, tooling config, and CI workflows. Includes the data-layout revision made on this branch: the local data root defaults to ~/.penguin/data (PENGUIN_HOME still overrides; the installer keeps its binaries in ~/.penguin), and every Agent lives under <project>/agents/<agent>/ — path helpers, the three agent-enumeration scans, the system prompt, built-in Skills, tests and docs all follow the new layout. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_018ihk8iQuo3kv2aPjAYEPuR
This commit is contained in:
@@ -0,0 +1,116 @@
|
||||
/**
|
||||
* Tool call approval: ApproveFn factory + pending
|
||||
* approval registry.
|
||||
*
|
||||
* Approval mode semantics match the CLI (packages/cli/src/approval.ts):
|
||||
* allow-all auto-approves; deny-all auto-denies; read-only allows read-only tools
|
||||
* (permission==="r") and routes the rest to manual approval; always-ask routes
|
||||
* everything to manual approval. Routing to manual approval registers a pending entry
|
||||
* and pushes an `approval_request` via SSE, suspending until the frontend decides via
|
||||
* `POST /approvals/:toolCallId`; no timeout — pending approvals are resolved to deny
|
||||
* when the Task is interrupted (then proceeds through the abort flow).
|
||||
*
|
||||
* Every approval decision re-reads the current approval_mode (`getMode` reads the DB),
|
||||
* so mode changes take effect immediately.
|
||||
* Docs: /docs/tools § "Approval".
|
||||
*/
|
||||
import type { ApprovalMode } from "../api/types.js";
|
||||
import type {
|
||||
ApprovalDecision,
|
||||
ApproveFn,
|
||||
OmniMessage,
|
||||
ToolCallPayload,
|
||||
} from "@prismshadow/penguin-core";
|
||||
|
||||
export interface PendingApproval {
|
||||
toolCall: OmniMessage<ToolCallPayload>;
|
||||
origin?: string[];
|
||||
}
|
||||
|
||||
interface PendingEntry extends PendingApproval {
|
||||
resolve: (decision: ApprovalDecision) => void;
|
||||
}
|
||||
|
||||
/** Pending approval registry (key = tool_call_id), one per Session runtime. */
|
||||
export class ApprovalRegistry {
|
||||
private readonly pending = new Map<string, PendingEntry>();
|
||||
|
||||
get size(): number {
|
||||
return this.pending.size;
|
||||
}
|
||||
|
||||
/** All currently pending approvals (for subscription replay). */
|
||||
list(): PendingApproval[] {
|
||||
return [...this.pending.values()].map(({ toolCall, origin }) => ({
|
||||
toolCall,
|
||||
...(origin !== undefined ? { origin } : {}),
|
||||
}));
|
||||
}
|
||||
|
||||
/** Register and wait for a decision. Re-registering the same id (defensive) resolves the old entry as deny. */
|
||||
wait(toolCall: OmniMessage<ToolCallPayload>): Promise<ApprovalDecision> {
|
||||
const id = toolCall.payload.tool_call_id;
|
||||
this.pending.get(id)?.resolve("deny");
|
||||
return new Promise<ApprovalDecision>((resolve) => {
|
||||
const entry: PendingEntry = {
|
||||
toolCall,
|
||||
...(toolCall.origin !== undefined ? { origin: toolCall.origin } : {}),
|
||||
resolve: (decision) => {
|
||||
this.pending.delete(id);
|
||||
resolve(decision);
|
||||
},
|
||||
};
|
||||
this.pending.set(id, entry);
|
||||
});
|
||||
}
|
||||
|
||||
/** Submit a decision; returns false if not found (already decided/unknown). */
|
||||
decide(toolCallId: string, decision: ApprovalDecision): boolean {
|
||||
const entry = this.pending.get(toolCallId);
|
||||
if (!entry) return false;
|
||||
entry.resolve(decision);
|
||||
return true;
|
||||
}
|
||||
|
||||
/** Interruption convergence: resolve all pending approvals as deny. */
|
||||
denyAll(): void {
|
||||
for (const entry of [...this.pending.values()]) entry.resolve("deny");
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Build the approve callback: re-reads the approval mode on every call; when routed to
|
||||
* manual approval, registers a pending entry and suspends after pushing an
|
||||
* `approval_request` server event via `publishRequest`.
|
||||
*/
|
||||
export function makeApprove(args: {
|
||||
getMode: () => ApprovalMode;
|
||||
toolPermission: (name: string) => "r" | "rw" | undefined;
|
||||
registry: ApprovalRegistry;
|
||||
publishRequest: (pending: PendingApproval) => void;
|
||||
}): ApproveFn {
|
||||
const { getMode, toolPermission, registry, publishRequest } = args;
|
||||
const manual = (toolCall: OmniMessage<ToolCallPayload>): Promise<ApprovalDecision> => {
|
||||
const promise = registry.wait(toolCall);
|
||||
publishRequest({
|
||||
toolCall,
|
||||
...(toolCall.origin !== undefined ? { origin: toolCall.origin } : {}),
|
||||
});
|
||||
return promise;
|
||||
};
|
||||
return async (toolCall) => {
|
||||
switch (getMode()) {
|
||||
case "allow-all":
|
||||
return "allow";
|
||||
case "deny-all":
|
||||
return "deny";
|
||||
case "read-only":
|
||||
// Auto-approve read-only tools; route read-write/unknown tools to manual approval (matches CLI semantics).
|
||||
if (toolPermission(toolCall.payload.name) === "r") return "allow";
|
||||
return manual(toolCall);
|
||||
case "always-ask":
|
||||
default:
|
||||
return manual(toolCall);
|
||||
}
|
||||
};
|
||||
}
|
||||
@@ -0,0 +1,194 @@
|
||||
/**
|
||||
* SSE event channel.
|
||||
*
|
||||
* The Session channel and the user channel share this implementation:
|
||||
* - Event id is an opaque string `<epoch>-<seq>`: epoch is a random short string
|
||||
* generated when each Channel instance is created, seq is a monotonically increasing
|
||||
* integer within the channel. epoch necessarily changes when the channel is
|
||||
* 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,
|
||||
* 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 /
|
||||
* resync / hello): it consumes a seq number but doesn't enter the buffer or get
|
||||
* broadcast — if that subscriber later reconnects with this id, the hit check is
|
||||
* still safe (seq is monotonic).
|
||||
*
|
||||
* This module only handles event numbering / buffering / dispatch, not HTTP — SSE
|
||||
* output is adapted at the routing layer.
|
||||
* Docs: /docs/server-api § "Delivery Guarantees".
|
||||
*/
|
||||
import { randomUUID } from "node:crypto";
|
||||
|
||||
/** A numbered channel event; `id` is `<epoch>-<seq>`, `data` is serialized single-line JSON. */
|
||||
export interface ChannelEvent {
|
||||
id: string;
|
||||
/** SSE event name; omitted (OmniMessage) means no `event:` line. */
|
||||
event?: string;
|
||||
data: string;
|
||||
}
|
||||
|
||||
export type ChannelListener = (evt: ChannelEvent) => void;
|
||||
|
||||
export interface ChannelOptions {
|
||||
maxBufferCount?: number;
|
||||
maxBufferBytes?: number;
|
||||
}
|
||||
|
||||
const DEFAULT_MAX_COUNT = 1000;
|
||||
const DEFAULT_MAX_BYTES = 2 * 1024 * 1024;
|
||||
|
||||
/** Buffered entry: seq is stored separately so hit checks never need to parse the string id. */
|
||||
interface BufferedEvent {
|
||||
seq: number;
|
||||
evt: ChannelEvent;
|
||||
}
|
||||
|
||||
export class Channel {
|
||||
/** Channel epoch: generated at instance creation, prefixed onto event ids (necessarily changes after recycle/recreate or restart). */
|
||||
readonly epoch: string = randomUUID().slice(0, 8);
|
||||
private nextSeq = 1;
|
||||
private buffer: BufferedEvent[] = [];
|
||||
private bufferBytes = 0;
|
||||
/** Max seq among evicted events (0 means never evicted): lower bound for hit checks. */
|
||||
private lastEvictedSeq = 0;
|
||||
private readonly listeners = new Set<ChannelListener>();
|
||||
private readonly maxCount: number;
|
||||
private readonly maxBytes: number;
|
||||
/** Timestamp of last activity (publish/subscription change), used for idle-reclaim checks. */
|
||||
lastActivityMs = Date.now();
|
||||
|
||||
constructor(opts: ChannelOptions = {}) {
|
||||
this.maxCount = opts.maxBufferCount ?? DEFAULT_MAX_COUNT;
|
||||
this.maxBytes = opts.maxBufferBytes ?? DEFAULT_MAX_BYTES;
|
||||
}
|
||||
|
||||
get subscriberCount(): number {
|
||||
return this.listeners.size;
|
||||
}
|
||||
|
||||
/** Broadcast an event: number it, buffer it (evicting the oldest), notify all subscribers. */
|
||||
publish(data: unknown, event?: string): ChannelEvent {
|
||||
const entry = this.makeEvent(data, event);
|
||||
this.buffer.push(entry);
|
||||
this.bufferBytes += entry.evt.data.length;
|
||||
while (
|
||||
this.buffer.length > 0 &&
|
||||
(this.buffer.length > this.maxCount || this.bufferBytes > this.maxBytes)
|
||||
) {
|
||||
const evicted = this.buffer.shift()!;
|
||||
this.bufferBytes -= evicted.evt.data.length;
|
||||
this.lastEvictedSeq = Math.max(this.lastEvictedSeq, evicted.seq);
|
||||
}
|
||||
for (const listener of this.listeners) listener(entry.evt);
|
||||
return entry.evt;
|
||||
}
|
||||
|
||||
/** Unicast an event to a single subscriber: consumes a seq but doesn't buffer or broadcast (used for replay at subscribe time). */
|
||||
sendTo(listener: ChannelListener, data: unknown, event?: string): ChannelEvent {
|
||||
const entry = this.makeEvent(data, event);
|
||||
listener(entry.evt);
|
||||
return entry.evt;
|
||||
}
|
||||
|
||||
subscribe(listener: ChannelListener): () => void {
|
||||
this.listeners.add(listener);
|
||||
this.lastActivityMs = Date.now();
|
||||
return () => {
|
||||
this.listeners.delete(listener);
|
||||
this.lastActivityMs = Date.now();
|
||||
};
|
||||
}
|
||||
|
||||
/**
|
||||
* Compute replay from a Last-Event-ID (`<epoch>-<seq>`): a mismatched epoch (channel
|
||||
* recycled/recreated, process restarted, or malformed id) always misses; a matching
|
||||
* epoch hits the buffer (if no events after that seq have been evicted and the seq was
|
||||
* indeed assigned by this channel) and returns the buffered events after it; otherwise
|
||||
* miss (the caller should send `resync_required` first).
|
||||
*/
|
||||
replayAfter(lastEventId: string): { hit: boolean; events: ChannelEvent[] } {
|
||||
const sep = lastEventId.lastIndexOf("-");
|
||||
if (sep <= 0) return { hit: false, events: [] };
|
||||
const epoch = lastEventId.slice(0, sep);
|
||||
const seq = Number.parseInt(lastEventId.slice(sep + 1), 10);
|
||||
if (epoch !== this.epoch || !Number.isInteger(seq) || seq < 0) {
|
||||
return { hit: false, events: [] };
|
||||
}
|
||||
const hit = seq >= this.lastEvictedSeq && seq < this.nextSeq;
|
||||
if (!hit) return { hit: false, events: [] };
|
||||
return { hit: true, events: this.buffer.filter((e) => e.seq > seq).map((e) => e.evt) };
|
||||
}
|
||||
|
||||
private makeEvent(data: unknown, event?: string): BufferedEvent {
|
||||
this.lastActivityMs = Date.now();
|
||||
const serialized = typeof data === "string" ? data : JSON.stringify(data);
|
||||
const seq = this.nextSeq++;
|
||||
const evt: ChannelEvent = { id: `${this.epoch}-${seq}`, data: serialized };
|
||||
if (event !== undefined) evt.event = event;
|
||||
return { seq, evt };
|
||||
}
|
||||
}
|
||||
|
||||
const DEFAULT_IDLE_MS = 30 * 60 * 1000;
|
||||
const SWEEP_INTERVAL_MS = 60 * 1000;
|
||||
|
||||
export interface ChannelHubOptions {
|
||||
idleMs?: number;
|
||||
/**
|
||||
* Active check: keys for which this returns true are excluded from idle reclaim (app
|
||||
* assembly injects `manager.statusOf(key) !== "idle"`, so a running/compacting
|
||||
* Session channel is never reclaimed no matter how long since its last publish; a
|
||||
* user channel key looks like `user:<id>` and is always considered active).
|
||||
*/
|
||||
isActive?: (key: string) => boolean;
|
||||
}
|
||||
|
||||
/**
|
||||
* Channel collection: lazily created by key (Session id or `user:<user_id>`);
|
||||
* a channel whose Session is idle and has had no subscribers for over 30 minutes is
|
||||
* reclaimed, releasing its buffer as well.
|
||||
*/
|
||||
export class ChannelHub {
|
||||
private readonly channels = new Map<string, Channel>();
|
||||
private readonly timer: NodeJS.Timeout;
|
||||
private readonly idleMs: number;
|
||||
private readonly isActive: (key: string) => boolean;
|
||||
|
||||
constructor(opts: ChannelHubOptions = {}) {
|
||||
this.idleMs = opts.idleMs ?? DEFAULT_IDLE_MS;
|
||||
this.isActive = opts.isActive ?? (() => false);
|
||||
this.timer = setInterval(() => this.sweep(), SWEEP_INTERVAL_MS);
|
||||
this.timer.unref?.();
|
||||
}
|
||||
|
||||
get(key: string): Channel {
|
||||
let ch = this.channels.get(key);
|
||||
if (!ch) {
|
||||
ch = new Channel();
|
||||
this.channels.set(key, ch);
|
||||
}
|
||||
return ch;
|
||||
}
|
||||
|
||||
peek(key: string): Channel | undefined {
|
||||
return this.channels.get(key);
|
||||
}
|
||||
|
||||
/** Reclaim idle channels (skips active Sessions: no reclaim even without a publish while awaiting approval); `now` is injectable for tests. */
|
||||
sweep(now: number = Date.now()): void {
|
||||
for (const [key, ch] of this.channels) {
|
||||
if (this.isActive(key)) continue;
|
||||
if (ch.subscriberCount === 0 && now - ch.lastActivityMs > this.idleMs) {
|
||||
this.channels.delete(key);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
dispose(): void {
|
||||
clearInterval(this.timer);
|
||||
this.channels.clear();
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,161 @@
|
||||
/**
|
||||
* Error persistence: errors caught on the server are all
|
||||
* written to error_records through here, for display on the stats dashboard. Shape
|
||||
* mirrors usage-recorder — persist only raw facts, leave aggregation to query time.
|
||||
*
|
||||
* **The classification (kind) criterion is "does a human need to step in"**, not where
|
||||
* the error originated:
|
||||
*
|
||||
* - `expected`: anticipated by the system, has a defined handling path, part of normal
|
||||
* operation, no human needed — HTTP business errors (`HttpError`, mostly 4xx); LLM
|
||||
* `timeout` / `malformed` (the engine already reconnects and retries); tool execution
|
||||
* `failed` / `timeout` (the error is fed back to the model, and the Agent adjusts on
|
||||
* its own).
|
||||
* - `unexpected`: shouldn't happen, usually a bug or a config/environment fault,
|
||||
* **needs a human** — internal errors converged to 500; process crashes; runtime
|
||||
* errors escaping from background tasks (Session drive / usage persistence / title
|
||||
* generation / subagent registration); LLM `failed` (not retryable: auth failure,
|
||||
* invalid params, etc.).
|
||||
* - User-initiated actions **are not errors** and are never recorded: request/tool
|
||||
* `aborted` (user clicked "stop", or denied a tool).
|
||||
*
|
||||
* Determination: HTTP sources are inferred automatically from `HttpError` (preserving
|
||||
* existing behavior); other sources must pass `kind` explicitly at the capture site.
|
||||
* The frontend highlights unexpected by default; expected is still recorded without
|
||||
* losing information.
|
||||
*
|
||||
* Sources cover HTTP, Session drive, LLM requests, Environment (tool execution), usage
|
||||
* persistence, title generation, subagent registration, and process-level fallback;
|
||||
* among these, `llm` / `environment` errors are not expressed via throw (core converges
|
||||
* them into the message stream instead), and are fished out by stream-error-watcher from
|
||||
* the Session output stream.
|
||||
*
|
||||
* **This recorder never throws**: it's hooked onto app.onError, and throwing from
|
||||
* within it would turn error handling into infinite recursion; if persistence itself
|
||||
* fails (disk full / DB already closed, etc.), it's fine to drop that one record.
|
||||
*
|
||||
* **Short-window dedup (DEDUP_WINDOW_MS)**: error storms are the norm — someone scanning
|
||||
* the API produces a wall of 404s, or a tool fails repeatedly in a loop. Persisting each
|
||||
* one both write-amplifies and floods the table, and makes the dashboard's "most recent
|
||||
* 20" all the same error. So the same `(source, code, Project)` is persisted at most
|
||||
* once per window; repeats within the window are **dropped outright** (not persisted);
|
||||
* only an actual persist refreshes the timestamp, so a sustained storm leaves a steady
|
||||
* one record per window instead of being suppressed indefinitely.
|
||||
* **Tradeoff**: aggregate counts therefore **underestimate** — a storm of the same error
|
||||
* only counts once, check the logs for true frequency; in exchange, a single error storm
|
||||
* doesn't drown out error_records or the stats dashboard. The second line of defense is
|
||||
* ErrorsRepo's capacity cap. The dedup table (lastSeen) must stay bounded: past
|
||||
* DEDUP_KEYS_MAX, expired entries are cleared first, and if still over the limit the
|
||||
* whole table is cleared — better to miss some dedup than let it grow unbounded across
|
||||
* different error codes.
|
||||
*/
|
||||
import { formatLocalDate } from "../internal/dates.js";
|
||||
import { HttpError } from "../http/errors.js";
|
||||
import type { ErrorsRepo } from "../db/repos/errors.js";
|
||||
|
||||
/** Capture-site source (maps one-to-one to error_records.source). */
|
||||
export type ErrorSource =
|
||||
| "http"
|
||||
| "session"
|
||||
| "llm"
|
||||
| "environment"
|
||||
| "usage"
|
||||
| "title"
|
||||
| "subagent"
|
||||
| "process"
|
||||
| "schedule";
|
||||
|
||||
/** Error classification: see file header — the criterion is "does a human need to step in". */
|
||||
export type ErrorKind = "expected" | "unexpected";
|
||||
|
||||
/** Attribution context (all optional: the login endpoint has no Project, and process-level fallback has no request at all). */
|
||||
export interface ErrorContext {
|
||||
projectId?: string;
|
||||
agentId?: string;
|
||||
sessionId?: string;
|
||||
}
|
||||
|
||||
export interface ErrorRecordArgs {
|
||||
source: ErrorSource;
|
||||
/** The caught error (unknown: the value caught may not be an Error; failures from the message stream pass the reason text directly). */
|
||||
err: unknown;
|
||||
ctx?: ErrorContext;
|
||||
/** Semantic code (required for non-HTTP sources, e.g. session_run_failed); defaults to HttpError.code. */
|
||||
code?: string;
|
||||
/** HTTP status code; leave empty for non-HTTP sources. */
|
||||
status?: number;
|
||||
/** Explicit classification (see file header); defaults to inferring from `HttpError` — HTTP sources rely on this, other sources should pass it explicitly. */
|
||||
kind?: ErrorKind;
|
||||
}
|
||||
|
||||
/** Message truncation length (keep only a readable summary; the full stack is still logged). */
|
||||
export const MESSAGE_MAX = 500;
|
||||
|
||||
/** Short-window dedup window: the same (source, code, Project) is persisted at most once per window (see the file header's tradeoff). */
|
||||
export const DEDUP_WINDOW_MS = 2000;
|
||||
|
||||
/** Cap on dedup table keys (bounded; over the limit, expired entries are cleared first, and if still over, the whole table is cleared). */
|
||||
export const DEDUP_KEYS_MAX = 1000;
|
||||
|
||||
function messageOf(err: unknown): string {
|
||||
const raw = err instanceof Error ? err.message : String(err);
|
||||
return raw.length > MESSAGE_MAX ? raw.slice(0, MESSAGE_MAX) : raw;
|
||||
}
|
||||
|
||||
export class ErrorRecorder {
|
||||
/** Dedup table: `source \0 code \0 projectId` → timestamp of the last **persist** (see file header). */
|
||||
private readonly lastSeen = new Map<string, number>();
|
||||
|
||||
constructor(
|
||||
private readonly errors: ErrorsRepo,
|
||||
private readonly now: () => Date = () => new Date(),
|
||||
) {}
|
||||
|
||||
/** Record an error (synchronous, fails silently; same-window duplicates are dropped outright, see file header). */
|
||||
record(args: ErrorRecordArgs): void {
|
||||
try {
|
||||
const http = args.err instanceof HttpError ? args.err : null;
|
||||
const now = this.now();
|
||||
const projectId = args.ctx?.projectId ?? null;
|
||||
const code = args.code ?? http?.code ?? "internal";
|
||||
// Short-window dedup: coarse-grained to "same kind of error for the same Project"; repeats within the window aren't persisted.
|
||||
if (this.deduped(`${args.source}\0${code}\0${projectId ?? ""}`, now.getTime())) return;
|
||||
this.errors.insert({
|
||||
ts: now.toISOString(),
|
||||
date: formatLocalDate(now),
|
||||
projectId,
|
||||
agentId: args.ctx?.agentId ?? null,
|
||||
sessionId: args.ctx?.sessionId ?? null,
|
||||
source: args.source,
|
||||
// Explicit classification takes priority; otherwise infer from HttpError (business error = expected, else unexpected).
|
||||
kind: args.kind ?? (http ? "expected" : "unexpected"),
|
||||
code,
|
||||
// Unexpected errors from HTTP sources are converged to 500 externally (matches handleError's response).
|
||||
status: args.status ?? http?.status ?? (args.source === "http" ? 500 : null),
|
||||
message: messageOf(args.err),
|
||||
});
|
||||
} catch {
|
||||
// See file header: if the recorder itself errors, dropping this one record is the only option — never rethrow.
|
||||
}
|
||||
}
|
||||
|
||||
/** true if a same-kind error was already recorded within the window (drop it); otherwise register this persist timestamp and keep the dedup table bounded. */
|
||||
private deduped(key: string, nowMs: number): boolean {
|
||||
const last = this.lastSeen.get(key);
|
||||
if (last !== undefined && nowMs - last < DEDUP_WINDOW_MS) return true;
|
||||
this.lastSeen.set(key, nowMs);
|
||||
if (this.lastSeen.size > DEDUP_KEYS_MAX) this.evict(nowMs);
|
||||
return false;
|
||||
}
|
||||
|
||||
/** Keep the dedup table bounded (see file header): clear expired entries first; if still over the limit (hundreds/thousands of distinct error codes erupting at once), clear it entirely. */
|
||||
private evict(nowMs: number): void {
|
||||
for (const [key, at] of this.lastSeen) {
|
||||
if (nowMs - at >= DEDUP_WINDOW_MS) this.lastSeen.delete(key);
|
||||
}
|
||||
if (this.lastSeen.size > DEDUP_KEYS_MAX) this.lastSeen.clear();
|
||||
}
|
||||
}
|
||||
|
||||
/** Minimal dependency a capture site needs on the recorder (tests inject a fake; structurally matches SessionManager's UsageRecorderLike). */
|
||||
export type ErrorSink = Pick<ErrorRecorder, "record">;
|
||||
@@ -0,0 +1,214 @@
|
||||
/**
|
||||
* Schedule file parsing, validation, and trigger-time computation.
|
||||
*
|
||||
* `agent_state/schedule/<name>.toml` is declarative intent; the system never writes it
|
||||
* back. This module does pure parsing and pure time math only: an invalid file returns
|
||||
* an error (the scheduler skips it and records the error); runtime state (fired /
|
||||
* missed / disabled) doesn't live here — it belongs to SQLite (db/repos/schedules.ts).
|
||||
* Docs: /docs/configuration § "Schedules".
|
||||
*/
|
||||
import { parse as parseToml } from "smol-toml";
|
||||
|
||||
/** `period` lower bound: below 5 minutes is treated as an invalid file (guards against runaway high-frequency tasks). */
|
||||
export const MIN_PERIOD_MS = 5 * 60_000;
|
||||
|
||||
/** A parsed schedule definition (the filename minus `.toml` is its identity). */
|
||||
export interface ScheduleDefinition {
|
||||
name: string;
|
||||
/** The Prompt to send (required). */
|
||||
prompt: string;
|
||||
/** Enabled switch; disabled by default. */
|
||||
enabled: boolean;
|
||||
/** Original text of the first trigger time (for API echo, preserving the written form). */
|
||||
startAt: string;
|
||||
/** First trigger time (epoch ms). */
|
||||
startAtMs: number;
|
||||
/** Original text of the end time. */
|
||||
endAt?: string;
|
||||
/** Original text of the trigger period (e.g. `30m`, for API echo); undefined means a one-shot task. */
|
||||
period?: string;
|
||||
/** Trigger period (ms); undefined means a one-shot task. */
|
||||
periodMs?: number;
|
||||
/** End time (epoch ms); no more triggers once past it. */
|
||||
endAtMs?: number;
|
||||
/** The target Session to bind to; defaults to creating a new Session each time. */
|
||||
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). */
|
||||
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).
|
||||
*/
|
||||
provider?: string;
|
||||
}
|
||||
|
||||
export type ScheduleParseResult =
|
||||
{ ok: true; def: ScheduleDefinition } | { ok: false; error: string };
|
||||
|
||||
/** Parse a fixed interval in `30m` / `12h` / `7d` form; returns null if invalid. */
|
||||
export function parsePeriod(raw: string): number | null {
|
||||
const m = /^(\d+)([mhd])$/.exec(raw.trim());
|
||||
if (!m) return null;
|
||||
const n = Number(m[1]);
|
||||
if (!Number.isInteger(n) || n <= 0) return null;
|
||||
const unit = m[2] === "m" ? 60_000 : m[2] === "h" ? 3_600_000 : 86_400_000;
|
||||
return n * unit;
|
||||
}
|
||||
|
||||
/** Parse an ISO 8601 instant into epoch ms plus the original text for echo; returns null if invalid (smol-toml's date values are also accepted). */
|
||||
function parseInstant(value: unknown): { ms: number; raw: string } | null {
|
||||
if (value instanceof Date) {
|
||||
const ms = value.getTime();
|
||||
return Number.isNaN(ms) ? null : { ms, raw: value.toISOString() };
|
||||
}
|
||||
if (typeof value !== "string") return null;
|
||||
const ms = Date.parse(value);
|
||||
return Number.isNaN(ms) ? null : { ms, raw: value };
|
||||
}
|
||||
|
||||
/**
|
||||
* Parse and validate a schedule file. A field with the wrong type invalidates the whole
|
||||
* file (the baseline for hand-edit tolerance is to never let bad config reach the
|
||||
* scheduler); unknown keys are ignored (forward compatibility).
|
||||
*/
|
||||
export function parseScheduleFile(name: string, raw: string): ScheduleParseResult {
|
||||
let parsed: unknown;
|
||||
try {
|
||||
parsed = parseToml(raw);
|
||||
} catch (err) {
|
||||
return {
|
||||
ok: false,
|
||||
error: `TOML 解析失败:${err instanceof Error ? err.message : String(err)}`,
|
||||
};
|
||||
}
|
||||
if (parsed === null || typeof parsed !== "object")
|
||||
return { ok: false, error: "内容不是 TOML 表" };
|
||||
const t = parsed as Record<string, unknown>;
|
||||
|
||||
const prompt = t["prompt"];
|
||||
if (typeof prompt !== "string" || prompt.trim() === "") {
|
||||
return { ok: false, error: "缺少必填字段 prompt" };
|
||||
}
|
||||
const enabled = t["enabled"] === undefined ? false : t["enabled"];
|
||||
if (typeof enabled !== "boolean") return { ok: false, error: "enabled 必须是布尔值" };
|
||||
|
||||
const startAt = parseInstant(t["start_at"]);
|
||||
if (startAt === null) return { ok: false, error: "start_at 缺失或不是合法的 ISO 8601 时刻" };
|
||||
|
||||
let period: string | undefined;
|
||||
let periodMs: number | undefined;
|
||||
if (t["period"] !== undefined) {
|
||||
if (typeof t["period"] !== "string") return { ok: false, error: "period 必须是字符串" };
|
||||
const ms = parsePeriod(t["period"]);
|
||||
if (ms === null) return { ok: false, error: "period 必须形如 30m / 12h / 7d" };
|
||||
if (ms < MIN_PERIOD_MS) return { ok: false, error: "period 低于下限 5m" };
|
||||
period = t["period"].trim();
|
||||
periodMs = ms;
|
||||
}
|
||||
|
||||
let endAt: { ms: number; raw: string } | undefined;
|
||||
if (t["end_at"] !== undefined) {
|
||||
const parsedEnd = parseInstant(t["end_at"]);
|
||||
if (parsedEnd === null) return { ok: false, error: "end_at 不是合法的 ISO 8601 时刻" };
|
||||
if (parsedEnd.ms <= startAt.ms) return { ok: false, error: "end_at 必须晚于 start_at" };
|
||||
endAt = parsedEnd;
|
||||
}
|
||||
|
||||
let sessionId: string | undefined;
|
||||
if (t["session_id"] !== undefined) {
|
||||
if (typeof t["session_id"] !== "string" || t["session_id"] === "") {
|
||||
return { ok: false, error: "session_id 必须是非空字符串" };
|
||||
}
|
||||
sessionId = t["session_id"];
|
||||
}
|
||||
let workspace: string | undefined;
|
||||
if (t["workspace"] !== undefined) {
|
||||
if (typeof t["workspace"] !== "string" || t["workspace"] === "") {
|
||||
return { ok: false, error: "workspace 必须是非空字符串" };
|
||||
}
|
||||
workspace = t["workspace"];
|
||||
}
|
||||
let modelId: string | undefined;
|
||||
if (t["model_id"] !== undefined) {
|
||||
if (typeof t["model_id"] !== "string" || t["model_id"] === "") {
|
||||
return { ok: false, error: "model_id 必须是非空字符串" };
|
||||
}
|
||||
modelId = t["model_id"];
|
||||
}
|
||||
let provider: string | undefined;
|
||||
if (t["provider"] !== undefined) {
|
||||
if (typeof t["provider"] !== "string" || t["provider"] === "") {
|
||||
return { ok: false, error: "provider 必须是非空字符串" };
|
||||
}
|
||||
provider = t["provider"];
|
||||
}
|
||||
if (provider !== undefined && modelId === undefined) {
|
||||
return { ok: false, error: "provider 仅与 model_id 成对使用(模型引用须成对给出)" };
|
||||
}
|
||||
if (
|
||||
sessionId !== undefined &&
|
||||
(workspace !== undefined || modelId !== undefined || provider !== undefined)
|
||||
) {
|
||||
return {
|
||||
ok: false,
|
||||
error: "目标二选一:workspace 与 provider / model_id 仅用于新建 Session 模式",
|
||||
};
|
||||
}
|
||||
|
||||
return {
|
||||
ok: true,
|
||||
def: {
|
||||
name,
|
||||
prompt,
|
||||
enabled,
|
||||
startAt: startAt.raw,
|
||||
startAtMs: startAt.ms,
|
||||
...(period !== undefined ? { period } : {}),
|
||||
...(periodMs !== undefined ? { periodMs } : {}),
|
||||
...(endAt !== undefined ? { endAt: endAt.raw, endAtMs: endAt.ms } : {}),
|
||||
...(sessionId !== undefined ? { sessionId } : {}),
|
||||
...(workspace !== undefined ? { workspace } : {}),
|
||||
...(modelId !== undefined ? { modelId } : {}),
|
||||
...(provider !== undefined ? { provider } : {}),
|
||||
},
|
||||
};
|
||||
}
|
||||
|
||||
/**
|
||||
* Step from `start_at` by `period` and return the most recent scheduled time not later
|
||||
* than `nowMs`; null if `start_at` hasn't been reached yet. A one-shot task's only slot
|
||||
* is `start_at` itself.
|
||||
*/
|
||||
export function latestSlotAt(def: ScheduleDefinition, nowMs: number): number | null {
|
||||
if (nowMs < def.startAtMs) return null;
|
||||
if (def.periodMs === undefined) return def.startAtMs;
|
||||
const k = Math.floor((nowMs - def.startAtMs) / def.periodMs);
|
||||
return def.startAtMs + k * def.periodMs;
|
||||
}
|
||||
|
||||
/** Whether a scheduled slot still falls within the `[start_at, end_at]` window (always true if there's no end_at). */
|
||||
export function slotInWindow(def: ScheduleDefinition, slotMs: number): boolean {
|
||||
return def.endAtMs === undefined || slotMs <= def.endAtMs;
|
||||
}
|
||||
|
||||
/**
|
||||
* The next scheduled time strictly after `nowMs` (used to display "next trigger");
|
||||
* for a one-shot task this only has a value while start_at hasn't been reached, and
|
||||
* returns null once past end_at.
|
||||
*/
|
||||
export function nextSlotAfter(def: ScheduleDefinition, nowMs: number): number | null {
|
||||
let next: number;
|
||||
if (nowMs < def.startAtMs) {
|
||||
next = def.startAtMs;
|
||||
} else if (def.periodMs === undefined) {
|
||||
return null;
|
||||
} else {
|
||||
const k = Math.floor((nowMs - def.startAtMs) / def.periodMs) + 1;
|
||||
next = def.startAtMs + k * def.periodMs;
|
||||
}
|
||||
return slotInWindow(def, next) ? next : null;
|
||||
}
|
||||
@@ -0,0 +1,140 @@
|
||||
/**
|
||||
* Schedule file access: `agent_state/schedule/<name>.toml`, where
|
||||
* the filename (a semantic name) is the identity. Reads are fault-tolerant (an invalid
|
||||
* file is recorded as an error and skipped by the caller); writes only go through the
|
||||
* API routes (the system never rewrites existing file content — PUT is a full-file
|
||||
* replacement expressing user intent).
|
||||
*/
|
||||
import fs from "node:fs/promises";
|
||||
import path from "node:path";
|
||||
import { stringify as stringifyToml } from "smol-toml";
|
||||
import { loadProjectConfig, resolveModelRef, scheduleDir } from "@prismshadow/penguin-core";
|
||||
import type { ScheduleDefinition } from "./schedule-file.js";
|
||||
import { parseScheduleFile, type ScheduleParseResult } from "./schedule-file.js";
|
||||
|
||||
export interface ScheduleFileEntry {
|
||||
name: string;
|
||||
raw: string;
|
||||
parsed: ScheduleParseResult;
|
||||
}
|
||||
|
||||
/** List all schedule files for this Agent (a missing directory is treated as empty). */
|
||||
export async function listScheduleFiles(
|
||||
root: string,
|
||||
projectId: string,
|
||||
agentId: string,
|
||||
): Promise<ScheduleFileEntry[]> {
|
||||
const dir = scheduleDir(root, projectId, agentId);
|
||||
let names: string[];
|
||||
try {
|
||||
names = await fs.readdir(dir);
|
||||
} catch {
|
||||
return [];
|
||||
}
|
||||
const entries: ScheduleFileEntry[] = [];
|
||||
for (const file of names.sort()) {
|
||||
if (!file.endsWith(".toml")) continue;
|
||||
const name = file.slice(0, -".toml".length);
|
||||
let raw: string;
|
||||
try {
|
||||
raw = await fs.readFile(path.join(dir, file), "utf8");
|
||||
} catch {
|
||||
continue; // Deleted during reconciliation: revisit next round.
|
||||
}
|
||||
entries.push({ name, raw, parsed: parseScheduleFile(name, raw) });
|
||||
}
|
||||
return entries;
|
||||
}
|
||||
|
||||
export async function readScheduleFile(
|
||||
root: string,
|
||||
projectId: string,
|
||||
agentId: string,
|
||||
name: string,
|
||||
): Promise<ScheduleFileEntry | null> {
|
||||
const file = path.join(scheduleDir(root, projectId, agentId), `${name}.toml`);
|
||||
try {
|
||||
const raw = await fs.readFile(file, "utf8");
|
||||
return { name, raw, parsed: parseScheduleFile(name, raw) };
|
||||
} catch {
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
||||
/** Serialize API fields into file content (validation uniformly goes through parseScheduleFile, avoiding two sets of rules). */
|
||||
export function serializeSchedule(fields: {
|
||||
prompt: string;
|
||||
enabled: boolean;
|
||||
startAt: string;
|
||||
period?: string;
|
||||
endAt?: string;
|
||||
sessionId?: string;
|
||||
workspace?: string;
|
||||
modelId?: string;
|
||||
provider?: string;
|
||||
}): string {
|
||||
const table: Record<string, unknown> = {
|
||||
prompt: fields.prompt,
|
||||
enabled: fields.enabled,
|
||||
start_at: fields.startAt,
|
||||
...(fields.period !== undefined ? { period: fields.period } : {}),
|
||||
...(fields.endAt !== undefined ? { end_at: fields.endAt } : {}),
|
||||
...(fields.sessionId !== undefined ? { session_id: fields.sessionId } : {}),
|
||||
...(fields.workspace !== undefined ? { workspace: fields.workspace } : {}),
|
||||
...(fields.provider !== undefined ? { provider: fields.provider } : {}),
|
||||
...(fields.modelId !== undefined ? { model_id: fields.modelId } : {}),
|
||||
};
|
||||
return `${stringifyToml(table)}\n`;
|
||||
}
|
||||
|
||||
/**
|
||||
* 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).
|
||||
*/
|
||||
export async function validateScheduleModelRef(
|
||||
root: string,
|
||||
projectId: string,
|
||||
def: Pick<ScheduleDefinition, "modelId" | "provider">,
|
||||
): Promise<string | null> {
|
||||
if (def.modelId === undefined) return null;
|
||||
try {
|
||||
const cfg = await loadProjectConfig(root, projectId);
|
||||
resolveModelRef(cfg, def.modelId, def.provider);
|
||||
return null;
|
||||
} catch (err) {
|
||||
return err instanceof Error ? err.message : String(err);
|
||||
}
|
||||
}
|
||||
|
||||
/** Write a schedule file to disk (full-file replacement for POST/PUT). */
|
||||
export async function writeScheduleFile(
|
||||
root: string,
|
||||
projectId: string,
|
||||
agentId: string,
|
||||
name: string,
|
||||
raw: string,
|
||||
): Promise<void> {
|
||||
const dir = scheduleDir(root, projectId, agentId);
|
||||
await fs.mkdir(dir, { recursive: true });
|
||||
await fs.writeFile(path.join(dir, `${name}.toml`), raw, "utf8");
|
||||
}
|
||||
|
||||
/** Delete a schedule file; returns false if it doesn't exist. */
|
||||
export async function deleteScheduleFile(
|
||||
root: string,
|
||||
projectId: string,
|
||||
agentId: string,
|
||||
name: string,
|
||||
): Promise<boolean> {
|
||||
const file = path.join(scheduleDir(root, projectId, agentId), `${name}.toml`);
|
||||
try {
|
||||
await fs.unlink(file);
|
||||
return true;
|
||||
} catch {
|
||||
return false;
|
||||
}
|
||||
}
|
||||
Binary file not shown.
@@ -0,0 +1,841 @@
|
||||
/**
|
||||
* Active Session runtime.
|
||||
*
|
||||
* Responsibilities:
|
||||
* - get-or-resume-or-heal: use it directly on an active-table hit; with a Trace,
|
||||
* recover via `agent.resumeSession`; a stale Session that was created but never run
|
||||
* and survived a process restart (no Trace) **self-heals** — recreated via
|
||||
* 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;
|
||||
* - 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
|
||||
* message to the SSE channel and handing it to usage-recorder for persistence;
|
||||
* on completion (including errors) resets to idle and pushes a `task_state` server
|
||||
* event;
|
||||
* - Approval registration and interrupt convergence: each approval decision re-reads
|
||||
* approval_mode from the DB (takes effect immediately); an interrupt first
|
||||
* converges pending approvals to deny, then aborts.
|
||||
*
|
||||
* The underlying implementation of get-or-resume-or-heal is injected via
|
||||
* `SessionLoader`: production uses the core SDK (createCoreSessionLoader), tests inject
|
||||
* a fake Session (issuing no real LLM requests).
|
||||
*/
|
||||
import fs from "node:fs/promises";
|
||||
import path from "node:path";
|
||||
import {
|
||||
createAgent,
|
||||
findLatestTraceFile,
|
||||
isSessionMeta,
|
||||
tracesDir,
|
||||
} from "@prismshadow/penguin-core";
|
||||
import type {
|
||||
ApproveFn,
|
||||
CompactAvailability,
|
||||
OmniMessage,
|
||||
SessionMetaPayload,
|
||||
SessionTitleResult,
|
||||
TextPayload,
|
||||
} from "@prismshadow/penguin-core";
|
||||
import type { ServerEvent, SessionStatus } from "../api/types.js";
|
||||
import { HttpError, isMissingCredential, modelCredentialMissing } from "../http/errors.js";
|
||||
import type { SessionRow, SessionsRepo } from "../db/repos/sessions.js";
|
||||
import { ApprovalRegistry, makeApprove } from "./approvals.js";
|
||||
import type { PendingApproval } from "./approvals.js";
|
||||
import type { ChannelHub } from "./channel.js";
|
||||
import type { ErrorSink } from "./error-recorder.js";
|
||||
import { StreamErrorWatcher } from "./stream-error-watcher.js";
|
||||
import type { TitleNotifier } from "./title-generator.js";
|
||||
import type { UsageContext } from "./usage-recorder.js";
|
||||
|
||||
/** 409 for when there's nothing to compact: give the specific reason rather than a one-size-fits-none message. */
|
||||
function compactUnavailable(why: Exclude<CompactAvailability, "ok">): HttpError {
|
||||
const messages: Record<typeof why, string> = {
|
||||
unsupported: "该 Agent 未配置上下文压缩能力。",
|
||||
empty: "当前上下文没有可压缩的内容(尚无已完成的对话轮次)。",
|
||||
just_compacted: "上下文刚压缩过,此后还没有新的对话,无需再次压缩。",
|
||||
};
|
||||
return new HttpError(409, "nothing_to_compact", messages[why]);
|
||||
}
|
||||
|
||||
/** Minimal interface for a runtime Session (satisfied by core Session; tests may inject a fake implementation). */
|
||||
export interface RuntimeSession {
|
||||
readonly sessionId: string;
|
||||
run(
|
||||
newMessages: OmniMessage[],
|
||||
opts: { approve: ApproveFn; signal: AbortSignal },
|
||||
): AsyncGenerator<OmniMessage>;
|
||||
compact(opts: { signal: AbortSignal }): AsyncGenerator<OmniMessage>;
|
||||
/** Whether compaction is possible and why; when not ok, compact() yields no messages (see core ContextEngine.compactability). */
|
||||
compactability(): CompactAvailability;
|
||||
toolPermission(name: string): "r" | "rw" | undefined;
|
||||
/**
|
||||
* Out-of-band one-shot request for title generation (core `Session.generateTitle`,
|
||||
* writes no history/Trace). Material defaults to what the Session collects itself
|
||||
* (the first Task's text gathered during run); `material` overrides this for
|
||||
* subagents.
|
||||
*/
|
||||
generateTitle(args?: {
|
||||
material?: { userText: string; assistantText: string };
|
||||
signal?: AbortSignal;
|
||||
}): Promise<SessionTitleResult>;
|
||||
}
|
||||
|
||||
/** The underlying loader behind get-or-resume-or-heal. */
|
||||
export interface SessionLoader {
|
||||
/**
|
||||
* Load a runtime Session from an index row: recover (with a Trace) or self-heal
|
||||
* rebuild (no Trace, session_id will change). Throws HttpError(409) for unrecoverable
|
||||
* cases such as a missing Workspace.
|
||||
*/
|
||||
load(row: SessionRow): Promise<RuntimeSession>;
|
||||
}
|
||||
|
||||
/** Production loader: the core SDK's resumeSession / createSession. */
|
||||
export function createCoreSessionLoader(root: string): SessionLoader {
|
||||
return {
|
||||
async load(row: SessionRow): Promise<RuntimeSession> {
|
||||
const agent = await createAgent({
|
||||
root,
|
||||
projectId: row.projectId,
|
||||
agentId: row.agentId,
|
||||
});
|
||||
const located = await findLatestTraceFile(
|
||||
tracesDir(root, row.projectId, row.agentId),
|
||||
row.sessionId,
|
||||
);
|
||||
if (located) {
|
||||
// With a Trace: rebuild via "Session Recovery" (history injected via setHistory,
|
||||
// carrying over any residual state).
|
||||
// core's recognizable recovery failures (Workspace deleted / Model removed from
|
||||
// config / Trace missing session_meta, etc.) are converged to 409, preserving
|
||||
// the original message rather than bubbling up as 500.
|
||||
try {
|
||||
return await agent.resumeSession({ sessionId: row.sessionId });
|
||||
} catch (err) {
|
||||
// The credential key was deleted after the Session was created: only caught
|
||||
// here at recovery time; give the same actionable message.
|
||||
if (isMissingCredential(err)) throw modelCredentialMissing(row.modelId);
|
||||
throw toUnrecoverableError(err);
|
||||
}
|
||||
}
|
||||
// No Trace (created but never run, and the process has restarted since): self-heal
|
||||
// rebuild. A missing Workspace → 409.
|
||||
try {
|
||||
const stat = await fs.stat(row.workspace);
|
||||
if (!stat.isDirectory()) throw new Error("not a directory");
|
||||
} catch {
|
||||
throw new HttpError(
|
||||
409,
|
||||
"workspace_missing",
|
||||
`该 Session 的 Workspace 已不存在:${row.workspace},无法继续。请新建 Session。`,
|
||||
);
|
||||
}
|
||||
try {
|
||||
return await agent.createSession({
|
||||
workspaceDir: row.workspace,
|
||||
modelId: row.modelId,
|
||||
provider: row.provider,
|
||||
});
|
||||
} catch (err) {
|
||||
if (isMissingCredential(err)) throw modelCredentialMissing(row.modelId);
|
||||
throw toUnrecoverableError(err);
|
||||
}
|
||||
},
|
||||
};
|
||||
}
|
||||
|
||||
/** A plain Error thrown by core recovery/self-heal rebuild → 409 (preserving the original, actionable message). */
|
||||
function toUnrecoverableError(err: unknown): HttpError {
|
||||
if (err instanceof HttpError) return err;
|
||||
return new HttpError(
|
||||
409,
|
||||
"session_unrecoverable",
|
||||
err instanceof Error ? err.message : String(err),
|
||||
);
|
||||
}
|
||||
|
||||
export interface UsageRecorderLike {
|
||||
record(ctx: UsageContext, msg: OmniMessage): Promise<void>;
|
||||
}
|
||||
|
||||
export interface SessionManagerDeps {
|
||||
sessions: SessionsRepo;
|
||||
channels: ChannelHub;
|
||||
loader: SessionLoader;
|
||||
recorder: UsageRecorderLike;
|
||||
/** Automatic Session title generation (optional: not injected in tests or when disabled). */
|
||||
titles?: TitleNotifier;
|
||||
/** Error persistence (optional: without it, only logs — same as before this was wired up). */
|
||||
errors?: ErrorSink;
|
||||
log?: (line: string) => void;
|
||||
}
|
||||
|
||||
/** Active-table entry: a loaded runtime Session plus its running state. */
|
||||
interface RuntimeEntry {
|
||||
sessionId: string;
|
||||
projectId: string;
|
||||
agentId: string;
|
||||
/** Vendor grouping for the Session's model (paired with modelId to form a model reference). */
|
||||
provider: string;
|
||||
modelId: string;
|
||||
session: RuntimeSession;
|
||||
status: SessionStatus;
|
||||
approvals: ApprovalRegistry;
|
||||
abort: AbortController | null;
|
||||
/** The in-flight drive Promise (awaited during graceful shutdown). */
|
||||
running: Promise<void> | null;
|
||||
/** Timestamp of last activity (refreshed on load / status flip / drive completion), used for idle-eviction checks. */
|
||||
lastActivityMs: number;
|
||||
}
|
||||
|
||||
/** Active-table idle eviction: same convention as the SSE channel (an idle entry with no activity for 30 minutes releases its memory). */
|
||||
const ENTRY_IDLE_MS = 30 * 60 * 1000;
|
||||
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;
|
||||
|
||||
/** Composite Agent key (used as a Set key, avoiding projectId/agentId concatenation ambiguity). */
|
||||
function agentKey(projectId: string, agentId: string): string {
|
||||
return `${projectId}\0${agentId}`;
|
||||
}
|
||||
|
||||
/** If msg is a run_subagent tool call carrying a `prompt`, return its id and prompt (for use as the subagent's title); otherwise null. */
|
||||
function runSubagentCall(msg: OmniMessage): { toolCallId: string; prompt: string } | null {
|
||||
const p = msg.payload as {
|
||||
type?: string;
|
||||
name?: string;
|
||||
arguments?: string;
|
||||
tool_call_id?: string;
|
||||
};
|
||||
if (msg.type !== "model_msg" || p.type !== "tool_call" || p.name !== "run_subagent") return null;
|
||||
if (typeof p.arguments !== "string" || typeof p.tool_call_id !== "string") return null;
|
||||
try {
|
||||
const args = JSON.parse(p.arguments) as { prompt?: unknown };
|
||||
if (typeof args.prompt !== "string" || !args.prompt.trim()) return null;
|
||||
return { toolCallId: p.tool_call_id, prompt: args.prompt };
|
||||
} catch {
|
||||
return null; // Arguments were truncated/malformed: this call is doomed, no subagent will result
|
||||
}
|
||||
}
|
||||
|
||||
/** The denied tool_call_id (approval_decision with decision ≠ allow); otherwise null. */
|
||||
function deniedToolCallId(msg: OmniMessage): string | null {
|
||||
const p = msg.payload as { type?: string; decision?: string; tool_call_id?: string };
|
||||
if (msg.type !== "event_msg" || p.type !== "approval_decision") return null;
|
||||
if (p.decision === "allow" || typeof p.tool_call_id !== "string") return null;
|
||||
return p.tool_call_id;
|
||||
}
|
||||
|
||||
/** The tool_call_id of a parent-level tool call that has settled (a complete tool_call_output); otherwise null. */
|
||||
function settledToolCallId(msg: OmniMessage): string | null {
|
||||
const p = msg.payload as { type?: string; tool_call_id?: string };
|
||||
if (msg.type !== "model_msg" || p.type !== "tool_call_output") return null;
|
||||
return typeof p.tool_call_id === "string" ? p.tool_call_id : null;
|
||||
}
|
||||
|
||||
/** A subagent registered during this run, plus its title material. */
|
||||
interface ChildSession {
|
||||
sessionId: string;
|
||||
agentId: string;
|
||||
modelId: string;
|
||||
/** The prompt of the run_subagent call that spawned it (user material for title generation, and the fallback title). */
|
||||
prompt: string;
|
||||
/** The model text the subagent itself produced (assistant material for title generation). */
|
||||
assistantExcerpt: string;
|
||||
}
|
||||
|
||||
/** Predicate for a plain-text message on the main session (no origin): title material is drawn only from user/model text. */
|
||||
function isPlainText(role: "user" | "assistant") {
|
||||
return (msg: OmniMessage): msg is OmniMessage<TextPayload> => {
|
||||
const payload = msg.payload as { type?: string; role?: string };
|
||||
return (
|
||||
msg.type === "model_msg" &&
|
||||
payload.type === "text" &&
|
||||
payload.role === role &&
|
||||
(!msg.origin || msg.origin.length === 0)
|
||||
);
|
||||
};
|
||||
}
|
||||
|
||||
/** For a nested message, the owning Session (end of the origin chain) and text of the model reply; null if it isn't model text. */
|
||||
function nestedAssistantText(msg: OmniMessage): { sessionId: string; text: string } | null {
|
||||
const p = msg.payload as { type?: string; role?: string; text?: string };
|
||||
if (msg.type !== "model_msg" || p.type !== "text" || p.role !== "assistant") return null;
|
||||
if (!msg.origin || msg.origin.length === 0 || typeof p.text !== "string") return null;
|
||||
return { sessionId: msg.origin[msg.origin.length - 1]!, text: p.text };
|
||||
}
|
||||
|
||||
export class SessionManager {
|
||||
private readonly entries = new Map<string, RuntimeEntry>();
|
||||
/** Per-Session mutex (serializes get-or-load and status flips); auto-cleaned once the chain drains. */
|
||||
private readonly locks = new Map<string, Promise<unknown>>();
|
||||
private readonly log: (line: string) => void;
|
||||
/** Graceful-shutdown flag: once set, new Tasks/compactions are rejected (503). */
|
||||
private closed = false;
|
||||
/** Agents currently being deleted (key = agentKey): new Tasks/compactions are always rejected with 409 during this window. */
|
||||
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>();
|
||||
private readonly sweepTimer: NodeJS.Timeout;
|
||||
|
||||
constructor(private readonly deps: SessionManagerDeps) {
|
||||
this.log = deps.log ?? ((line) => console.error(line));
|
||||
this.sweepTimer = setInterval(() => this.sweepIdle(), ENTRY_SWEEP_INTERVAL_MS);
|
||||
this.sweepTimer.unref?.();
|
||||
}
|
||||
|
||||
// —— Query surface (used by Session listing / Agent active-count / SSE subscription replay) ——
|
||||
|
||||
statusOf(sessionId: string): SessionStatus {
|
||||
return this.entries.get(sessionId)?.status ?? "idle";
|
||||
}
|
||||
|
||||
pendingApprovalCount(sessionId: string): number {
|
||||
return this.entries.get(sessionId)?.approvals.size ?? 0;
|
||||
}
|
||||
|
||||
pendingApprovals(sessionId: string): PendingApproval[] {
|
||||
return this.entries.get(sessionId)?.approvals.list() ?? [];
|
||||
}
|
||||
|
||||
/** Number of Sessions for this Agent that are currently running / compacting. */
|
||||
activeCountForAgent(projectId: string, agentId: string): number {
|
||||
let n = 0;
|
||||
for (const e of this.entries.values()) {
|
||||
if (e.projectId === projectId && e.agentId === agentId && e.status !== "idle") n++;
|
||||
}
|
||||
return n;
|
||||
}
|
||||
|
||||
/** Add a newly created Session to the active table (status idle), avoiding a redundant load on the next Task. */
|
||||
adopt(row: SessionRow, session: RuntimeSession): void {
|
||||
this.entries.set(row.sessionId, {
|
||||
sessionId: row.sessionId,
|
||||
projectId: row.projectId,
|
||||
agentId: row.agentId,
|
||||
provider: row.provider,
|
||||
modelId: row.modelId,
|
||||
session,
|
||||
status: "idle",
|
||||
approvals: new ApprovalRegistry(),
|
||||
abort: null,
|
||||
running: null,
|
||||
lastActivityMs: Date.now(),
|
||||
});
|
||||
}
|
||||
|
||||
// —— Task / compaction drive ——
|
||||
|
||||
/**
|
||||
* Start a Task: get-or-load → 409
|
||||
* mutual-exclusion check → publish the input messages first → drive run in the
|
||||
* background. Returns the current actual session_id (the new id after self-heal).
|
||||
*/
|
||||
async startTask(sessionId: string, input: OmniMessage[]): Promise<{ sessionId: string }> {
|
||||
return this.withLock(sessionId, async () => {
|
||||
this.assertOpen();
|
||||
this.assertAgentNotDeleting(sessionId);
|
||||
this.assertSessionNotDeleting(sessionId);
|
||||
const entry = await this.ensureEntry(sessionId);
|
||||
this.assertIdle(entry);
|
||||
const channel = this.deps.channels.get(entry.sessionId);
|
||||
const ac = new AbortController();
|
||||
entry.status = "running";
|
||||
entry.abort = ac;
|
||||
entry.lastActivityMs = Date.now();
|
||||
// Publish the input messages first (visible to other subscribers; the Trace is
|
||||
// persisted by the SDK), then flip the running status.
|
||||
for (const msg of input) channel.publish(msg);
|
||||
this.publishState(entry, "running");
|
||||
|
||||
const approve = makeApprove({
|
||||
// Re-reads approval_mode from the DB on every decision (a PATCH takes effect immediately).
|
||||
getMode: () => this.deps.sessions.findById(entry.sessionId)?.approvalMode ?? "always-ask",
|
||||
toolPermission: (name) => entry.session.toolPermission(name),
|
||||
registry: entry.approvals,
|
||||
publishRequest: (pending) =>
|
||||
this.publishEvent(entry, {
|
||||
type: "approval_request",
|
||||
toolCall: pending.toolCall,
|
||||
...(pending.origin !== undefined ? { origin: pending.origin } : {}),
|
||||
}),
|
||||
});
|
||||
const gen = entry.session.run(input, { approve, signal: ac.signal });
|
||||
// Title material is collected by the core Session itself during run; here we only
|
||||
// keep this call's input user text, used both as the "material present → attempt
|
||||
// generation" criterion and as the fallback title source if the LLM call fails.
|
||||
const userExcerpt = input
|
||||
.filter(isPlainText("user"))
|
||||
.map((m) => m.payload.text)
|
||||
.join("\n");
|
||||
entry.running = this.drive(entry, gen, { userExcerpt });
|
||||
return { sessionId: entry.sessionId };
|
||||
});
|
||||
}
|
||||
|
||||
/** Manually compact the context: 409 if already running; compaction output also flows into the SSE channel. */
|
||||
async startCompact(sessionId: string): Promise<{ sessionId: string }> {
|
||||
return this.withLock(sessionId, async () => {
|
||||
this.assertOpen();
|
||||
this.assertAgentNotDeleting(sessionId);
|
||||
this.assertSessionNotDeleting(sessionId);
|
||||
const entry = await this.ensureEntry(sessionId);
|
||||
this.assertIdle(entry);
|
||||
// When there's nothing to compact, core's compact() yields no messages at all: we
|
||||
// can't just return 202 and walk away, or the frontend would wait forever for a
|
||||
// compaction banner that never comes (this is exactly the "/compact does nothing
|
||||
// after an interrupt" complaint). Reject explicitly, and **say why** clearly —
|
||||
// "just compacted" and "haven't talked yet" share the same internal state
|
||||
// (sessionTurns === 0), but are two completely different messages to the user:
|
||||
// telling someone who just compacted that there's "no completed conversation turn
|
||||
// yet" tells them nothing.
|
||||
const why = entry.session.compactability();
|
||||
if (why !== "ok") throw compactUnavailable(why);
|
||||
const ac = new AbortController();
|
||||
entry.status = "compacting";
|
||||
entry.abort = ac;
|
||||
entry.lastActivityMs = Date.now();
|
||||
this.publishState(entry, "compacting");
|
||||
const gen = entry.session.compact({ signal: ac.signal });
|
||||
entry.running = this.drive(entry, gen);
|
||||
return { sessionId: entry.sessionId };
|
||||
});
|
||||
}
|
||||
|
||||
/** Submit an approval decision; returns false if the pending approval doesn't exist (already decided/unknown). */
|
||||
decideApproval(sessionId: string, toolCallId: string, decision: "allow" | "deny"): boolean {
|
||||
const entry = this.entries.get(sessionId);
|
||||
if (!entry) return false;
|
||||
return entry.approvals.decide(toolCallId, decision);
|
||||
}
|
||||
|
||||
/**
|
||||
* Interrupt the current Task/compaction: pending approvals converge to deny first,
|
||||
* then the AbortSignal fires. Returns false if nothing is in progress (the route
|
||||
* treats this as a 204 no-op).
|
||||
*/
|
||||
abortTask(sessionId: string): boolean {
|
||||
const entry = this.entries.get(sessionId);
|
||||
if (!entry || !entry.abort) return false;
|
||||
entry.approvals.denyAll();
|
||||
entry.abort.abort();
|
||||
return true;
|
||||
}
|
||||
|
||||
/**
|
||||
* Before deleting a Project, converge all its active runs and clear them out of the
|
||||
* active table. Returns the in-flight drive Promises of the affected entries: the
|
||||
* caller (deleteProject) should await them before removing the directory, so that
|
||||
* interrupt-cleanup Trace writes don't recreate the directory after deletion.
|
||||
*/
|
||||
abortProject(projectId: string): Promise<void>[] {
|
||||
const runnings: Promise<void>[] = [];
|
||||
for (const [key, entry] of [...this.entries]) {
|
||||
if (entry.projectId !== projectId) continue;
|
||||
entry.approvals.denyAll();
|
||||
entry.abort?.abort();
|
||||
if (entry.running) runnings.push(entry.running);
|
||||
this.entries.delete(key);
|
||||
}
|
||||
return runnings;
|
||||
}
|
||||
|
||||
/**
|
||||
* Before deleting an Agent, converge all its active runs and clear them out of the
|
||||
* active table (same semantics as abortProject). Also marks this Agent as "being
|
||||
* deleted": new Tasks/compactions entering during the deletion process are always
|
||||
* rejected with 409 (assertAgentNotDeleting), closing the race window where a new
|
||||
* task recreates the directory and revives an already-deleted Agent between the
|
||||
* abortAgent snapshot and the directory removal. The caller must call
|
||||
* endAgentDeletion once deletion finishes (success or failure).
|
||||
*/
|
||||
beginAgentDeletion(projectId: string, agentId: string): Promise<void>[] {
|
||||
this.deletingAgents.add(agentKey(projectId, agentId));
|
||||
const runnings: Promise<void>[] = [];
|
||||
for (const [key, entry] of [...this.entries]) {
|
||||
if (entry.projectId !== projectId || entry.agentId !== agentId) continue;
|
||||
entry.approvals.denyAll();
|
||||
entry.abort?.abort();
|
||||
if (entry.running) runnings.push(entry.running);
|
||||
this.entries.delete(key);
|
||||
}
|
||||
return runnings;
|
||||
}
|
||||
|
||||
endAgentDeletion(projectId: string, agentId: string): void {
|
||||
this.deletingAgents.delete(agentKey(projectId, agentId));
|
||||
}
|
||||
|
||||
/**
|
||||
* Before deleting a single Session, converge its active run and clear it out of the
|
||||
* active table (same semantics as beginAgentDeletion). Also marks this Session as
|
||||
* "being deleted": new Tasks/compactions entering during the deletion process are
|
||||
* always rejected with 409 (assertSessionNotDeleting), closing the race window where
|
||||
* a new task recreates the entry and Trace file, reviving an already-deleted Session
|
||||
* between the abort snapshot and the file removal. The caller must call
|
||||
* endSessionDeletion once deletion finishes (success or failure). Returns the
|
||||
* in-flight drive Promise: the caller should await it before deleting the Trace file,
|
||||
* so cleanup writes don't recreate the file.
|
||||
*/
|
||||
beginSessionDeletion(sessionId: string): Promise<void>[] {
|
||||
this.deletingSessions.add(sessionId);
|
||||
const entry = this.entries.get(sessionId);
|
||||
if (!entry) return [];
|
||||
entry.approvals.denyAll();
|
||||
entry.abort?.abort();
|
||||
this.entries.delete(sessionId);
|
||||
return entry.running ? [entry.running] : [];
|
||||
}
|
||||
|
||||
endSessionDeletion(sessionId: string): void {
|
||||
this.deletingSessions.delete(sessionId);
|
||||
}
|
||||
|
||||
/** Graceful shutdown: reject new tasks (503), interrupt all active runs, and wait for them to finish (default ≤5s). */
|
||||
async shutdown(timeoutMs = 5000): Promise<void> {
|
||||
this.closed = true;
|
||||
clearInterval(this.sweepTimer);
|
||||
const pending: Promise<void>[] = [];
|
||||
for (const entry of this.entries.values()) {
|
||||
if (!entry.abort) continue;
|
||||
entry.approvals.denyAll();
|
||||
entry.abort.abort();
|
||||
if (entry.running) pending.push(entry.running);
|
||||
}
|
||||
if (pending.length === 0) return;
|
||||
await Promise.race([
|
||||
Promise.allSettled(pending).then(() => undefined),
|
||||
new Promise<void>((resolve) => setTimeout(resolve, timeoutMs).unref?.()),
|
||||
]);
|
||||
}
|
||||
|
||||
/**
|
||||
* Active-table idle eviction: removes entries that are idle (idle status, no pending
|
||||
* approvals, no in-flight drive) and have been inactive past the timeout, releasing
|
||||
* the core Session's full in-memory history. This is purely memory reclamation: the
|
||||
* next access re-resumes via the loader, so correctness is unaffected. Lock-table
|
||||
* entries are auto-cleaned by withLock once their chain drains (including leftover
|
||||
* entries under the old id after self-heal). `now` / `idleMs` are injectable for
|
||||
* tests and timers.
|
||||
*/
|
||||
sweepIdle(now: number = Date.now(), idleMs: number = ENTRY_IDLE_MS): void {
|
||||
for (const [key, entry] of this.entries) {
|
||||
if (entry.status !== "idle" || entry.approvals.size !== 0 || entry.running !== null) continue;
|
||||
if (now - entry.lastActivityMs <= idleMs) continue;
|
||||
this.entries.delete(key);
|
||||
}
|
||||
}
|
||||
|
||||
// —— Internal ——
|
||||
|
||||
private assertOpen(): void {
|
||||
if (this.closed) {
|
||||
throw new HttpError(503, "shutting_down", "服务端正在关停,暂不接受新任务。");
|
||||
}
|
||||
}
|
||||
|
||||
/** The Agent owning this Session is being deleted → 409 (guards against directory recreation inside the deletion race window). */
|
||||
private assertAgentNotDeleting(sessionId: string): void {
|
||||
const row = this.deps.sessions.findById(sessionId);
|
||||
if (row && this.deletingAgents.has(agentKey(row.projectId, row.agentId))) {
|
||||
throw new HttpError(409, "agent_deleting", "该 Agent 正在删除,暂不接受新任务。");
|
||||
}
|
||||
}
|
||||
|
||||
/** This Session is being deleted → 409 (guards against the entry/Trace being rebuilt and reviving it inside the deletion race window). */
|
||||
private assertSessionNotDeleting(sessionId: string): void {
|
||||
if (this.deletingSessions.has(sessionId)) {
|
||||
throw new HttpError(409, "session_deleting", "该 Session 正在删除,暂不接受新任务。");
|
||||
}
|
||||
}
|
||||
|
||||
private assertIdle(entry: RuntimeEntry): void {
|
||||
if (entry.status === "running") {
|
||||
throw new HttpError(409, "task_in_progress", "该 Session 已有进行中的 Task。");
|
||||
}
|
||||
if (entry.status === "compacting") {
|
||||
throw new HttpError(409, "compacting", "该 Session 正在压缩上下文,暂不接受新输入。");
|
||||
}
|
||||
}
|
||||
|
||||
/** 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;
|
||||
const row = this.deps.sessions.findById(sessionId);
|
||||
if (!row) {
|
||||
throw new HttpError(404, "session_not_found", "Session 不存在或无权访问。");
|
||||
}
|
||||
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).
|
||||
this.assertSessionNotDeleting(row.sessionId);
|
||||
this.assertAgentNotDeleting(row.sessionId);
|
||||
let currentId = row.sessionId;
|
||||
if (session.sessionId !== row.sessionId) {
|
||||
// Self-heal produced a new session_id: update the index's primary key; the SSE
|
||||
// channel and pending state are naturally empty for it.
|
||||
this.deps.sessions.replaceId(row.sessionId, session.sessionId);
|
||||
currentId = session.sessionId;
|
||||
}
|
||||
const entry: RuntimeEntry = {
|
||||
sessionId: currentId,
|
||||
projectId: row.projectId,
|
||||
agentId: row.agentId,
|
||||
provider: row.provider,
|
||||
modelId: row.modelId,
|
||||
session,
|
||||
status: "idle",
|
||||
approvals: new ApprovalRegistry(),
|
||||
abort: null,
|
||||
running: null,
|
||||
lastActivityMs: Date.now(),
|
||||
};
|
||||
this.entries.set(currentId, entry);
|
||||
return entry;
|
||||
}
|
||||
|
||||
/**
|
||||
* Drive the output stream in the background: publish each message + persist usage +
|
||||
* persist LLM/tool errors; on completion (including errors) resets to idle and pushes
|
||||
* the status. `titleSource` is passed only for Task runs (compaction doesn't generate
|
||||
* a title): it collects model text for automatic title generation.
|
||||
*/
|
||||
private async drive(
|
||||
entry: RuntimeEntry,
|
||||
gen: AsyncGenerator<OmniMessage>,
|
||||
titleSource?: { userExcerpt: string },
|
||||
): Promise<void> {
|
||||
const ctx: UsageContext = {
|
||||
projectId: entry.projectId,
|
||||
agentId: entry.agentId,
|
||||
sessionId: entry.sessionId,
|
||||
provider: entry.provider,
|
||||
modelId: entry.modelId,
|
||||
};
|
||||
// LLM request failures and tool execution failures aren't expressed via throw (core
|
||||
// converges them into the message stream), so the try/catch below can't catch them:
|
||||
// the watcher inspects messages one by one and fishes them out for persistence
|
||||
// (subagent failures flow through this same stream too; see stream-error-watcher).
|
||||
const watcher = this.deps.errors
|
||||
? new StreamErrorWatcher(this.deps.errors, {
|
||||
projectId: entry.projectId,
|
||||
agentId: entry.agentId,
|
||||
sessionId: entry.sessionId,
|
||||
})
|
||||
: null;
|
||||
// Subagent (origin) registration: as soon as session_meta arrives, the child Session
|
||||
// is persisted so it appears immediately in the sidebar (the frontend picks it up
|
||||
// when it refreshes the list at task completion). The title material is "the prompt
|
||||
// of the run_subagent call that spawned this subagent" — the subagent's user input
|
||||
// is never replayed onto the parent stream (ContextEngine writes the Trace but never
|
||||
// yields it), so we can't rely on the subagent's first user message; instead we use
|
||||
// the run_subagent tool_call arguments immediately preceding it on the parent stream
|
||||
// (depth limited to 1, spawned in order, so taking the most recent one suffices).
|
||||
/** Subagents registered during this run (keyed by session id); titles are generated for each on completion. */
|
||||
const children = new Map<string, ChildSession>();
|
||||
// Unclaimed run_subagent prompts, queued in call order: a single round may spawn
|
||||
// multiple subagents in parallel, and a subagent's session_meta only carries the
|
||||
// session id (no tool_call_id), so pairing can only be approximated via FIFO (when
|
||||
// spawned in parallel and session_meta arrives out of order, two subagents' titles
|
||||
// may end up swapped — this only affects the displayed title). A call that will
|
||||
// never produce a subagent must be dequeued, or its prompt would be mismatched onto
|
||||
// the next subagent: this covers denied calls (approval_decision ≠ allow), and calls
|
||||
// that were approved but failed before spawning the subagent (e.g. agent_id doesn't
|
||||
// exist) — the latter is cleaned up when the parent-level tool_call_output settles;
|
||||
// if the call is still in the queue at that point, it never produced a session_meta.
|
||||
const subagentPrompts = new Map<string, string>();
|
||||
try {
|
||||
for await (const msg of gen) {
|
||||
// A parent-level (no origin) run_subagent call: record its prompt for the child
|
||||
// session_meta that arrives later to use as its title.
|
||||
if (!msg.origin || msg.origin.length === 0) {
|
||||
const call = runSubagentCall(msg);
|
||||
if (call) subagentPrompts.set(call.toolCallId, call.prompt);
|
||||
const denied = deniedToolCallId(msg);
|
||||
if (denied) subagentPrompts.delete(denied);
|
||||
const settled = settledToolCallId(msg);
|
||||
if (settled) subagentPrompts.delete(settled);
|
||||
} else if (isSessionMeta(msg)) {
|
||||
// Subagent registration is only a "side effect" — it must never interrupt the
|
||||
// main run flow on error: wrap the whole thing in a defensive try/catch.
|
||||
try {
|
||||
const child = this.registerChildSession(entry, msg, children);
|
||||
// Only a **direct** subagent (origin length 1) claims a queued parent-level
|
||||
// run_subagent prompt; deeper sessions are spawned by their own parent and
|
||||
// shouldn't consume from this queue.
|
||||
if (child && msg.origin!.length === 1) {
|
||||
const [pendingId] = subagentPrompts.keys();
|
||||
if (pendingId !== undefined) {
|
||||
child.prompt = subagentPrompts.get(pendingId) ?? "";
|
||||
subagentPrompts.delete(pendingId); // Consumed by this session_meta
|
||||
}
|
||||
}
|
||||
} catch (err) {
|
||||
this.log(
|
||||
`[subagent] 子会话登记失败: ${err instanceof Error ? err.message : String(err)}`,
|
||||
);
|
||||
this.deps.errors?.record({
|
||||
source: "subagent",
|
||||
err,
|
||||
ctx,
|
||||
code: "subagent_register_failed",
|
||||
});
|
||||
}
|
||||
} else {
|
||||
// A subagent's model text: its title is generated from the subagent's **own
|
||||
// conversation**, so the material is accumulated here.
|
||||
const nested = nestedAssistantText(msg);
|
||||
const child = nested ? children.get(nested.sessionId) : undefined;
|
||||
if (nested && child && child.assistantExcerpt.length < TITLE_EXCERPT_LIMIT) {
|
||||
child.assistantExcerpt += (child.assistantExcerpt ? "\n" : "") + nested.text;
|
||||
}
|
||||
}
|
||||
// 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.
|
||||
this.deps.channels.get(entry.sessionId).publish(msg);
|
||||
watcher?.observe(msg);
|
||||
try {
|
||||
await this.deps.recorder.record(ctx, msg);
|
||||
} catch (err) {
|
||||
this.log(`[usage] 落库失败: ${err instanceof Error ? err.message : String(err)}`);
|
||||
this.deps.errors?.record({ source: "usage", err, ctx, code: "usage_insert_failed" });
|
||||
}
|
||||
}
|
||||
} catch (err) {
|
||||
// The SDK doesn't normally throw (errors are converged into the message stream);
|
||||
// this is a defensive record here to avoid crashing the runtime.
|
||||
this.log(
|
||||
`[session] 运行异常: ${err instanceof Error ? (err.stack ?? err.message) : String(err)}`,
|
||||
);
|
||||
this.deps.errors?.record({ source: "session", err, ctx, code: "session_run_failed" });
|
||||
} finally {
|
||||
// Wrap-up: persist any still-pending LLM failure and clear the tool-name cache (the watcher's state doesn't carry across runs).
|
||||
watcher?.close();
|
||||
entry.approvals.denyAll();
|
||||
entry.status = "idle";
|
||||
entry.abort = null;
|
||||
entry.running = null;
|
||||
entry.lastActivityMs = Date.now();
|
||||
this.publishState(entry, "idle");
|
||||
if (titleSource && titleSource.userExcerpt.trim()) {
|
||||
// Attempt generation whenever there's user material; whether generation is
|
||||
// actually needed (title still NULL, etc.) is decided by the generator itself.
|
||||
// Material is collected by the core Session during run; here we only pass the
|
||||
// fallback text.
|
||||
this.deps.titles?.maybeGenerate(ctx, entry.session, {
|
||||
fallbackText: titleSource.userExcerpt,
|
||||
});
|
||||
}
|
||||
// A subagent's title is likewise generated by the model, with material being the
|
||||
// subagent's **own conversation**: the prompt that spawned it plus its own reply
|
||||
// (material the parent Session collects belongs to the parent, hence the explicit
|
||||
// override here). It piggybacks a one-shot request on the parent Session's bare
|
||||
// LLM (the child Session object never leaves the SDK); on failure/empty result the
|
||||
// generator falls back to the prompt's first line.
|
||||
for (const child of children.values()) {
|
||||
if (!child.prompt.trim()) continue;
|
||||
this.deps.titles?.maybeGenerate(
|
||||
// Bookkeeping: Session/Agent record the subagent (the title belongs to it),
|
||||
// but the model reference still uses ctx's **parent-Session** pair
|
||||
// (provider, modelId) — this one-shot request really does run on the parent
|
||||
// Session's bare LLM (a subagent may switch models via run_subagent's
|
||||
// model_id).
|
||||
{ ...ctx, agentId: child.agentId, sessionId: child.sessionId },
|
||||
entry.session,
|
||||
{
|
||||
fallbackText: child.prompt,
|
||||
material: { userText: child.prompt, assistantText: child.assistantExcerpt },
|
||||
notifyOn: entry.sessionId, // Notify the frontend via the parent Session's SSE channel
|
||||
},
|
||||
);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Register a subagent: persisted only when the origin message is session_meta
|
||||
* (agentId is derived from the agent_state path: `<…>/<agentId>/agent_state`).
|
||||
* **The title is left blank** — it's generated at the end of this run by the model
|
||||
* from the subagent's own conversation (see drive's finally), falling back to the
|
||||
* first line of the run_subagent prompt if generation fails. Idempotent (children
|
||||
* dedup + insertOrIgnore); a subagent has its own Trace, so it's visible in both the
|
||||
* list and the trace view. On successful registration, the entry is put into
|
||||
* `children` and returned; a duplicate session_meta returns null.
|
||||
*/
|
||||
private registerChildSession(
|
||||
entry: RuntimeEntry,
|
||||
msg: OmniMessage,
|
||||
children: Map<string, ChildSession>,
|
||||
): ChildSession | null {
|
||||
if (!isSessionMeta(msg)) return null;
|
||||
const childSid = msg.origin![msg.origin!.length - 1]!;
|
||||
if (children.has(childSid)) return null;
|
||||
const p = msg.payload as SessionMetaPayload;
|
||||
const agentId = path.basename(path.dirname(p.agent_state));
|
||||
if (!agentId || agentId === "." || agentId === "..") return null;
|
||||
this.deps.sessions.insertOrIgnore({
|
||||
sessionId: childSid,
|
||||
projectId: entry.projectId,
|
||||
agentId,
|
||||
provider: p.provider,
|
||||
modelId: p.model_id,
|
||||
workspace: p.workspace,
|
||||
// A subagent's approvals are inherited from the parent Session; the index row is
|
||||
// inserted with defaults (matches the convention for Sessions discovered by the CLI).
|
||||
approvalMode: "allow-all",
|
||||
title: null,
|
||||
source: "subagent",
|
||||
createdAt: new Date().toISOString(),
|
||||
});
|
||||
// Make the subagent appear immediately in the sidebar: notify via the parent
|
||||
// Session's channel (a frontend currently watching the parent run refreshes its list in place).
|
||||
this.publishEvent(entry, {
|
||||
type: "session_created",
|
||||
projectId: entry.projectId,
|
||||
agentId,
|
||||
sessionId: childSid,
|
||||
source: "subagent",
|
||||
});
|
||||
const child: ChildSession = {
|
||||
sessionId: childSid,
|
||||
agentId,
|
||||
modelId: p.model_id,
|
||||
prompt: "",
|
||||
assistantExcerpt: "",
|
||||
};
|
||||
children.set(childSid, child);
|
||||
return child;
|
||||
}
|
||||
|
||||
private publishState(entry: RuntimeEntry, state: SessionStatus): void {
|
||||
this.publishEvent(entry, { type: "task_state", state });
|
||||
}
|
||||
|
||||
private publishEvent(entry: RuntimeEntry, event: ServerEvent): void {
|
||||
this.deps.channels.get(entry.sessionId).publish(event, "server_event");
|
||||
}
|
||||
|
||||
/** Serialize (mutually exclude) execution by sessionId; cleans up the lock-table entry once its chain drains (avoids unbounded growth). */
|
||||
private async withLock<T>(sessionId: string, fn: () => Promise<T>): Promise<T> {
|
||||
// What's stored in the chain is the already-caught version (used only for
|
||||
// sequencing, never propagates errors); the caller gets the original result from `next`.
|
||||
const prev = this.locks.get(sessionId) ?? Promise.resolve();
|
||||
const next = prev.then(fn);
|
||||
const settled: Promise<void> = next
|
||||
.then(
|
||||
() => undefined,
|
||||
() => undefined,
|
||||
)
|
||||
.then(() => {
|
||||
// Only delete if still the tail of the chain (no later waiter): preserves mutual-exclusion semantics.
|
||||
if (this.locks.get(sessionId) === settled) this.locks.delete(sessionId);
|
||||
});
|
||||
this.locks.set(sessionId, settled);
|
||||
return next;
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,274 @@
|
||||
/**
|
||||
* Error capture within the message stream: LLM request
|
||||
* failures and tool execution failures are **both never expressed via throw** — core
|
||||
* converges them into the message stream (LLM and Environment handle errors
|
||||
* internally and never throw), so a try/catch can't catch a single one. This watcher
|
||||
* hooks onto SessionManager's drive, inspects messages one by one, and fishes them out
|
||||
* into error_records (source = `llm` / `environment`), matching usage-recorder's shape:
|
||||
* recognizes only a few payload types, no-op on the rest. **One instance per run/compact**
|
||||
* (its state wraps up accordingly, see close).
|
||||
*
|
||||
* LLM (source = `llm`): reads the status of `request_end` —
|
||||
* - `failed` → unexpected (not retryable: auth failure, invalid params, etc., needs a human);
|
||||
* - `timeout` / `malformed` → expected (the engine already reconnects and retries, part
|
||||
* of normal operation);
|
||||
* - `aborted` / `completed` are not recorded (the former is a user-initiated interrupt,
|
||||
* not an error).
|
||||
*
|
||||
* The message uses the real reason: `request_end` only carries status, and **the only
|
||||
* place core carries the actual failure-reason text is the `abort` event's reason**
|
||||
* (e.g. `llm request error: 401 …` / `malformed response failed after N retries`). So a
|
||||
* `request_end` failure is first held pending, not persisted immediately, and is
|
||||
* resolved at the next request boundary:
|
||||
* - Immediately followed by `abort` → use its reason as the message (the real reason);
|
||||
* - Immediately followed by `request_begin` (the engine is retrying) → no reason text
|
||||
* left to wait for, use the status text;
|
||||
* - Still unresolved when the run ends → close persists it as a fallback.
|
||||
* Exception: when reason is a user-interrupt message (`aborted …`), it's not trusted —
|
||||
* "the user clicked stop during backoff" isn't the reason for this timeout, so the
|
||||
* status text is used instead (that timeout is a genuine failure and is still recorded).
|
||||
* Pending state is bucketed by origin: subagent messages interleave with the parent
|
||||
* session's (even more so with parallel subagents), and mixing them up would misattribute.
|
||||
*
|
||||
* Environment (source = `environment`): reads `tool_call_output`'s stop_reason ∈
|
||||
* {failed, timeout} → expected (the error is fed back to the model, and the Agent
|
||||
* adjusts on its own; `aborted` is denial/interruption, not recorded).
|
||||
* `tool_call_output` only has tool_call_id, no tool name, so `tool_call_id → tool name`
|
||||
* is cached (tool_call always arrives before its output), and the tool name is written
|
||||
* into code (`tool_failed:exec_command`) — so the stats dashboard's "most common error
|
||||
* code" and the error table can show at a glance which tool failed.
|
||||
*
|
||||
* **Attribution (ctx) is recorded against the session that actually produced the error,
|
||||
* not always the parent Session**: a subagent's LLM failures and tool failures also flow
|
||||
* through this same stream (carrying origin); if we simply reused the parent ctx passed
|
||||
* in at construction, filtering errors by Agent would always show 0 for the child Agent
|
||||
* and an inflated count for the parent — both attribution stats and the troubleshooting
|
||||
* target would be wrong. So we recognize `session_meta` carrying origin (a subagent's
|
||||
* first message, always arriving before any of its failures), registering
|
||||
* `origin → {agentId, sessionId}` (agentId derived from the agent_state path, matching
|
||||
* SessionManager.registerChildSession's convention); at persist time we look up the
|
||||
* message's origin: a hit records the subagent, a miss (a main-session message, or
|
||||
* session_meta hasn't arrived yet) falls back to the parent ctx. projectId is always
|
||||
* taken from the parent — a subagent is necessarily in the same Project.
|
||||
*/
|
||||
import { isEventMessage, isModelMessage, isSessionMeta } from "@prismshadow/penguin-core";
|
||||
import path from "node:path";
|
||||
import type { OmniMessage, SessionMetaMessage, StopReason } from "@prismshadow/penguin-core";
|
||||
import { MESSAGE_MAX } from "./error-recorder.js";
|
||||
import type { ErrorContext, ErrorKind, ErrorSink } from "./error-recorder.js";
|
||||
|
||||
/** Cap on the tool-name cache (bounded, to prevent unbounded growth over a long run; over the limit, evicts the oldest by registration order). */
|
||||
export const TOOL_NAMES_MAX = 1000;
|
||||
|
||||
/**
|
||||
* Cap on the subagent-identity cache (bounded, same reasoning as TOOL_NAMES_MAX: prevent
|
||||
* unbounded growth over a long run). The number of in-flight subagents is naturally
|
||||
* bounded by the subagent concurrency limit and falls far short of this value; over the
|
||||
* limit, evicts the oldest by registration order (those subagents have long since
|
||||
* settled, so even if a failure still arrives, it just falls back to the parent ctx —
|
||||
* i.e., the pre-fix behavior).
|
||||
*/
|
||||
export const ORIGIN_CTX_MAX = 200;
|
||||
|
||||
/** Recorded LLM failure states (`aborted` / `completed` are not errors and aren't included here). */
|
||||
type LlmFailure = "failed" | "timeout" | "malformed";
|
||||
|
||||
/** LLM failure state → error code, classification, and fallback message (used when the abort reason isn't available). */
|
||||
const LLM_FAILURES: Record<LlmFailure, { code: string; kind: ErrorKind; text: string }> = {
|
||||
failed: { code: "llm_failed", kind: "unexpected", text: "LLM 请求失败(不可重试)。" },
|
||||
timeout: { code: "llm_timeout", kind: "expected", text: "LLM 请求超时(引擎重连重试)。" },
|
||||
malformed: {
|
||||
code: "llm_malformed",
|
||||
kind: "expected",
|
||||
text: "LLM 响应无法解析(引擎重连重试)。",
|
||||
},
|
||||
};
|
||||
|
||||
/** Recorded tool failure states (`aborted` = denial/interruption, not an error). */
|
||||
type ToolFailure = "failed" | "timeout";
|
||||
|
||||
function isLlmFailure(s: unknown): s is LlmFailure {
|
||||
return s === "failed" || s === "timeout" || s === "malformed";
|
||||
}
|
||||
|
||||
function isToolFailure(s: unknown): s is ToolFailure {
|
||||
return s === "failed" || s === "timeout";
|
||||
}
|
||||
|
||||
/** A user-interrupt abort message (core's `aborted by user` / `aborted during …`): not a failure reason. */
|
||||
function isUserAbortReason(reason: string): boolean {
|
||||
return /^aborted\b/i.test(reason);
|
||||
}
|
||||
|
||||
/** The session a message belongs to (last origin element; empty string for the main session) — both pending state and the tool-name cache are bucketed by it. */
|
||||
function originKey(msg: OmniMessage): string {
|
||||
const origin = msg.origin;
|
||||
return origin && origin.length > 0 ? origin[origin.length - 1]! : "";
|
||||
}
|
||||
|
||||
/**
|
||||
* Take the **tail** of the tool output (not the head) as the message: core appends the
|
||||
* failure reason (`[tool error] …` / `[tool timeout: …]` / exit code) at the end of the
|
||||
* output, so truncating from the head would leave only a chunk of normal stdout and
|
||||
* drop the reason.
|
||||
*/
|
||||
function toolFailureText(output: string): string {
|
||||
if (!output) return "工具执行失败(无输出)。";
|
||||
if (output.length <= MESSAGE_MAX) return output;
|
||||
return `…${output.slice(output.length - (MESSAGE_MAX - 1))}`;
|
||||
}
|
||||
|
||||
export class StreamErrorWatcher {
|
||||
/** LLM failures awaiting a real reason: origin → failure state (see file header; each session has at most one in-flight Request). */
|
||||
private readonly pending = new Map<string, LlmFailure>();
|
||||
/** Names of in-flight tool calls: `origin \0 tool_call_id` → tool name (dequeued once output arrives). */
|
||||
private readonly toolNames = new Map<string, string>();
|
||||
/** Subagent identity: origin → that subagent's `{agentId, sessionId}` (see the file header's attribution section). */
|
||||
private readonly originCtx = new Map<string, { agentId: string; sessionId: string }>();
|
||||
|
||||
constructor(
|
||||
private readonly errors: ErrorSink,
|
||||
private readonly ctx: ErrorContext,
|
||||
) {}
|
||||
|
||||
/** Consume one outgoing message; messages irrelevant to this watcher are a no-op. */
|
||||
observe(msg: OmniMessage): void {
|
||||
if (isSessionMeta(msg)) {
|
||||
this.registerOrigin(msg);
|
||||
return;
|
||||
}
|
||||
if (isModelMessage(msg)) {
|
||||
this.observeTool(msg);
|
||||
return;
|
||||
}
|
||||
if (isEventMessage(msg)) this.observeLlm(msg);
|
||||
}
|
||||
|
||||
/** run/compact wrap-up: persist any still-pending failure (that never got its abort), and clear caches (prevents leaks). */
|
||||
close(): void {
|
||||
for (const key of [...this.pending.keys()]) this.flush(key);
|
||||
this.pending.clear();
|
||||
this.toolNames.clear();
|
||||
this.originCtx.clear();
|
||||
}
|
||||
|
||||
// —— Attribution (see file header) ——
|
||||
|
||||
/**
|
||||
* Register a subagent's identity: `session_meta` carrying origin is the subagent's
|
||||
* first message, always arriving before any of its failures. agentId is derived from
|
||||
* the absolute agent_state path (`<…>/<agentId>/agent_state`) — matching
|
||||
* SessionManager.registerChildSession's convention; not registered if the path is
|
||||
* malformed (that subagent's failures fall back to the parent ctx — better to
|
||||
* misattribute than write into a nonexistent agentId). The main session's session_meta
|
||||
* (no origin) is already the parent ctx and isn't registered.
|
||||
*/
|
||||
private registerOrigin(msg: SessionMetaMessage): void {
|
||||
const key = originKey(msg);
|
||||
if (!key) return; // Main session
|
||||
const agentId = path.basename(path.dirname(msg.payload.agent_state));
|
||||
if (!agentId || agentId === "." || agentId === "..") return;
|
||||
// Bounded (re-registering refreshes registration order; over the limit, evicts the oldest).
|
||||
this.originCtx.delete(key);
|
||||
this.originCtx.set(key, { agentId, sessionId: key });
|
||||
if (this.originCtx.size > ORIGIN_CTX_MAX) {
|
||||
const oldest = this.originCtx.keys().next().value;
|
||||
if (oldest !== undefined) this.originCtx.delete(oldest);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* The attribution to persist for this origin: a registered subagent hit → record its
|
||||
* own Agent/Session; a miss (a main-session message, or session_meta hasn't arrived
|
||||
* yet) → fall back to the parent ctx passed at construction. projectId is always taken
|
||||
* from the parent (a subagent is necessarily in the same Project).
|
||||
*/
|
||||
private ctxFor(key: string): ErrorContext {
|
||||
const child = this.originCtx.get(key);
|
||||
if (!child) return this.ctx;
|
||||
return { projectId: this.ctx.projectId, agentId: child.agentId, sessionId: child.sessionId };
|
||||
}
|
||||
|
||||
// —— LLM ——
|
||||
|
||||
private observeLlm(msg: OmniMessage): void {
|
||||
const p = msg.payload as { type?: string; status?: StopReason; reason?: string | null };
|
||||
const key = originKey(msg);
|
||||
if (p.type === "request_end") {
|
||||
this.flush(key); // Defensive: if a previous failure is still pending (normally resolved by request_begin), persist it first
|
||||
if (isLlmFailure(p.status)) this.pending.set(key, p.status);
|
||||
return;
|
||||
}
|
||||
// A new attempt begins (the engine is retrying): no reason text left to wait for the previous failure, persist using the status text.
|
||||
if (p.type === "request_begin") {
|
||||
this.flush(key);
|
||||
return;
|
||||
}
|
||||
// Interrupted/failed exit: reason is core's only failure-reason text.
|
||||
if (p.type === "abort") {
|
||||
this.flush(key, typeof p.reason === "string" ? p.reason : null);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Persist a pending LLM failure (no-op if none is pending); `reason` is the abort
|
||||
* message that arrived afterward. Pending state is already bucketed by origin, so
|
||||
* `key` is exactly "the session that produced this failure" — attribution is looked
|
||||
* up from it (see file header).
|
||||
*/
|
||||
private flush(key: string, reason?: string | null): void {
|
||||
const status = this.pending.get(key);
|
||||
if (status === undefined) return;
|
||||
this.pending.delete(key);
|
||||
const spec = LLM_FAILURES[status];
|
||||
const trimmed = reason?.trim();
|
||||
// A user-interrupt message isn't a failure reason (see file header); fall back to the status text.
|
||||
const message = trimmed && !isUserAbortReason(trimmed) ? trimmed : spec.text;
|
||||
this.errors.record({
|
||||
source: "llm",
|
||||
err: message,
|
||||
ctx: this.ctxFor(key),
|
||||
code: spec.code,
|
||||
kind: spec.kind,
|
||||
});
|
||||
}
|
||||
|
||||
// —— Environment (tool execution) ——
|
||||
|
||||
private observeTool(msg: OmniMessage): void {
|
||||
const p = msg.payload as {
|
||||
type?: string;
|
||||
name?: string;
|
||||
output?: string;
|
||||
tool_call_id?: string;
|
||||
stop_reason?: StopReason;
|
||||
};
|
||||
if (typeof p.tool_call_id !== "string") return;
|
||||
const origin = originKey(msg); // The session that made this call (both attribution and the tool-name cache are bucketed by it)
|
||||
const key = `${origin}\0${p.tool_call_id}`;
|
||||
|
||||
if (p.type === "tool_call" && typeof p.name === "string") {
|
||||
// tool_call arrives before its output: record the tool name (bounded, re-registering refreshes registration order).
|
||||
this.toolNames.delete(key);
|
||||
this.toolNames.set(key, p.name);
|
||||
if (this.toolNames.size > TOOL_NAMES_MAX) {
|
||||
const oldest = this.toolNames.keys().next().value;
|
||||
if (oldest !== undefined) this.toolNames.delete(oldest);
|
||||
}
|
||||
return;
|
||||
}
|
||||
if (p.type !== "tool_call_output") return;
|
||||
|
||||
const name = this.toolNames.get(key);
|
||||
this.toolNames.delete(key); // This call has settled: dequeue it, the cache only keeps in-flight calls
|
||||
if (!isToolFailure(p.stop_reason)) return; // completed / aborted (denial, user interrupt) are not errors
|
||||
this.errors.record({
|
||||
source: "environment",
|
||||
err: toolFailureText(p.output ?? ""),
|
||||
ctx: this.ctxFor(origin),
|
||||
// The tool name goes into code: so the stats dashboard's "most common error code" and table can show which tool failed.
|
||||
code: `tool_${p.stop_reason}:${name ?? "unknown"}`,
|
||||
kind: "expected",
|
||||
});
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,137 @@
|
||||
/**
|
||||
* The **policy layer** for automatic Session title generation (conversation-page
|
||||
* extension).
|
||||
*
|
||||
* "How to generate" lives in the core SDK (`session.generateTitle`: an out-of-band
|
||||
* one-shot request on the session's own Model, no tools, thinking disabled, writes no
|
||||
* history/Trace); this module is only responsible for host-side policy:
|
||||
* - When to generate: after a Task completes and the DB row's title is still NULL
|
||||
* (i.e. after the first successful conversation);
|
||||
* - Persistence and notification: writes sessions.title and pushes a `session_title`
|
||||
* server event to the Session channel;
|
||||
* - Bookkeeping: the one-shot request's token consumption is converted to token_usage
|
||||
* and handed to usage-recorder for persistence;
|
||||
* - Silent failure (logged): the title stays NULL and naturally retries after the next
|
||||
* Task completes.
|
||||
*/
|
||||
import { emptyTokenCounts, sanitizeTitle, 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";
|
||||
import type { RuntimeSession } from "./session-manager.js";
|
||||
import type { UsageContext, UsageRecorder } from "./usage-recorder.js";
|
||||
|
||||
export interface TitleGeneratorDeps {
|
||||
sessions: SessionsRepo;
|
||||
channels: ChannelHub;
|
||||
recorder: Pick<UsageRecorder, "record">;
|
||||
/** Error persistence (optional: without it, only logs — same as before this was wired up). */
|
||||
errors?: ErrorSink;
|
||||
log?: (line: string) => void;
|
||||
}
|
||||
|
||||
/** Host-side parameters for one title-generation request. */
|
||||
export interface TitleRequest {
|
||||
/** Fallback material for when the LLM fails or returns an empty result (cleaned and truncated from the first non-empty line). */
|
||||
fallbackText: string;
|
||||
/** Material override (for subagents — the material is the subagent's own conversation); defaults to the first Task's material self-collected by the core Session. */
|
||||
material?: { userText: string; assistantText: string };
|
||||
/** The channel to push the `session_title` event to; defaults to `ctx.sessionId`. A
|
||||
* subagent has no SSE channel of its own, so its title must reach the frontend via
|
||||
* the **parent Session's** channel (the list updates in place by sessionId). */
|
||||
notifyOn?: string;
|
||||
}
|
||||
|
||||
/** session-manager's minimal dependency on the title generator (tests inject a fake implementation). */
|
||||
export interface TitleNotifier {
|
||||
maybeGenerate(
|
||||
ctx: UsageContext,
|
||||
session: Pick<RuntimeSession, "generateTitle">,
|
||||
req: TitleRequest,
|
||||
): void;
|
||||
}
|
||||
|
||||
export class TitleGenerator implements TitleNotifier {
|
||||
private readonly inflight = new Set<string>();
|
||||
private readonly log: (line: string) => void;
|
||||
|
||||
constructor(private readonly deps: TitleGeneratorDeps) {
|
||||
this.log = deps.log ?? ((line) => console.error(line));
|
||||
}
|
||||
|
||||
/** Generate a title in the background (fire-and-forget) when conditions are met: the row exists, title is still NULL, and no generation is already in flight. */
|
||||
maybeGenerate(
|
||||
ctx: UsageContext,
|
||||
session: Pick<RuntimeSession, "generateTitle">,
|
||||
req: TitleRequest,
|
||||
): void {
|
||||
const row = this.deps.sessions.findById(ctx.sessionId);
|
||||
if (!row || row.title !== null) return;
|
||||
if (this.inflight.has(ctx.sessionId)) return;
|
||||
this.inflight.add(ctx.sessionId);
|
||||
void this.generate(ctx, session, req)
|
||||
.catch((err: unknown) => {
|
||||
this.log(`[title] 生成失败: ${err instanceof Error ? err.message : String(err)}`);
|
||||
this.deps.errors?.record({ source: "title", err, ctx, code: "title_failed" });
|
||||
})
|
||||
.finally(() => {
|
||||
this.inflight.delete(ctx.sessionId);
|
||||
});
|
||||
}
|
||||
|
||||
private async generate(
|
||||
ctx: UsageContext,
|
||||
session: Pick<RuntimeSession, "generateTitle">,
|
||||
req: TitleRequest,
|
||||
): Promise<void> {
|
||||
let title: string | null = null;
|
||||
try {
|
||||
// Material defaults to what the core Session self-collects during run; it's only
|
||||
// overridden in scenarios like subagents where the material isn't on that Session.
|
||||
const res = await session.generateTitle(
|
||||
req.material ? { material: req.material } : undefined,
|
||||
);
|
||||
title = res.title;
|
||||
// The one-shot request's real consumption is metered as usual (converted to token_usage and handed to recorder, attributed to this Session).
|
||||
if (res.usage) {
|
||||
try {
|
||||
await this.deps.recorder.record(ctx, tokenUsage(emptyTokenCounts(), res.usage));
|
||||
} catch (err) {
|
||||
this.log(`[title] 用量落库失败: ${err instanceof Error ? err.message : String(err)}`);
|
||||
this.deps.errors?.record({
|
||||
source: "title",
|
||||
err,
|
||||
ctx,
|
||||
code: "title_usage_insert_failed",
|
||||
});
|
||||
}
|
||||
}
|
||||
} catch (err) {
|
||||
// A model request error (rate limit / timeout / network, etc.) shouldn't leave the
|
||||
// title permanently missing: log it and fall through to the fallback.
|
||||
this.log(`[title] 模型请求失败: ${err instanceof Error ? err.message : String(err)}`);
|
||||
this.deps.errors?.record({ source: "title", err, ctx, code: "title_llm_failed" });
|
||||
}
|
||||
// When the LLM produces no usable title (failure / empty result), truncate the fallback material's first line — this guarantees a title is always generated.
|
||||
const finalTitle = title ?? fallbackTitle(req.fallbackText);
|
||||
if (finalTitle === null) return;
|
||||
// There may already be a concurrent write during generation (e.g. a future manual rename): only persist if still NULL.
|
||||
const latest = this.deps.sessions.findById(ctx.sessionId);
|
||||
if (!latest || latest.title !== null) return;
|
||||
this.deps.sessions.updateTitle(ctx.sessionId, finalTitle);
|
||||
this.deps.channels
|
||||
.get(req.notifyOn ?? ctx.sessionId)
|
||||
.publish(
|
||||
{ type: "session_title", sessionId: ctx.sessionId, title: finalTitle },
|
||||
"server_event",
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
/** 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);
|
||||
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);
|
||||
}
|
||||
@@ -0,0 +1,116 @@
|
||||
/**
|
||||
* Usage persistence.
|
||||
*
|
||||
* Consumes the Session output stream:
|
||||
* - A subagent's `session_meta` (carrying origin) → registers the mapping "origin's
|
||||
* last session_id → (provider, model_id)" (a subagent may use a different Model, so
|
||||
* cost is priced against the actual Model used);
|
||||
* - `token_usage` → inserts one usage_records row: the four token fields come from
|
||||
* `payload.request` (per-Request increment; subagents are recorded one row at a
|
||||
* time, with `session_id` attributed to their owning main Session). Only tokens are
|
||||
* persisted, not cost — pricing may be added later, and cost is computed on the fly
|
||||
* against current pricing when usage-service queries.
|
||||
* The attribution key is always the paired reference `(provider, model_id)` (the same
|
||||
* model_id name across different vendors is attributed separately).
|
||||
*/
|
||||
import { isEventMessage, isSessionMeta } from "@prismshadow/penguin-core";
|
||||
import type { OmniMessage } from "@prismshadow/penguin-core";
|
||||
import { formatLocalDate } from "../internal/dates.js";
|
||||
import type { UsageRepo } from "../db/repos/usage.js";
|
||||
|
||||
/** Attribution context for one record (top-level Session scope). */
|
||||
export interface UsageContext {
|
||||
projectId: string;
|
||||
agentId: string;
|
||||
/** Top-level Session id (the current actual id after self-heal). */
|
||||
sessionId: string;
|
||||
/** Vendor grouping for the top-level Session's model (paired with modelId; the fallback attribution when the origin mapping has no hit). */
|
||||
provider: string;
|
||||
/** Upstream model_id of the top-level Session (paired with provider). */
|
||||
modelId: string;
|
||||
}
|
||||
|
||||
/** Cap on the subagent attribution mapping: over the limit, evicts the oldest by insertion order (an evicted entry falls back to the main Session's Model attribution). */
|
||||
export const ORIGIN_MODELS_MAX = 1000;
|
||||
|
||||
export class UsageRecorder {
|
||||
/** Subagent model attribution mapping: origin's last session_id → paired reference (session_id is globally unique). */
|
||||
private readonly originModels = new Map<string, { provider: string; modelId: string }>();
|
||||
|
||||
constructor(
|
||||
private readonly usage: UsageRepo,
|
||||
private readonly now: () => Date = () => new Date(),
|
||||
) {}
|
||||
|
||||
/** Consume one outgoing message; messages other than session_meta / token_usage are a no-op. */
|
||||
async record(ctx: UsageContext, msg: OmniMessage): Promise<void> {
|
||||
if (isSessionMeta(msg) && msg.origin && msg.origin.length > 0) {
|
||||
const originSessionId = msg.origin[msg.origin.length - 1]!;
|
||||
// Bounded mapping (avoids unbounded growth over a long-running process): re-inserting refreshes insertion order, over the limit evicts the oldest.
|
||||
this.originModels.delete(originSessionId);
|
||||
this.originModels.set(originSessionId, {
|
||||
provider: msg.payload.provider,
|
||||
modelId: msg.payload.model_id,
|
||||
});
|
||||
if (this.originModels.size > ORIGIN_MODELS_MAX) {
|
||||
const oldest = this.originModels.keys().next().value;
|
||||
if (oldest !== undefined) this.originModels.delete(oldest);
|
||||
}
|
||||
return;
|
||||
}
|
||||
if (!isEventMessage(msg)) return;
|
||||
const payload = msg.payload as {
|
||||
type?: string;
|
||||
request?: { cache_read: number; cache_write: number; output: number; total: number };
|
||||
status?: string;
|
||||
};
|
||||
|
||||
const originSessionId =
|
||||
msg.origin && msg.origin.length > 0 ? msg.origin[msg.origin.length - 1]! : null;
|
||||
// Empty origin → main session's Model; otherwise look up the mapping, falling back to the main Session's Model (paired) on a miss.
|
||||
const ref =
|
||||
originSessionId === null
|
||||
? { provider: ctx.provider, modelId: ctx.modelId }
|
||||
: (this.originModels.get(originSessionId) ?? {
|
||||
provider: ctx.provider,
|
||||
modelId: ctx.modelId,
|
||||
});
|
||||
const now = this.now();
|
||||
const base = {
|
||||
ts: now.toISOString(),
|
||||
date: formatLocalDate(now),
|
||||
projectId: ctx.projectId,
|
||||
agentId: ctx.agentId,
|
||||
sessionId: ctx.sessionId,
|
||||
originSessionId,
|
||||
provider: ref.provider,
|
||||
modelId: ref.modelId,
|
||||
};
|
||||
|
||||
if (payload.type === "token_usage" && payload.request) {
|
||||
// A successful request: persist along with tokens (status defaults to completed).
|
||||
const r = payload.request;
|
||||
this.usage.insert({
|
||||
...base,
|
||||
cacheRead: r.cache_read,
|
||||
cacheWrite: r.cache_write,
|
||||
output: r.output,
|
||||
total: r.total,
|
||||
});
|
||||
return;
|
||||
}
|
||||
// A failed request (request_end and not completed, usually with no token_usage):
|
||||
// persist 0 tokens + status, feeding the "model success rate" stat; a successful
|
||||
// request is already counted once via the token_usage branch above, not repeated here.
|
||||
if (payload.type === "request_end" && payload.status && payload.status !== "completed") {
|
||||
this.usage.insert({
|
||||
...base,
|
||||
cacheRead: 0,
|
||||
cacheWrite: 0,
|
||||
output: 0,
|
||||
total: 0,
|
||||
status: payload.status,
|
||||
});
|
||||
}
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user