@crow-agent/agent-core
v0.2.0
Published
Agentic runtime for Crow: tool calling, session state, and pluggable transports.
Maintainers
Readme
@crow-agent/agent-core
@crow-agent/agent-core is the runtime that turns a language model into a working
agent: it manages the message transcript, calls the model in a loop, executes
tool calls (in parallel or one at a time), streams progress events to your UI,
and offers hooks — steering, follow-up queues, per-request context rebuilds —
for steering a run while it is live. It is built on @crow-agent/ai for model
access and stays out of the persistence business: session storage lives in
@crow-agent/session-backend-sqlite-node and friends.
npm install @crow-agent/agent-coreA first run
import { Agent } from "@crow-agent/agent-core";
import { createModels } from "@crow-agent/ai";
import { anthropicProvider } from "@crow-agent/ai/vendors/anthropic";
const models = createModels();
const anthropic = anthropicProvider();
models.setProvider(anthropic);
const model = models.getModel("anthropic", "claude-opus-4-6");
if (!model) throw new Error("No such model");
const agent = new Agent({
initialState: {
systemPrompt: "You are Crow, a helpful coding assistant.",
model,
},
streamFn: (m, c, o) => models.streamSimple(m, c, o),
});
// Render streamed text as it arrives.
agent.subscribe((event) => {
if (event.type !== "message_update") return;
if (event.assistantMessageEvent.type !== "text_delta") return;
process.stdout.write(event.assistantMessageEvent.delta);
});
await agent.prompt("Hello!");prompt() resolves once the run settles — the final assistant message is in,
every tool call finished, and all awaited subscribers ran.
Two message worlds
The agent keeps an AgentMessage[] transcript. That type is deliberately
wider than what a model accepts:
user,assistant, andtoolResultmessages go to the model as-is.- Anything else — UI markers, app-specific kinds declared through module augmentation — lives only in the transcript.
Two functions bridge the worlds, and they run in a fixed order before every model call:
AgentMessage[] → transformContext() → AgentMessage[] → convertToLlm() → Message[] → model
(optional) (required)transformContext reshapes the transcript: prune old turns, pull in fresh
context, compact. convertToLlm then maps the result onto model-shaped
messages, dropping whatever the model must never see. If you add a custom
message kind, convertToLlm is where you teach the agent how to render it.
What a run looks like, event by event
Subscribe with agent.subscribe() and you get a deterministic event stream.
Handlers run in registration order and are awaited — prompt() and
waitForIdle() do not resolve until awaited agent_end handlers finish.
A plain prompt("Hello") produces:
prompt("Hello")
├─ agent_start run begins
├─ turn_start one model call + its tool batch
├─ message_start / message_end your user message
├─ message_start assistant starts
├─ message_update (×n) streamed partials
├─ message_end full assistant message
├─ turn_end turn's message + tool results
└─ agent_end nothing else followsWhen the assistant calls tools, the loop keeps going instead of ending the turn:
prompt("Read config.json")
├─ agent_start
├─ turn_start
├─ message_start / message_end user message
├─ message_start assistant message with tool calls
├─ message_update… / message_end
├─ tool_execution_start (toolCallId, toolName, args)
├─ tool_execution_update { partialResult } only when the tool streams
├─ tool_execution_end (toolCallId, result)
├─ message_start / message_end toolResult message
├─ turn_end { toolResults: [...] }
│
├─ turn_start next turn
├─ message_start / message_update… / message_end model answers the results
├─ turn_end
└─ agent_endParallel and sequential tool batches
toolExecution (global) and executionMode (per tool) control batching:
- parallel (default) — validate every call in order, run the permitted
ones together.
tool_execution_endfires the instant each call finishes, so completion order follows wall-clock; the stored toolResult messages still keep the assistant's original call order. - sequential — run calls strictly one after another. A single
sequentialtool in a batch drags the whole batch into sequential mode, regardless of the global setting.
Hooks around each turn
| Hook | When it runs | What it can do |
|---|---|---|
| beforeToolCall | After tool_execution_start and argument validation | Block the call ({ block: true, reason }), optionally with terminate: true |
| afterToolCall | After the tool settles, before tool_execution_end | Rewrite the result, or set terminate: true |
| prepareRequest | Before every model request, first included | Rebuild the context — e.g. install canonical persisted messages after flushing pending input. Never drains the steering queue; queued steering waits for the next regular poll |
| finishTurn | After the assistant message and all tool results are final, before turn_end | Return { action: "continue" } to force one more request, { action: "end" } to stop after turn_end, or undefined for default scheduling |
terminate: true (from a tool, a blocked beforeToolCall, or
afterToolCall) proposes skipping the follow-up model call — but it only
lands when every finalized result in the batch agrees. Mixed batches
carry on as normal.
finishTurn runs on normal, error, and aborted turns alike. On a normal
response, { action: "continue" } guarantees one more request; if tool
results, steering, or follow-up would already trigger one, that satisfies the
decision and nothing extra is issued — otherwise the loop sends a
context-only request. Error and aborted responses exit hard no matter what
the hook returns, and the hook runs again after the next request, so an
unconditional { action: "continue" } loops forever.
With the Agent class, the assistant's message_end is a barrier ahead of
tool preflight, so beforeToolCall observes agent state that already
contains the requesting assistant message.
Interrupting and resuming
Two queues let you talk to a running agent:
agent.steeringMode = "one-at-a-time"; // or "all"
agent.followUpMode = "one-at-a-time"; // or "all"
agent.steer({ role: "user", content: "Stop! Do this instead.", timestamp: Date.now() });
agent.followUp({ role: "user", content: "Also summarize the result.", timestamp: Date.now() });
await agent.continue();- Steering interrupts mid-tool-run. After the in-flight tool calls finish, steering messages are injected and the model responds next turn.
- Follow-up queues work for when the run would otherwise stop; it is only consulted once no tool calls and no steering remain.
continue() replays the queue semantics: an empty or system-only transcript
rejects without touching the queues; a non-assistant tail resumes from the
current context, polling steering at startup and waiting for a natural stop
on follow-up. An assistant tail cannot be sent to the model directly, so
continue() consumes one steering batch and then one follow-up batch, with
the queue mode deciding whether each batch carries a single message or all
of them.
peekQueuedMessages() previews the queues without consuming;
clearSteeringQueue(), clearFollowUpQueue(), and clearAllQueues() drop
them.
Configuring an agent
Every option below is settable at construction via initialState (state
fields) or top-level config (behavior fields), and most are adjustable
later through agent.state or direct assignment.
const agent = new Agent({
initialState: {
systemPrompt: string, // seeds the leading system message
model: Model<any>,
thinkingLevel:
| "off" | "minimal" | "low" | "medium" | "high" | "xhigh" | "max",
tools: AgentTool<any>[],
messages: AgentMessage[],
},
convertToLlm: (messages) => messages.filter((m) => m.role !== "system"),
transformContext: async (messages) => dropOldTurns(messages),
steeringMode: "one-at-a-time", // or "all"
followUpMode: "one-at-a-time", // or "all"
streamFn: models.streamSimple.bind(models), // required
sessionId: "session-123", // provider-side caching
getApiKey: async (provider) => refreshToken(), // dynamic keys
toolExecution: "parallel", // or "sequential"
beforeToolCall: async ({ toolCall, args }) => {
if (toolCall.name === "bash") {
return { block: true, reason: "bash is disabled here", terminate: true };
}
},
afterToolCall: async ({ toolCall, result, isError }) => {
if (!isError) return { details: { ...result.details, audited: true } };
},
prepareRequest: async ({ context }) => {
return { context: { ...context, messages: await loadCanonicalMessages() } };
},
finishTurn: async ({ message, toolResults }) => {
return wantsAnotherTurn(message, toolResults) ? { action: "continue" } : undefined;
},
thinkingBudgets: { minimal: 128, low: 512, medium: 1024, high: 2048 },
});Reading and writing state
agent.state.model = getModel("openai", "gpt-4o-mini");
agent.state.thinkingLevel = "medium";
agent.state.tools = [myTool];
agent.state.messages = newMessages; // the top-level array is copied on write
agent.state.messages.push(message); // …but pushing mutates live state
agent.toolExecution = "sequential";
agent.prepareRequest = async ({ context, signal }) => ({
context: { ...context, messages: await reloadPersistedMessages(signal) },
});
agent.abort(); // cancel the active operation
await agent.waitForIdle(); // wait until everything settles
agent.reset();The transcript owns the system prompt and the declared tools: the first
system message is the prompt, and later system messages patch it (see
SystemMessage in @crow-agent/ai). agent.state.systemPrompt is read-only and
replays the transcript. Before each request the loop diffs agent.state.tools
against the transcript's declared tools; on mismatch it announces the delta
in a system message (folded into a pending one when there is one), and
@crow-agent/ai's getCurrentSystemMessage(messages) replays the head — declared
tools included — for any message array.
To change the prompt mid-run, append a system message with content (extra
instructions) or sections (named replacements):
await agent.prompt([
{
role: "system", content: "",
sections: { skills: "<skills>...</skills>" }, timestamp: Date.now(),
},
{ role: "user", content: "Pick up where you left off", timestamp: Date.now() },
]);While streaming, agent.state.streamingMessage holds the partial assistant
message, and agent.state.isStreaming stays true until the run fully
settles — awaited agent_end subscribers included.
Prompting shapes
await agent.prompt("Hello"); // plain text
await agent.prompt("What's in this image?", [ // with images
{ type: "image", data: base64Data, mimeType: "image/jpeg" },
]);
await agent.prompt({ role: "user", content: "Hello", timestamp: Date.now() }); // raw message
await agent.continue(); // resume queued / unfinished inputCustom message kinds
Extend the transcript through declaration merging, then teach convertToLlm
what the model should see:
declare module "@crow-agent/agent-core" {
interface CustomAgentMessages {
notification: { role: "notification"; body: string; timestamp: number };
}
}
const agent = new Agent({
streamFn: (m, c, o) => models.streamSimple(m, c, o),
convertToLlm: (messages) => messages.flatMap((m) =>
m.role === "notification" ? [] : [m], // UI-only kinds never reach the model
),
});Writing tools
import { readFile } from "node:fs/promises";
import { Type } from "typebox";
const readFileTool: AgentTool = {
name: "read_file",
label: "Read File", // shown in UIs
description: "Read a file's contents",
parameters: Type.Object({
path: Type.String({ description: "Path of the file to read" }),
}),
executionMode: "sequential", // optional per-tool override
execute: async (callId, params, signal, onUpdate) => {
const content = await readFile(params.path, "utf-8");
onUpdate?.({ content: [{ type: "text", text: "Reading file…" }], details: {} });
return {
content: [{ type: "text", text: content }],
details: { path: params.path, bytes: content.length },
};
},
};
agent.state.tools = [readFileTool];Throw on failure — never return error text as content. A thrown error
becomes a tool error surfaced to the model with isError: true.
@crow-agent/mcp exposes MCP servers as AgentTools and @crow-agent/codemode runs
model-authored JavaScript that can invoke tools; the
examples/mcp-agent directory wires both up, and
beforeToolCall/afterToolCall cover codemode's tool invocations too.
Browser apps: proxying the model stream
import { Agent, streamProxy } from "@crow-agent/agent-core";
const agent = new Agent({
streamFn: (model, context, options) =>
streamProxy(model, context, {
...options,
authToken: "...",
proxyUrl: "https://your-server.com",
}),
});The low-level loop
agentLoop / agentLoopContinue give you the raw event stream without the
Agent class. They are observational — event order is preserved, but async
handlers are not awaited before later producer phases proceed. If you need
barrier semantics (message handling settled before tool preflight), use
Agent.
import { agentLoop, agentLoopContinue } from "@crow-agent/agent-core";
const context: AgentContext = {
messages: [{ role: "system", content: "Be concise.", timestamp: Date.now() }],
tools: [],
};
const config: AgentLoopConfig = {
model: getModel("openai", "gpt-4o"),
convertToLlm: (msgs) => msgs.filter((m) => m.role !== "system"),
toolExecution: "parallel",
beforeToolCall: async () => undefined,
afterToolCall: async () => undefined,
};
const stream = models.streamSimple.bind(models);
const userMessage = { role: "user", content: "Hello there", timestamp: Date.now() };
for await (const event of agentLoop([userMessage], context, config, undefined, stream)) {
console.log(event.type);
}
for await (const event of agentLoopContinue(context, config, undefined, stream)) {
console.log(event.type);
}Subpath exports
@crow-agent/agent-core— theAgentclass, the low-level loop, proxying@crow-agent/agent-core/node— everything above plus the Node.js execution environment@crow-agent/agent-core/harness/context— chord-backed context helpers (Context,withTelemetryContext, …)@crow-agent/agent-core/experimental/pico3— the experimental pico3 harness surface@crow-agent/agent-core/harness/env/nodejs— the Node.js execution environment on its own@crow-agent/agent-core/harness/runtime/reducer— the lane-snapshot reducer@crow-agent/agent-core/harness/session— session, storage, and lane types@crow-agent/agent-core/harness/session/testing— fixtures and seeded benchmark datasets
License
MIT
