@m6d/cortex-server
v0.1.0
Published
Reusable AI agent chat server library for Hono + Bun
Readme
@m6d/cortex-server
Multi-agent AI chat server built on Hono + Bun. Supports MSSQL or Postgres persistence, MinIO file storage, Neo4j knowledge graphs, and AI streaming with tool execution.
This package is a runtime library: it serves requests, resolves against the knowledge graph on every prompt, and exports the define* helpers you author domains with. Generation and seeding are not here — seeding the graph and generating endpoint schemas from Swagger are @m6d/cortex-cli commands (cortex graph seed, cortex swagger sync). What the two share is the graph contract — GRAPH_SCHEMA, GRAPH_SCHEMA_VERSION, the define* helpers and their types, the Neo4j client and the embedder — which lives in the private @cortex/contracts workspace package (internal/contracts) and is vendored into this package's tarball at pack time. It is all re-exported from here, so authoring your domains stays one import. The CLI stamps the schema version it seeded onto the graph; this package reads that stamp on its first resolve and refuses a graph it cannot read, instead of quietly retrieving nothing.
Usage
bun add @m6d/cortex-serverimport { createCortex } from "@m6d/cortex-server";
const cortex = createCortex({
database: {
// Or `type: "mssql"` with an ADO.NET connection string.
type: "postgres",
connectionString: "Server=localhost;Database=cortex;...",
},
storage: {
endPoint: "localhost",
port: 9000,
useSSL: false,
accessKey: "minio",
secretKey: "minio-secret-key",
},
auth: {
kind: "jwks",
jwksUri: "https://auth.example.com/.well-known/jwks.json",
issuer: "https://auth.example.com",
},
model: {
baseURL: "https://api.openai.com/v1",
apiKey: "sk-...",
modelName: "gpt-4o",
},
embedding: {
baseURL: "https://api.openai.com/v1",
apiKey: "sk-...",
modelName: "text-embedding-3-small",
dimension: 1536,
},
neo4j: {
url: "http://localhost:7474",
user: "neo4j",
password: "password",
},
agents: {
assistant: {
systemPrompt: "You are a helpful assistant.",
},
},
});
// Start the server
const server = await cortex.serve();
export default {
fetch: server.fetch,
websocket: server.websocket,
};cortex.serve() is async — it runs database migrations on startup, from the migration history belonging to the configured dialect.
Configuration
createCortex(config) takes a CortexConfig. Only database and model are required; every provider declared at server level is inherited by each agent unless the agent overrides it.
CortexConfig
| Field | Type | Required | Description |
| ----------------------- | -------------------------------------------------------------- | -------- | -------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- |
| port | number | No | The port serve() reports back; the caller binds it. Defaults to 3331. |
| database | { type: "mssql" \| "postgres"; connectionString } | Yes | The application database. Each dialect keeps its own migration history, run on serve(). |
| storage | StorageConfig | No | MinIO / S3 storage for composer attachments and attachment#<id> interpolation. |
| cors | { origins: string[] } | No | Browser origins allowed to call this server directly, for a widget hosted elsewhere without a same-origin proxy (the Control Center's assistant panel). Unset, cross-origin browsers are refused. |
| redis | { url; streamTtlSeconds? } | No | Redis-backed resumable streams that survive restarts and reach every instance. Omitted, streams fall back to in-memory (single process only). TTL defaults to 24 h. |
| auth | AuthConfig & { tokenExtractor?; cookieName? } | No | End-user token verification for the whole server (see Authentication), plus an optional extractor and cookie name for tokens that do not travel as a bearer. |
| model | ModelConfig | Yes | The chat model: { baseURL, apiKey, modelName, providerName? } for any OpenAI-compatible endpoint. |
| fastModel | ModelConfig | No | A cheaper model for auxiliary passes (summaries, descriptions). Same shape as model. |
| vision | ModelConfig | No | An image-capable model that describes uploaded attachments and enables the built-in readAttachment tool. A PDF's text pages are read as text and its scanned or chart pages as images (the note names any the image budget left out); a mostly scanned PDF as images only. Needs storage. |
| embedding | { baseURL; apiKey; modelName; dimension } | No | The embeddings endpoint; required with neo4j. dimension is what cortex graph seed sizes the vector index to. |
| neo4j | { url; user; password } | No | The knowledge graph over Neo4j's HTTP transaction API. Requires embedding. |
| reranker | { url; apiKey } | No | Re-ranks graph retrieval when it comes back noisy. |
| controlCenter | ControlCenterConfig | No | The Control Center this server belongs to: { url, apiKey, playground?, playgroundIssuer? }. Server-level only; which agents it backs is decided by publishing them there. |
| context | Partial<ContextConfig> | No | Context-window tuning: maxContextTokens, summarizationThreshold, toolResultMaxTokens, recentMessagesToKeep, all defaulted, and summarizationModel, which falls back to model. A thread compacts into a checkpoint summary once the messages since the last checkpoint pass the threshold; tool results are capped once, when recorded. |
| mcpServers | McpServerConfig[] | No | MCP servers every agent gets; see MCP servers. |
| knowledge | { swagger?: { url }; domains? } | No | The knowledge graph's inputs: the Swagger spec cortex swagger sync reads and the domains cortex graph seed loads. |
| variables | Array<{ name; description; example }> | No | The {{variables}} this server fills per request for the agents the Control Center creates on it. Declared to cc at boot. |
| loadSessionData | (token) => Promise<Record<string, unknown>> | No | Fetches session data on first request; becomes context.session in a functional systemPrompt and fills variables. |
| resolveRequestContext | (request, identity) => Record<string, unknown> \| Promise<…> | No | Per-request metadata; becomes context.requestContext and fills variables. identity is { userId, claims } from the verifier. |
| triggerWaitMaxMs | number | No | The longest a ?wait= fire holds its request for the run; longer waits are capped to it. Default five minutes. |
| agents | Record<string, CortexAgentDefinition> | No | The agents this server's own code runs, keyed by slug (assistant → /agents/assistant/...). A slug with no entry is served when cc created an agent by that slug here. |
StorageConfig
| Field | Type | Required | Description |
| ------------ | --------- | -------- | -------------------------------------- |
| endPoint | string | Yes | MinIO server hostname. |
| port | number | Yes | MinIO server port. |
| useSSL | boolean | Yes | Whether to use HTTPS. |
| accessKey | string | Yes | MinIO access key. |
| secretKey | string | Yes | MinIO secret key. |
| bucketName | string | No | Bucket name; defaults to ai-storage. |
storage is required for attachments. The widget uploads files from the composer to
POST /threads/:threadId/files; those thread-owned pending rows are claimed by the
next user turn and listed in the turn's context message. With backendFetch, the
agent can pass attachment#<id> anywhere in a request body and the default request
interceptor replaces it with the stored file before forwarding the request.
A custom backendFetch.transformRequestBody replaces that default interceptor and
must resolve the attachment#<id> sentinel itself if it advertises or accepts it.
CortexAgentDefinition
Everything an agent does, plus its own copy of any provider it wants to differ on.
| Field | Type | Required | Description |
| ---------------------------------------------------------------- | -------------------------------------------------------------------------------- | -------- | ------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------ |
| systemPrompt | string \| (context) => string \| Promise<string> | No | Static, or a function of { session, requestContext }. Defaults to an empty prompt. |
| tools | ToolSet \| (context) => ToolSet | No | The agent's tools, as chat() takes them: one with execute runs on the server, one without is a static client tool. A factory sees the request context. |
| mcpServers | McpServerConfig[] | No | MCP servers whose tools this agent's scripts can call each turn; an entry replaces a server-level one of the same name. See MCP servers. |
| skills | CodeSkill[] | No | On-demand instructions; see Skills. |
| triggers | CodeTrigger[] | No | Entry points external systems fire; see Triggers. |
| messageTrigger | boolean | No | The built-in message trigger, on by default. false makes the agent headless. |
| suggestions | Suggestions | No | The widget's landing-page chips per locale, at most six { title, prompt } each: the chip shows the title and sends the prompt. Wins over the Control Center's list when set. |
| greeting | Greeting | No | The landing-page headline per locale; "" keeps the widget's default. Same precedence. |
| backendFetch | { baseUrl; apiKey; headers?; userToken?; transformRequestBody?; interceptor? } | No | Where the sandbox api helper points. userToken: { header, prefix? } says where the end user's token goes (default X-Service-Token; { header: "Authorization", prefix: "Bearer " } hands it over as the user's own bearer). Without transformRequestBody the default interceptor swaps attachment#<id> sentinels for the stored file. |
| loadSessionData | (token) => Promise<TSession> | No | Per-agent override of the server-level loader. |
| resolveRequestContext | (request, identity) => TRequestContext \| Promise<TRequestContext> | No | Per-agent override of the server-level resolver. |
| onToolCall | ({ toolName, toolCallId, args }) => void | No | Fired on each tool call the model makes. Tools called from inside executeCode are not model tool calls and do not fire it. |
| onStreamFinish | ({ messages, isAborted }) => void | No | Fired when the stream finishes. |
| auth | AuthConfig | No | Agent-level token verification; see Authentication. |
| knowledge | KnowledgeConfig \| null | No | Replaces the server-level knowledge config for this agent; null opts out of the graph entirely. |
| context | Partial<ContextConfig> | No | Per-agent context-window tuning. |
| attachments | boolean | No | false refuses end-user uploads on this agent even with server-level storage. |
| imageFetch | { allowedHostnames } | No | Enables the built-in fetchImage tool for these exact hosts: the image is fetched server-side under the egress policy, stored as an attachment, and shown inline. Needs storage. |
| model, fastModel, vision, embedding, neo4j, reranker | same as server level | No | Per-agent provider overrides; an unset one inherits. |
Knowledge Graph
Cortex uses a Neo4j knowledge graph to give agents structured understanding of your domain. The graph is built from declarative domain definitions.
Domain Model
- Domains -- Top-level groupings (e.g. "Leaves", "Attendance")
- Concepts -- Things the AI understands (e.g. "Leave", "Employee"). Support aliases for better vector search matching.
- Endpoints -- API endpoints with auto-generated params, body, and response schemas from Swagger.
- Services -- Actionable workflows tied to concepts (e.g. "Request Leave").
- Rules -- Business constraints that govern concepts, endpoints, or services.
Helper Functions
import {
defineDomain,
defineConcept,
defineEndpoint,
defineRule,
defineService,
} from "@m6d/cortex-server";
const leaveConcept = defineConcept({
name: "Leave",
description: "An employee leave request (annual, sick, etc.)",
aliases: ["time off", "vacation", "absence"],
});
const paginationRule = defineRule({
name: "Pagination Limit",
description: "List endpoints must use bounded pageSize (max 100)",
});
const listLeavesEndpoint = defineEndpoint({
name: "List Leaves",
path: "/leaves/leaves",
method: "GET",
autoGenerated: {
params: [],
body: [],
response: [],
successStatus: 200,
errorStatuses: [400],
},
queries: [leaveConcept],
governedBy: [paginationRule],
});
const requestLeaveService = defineService({
name: "Request Leave",
description: "Workflow to submit a new leave request",
builtInId: "services:annual_leave",
belongsTo: leaveConcept,
});
const leavesDomain = defineDomain({
name: "Leaves",
description: "Leave management and balances",
concepts: [leaveConcept],
endpoints: [listLeavesEndpoint],
rules: [paginationRule],
services: [requestLeaveService],
});File Structure Convention
src/domains/
leaves/
index.ts
concepts/leave.concept.ts
concepts/leaveBalance.concept.ts
endpoints/listLeaves.endpoint.ts
endpoints/getLeaveById.endpoint.ts
services/leaveRequests.service.ts
attendance/
index.ts
concepts/punch.concept.ts
endpoints/listPunches.endpoint.ts
...Wiring Domains into Config
import { createCortex } from "@m6d/cortex-server";
import { employeesDomain, leavesDomain, attendanceDomain } from "./domains";
const cortex = createCortex({
// ...other config
knowledge: {
swagger: { url: "https://api.example.com/swagger/v1/swagger.json" },
domains: {
employees: employeesDomain,
leaves: leavesDomain,
attendance: attendanceDomain,
},
},
// ...
});The CLI
@m6d/cortex-cli scaffolds new Cortex servers (bunx @m6d/cortex-cli new my-app) and drives the
graph generators, so there are no seed scripts to write. The full command surface lives in
apps/cortex-cli/README.md.
bun add -d @m6d/cortex-cliEvery command loads cortex.config.ts from the directory you run in — that exact name, no upward search. --config <path> points it elsewhere.
Seed the Knowledge Graph
Reads your domain definitions and seeds the Neo4j knowledge graph with concepts, endpoints, services, rules, and their embeddings. No flags. Exits 1 if any statement failed, so a half-seeded graph doesn't pass as success.
The CLI runs standalone against your cortex.config.ts; it never constructs or serves your server. Both sides take the graph vocabulary from the private @cortex/contracts package they each vendor, and the seeder stamps the schema version it wrote onto the graph. If your server later reads a graph stamped with a different version it says so and stops, rather than resolving nothing and looking like a bad prompt.
bunx cortex graph seedExtract Endpoint Schemas
Fetches the Swagger spec(s) named in your config and updates .endpoint.ts files with params, body, and response schemas. Prints the files it changed. Endpoint files with no match in the spec are listed but don't fail the run.
bunx cortex swagger synccheck is the read-only sibling — same inputs, writes nothing, exits 2 when it finds drift or an unmatched endpoint file. That makes it a CI step:
bunx cortex swagger checkBoth read src/domains, resolved next to the config file. --domains-dir <path> overrides:
bunx cortex swagger sync --domains-dir ./server/src/domainsSkills
Instructions the model pulls in on demand. Every skill's name: description line is in the system prompt on each turn under ## Skills; the body enters context only when the model calls the built-in loadSkill tool with that name. The description is the whole trigger, so keep it to one or two specific sentences.
agents: {
assistant: defineAgent({
skills: [
{
name: "refund-policy",
description: "Rules and tone for handling refund requests.",
body: "# Refunds\nRefund within 14 days, no questions asked. Address {{userName}} by name once.",
},
],
}),
},Names are lowercase slugs, unique per agent. Bodies are markdown; {{variables}} are filled per request the same way the prompt's are. Skills published in the Control Center join the same list; a code skill wins a name collision.
Triggers
A trigger is a named entry point external systems fire on an agent. Each firing validates its payload, lands in a thread the trigger owns, and starts a run with no user attached. The first message of that run is the trigger's instructions rendered with {{variables}} and {{/json/pointer}} payload fields, followed by the payload as a JSON block. The end user's own message is the built-in message trigger, on by default; messageTrigger: false on the agent makes it headless.
import { defineAgent, defineTrigger } from "@m6d/cortex-server";
import { z } from "zod";
agents: {
reviewer: defineAgent({
triggers: [
defineTrigger({
name: "pr-opened",
description: "A pull request was opened.",
input: z.object({ pull_request: z.object({ number: z.number(), url: z.string() }) }),
instructions: "Review pull request #{{/pull_request/number}} at {{/pull_request/url}}.",
output: z.object({
verdict: z.enum(["approve", "request_changes"]),
summary: z.string(),
}),
threadPolicy: { policy: "keyed", key: "/pull_request/number" },
secret: process.env.PR_TRIGGER_SECRET!,
credential: process.env.PR_REVIEW_TOKEN,
budgetPerDay: 200,
}),
],
}),
},output makes the run answer in that shape: the tool loop runs as usual, then one final call asks the model for the shape over the finished conversation, the answer is validated against it, and it is stored as the firing's output. A run whose answer does not parse fails with reason invalid_output. Without output the run is free text and output stays null.
Fire it with the secret as the bearer. Ask for text/event-stream to follow the run live, or ?wait=<ms> to hold the request for the run and get the outcome with its output; a run that outlasts the wait (capped by the server's triggerWaitMaxMs, five minutes by default) still answers started, as does one whose outcome could not be written to the firing, and GET .../firings/<firingId> returns the same record once it ends. A reverse proxy in front needs a read timeout above the wait.
curl -X POST 'https://your-server/agents/reviewer/triggers/pr-opened?wait=120000' \
-H "Authorization: Bearer $PR_TRIGGER_SECRET" \
-H "Content-Type: application/json" \
-d '{"pull_request":{"number":42,"url":"https://github.com/org/repo/pull/42"}}'| Status | Body | Meaning |
| ------ | ------------------------------------------------------------------------ | ---------------------------------------------------------------------------------------------------------------- |
| 202 | { firingId, outcome: "started", threadId, runId } | The run is going, or outlasted wait. |
| 202 | { firingId, outcome: "awaiting_approval", threadId, runId, approvals } | With wait: the run parked on tools needing approval; answer through the firing route. |
| 200 | { firingId, outcome: "completed", threadId, runId, output } | With wait: the run ended; output is the validated answer, or null. |
| 200 | { firingId, outcome: "failed", threadId, runId, reason } | With wait: the run failed; invalid_output is an answer the schema refused, no_output a run that gave none. |
| 202 | { firingId, outcome: "dropped", reason: "busy" } | The thread is mid-run. Nothing is queued. |
| 202 | { firingId, outcome: "dropped", reason: "budget" } | Today's budgetPerDay is spent. |
| 202 | { firingId, outcome: "disabled" } | Auto-disabled after disableAfterFailures (5). |
| 422 | { firingId, outcome: "rejected", issues } | The payload failed input. |
| 401 | | Wrong or missing secret. |
| 404 | | No such trigger. |
threadPolicy picks where firings land: single (default) shares one thread, per-firing opens a new one each time, ephemeral opens a new one and deletes it once the run completes (a one-shot call: the answer reaches the caller as the firing's output or through the live stream, and nothing else outlives the run, so a caller that neither streams nor has an output schema gets no answer; a failed run keeps its thread, and the threadId a completed ?wait= answers is already gone), keyed groups by a payload value. A busy thread drops the newcomer for every trigger, message included; nothing aborts a run except the abort endpoint. A firing runs with nobody at a widget, so the agent's client tools are not declared to it. GET .../triggers/<name>/firings lists recent firings (each with its output) and GET .../triggers/<name>/firings/<firingId> reads one, with the same bearer.
Cortex ships no scheduler or webhook receiver: cron, one-shot timers, and webhooks are whatever fires the endpoint. Triggers published in the Control Center join the same list, verified through cc; a code trigger wins a name collision.
Authentication
The top-level auth verifies end-user tokens for the whole server; the verified subject becomes the user id and the claims the verifier read travel with the request. Without it, every request shares one anonymous identity — fine for local development, wrong for anything multi-user. auth is one of three kinds:
| kind | Fields | Verifies |
| -------------- | --------------------------------------------- | ------------------------------------------------------------------------------------------------------------------------------------------------------------ |
| "jwks" | jwksUri, issuer, audience? | Asymmetric JWTs (RS256/ES256/EdDSA, with a kid header) against the keys the issuer publishes at the JWKS URL. |
| "jwt-secret" | secret, algorithms, issuer, audience? | Symmetric JWTs (HS256/HS384/HS512) against a shared secret of at least 32 characters; algorithms pins what is accepted. |
| "custom" | verify(token, request) | Whatever the host says: return { sub, claims } for a caller, null for nobody. For opaque tokens, a session store, anything that is not a JWT. Code only. |
An agent can carry its own auth, either in its agents entry or configured on the agent in the Control Center (the code entry wins; the cc value reaches the server through the same config sync as the rest of the agent, with a jwt-secret read from cc's secret store at publish time). On such an agent, identity comes only from tokens that verifier accepts — the server-level default and the anonymous fallback no longer apply, and requests without a verifiable token get a 401. cc playground tokens remain valid so the Control Center can test auth-enabled agents.
Set the optional audience when the issuer mints tokens for more than one application: the token's aud claim must then contain it, so a token issued for a different audience under the same issuer is rejected. Tokens need exact iss and exp claims. An owner-supplied jwksUri is fetched under an SSRF policy that refuses loopback, private, link-local, and cloud-metadata destinations, pins DNS against rebinding, and rejects redirects; a self-hosted IdP on a private network is reached by setting EGRESS_ALLOW_PRIVATE_NETWORK=true (development only).
Tokens travel as Authorization: Bearer <token>; WebSocket upgrades use ?token=<token> since browsers cannot set headers there. resolveRequestContext(request, { userId, claims }) sees the verified identity on every request, so a host that gates tools on a role puts the claim into the request context there; the tools factory reads it from requestContext on every turn, unlike session, which is loaded once per thread.
Client tools
A tool in an agent's tools with an execute function runs on the server. One
without is a static client tool: it is declared to the model, the call streams
to the widget, and the host app answers it via hooks.onToolCall (or a
toolComponents entry calling setOutput).
const getBrowserTimezone = toolDefinition({
name: "getBrowserTimezone",
description: "Read the user's IANA time zone from their browser.",
inputSchema: z.object({}),
}).client();
// agents.sample.tools = [getBrowserTimezone]A third kind sits between the two: a server tool whose call needs the user's
decision first. Declare it with the SDK's needsApproval: true and, spread
onto the tool, an approval sentence per locale; the model calls it directly
(never from executeCode), the run parks with the call marked
approval-requested, the widget asks with that sentence (never the function's
name or its arguments, which only the debug view shows) and offers approve and
deny, and the decision comes back as a continuation: approved runs the tool's
execute server-side, denied records a refusal the model reads. A tool with
the flag and no sentence is asked about generically; an approval on a tool
that never parks is a configuration error.
const escalateToHuman = {
...toolDefinition({ name: "escalateToHuman", inputSchema, needsApproval: true }).server(run),
approval: {
en: "Hand this conversation to a person?",
ar: "تحويل هذه المحادثة إلى موظف بشري؟",
},
};A Control Center http or soap tool flagged "Requires approval" behaves the
same and runs through cc once approved. The
agent's onApprovalRequested hook fires when a run parks, for a host that
decides elsewhere; a headless firing reports awaiting_approval and is resumed
through POST .../triggers/<name>/firings/<firingId>/answer with
{ toolCallId, approved }, the same bearer and ?wait= semantics as a fire.
Of two decisions on one firing, one resumes it and the other gets 409; a
decision that finds the thread busy with another run answers dropped busy
and leaves the firing parked, the call undecided, for a retry.
Like client tools, a tool needing approval is declared at the start of a turn,
from the agent's own tools, a cc trigger's mentions, and what retrieval finds
for the message. One that a mid-turn search or a loaded skill surfaces is not
callable in that turn; searches leave it out, and a skill's prose that names it
reaches the model as-is, so it runs on a later turn only when retrieval or a
mention declares it.
| Tool | Declared to the model | Runs | Answered by |
| --------------------------------- | --------------------- | ------------------------------------- | ------------------------------------------------------ |
| Server (.server(execute)) | directly | on the server, at once | itself |
| Server with needsApproval | directly | on the server, once approved | the widget, onApprovalRequested, or the firing route |
| Client (.client(), cc embedded) | directly | in the browser | hooks.onToolCall, a tool component, the embed |
| cc http/soap, MCP | via executeCode | through cc or the MCP server, at once | the sandbox |
A server declares the {{variables}} it fills per request for the agents the
Control Center creates on it, in a server-level variables (name, description,
example); the values come from the server-level loadSessionData and
resolveRequestContext, and the SDK adds utcTime, locale and timezone (the
browser's IANA zone, sent by the widget as X-Cortex-Timezone; a consumer's own
value wins). The turn states the user's local time beside the UTC one, and
CURRENT_DATE in a dataset query is today in that zone. At boot the
server sends that declaration to the Control Center, which offers the variables
to every cc-created agent on this server. Agents in this server's own agents
entries are served from code alone and never reach the Control Center. In cc a
prompt may reference a declared variable, and a tool input field may take its
value from one instead of the model ({ emiratesId: "784…" } fills a field
sourced from emiratesId); the server sends only the variables a tool reads and
the model never sees the field. Source identity that way rather than forwarding
the end-user token when the upstream cannot verify it.
Control-Center tools published with type "embedded" are the dynamic client
tools: declared the same way, but the widget first calls
POST /chat/:chatId/tools/:toolCallId/initiate, which runs the tool's endpoint
server-to-server and hands the widget its embed payload (never the model).
The result the embedded page reports back is relayed to the agent as-is —
verifying it is the integrating backend's job. See
apps/cortex-cc/docs/client-tools.md for authoring them.
Dataset hosts
A dataset is a table an application exposes for the agent to ask grouped and joined questions of, with the application's own access rules applied before the query. The application becomes a dataset host by serving two routes; Cortex discovers the manifest, validates every query against it, caps it, forwards it under the user's identity, and logs the read. Cortex never decides who may see which rows.
GET {host}/ -> { version, datasets: [{ name, description, columns, relations, maxRows }] }
POST {host}/{name} <- { query, context: { threadId } }
-> { rows, rowCount, truncated } | { error: { kind, detail } }The wire query names a dataset, joins related datasets by relation name (any
depth, through for many-to-many) and reaches their columns as
relation.column. Its JSON nests at most 32 objects and arrays deep; every
reader refuses a deeper query before parsing the rest. Every expression a host
translates (computed values, conditions, select values, group keys, aggregates
and windows) grows to at most 2000 nodes, each part counted as often as a host
writes it out, once the computed aliases and aggregates it names are written
out in place, as hosts do; a window also counts the row it carries. Everything
else is built from two shapes:
- An operand is a column reference, a number, a text or boolean literal
{ value }, arithmetic{ op, left, right }(division is exact), a date truncated to a period{ fn: "dateTrunc", unit, of }(answered as the ISO date of the period's first day, weeks starting Monday) or one of its parts{ fn: "datePart", unit, of }(an integer;dayOfWeekis 1 for Monday to 7 for Sunday), a scalar function{ fn, args }(lower,upper,trim,length,abs,round,coalesce), or{ case: [{ when, then }], else }whose branches share one type. - A filter is a typed condition over an operand (the right side a literal or
{ col: operand }), or{ any },{ all },{ not }over filters.notkeeps SQL's null semantics: a row whose comparison is unknown is left out either way.
With those, a query selects columns or named values { expr, as }, filters
rows, groupBys columns or named values, aggregates (count,
countDistinct, sum, avg, min, max over an operand, each with an
optional filter over its own rows), computes named values per row from the
aggregate aliases, group keys and window aliases, filters the grouped rows with
having (which sees no window), adds windows (an aggregate over the answered
rows sharing a partitionBy, answered on each of them: SQL's fn(col) OVER
(PARTITION BY ...), on a plain or a grouped query, over the row's own names),
filters the finished rows with qualify (window and computed aliases included;
the host wraps the select in a subquery, since none of its dialects has
QUALIFY), keeps the rows
ranked 1 to n in each partition with rank (SQL's RANK() OVER (PARTITION BY
... ORDER BY ...) <= n, ranked after qualify; rows tied on every sort key
share a rank and all stay), then sorts, offsets
(only with a sort) and limits. Row keys are the column reference
(department.name), the group key's or select value's alias, the aggregate
alias or the computed alias. The SQL layer then reshapes each row the way the
statement wrote it, nesting departments.name as { departments: { name } }
and keeping aliases at the top level. The model never writes the wire shape:
it writes one SQL SELECT with JOIN dataset ON a.key = b.key as
describeDataset lists the joins (both JOINs of a many-to-many the same kind,
since the link table's own columns are never read), and the server compiles it
(parseDatasetSql in the contracts). The conformance fixtures in
internal/contracts/src/runtime/datasets.fixtures.json pin every shape's rows;
each host runs them. Error kinds: invalid_query, unsupported, forbidden,
timeout, internal, each with a fixed HTTP status. The end user's token
travels in the header the agent's datasetHosts entry names in userToken
and the host reads from its own userToken (both default to
X-End-User-Token); a headless run sends the trigger's credential instead.
A TypeScript host needs one line per dataset over drizzle, plus its rule:
import { defineDatasetHost, drizzleDataset } from "@m6d/cortex-server/datasets";
const orders = drizzleDataset<User>({
name: "orders",
description: "One row per order: who bought what, for how much, and when.",
table: demoOrders,
dialect: "postgres",
db,
// The rule the application already has; undefined means every row.
scoped: (user) => eq(demoOrders.customerId, user.customerId),
columns: { exclude: ["internalNotes"] },
descriptions: { total: "Order total after quantity" },
relations: () => ({ pet: { dataset: pets, on: ["petId", "id"] } }),
});
app.route(
"/api/datasets",
defineDatasetHost<User>({
expose: [orders, pets],
resolveUser: ({ token }) => verifyAndLoad(token), // or throw new DatasetHostError("forbidden", ...)
}),
);Columns and their types derive from the table (types overrides one the table
cannot describe, such as a date kept as ISO text); relations are a function so
datasets may refer to each other. resolveUser is required: what it returns
is the user every scoped rule receives. The host refuses a request body over
256 KiB before reading it, clamps every query to the dataset's maxRows, runs
it under timeoutMs (30 seconds by default), and answers truncated when
more rows matched. When timeoutMs passes the host answers timeout, but
drizzle cannot cancel the statement, so give db a server-side statement
timeout no longer than timeoutMs (Postgres statement_timeout, SQL Server
requestTimeout). Another query builder implements the DatasetSource
interface: a definition and a run(query, user, signal) that returns rows.
Other stacks implement the two routes over the same JSON; .NET hosts use the
M6d.Cortex package.
On the agent side, point at hosts; never redefine their datasets:
defineAgent({
datasetHosts: [
{
name: "store",
url: "https://store.example/api/datasets",
auth: { headers: { "X-Service-Api-Key": process.env["STORE_KEY"]! } },
expose: ["orders", "pets"], // every dataset when absent
pinned: ["orders"], // listed in every turn's context
},
],
});The manifest is read at the first turn and again when it is older than five
minutes; a host that cannot be read keeps serving its last manifest. A
credential the host echoes back (an auth.headers value, or the user's token
bare or with its prefix) reads [REDACTED] in the manifest, the answer and the
logged failure. The
server stops reading a manifest past 4 MiB and an answer past 8 KiB per row of
the dataset's maxRows, and refuses either. The model
sees one line per dataset under ## Datasets in the turn context and gets
queryDataset (one SQL SELECT), describeDataset, and a datasets global inside
executeCode (datasets.query(sql), datasets.describe()) for reads a script
joins or compares. Each of the three tools takes a status sentence the model
writes in the user's words ("Counting active KPIs per department"); the widgets
show it while the call runs and keep it as the call's mark, with the rows a read
answered. Every read, refused or answered, is mirrored to the Control Center as a
dataset.read event with the host, dataset, the SQL the model wrote, row count,
truncated, duration, outcome and, when it did not answer, the refusal or the
host's detail; never the rows.
MCP servers
An agent's mcpServers are discovered on every turn over Streamable HTTP, with
the server's tool names prefixed by the entry's name:
mcpServers: [
{
name: "cc",
url: "https://cc.example.com/api/mcp",
auth: "forwardUserToken",
exclude: ["secrets_set", "secrets_rotate"],
},
],MCP tools are never declared to the model. Like Control Center's tools they are
called from executeCode through the tools global, as
await tools.cc_secrets_list(input), and they are listed in the system
prompt's Dynamic Tools section with a signature rendered from the schemas the
server published, in the same form Control Center renders its own tools (the
input typed field by field, the response as the field names it carries). A
server that declares an outputSchema for a tool gives the model that tool's
response shape; without one the signature returns unknown and the script has
to inspect the result. A tool that answers with text content parts
instead of structured content reaches the script as that text, so a server
returning JSON that way hands over a string the script parses itself; a result
carrying anything else (an image, an embedded resource) arrives as the raw
array of parts. A tool whose name is not a JavaScript identifier
(jira_get-issue, files.read) is listed as tools["jira_get-issue"](input),
since dot access cannot reach it.
auth is what the MCP server sees on each request: "forwardUserToken" sends
the caller's own bearer unchanged, for a server that authorizes per end user;
{ headers } sends fixed headers, for an API key; a function of the tool
context computes them per request. A server with auth must be https
(plain http is allowed only on localhost), or it is reported unavailable
rather than handed a credential over the network. include and exclude
filter by the server's own (unprefixed) tool names.
Server-level entries reach every agent; an agent entry with the same name
replaces it. An MCP tool never replaces anything: one named like a built-in
(loadSkill, executeCode), like one of the agent's code tools, or like a
Control Center tool already in the tools global is dropped, since a remote
server must not be able to stand in for code the agent author wrote. Two
servers whose prefixes collide on a name (a + b_c and a_b + c) resolve
in configured order: the earlier server keeps the name, the later one's tool is
dropped. A server that cannot be reached contributes no
tools that turn and is logged; it never fails the turn. Connections open per turn (the
forwarded bearer belongs to that caller) and close when the run ends.
The URL is fetched under the same egress policy as auth.jwksUri: a private
host needs EGRESS_ALLOW_PRIVATE_NETWORK=true in development.
Requirements
- Bun >= 1.0.0
