Skip to content
Merged
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
248 changes: 244 additions & 4 deletions src/services/executors/__tests__/claude-sdk.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -26,7 +26,7 @@ import { afterEach, beforeEach, describe, expect, it, mock } from 'bun:test';
// — recordAuditEvent / session-capture route through safePgCall, which is
// mocked per-test via setSafePgCall() (or guarded null for no-bridge tests).

const { ClaudeSdkOmniExecutor } = await import('../claude-sdk.js');
const { ClaudeSdkOmniExecutor, createSendMessageOmniHook } = await import('../claude-sdk.js');
const directory = await import('../../../lib/agent-directory.js');

describe('ClaudeSdkOmniExecutor', () => {
Expand Down Expand Up @@ -474,7 +474,7 @@ describe('ClaudeSdkOmniExecutor', () => {
// Verify turn-based prompt content
expect(systemPrompt).toContain('WhatsApp');
expect(systemPrompt).toContain('Alice');
expect(systemPrompt).toContain('omni say');
expect(systemPrompt).toContain('SendMessage');
expect(systemPrompt).toContain('omni done');
expect(systemPrompt).toContain('inst-wb');
});
Expand All @@ -496,8 +496,248 @@ describe('ClaudeSdkOmniExecutor', () => {
const systemPrompt = callArgs.options?.systemPrompt ?? '';

expect(systemPrompt).not.toContain('WhatsApp');
expect(systemPrompt).not.toContain('omni say');
expect(systemPrompt).not.toContain('omni done');
expect(systemPrompt).not.toContain('SendMessage');
});
});

// ==========================================================================
// SendMessage(to: omni) NATS routing — issue #1088
// ==========================================================================

describe('SendMessage omni interception', () => {
const omniEnv = { OMNI_INSTANCE: 'inst-x', OMNI_CHAT: 'chat-x', OMNI_AGENT: 'eugenia' };
const dummyMeta = { signal: new AbortController().signal };

it('publishes to omni.reply.{instance}.{chat} and denies the tool when recipient is "omni"', async () => {
const publishCalls: Array<[string, string]> = [];
const natsPublish = (subject: string, payload: string) => {
publishCalls.push([subject, payload]);
};
const hook = createSendMessageOmniHook(omniEnv, natsPublish);

const result = await hook(
{
hook_event_name: 'PreToolUse',
tool_name: 'SendMessage',
tool_input: { recipient: 'omni', message: 'Olá Cezar!' },
tool_use_id: 'tu-1',
} as never,
'tu-1',
dummyMeta,
);

expect(publishCalls).toHaveLength(1);
expect(publishCalls[0][0]).toBe('omni.reply.inst-x.chat-x');
const payload = JSON.parse(publishCalls[0][1]);
expect(payload).toMatchObject({
agent: 'eugenia',
chat_id: 'chat-x',
instance_id: 'inst-x',
content: 'Olá Cezar!',
});
expect(result).not.toBeNull();
const out = result as { hookSpecificOutput?: { permissionDecision?: string; permissionDecisionReason?: string } };
expect(out.hookSpecificOutput?.permissionDecision).toBe('deny');
expect(out.hookSpecificOutput?.permissionDecisionReason).toMatch(/delivered.*omni bridge/i);
});

it('accepts the alternate "to" + "content" field shape', async () => {
const publishCalls: Array<[string, string]> = [];
const hook = createSendMessageOmniHook(omniEnv, (s, p) => {
publishCalls.push([s, p]);
});

await hook(
{
hook_event_name: 'PreToolUse',
tool_name: 'SendMessage',
tool_input: { to: 'omni', content: 'second shape' },
tool_use_id: 'tu-2',
} as never,
'tu-2',
dummyMeta,
);

expect(publishCalls).toHaveLength(1);
expect(JSON.parse(publishCalls[0][1]).content).toBe('second shape');
});

it('passes through (no decision) when recipient is not "omni"', async () => {
const publishCalls: Array<[string, string]> = [];
const hook = createSendMessageOmniHook(omniEnv, (s, p) => {
publishCalls.push([s, p]);
});

const result = await hook(
{
hook_event_name: 'PreToolUse',
tool_name: 'SendMessage',
tool_input: { recipient: 'reviewer', message: 'peer ping' },
tool_use_id: 'tu-3',
} as never,
'tu-3',
dummyMeta,
);

expect(publishCalls).toHaveLength(0);
expect((result as Record<string, unknown>).hookSpecificOutput).toBeUndefined();
});

it('passes through for non-SendMessage tool calls', async () => {
const publishCalls: Array<[string, string]> = [];
const hook = createSendMessageOmniHook(omniEnv, (s, p) => {
publishCalls.push([s, p]);
});

const result = await hook(
{
hook_event_name: 'PreToolUse',
tool_name: 'Read',
tool_input: { file_path: '/tmp/x' },
tool_use_id: 'tu-4',
} as never,
'tu-4',
dummyMeta,
);

expect(publishCalls).toHaveLength(0);
expect((result as Record<string, unknown>).hookSpecificOutput).toBeUndefined();
});

it('denies with bridge-unavailable reason when natsPublish is null but env is set', async () => {
const hook = createSendMessageOmniHook(omniEnv, null);

const result = await hook(
{
hook_event_name: 'PreToolUse',
tool_name: 'SendMessage',
tool_input: { recipient: 'omni', message: 'oops' },
tool_use_id: 'tu-5',
} as never,
'tu-5',
dummyMeta,
);

const out = result as { hookSpecificOutput?: { permissionDecision?: string; permissionDecisionReason?: string } };
expect(out.hookSpecificOutput?.permissionDecision).toBe('deny');
expect(out.hookSpecificOutput?.permissionDecisionReason).toMatch(/bridge unavailable/i);
});

it('denies with retry-prompting reason when message body is missing', async () => {
const publishCalls: Array<[string, string]> = [];
const hook = createSendMessageOmniHook(omniEnv, (s, p) => {
publishCalls.push([s, p]);
});

// Model emits the wrong field name (e.g. `text` instead of `message`/`content`)
const result = await hook(
{
hook_event_name: 'PreToolUse',
tool_name: 'SendMessage',
tool_input: { recipient: 'omni', text: 'wrong field' },
tool_use_id: 'tu-empty-1',
} as never,
'tu-empty-1',
dummyMeta,
);

expect(publishCalls).toHaveLength(0);
const out = result as { hookSpecificOutput?: { permissionDecision?: string; permissionDecisionReason?: string } };
expect(out.hookSpecificOutput?.permissionDecision).toBe('deny');
expect(out.hookSpecificOutput?.permissionDecisionReason).toMatch(/non-empty.*message/i);
});

it('denies with retry-prompting reason when message body is whitespace-only', async () => {
const publishCalls: Array<[string, string]> = [];
const hook = createSendMessageOmniHook(omniEnv, (s, p) => {
publishCalls.push([s, p]);
});

const result = await hook(
{
hook_event_name: 'PreToolUse',
tool_name: 'SendMessage',
tool_input: { recipient: 'omni', message: ' \n\t ' },
tool_use_id: 'tu-empty-2',
} as never,
'tu-empty-2',
dummyMeta,
);

expect(publishCalls).toHaveLength(0);
const out = result as { hookSpecificOutput?: { permissionDecision?: string; permissionDecisionReason?: string } };
expect(out.hookSpecificOutput?.permissionDecision).toBe('deny');
expect(out.hookSpecificOutput?.permissionDecisionReason).toMatch(/non-empty.*message/i);
});

it('passes through when OMNI_INSTANCE is missing (non-bridge SDK session)', async () => {
const publishCalls: Array<[string, string]> = [];
const hook = createSendMessageOmniHook({}, (s, p) => {
publishCalls.push([s, p]);
});

const result = await hook(
{
hook_event_name: 'PreToolUse',
tool_name: 'SendMessage',
tool_input: { recipient: 'omni', message: 'no bridge' },
tool_use_id: 'tu-6',
} as never,
'tu-6',
dummyMeta,
);

expect(publishCalls).toHaveLength(0);
expect((result as Record<string, unknown>).hookSpecificOutput).toBeUndefined();
});

it('wires the SendMessage hook into runQuery options when OMNI_INSTANCE is set', async () => {
const session = await executor.spawn('test-agent', 'chat-wire', {
OMNI_INSTANCE: 'inst-wire',
OMNI_CHAT: 'chat-wire',
OMNI_AGENT: 'wirebot',
});

await executor.deliver(session, {
content: 'Hi',
sender: 'Alice',
instanceId: 'inst-wire',
chatId: 'chat-wire',
agent: 'test-agent',
});
await executor.waitForDeliveries(session.id);

expect(queryMock).toHaveBeenCalled();
const callArgs = (
queryMock.mock.calls.at(-1) as unknown as [{ options?: { hooks?: Record<string, unknown[]> } }]
)[0];
const preToolUseHooks = callArgs.options?.hooks?.PreToolUse;
expect(Array.isArray(preToolUseHooks)).toBe(true);
// permission gate matcher (*) + SendMessage matcher
const matchers = (preToolUseHooks as Array<{ matcher?: string }>).map((h) => h.matcher);
expect(matchers).toContain('SendMessage');
});

it('does NOT wire the SendMessage hook when OMNI_INSTANCE is absent', async () => {
const session = await executor.spawn('test-agent', 'chat-nowire', {});

await executor.deliver(session, {
content: 'Hi',
sender: 'bob',
instanceId: 'inst-1',
chatId: 'chat-nowire',
agent: 'test-agent',
});
await executor.waitForDeliveries(session.id);

const callArgs = (
queryMock.mock.calls.at(-1) as unknown as [
{ options?: { hooks?: Record<string, Array<{ matcher?: string }>> } },
]
)[0];
const preToolUseHooks = callArgs.options?.hooks?.PreToolUse ?? [];
const matchers = preToolUseHooks.map((h) => h.matcher);
expect(matchers).not.toContain('SendMessage');
});
});
});
108 changes: 107 additions & 1 deletion src/services/executors/claude-sdk.ts
Original file line number Diff line number Diff line change
@@ -1,4 +1,10 @@
import type { Query } from '@anthropic-ai/claude-agent-sdk';
import type {
HookCallback,
HookCallbackMatcher,
PreToolUseHookInput,
Query,
SyncHookJSONOutput,
} from '@anthropic-ai/claude-agent-sdk';

import { z } from 'zod';
import * as directory from '../../lib/agent-directory.js';
Expand Down Expand Up @@ -153,6 +159,95 @@ export function handleDoneTool(
return action.label;
}

// ============================================================================
// SendMessage interception — route SendMessage(to: "omni") to NATS reply path
// ============================================================================

/**
* Extract `recipient` and `message content` from a SendMessage tool input.
*
* Both field name variants are supported defensively:
* - recipient: `recipient` (CC native) or `to` (genie internal)
* - body: `message` (CC native) or `content` (genie internal)
*
* Returns `{ recipient: undefined }` when the input is not parseable.
*/
function parseSendMessageInput(input: unknown): { recipient?: string; body?: string } {
if (!input || typeof input !== 'object') return {};
const obj = input as Record<string, unknown>;
const recipient = typeof obj.recipient === 'string' ? obj.recipient : typeof obj.to === 'string' ? obj.to : undefined;
const body =
typeof obj.message === 'string' ? obj.message : typeof obj.content === 'string' ? obj.content : undefined;
return { recipient, body };
}

/**
* Build a PreToolUse hook that intercepts `SendMessage` to recipient "omni"
* and routes the message through NATS to the omni reply path.
*
* Mirrors `handleDoneTool`'s text-action publish so agents can use a single
* messaging interface (`SendMessage`) regardless of executor transport.
*
* Behavior:
* - tool_name !== 'SendMessage' → no decision (handler chain continues)
* - recipient !== 'omni' → no decision (native delivery proceeds)
* - OMNI_INSTANCE missing in env → no decision (not bridge mode)
* - matched → publish + deny with success-equivalent reason
*/
export function createSendMessageOmniHook(
env: Record<string, string>,
natsPublish: NatsPublishFn | null,
): HookCallback {
return async (input): Promise<SyncHookJSONOutput> => {
const hookInput = input as PreToolUseHookInput;
if (hookInput.tool_name !== 'SendMessage') return {};

const { recipient, body } = parseSendMessageInput(hookInput.tool_input);
if (recipient !== 'omni') return {};

const instanceId = env.OMNI_INSTANCE ?? '';
const chatId = env.OMNI_CHAT ?? '';
const agent = env.OMNI_AGENT ?? '';

if (!instanceId || !chatId) return {};

// Reject empty/missing message bodies with an explicit error so the agent
// retries with a real payload — otherwise we'd publish an empty reply and
// the model would believe delivery succeeded.
if (!body || body.trim() === '') {
return {
hookSpecificOutput: {
hookEventName: 'PreToolUse',
permissionDecision: 'deny',
permissionDecisionReason:
'SendMessage(recipient: "omni") requires a non-empty `message` field. Retry with the reply text.',
},
};
}

if (!natsPublish) {
console.warn('[claude-sdk] SendMessage(to: omni) intercepted but NATS publish unavailable — message dropped');
return {
hookSpecificOutput: {
hookEventName: 'PreToolUse',
permissionDecision: 'deny',
permissionDecisionReason: 'Omni bridge unavailable — message could not be delivered.',
},
};
}

natsPublish(`omni.reply.${instanceId}.${chatId}`, buildReplyPayload(agent, chatId, instanceId, { content: body }));

return {
hookSpecificOutput: {
hookEventName: 'PreToolUse',
permissionDecision: 'deny',
permissionDecisionReason: 'Message delivered to user via omni bridge.',
},
};
};
}

async function createDoneMcpServer(env: Record<string, string>, natsPublish: NatsPublishFn | null) {
const { createSdkMcpServer, tool } = await import('@anthropic-ai/claude-agent-sdk');
return createSdkMcpServer({
Expand Down Expand Up @@ -350,9 +445,20 @@ export class ClaudeSdkOmniExecutor implements IExecutor {
}

const doneMcp = await createDoneMcpServer(state.env, this.natsPublish);
const sendMessageHooks: Partial<Record<string, HookCallbackMatcher[]>> | undefined = isTurnBased
? {
PreToolUse: [
{
matcher: 'SendMessage',
hooks: [createSendMessageOmniHook(state.env, this.natsPublish)],
},
],
}
: undefined;
const extraOptions: Record<string, unknown> = {
abortController: state.abortController,
mcpServers: { 'genie-omni-tools': doneMcp },
...(sendMessageHooks && { hooks: sendMessageHooks }),
};
if (state.claudeSessionId) {
extraOptions.resume = state.claudeSessionId;
Expand Down
Loading
Loading