Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion packages/mcp/src/bin.ts
Original file line number Diff line number Diff line change
Expand Up @@ -12,7 +12,7 @@ if (!baseUrl) {
process.exit(1);
}

const server = createOmnigraphMcpServer({
const server = await createOmnigraphMcpServer({
baseUrl,
token: process.env.OMNIGRAPH_TOKEN,
defaultBranch: process.env.OMNIGRAPH_DEFAULT_BRANCH,
Expand Down
133 changes: 131 additions & 2 deletions packages/mcp/src/server.ts
Original file line number Diff line number Diff line change
Expand Up @@ -15,9 +15,11 @@ import { McpServer } from '@modelcontextprotocol/sdk/server/mcp.js';
import {
Omnigraph,
type FetchLike,
type SavedQuery,
type SavedQueryParam,
SERVER_VERSION as SDK_SERVER_VERSION,
} from '@modernrelay/omnigraph';
import { z } from 'zod';
import { z, type ZodTypeAny } from 'zod';
import { COOKBOOK } from './best-practices.gen';

const INSTRUCTIONS = `Omnigraph is a versioned property graph. Reads are typed GQ queries; writes are server-orchestrated and branchable.
Expand All @@ -44,6 +46,8 @@ Workflow norms (violating these breaks things or silently corrupts data):

Date format: ISO strings on \`change\` params; integer days-since-epoch in ingest JSONL \`Date\` fields. \`DateTime\` is ISO on both.

Saved queries: tools prefixed \`q_\` (e.g. \`q_find_person\`) are user-authored .gq queries the operator has persisted via \`PUT /queries/{name}\`. They dispatch through \`read\` with the saved source. Call them like any other tool, passing the declared params as named args. The list at session start is a snapshot — saved queries added or deleted after startup are not visible until the MCP reconnects. Browse \`omnigraph://queries\` for the full source if a tool description is not enough.

If you see \`sync_branch()\` in an error message, it is server-internal text, NOT a tool. Retry once; on persistent failure, fall back to \`ingest\` on a branch.

Depth: https://github.com/ModernRelay/omnigraph-cookbooks/tree/main/skills/omnigraph-best-practices`;
Expand All @@ -67,7 +71,67 @@ function plainText(text: string) {
return [{ type: 'text' as const, text }];
}

export function createOmnigraphMcpServer(opts: CreateServerOptions): McpServer {
// Prefix saved-query tools so they cannot shadow built-in tool names. The
// server side also rejects reserved names like `read`, but the prefix is
// the defence-in-depth and keeps the catalogue obviously partitioned for
// human readers.
const SAVED_QUERY_TOOL_PREFIX = 'q_';

// Map an Omnigraph scalar/composite type name (as it appears in `.gq`
// source — `String`, `I32`, `Vector(3072)`, etc.) onto a permissive zod
// schema. The MCP boundary is only doing argument-shape validation; the
// server still parses and typechecks the query against the live schema,
// so any tighter checking here would be duplicative and could reject
// future scalars we have not seen yet.
function paramTypeToZod(typeName: string): ZodTypeAny {
switch (typeName) {
case 'String':
return z.string();
case 'Bool':
return z.boolean();
case 'I32':
case 'I64':
case 'U32':
case 'U64':
return z.number().int();
case 'F32':
case 'F64':
return z.number();
case 'Date':
case 'DateTime':
return z.string();
default:
// Vector(N), Blob, [String], future scalars — pass through opaque.
return z.unknown();
}
}

function savedQueryInputShape(params: SavedQueryParam[]): Record<string, ZodTypeAny> {
const shape: Record<string, ZodTypeAny> = {};
for (const p of params) {
const base = paramTypeToZod(p.typeName);
shape[p.name] = p.nullable ? base.optional() : base;
}
// Every saved query is dispatched through `/read`, so the caller may
// also pin a branch or a snapshot at invocation time without having
// to bake it into the saved source.
shape.branch = z.string().optional();

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2: Saved-query parameters named branch or snapshot are silently shadowed by transport-level options, making it impossible to pass those values as actual query parameters.

Prompt for AI agents
Check if this issue is valid — if so, understand the root cause and fix it. At packages/mcp/src/server.ts, line 118:

<comment>Saved-query parameters named `branch` or `snapshot` are silently shadowed by transport-level options, making it impossible to pass those values as actual query parameters.</comment>

<file context>
@@ -67,7 +71,67 @@ function plainText(text: string) {
+  // Every saved query is dispatched through `/read`, so the caller may
+  // also pin a branch or a snapshot at invocation time without having
+  // to bake it into the saved source.
+  shape.branch = z.string().optional();
+  shape.snapshot = z.string().optional();
+  return shape;
</file context>

shape.snapshot = z.string().optional();
return shape;
}

function savedQueryDescription(q: SavedQuery): string {
const params =
q.params.length === 0
? 'no parameters'
: q.params
.map((p) => `${p.name}: ${p.typeName}${p.nullable ? '?' : ''}`)
.join(', ');
const head = q.description ?? `Saved query \`${q.name}\``;
return `${head} — params: ${params}. Dispatched through \`read\`; pass values in the matching argument names, optionally override \`branch\` or pin \`snapshot\`.`;
}

export async function createOmnigraphMcpServer(opts: CreateServerOptions): Promise<McpServer> {
const og = new Omnigraph({ baseUrl: opts.baseUrl, token: opts.token, fetch: opts.fetch });
const defaultBranch = opts.defaultBranch ?? 'main';

Expand Down Expand Up @@ -373,6 +437,71 @@ export function createOmnigraphMcpServer(opts: CreateServerOptions): McpServer {
);
}

// ---------- Saved queries -----------------------------------------------
// Pull whatever the user has saved server-side via `PUT /queries/{name}`
// and register one MCP tool per saved query, plus a list resource. This
// is best-effort: if the server is older than 0.4.3 the endpoint will
// 404, and if the network is flaky the list call will throw. In either
// case we keep the rest of the server intact rather than failing
// startup — operators get a stderr line so the absence is loud.
let savedQueries: SavedQuery[] = [];
try {
const listed = await og.queries.list();
if (Array.isArray(listed)) {
savedQueries = listed;
}
} catch (err) {
const message = err instanceof Error ? err.message : String(err);
console.warn(
`omnigraph-mcp: failed to list saved queries (${message}); skipping per-query tool registration`,
);
}
for (const saved of savedQueries) {
const toolName = `${SAVED_QUERY_TOOL_PREFIX}${saved.name}`;
server.registerTool(
toolName,
{
title: saved.description ?? `Saved query: ${saved.name}`,
description: savedQueryDescription(saved),
inputSchema: savedQueryInputShape(saved.params),
annotations: { readOnlyHint: true, openWorldHint: false },
},
async (input: Record<string, unknown>) => {
const { branch, snapshot, ...params } = input;
const target =
typeof snapshot === 'string'
? { snapshot }
: { branch: (typeof branch === 'string' ? branch : undefined) ?? defaultBranch };
const r = await og.read({
querySource: saved.source,
queryName: saved.name,
params,
...target,
});
return { content: jsonText(r) };
},
);
}
server.registerResource(
'queries',
'omnigraph://queries',
{
title: 'Saved queries',
description:
'JSON list of every saved query (name, description, source, params). Mirrors what is exposed as `q_<name>` tools — read this if you want to pick a query by description rather than scrolling the tool catalogue.',
mimeType: 'application/json',
},
async (uri) => ({
contents: [
{
uri: uri.href,
mimeType: 'application/json',
text: JSON.stringify(savedQueries, null, 2),
},
],
}),
);

// Index resource: a single small markdown that lists every cookbook
// reference + its purpose. An agent that has not yet decided which
// reference it needs can read this one cheap entry to orient.
Expand Down
144 changes: 140 additions & 4 deletions packages/mcp/test/server.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -57,12 +57,15 @@ function fakeFetch(): typeof globalThis.fetch {
],
});
}
if (method === 'GET' && path === '/queries') {
return respond(200, { queries: [] });
}
return respond(404, { error: 'not found', code: 'not_found' });
}) as unknown as typeof globalThis.fetch;
}

async function setup() {
const server = createOmnigraphMcpServer({ baseUrl: 'http://x', fetch: fakeFetch() });
async function setup(fetchImpl: typeof globalThis.fetch = fakeFetch()) {
const server = await createOmnigraphMcpServer({ baseUrl: 'http://x', fetch: fetchImpl });
const client = new Client({ name: 'test-client', version: '0.0.0' });
const [clientT, serverT] = InMemoryTransport.createLinkedPair();
await Promise.all([server.connect(serverT), client.connect(clientT)]);
Expand Down Expand Up @@ -171,6 +174,7 @@ describe('omnigraph-mcp server', () => {
'omnigraph://best-practices/schema',
'omnigraph://best-practices/search',
'omnigraph://branches',
'omnigraph://queries',
'omnigraph://schema',
].sort(),
);
Expand Down Expand Up @@ -216,7 +220,7 @@ it('branches_create honours configured defaultBranch when `from` is omitted', as
return new Response('{}', { status: 200, headers: { 'content-type': 'application/json' } });
}) as unknown as typeof globalThis.fetch;

const server = createOmnigraphMcpServer({
const server = await createOmnigraphMcpServer({
baseUrl: 'http://x',
defaultBranch: 'review-2026',
fetch: recordingFetch,
Expand Down Expand Up @@ -251,7 +255,7 @@ it('branches_create honours configured defaultBranch when `from` is omitted', as
return new Response('{}', { status: 200, headers: { 'content-type': 'application/json' } });
}) as unknown as typeof globalThis.fetch;

const server = createOmnigraphMcpServer({
const server = await createOmnigraphMcpServer({
baseUrl: 'http://x',
defaultBranch: 'main',
fetch: recordingFetch,
Expand All @@ -274,4 +278,136 @@ it('branches_create honours configured defaultBranch when `from` is omitted', as
const r = await client.callTool({ name: 'read', arguments: {} });
expect(r.isError).toBe(true);
});

it('registers one q_<name> tool per saved query and dispatches through /read', async () => {
let readBody: Record<string, unknown> | undefined;
const fetchWithSavedQuery: typeof globalThis.fetch = (async (
input: RequestInfo | URL,
init?: RequestInit,
) => {
const url =
typeof input === 'string' ? input : input instanceof URL ? input.toString() : input.url;
const path = new URL(url).pathname;
const method = init?.method ?? 'GET';
if (method === 'GET' && path === '/queries') {
return new Response(
JSON.stringify({
queries: [
{
name: 'find_person',
description: 'by name',
source:
'query find_person($name: String) { match { $p: Person { name: $name } } return { $p.name } }',
params: [{ name: 'name', type_name: 'String', nullable: false }],
updated_at_us: '1747315200000000',
},
],
}),
{ status: 200, headers: { 'content-type': 'application/json' } },
);
}
if (method === 'POST' && path === '/read') {
readBody = JSON.parse(typeof init?.body === 'string' ? init.body : '{}');
return new Response(
JSON.stringify({
query_name: 'find_person',
target: { branch: 'main', snapshot: null },
row_count: 1,
columns: ['$p.name'],
rows: [{ '$p.name': 'Alice' }],
}),
{ status: 200, headers: { 'content-type': 'application/json' } },
);
}
return new Response('{}', { status: 200, headers: { 'content-type': 'application/json' } });
}) as unknown as typeof globalThis.fetch;

const { client } = await setup(fetchWithSavedQuery);
const tools = await client.listTools();
const names = tools.tools.map((t) => t.name);
expect(names).toContain('q_find_person');

const result = await client.callTool({
name: 'q_find_person',
arguments: { name: 'Alice' },
});
expect(result.isError).toBeFalsy();
expect(readBody?.query_source).toContain('query find_person');
expect(readBody?.params).toEqual({ name: 'Alice' });
expect(readBody?.branch).toBe('main');
});

it('q_<name> tool routes to a snapshot when one is passed', async () => {
let readBody: Record<string, unknown> | undefined;
const fetchImpl: typeof globalThis.fetch = (async (
input: RequestInfo | URL,
init?: RequestInit,
) => {
const url =
typeof input === 'string' ? input : input instanceof URL ? input.toString() : input.url;
const path = new URL(url).pathname;
const method = init?.method ?? 'GET';
if (method === 'GET' && path === '/queries') {
return new Response(
JSON.stringify({
queries: [
{
name: 'all_people',
description: null,
source: 'query all_people() { match { $p: Person } return { $p.name } }',
params: [],
updated_at_us: '1747315200000000',
},
],
}),
{ status: 200, headers: { 'content-type': 'application/json' } },
);
}
if (method === 'POST' && path === '/read') {
readBody = JSON.parse(typeof init?.body === 'string' ? init.body : '{}');
return new Response(
JSON.stringify({
query_name: 'all_people',
target: { branch: null, snapshot: 'snap-1' },
row_count: 0,
columns: [],
rows: [],
}),
{ status: 200, headers: { 'content-type': 'application/json' } },
);
}
return new Response('{}', { status: 200, headers: { 'content-type': 'application/json' } });
}) as unknown as typeof globalThis.fetch;

const { client } = await setup(fetchImpl);
await client.callTool({ name: 'q_all_people', arguments: { snapshot: 'snap-1' } });
expect(readBody?.snapshot).toBe('snap-1');
expect(readBody?.branch).toBeUndefined();
});

it('continues to start when /queries returns 404 (older server)', async () => {
const fetchOldServer: typeof globalThis.fetch = (async (
input: RequestInfo | URL,
init?: RequestInit,
) => {
const url =
typeof input === 'string' ? input : input instanceof URL ? input.toString() : input.url;
const path = new URL(url).pathname;
const method = init?.method ?? 'GET';
if (method === 'GET' && path === '/queries') {
return new Response(JSON.stringify({ error: 'not found', code: 'not_found' }), {
status: 404,
headers: { 'content-type': 'application/json' },
});
}
return new Response('{}', { status: 200, headers: { 'content-type': 'application/json' } });
}) as unknown as typeof globalThis.fetch;

// Should not throw; built-in tools still registered.
const { client } = await setup(fetchOldServer);
const tools = await client.listTools();
const names = tools.tools.map((t) => t.name);
expect(names).toContain('read');
expect(names.find((n) => n.startsWith('q_'))).toBeUndefined();
});
});
3 changes: 3 additions & 0 deletions packages/sdk/src/client.ts
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@ import type { FetchLike } from './transport';
import { BranchesResource } from './resources/branches';
import type { CallOptions } from './internals';
import { CommitsResource } from './resources/commits';
import { QueriesResource } from './resources/queries';
import { SchemaResource } from './resources/schema';
import type {
Change,
Expand Down Expand Up @@ -44,6 +45,7 @@ export interface SnapshotInput {
export default class Omnigraph {
readonly branches: BranchesResource;
readonly commits: CommitsResource;
readonly queries: QueriesResource;
readonly schema: SchemaResource;

private readonly t: Transport;
Expand All @@ -52,6 +54,7 @@ export default class Omnigraph {
this.t = new Transport(opts);
this.branches = new BranchesResource(this.t);
this.commits = new CommitsResource(this.t);
this.queries = new QueriesResource(this.t);
this.schema = new SchemaResource(this.t);
}

Expand Down
Loading