diff --git a/.github/workflows/preview.yml b/.github/workflows/preview.yml index b13f9977ef..d26db25dd4 100644 --- a/.github/workflows/preview.yml +++ b/.github/workflows/preview.yml @@ -199,9 +199,6 @@ jobs: ;; esac echo "${secret_env_key}=$MOCK_API_TOKEN" >> "$OVERRIDES_FILE" - if [ "$service_key" = "CLOUDFLARE" ]; then - echo "CLOUDFLARE_ACCOUNT_ID=cf_account_mock_123" >> "$OVERRIDES_FILE" - fi mock_summary_line="- ${service}: ${mock_url} (\`${mock_worker_name}\`)" mock_comment_summary_line="- ${service}: [${mock_url}/__mocks](${mock_url}/__mocks?token=${MOCK_API_TOKEN}) (\`${mock_worker_name}\`)" @@ -573,16 +570,14 @@ jobs: - name: ℹ️ Skip GitHub preview environment delete if: >- - always() && - env.PREVIEW_ENVIRONMENT_GITHUB_TOKEN == '' + always() && env.PREVIEW_ENVIRONMENT_GITHUB_TOKEN == '' run: > echo "Skipping GitHub preview environment delete; PREVIEW_ENVIRONMENT_GITHUB_TOKEN is not configured." - name: 🏷️ Delete GitHub preview environment if: >- - always() && - env.PREVIEW_ENVIRONMENT_GITHUB_TOKEN != '' + always() && env.PREVIEW_ENVIRONMENT_GITHUB_TOKEN != '' uses: actions/github-script@v8.0.0 env: EVENT_NAME: ${{ github.event_name }} diff --git a/docs/contributing/adding-capabilities.md b/docs/contributing/adding-capabilities.md index f6056e93d0..b66ace018c 100644 --- a/docs/contributing/adding-capabilities.md +++ b/docs/contributing/adding-capabilities.md @@ -47,6 +47,12 @@ To merge extra domains later (e.g. plugins), the seam is: `buildCapabilityRegistry([...builtinDomains, ...extraDomains])` with real `Capability` handlers (typical Workers model: snapshot at deploy). +**Remote connectors:** at runtime, `getCapabilityRegistryForContext` also merges +domains synthesized from outbound WebSocket connectors (see +[`architecture/remote-connectors.md`](./architecture/remote-connectors.md)). +Those domains are driven by MCP **`remoteConnectors`** / **`homeConnectorId`** +rather than by editing `builtinDomains` in-repo. + `defineCapability()` in `packages/worker/src/mcp/capabilities/define-capability.ts` is still what normalizes Zod → JSON Schema and wraps handlers with logging; domain helpers diff --git a/docs/contributing/architecture/home-connector.md b/docs/contributing/architecture/home-connector.md index 9d7a8a3615..de5cba7728 100644 --- a/docs/contributing/architecture/home-connector.md +++ b/docs/contributing/architecture/home-connector.md @@ -3,6 +3,10 @@ The local `packages/home-connector` process is the bridge between Kody's Cloudflare Worker and devices that are only reachable on the local network. +It is a **remote connector** with `kind: home`. The wire protocol, URL shapes, +and secret configuration for **any** outbound connector are documented in +[Remote connectors](./remote-connectors.md). + ## Current adapters The connector currently exposes three local-device families: diff --git a/docs/contributing/architecture/index.md b/docs/contributing/architecture/index.md index b6c9208fcd..5c9eccbc64 100644 --- a/docs/contributing/architecture/index.md +++ b/docs/contributing/architecture/index.md @@ -19,6 +19,8 @@ is trying to become. Objects. - [Home Connector](./home-connector.md): local device adapters, Samsung token persistence, and connector-specific discovery/runtime behavior. +- [Remote connectors](./remote-connectors.md): generic outbound WebSocket + protocol, URLs, secrets, and MCP caller context for any `kind` / instance. - [Local Agent Bridge Direction](./local-agent-bridge.md): proposed direction for securely reaching local-network systems through an outbound agent connection. diff --git a/docs/contributing/architecture/remote-connectors.md b/docs/contributing/architecture/remote-connectors.md new file mode 100644 index 0000000000..2ae8ac4eea --- /dev/null +++ b/docs/contributing/architecture/remote-connectors.md @@ -0,0 +1,140 @@ +# Remote connectors + +A **remote connector** is any service that opens an **outbound WebSocket** to +the Kody Worker and exposes **MCP-style tools** (`tools/list`, `tools/call`) +over that socket. The Worker’s `HomeConnectorSession` Durable Object (binding +name `HOME_CONNECTOR_SESSION`) holds one live session per **session key** and +proxies HTTP `fetch` from Worker code to JSON-RPC on the socket. + +The first shipped connector is **`packages/home-connector`** (`kind: home`). +Additional kinds use the same protocol and routing pattern described below. + +## URLs and session keys + +- **Home (legacy URL, still supported):** + `wss:///home/connectors/` + Session key = `` (unchanged from historical behavior). + +- **Generic:** + `wss:///connectors//` + Session key = `:` when `kind` is not `home` (lowercase + compared after trim). + +The Worker sets header **`X-Kody-Connector-Session-Key`** on requests forwarded +into the Durable Object. The connector’s **`connector.hello`** must declare a +**`connectorKind`** and **`connectorId`** (instance id) that match the session +key implied by the WebSocket URL; otherwise the session closes with a mismatch +error. + +## WebSocket message protocol + +All messages are **JSON objects** with a **`type`** field. + +### Client → Worker (connector) + +1. **`connector.hello`** (required first logical message after open) + - **`type`:** `"connector.hello"` + - **`connectorId`:** string — instance id (for example `default`, + `living-room`). + - **`sharedSecret`:** string — must match Worker configuration (see + [Environment variables](../environment-variables.md#remote-connector-secrets)). + - **`connectorKind`:** string (optional but **required for generic + `/connectors/...` URLs**). Omit or set to `"home"` for the home connector. + Lowercase values are normalized. + +2. **`connector.heartbeat`** + - **`type`:** `"connector.heartbeat"` + - Keeps `lastSeenAt` fresh in the session DO. + +3. **`connector.jsonrpc`** + - **`type`:** `"connector.jsonrpc"` + - **`message`:** a single JSON-RPC 2.0 object (request or response). + +### Worker → Client (connector process) + +- **`server.ping`** — Worker may send this; connector should stay connected. +- **`server.ack`** — Successful hello; includes **`connectorId`** echo. +- **`server.error`** — Human-readable **`message`**; connection may close. + +## JSON-RPC on the socket + +The Worker sends MCP-style requests over the WebSocket wrapped in +`connector.jsonrpc`: + +- **`tools/list`** — Return `{ tools: [...] }` where each tool has at least + **`name`**, and typically **`description`**, **`inputSchema`**, optional + **`title`**, **`outputSchema`**, **`annotations`** (same shape as MCP tools). + +- **`tools/call`** — Params: `{ name: string, arguments?: object }`. Return a + normal MCP **`CallToolResult`**-compatible payload (content, structured + content, `isError`, etc.). + +If the Worker forwards **`notifications/tools/list_changed`**, the connector +should re-list tools when it supports dynamic registration. Separately, the +reference implementation in `packages/home-connector` **proactively** sends +`notifications/tools/list_changed` **to** the Worker right after +**`server.ack`** so the session performs an initial tool snapshot refresh. + +## HTTP helper endpoints (same origin) + +The same Durable Object serves snapshot and RPC helpers on paths **under the +connector URL** (for example `/snapshot`, `/rpc/tools-list`). External connector +authors normally only need the **WebSocket**; Worker-internal code uses these +for bridging. + +## Worker-side attachment (MCP caller context) + +For capabilities to be synthesized from a connector, the MCP session must list +that connector: + +- **`remoteConnectors`:** optional array of `{ kind, instanceId }`. When present + (including empty), it fully defines the set of remote connectors for that + session. +- **`homeConnectorId`:** when `remoteConnectors` is omitted, a non-null value + maps to `{ kind: "home", instanceId: homeConnectorId }`. + +Source: `packages/shared/src/chat.ts`, +`packages/shared/src/remote-connectors.ts`. + +## Capability naming (search / execute) + +- Single **`home`** connector with instance id **`default`:** synthesized + capabilities stay on the builtin **`home`** domain with names like + **`home_`** (legacy stability). + +- Any other combination (multiple home instances, non-`home` kinds): the Worker + uses distinct **domain ids** (for example `remote::`) and + **prefixed capability names** so nothing collides in `search` / `execute`. + +## Compatibility checklist + +1. **Outbound WebSocket** to the correct path for your **`kind`** and + **`instanceId`**. +2. **Hello first** with matching **`connectorKind`** + **`connectorId`** and a + **valid `sharedSecret`** for that `kind:instanceId` pair. +3. Implement **`tools/list`** and **`tools/call`** on the socket via + **`connector.jsonrpc`** envelopes. +4. **Heartbeats** if the service stays connected for a long time. +5. **Operator config:** Worker `REMOTE_CONNECTOR_SECRETS` and/or + `HOME_CONNECTOR_SHARED_SECRET` for `home`; MCP clients must pass + **`remoteConnectors`** / **`homeConnectorId`** so the registry merges your + domain. + +## Reference implementation + +- Protocol types and parsing: `packages/worker/src/home/types.ts`, + `packages/worker/src/home/utils.ts` +- Session Durable Object: `packages/worker/src/home/session.ts` +- Ingress and session key: + `packages/worker/src/remote-connector/connector-session-key.ts` +- Home connector WebSocket client: + `packages/home-connector/src/transport/worker-connector.ts` + +## Related docs + +- [Home Connector](./home-connector.md) — the shipped `home` implementation + (Roku, Lutron, Samsung TV, Sonos). +- [Request lifecycle](./request-lifecycle.md) — where connector routes sit in + the Worker. +- [Environment variables](../environment-variables.md#remote-connector-secrets) + — secrets and optional JSON map. diff --git a/docs/contributing/architecture/request-lifecycle.md b/docs/contributing/architecture/request-lifecycle.md index 601352c6a7..a65a6a5646 100644 --- a/docs/contributing/architecture/request-lifecycle.md +++ b/docs/contributing/architecture/request-lifecycle.md @@ -36,10 +36,16 @@ Requests are handled in this order: - `/.well-known/oauth-protected-resource/mcp` 5. MCP endpoint: - `/mcp` (requires OAuth bearer token) -6. Home connector session endpoint: - - `/home/connectors/:connectorId...` (internal-only Worker route that proxies - websocket upgrades and JSON-RPC helper requests to the - `HomeConnectorSession` Durable Object) +6. Remote connector session endpoints (internal-only Worker routes that proxy + WebSocket upgrades and JSON-RPC helper requests to the `HomeConnectorSession` + Durable Object): + - `/home/connectors/:connectorId...` — legacy **`home`** connector URL + (session key equals `connectorId`) + - `/connectors/:kind/:instanceId...` — generic **`kind`** + instance (session + key `kind:instanceId` when `kind` is not `home`) + + See [Remote connectors](./remote-connectors.md). + 7. Internal chat agent endpoint: - `/chat-agent/:threadId...` (requires the app session cookie and routes to the per-thread chat Agent instance) @@ -104,8 +110,11 @@ The home automation flow adds two more Durable Objects: home connector tools when needed. The chat agent still attaches to the main compact MCP server (`kody`), but it -also attaches to `home` and the runtime capability registry synthesizes a `home` -domain for `search` / `execute` from the connected home connector tool surface. +also attaches to `home` and the runtime capability registry **merges** +synthesized domains from **remote connectors** listed in MCP caller context +(`remoteConnectors` or legacy `homeConnectorId`). A single **`home`** + +**`default`** instance keeps the builtin `home` domain name; other combinations +use distinct domain ids. See [Remote connectors](./remote-connectors.md). Shared options are built in `packages/worker/src/sentry-options.ts`: **release** comes from `APP_COMMIT_SHA` when set (deploy workflows pass it as a var), and diff --git a/docs/contributing/environment-variables.md b/docs/contributing/environment-variables.md index 1b1bbe132b..9a11f335cd 100644 --- a/docs/contributing/environment-variables.md +++ b/docs/contributing/environment-variables.md @@ -126,6 +126,26 @@ Optional Worker secret/var (see `packages/worker/src/env-schema.ts` and `packages/home-connector` service when it opens the outbound WebSocket session to the worker. When unset, the worker rejects home connector registration and the internal home MCP bridge cannot route `home` capabilities. + +### Remote connector secrets (Worker) + +See `packages/worker/src/env-schema.ts` and +`packages/worker/src/remote-connector/resolve-remote-connector-secret.ts`. + +- `REMOTE_CONNECTOR_SECRETS` — optional Worker **secret** (JSON string) whose + value is a JSON object mapping **`"kind:instanceId"`** keys (trimmed, kind + lowercased) to **shared secret strings** for **`connector.hello`**. When a key + is present, it overrides per-connector lookup before any kind-specific + fallback. At Worker boot, invalid JSON or malformed keys fail env validation + with a clear error. At runtime, if the value is a plain string in a test + harness, malformed JSON is logged and ignored for map lookup only. +- For **`kind: home`**, if a key is missing in the map, the worker still falls + back to **`HOME_CONNECTOR_SHARED_SECRET`**. Non-`home` kinds have **no** + legacy fallback; they must appear in the map (or hello is rejected). + +Authoring guide for outbound WebSocket services: +[`architecture/remote-connectors.md`](./architecture/remote-connectors.md). + - `HOME_CONNECTOR_*` — when you start the full local stack with `npm run dev`, any `HOME_CONNECTOR_`-prefixed variable is forwarded to the child connector process with the prefix removed. For example, `HOME_CONNECTOR_MOCKS=false` diff --git a/docs/use/troubleshooting.md b/docs/use/troubleshooting.md index 5d432658de..6c6381575b 100644 --- a/docs/use/troubleshooting.md +++ b/docs/use/troubleshooting.md @@ -24,5 +24,8 @@ not signed in, user-scoped results are empty. ## Home automation -If home-related tools appear missing, check home connector status with -**`meta_get_home_connector_status`** when that capability is available. +If home-related tools appear missing, check connector status with +**`meta_list_remote_connector_status`** (all attached remote connectors) or +**`meta_get_home_connector_status`** (first **`home`** connector only) when +those capabilities are available. For protocol and URL requirements, see +[Remote connectors](../contributing/architecture/remote-connectors.md). diff --git a/packages/home-connector/src/transport/worker-connector.ts b/packages/home-connector/src/transport/worker-connector.ts index 0955ca8ddc..e8e7d2a1dd 100644 --- a/packages/home-connector/src/transport/worker-connector.ts +++ b/packages/home-connector/src/transport/worker-connector.ts @@ -245,6 +245,7 @@ export function createWorkerConnector(input: { ) const hello: HomeConnectorHelloMessage = { type: 'connector.hello', + connectorKind: 'home', connectorId: input.config.homeConnectorId, sharedSecret: input.config.sharedSecret!, } diff --git a/packages/shared/src/chat.ts b/packages/shared/src/chat.ts index 3488fa4ba9..8879693286 100644 --- a/packages/shared/src/chat.ts +++ b/packages/shared/src/chat.ts @@ -1,4 +1,7 @@ import { + array, + createSchema, + fail, nullable, number, object, @@ -7,6 +10,28 @@ import { type InferOutput, } from 'remix/data-schema' +const remoteConnectorKindFieldSchema = createSchema( + (value, context) => { + if (typeof value !== 'string') return fail('Expected string', context.path) + const trimmed = value.trim().toLowerCase() + if (!trimmed) { + return fail('remote connector kind must not be empty', context.path) + } + return { value: trimmed } + }, +) + +const remoteConnectorInstanceIdFieldSchema = createSchema( + (value, context) => { + if (typeof value !== 'string') return fail('Expected string', context.path) + const trimmed = value.trim() + if (!trimmed) { + return fail('remote connector instanceId must not be empty', context.path) + } + return { value: trimmed } + }, +) + export const aiModeValues = ['mock', 'remote'] as const export type AiMode = (typeof aiModeValues)[number] @@ -21,10 +46,16 @@ export const mcpStorageContextSchema = object({ appId: optional(nullable(string())), }) +const remoteConnectorRefSchema = object({ + kind: remoteConnectorKindFieldSchema, + instanceId: remoteConnectorInstanceIdFieldSchema, +}) + export const mcpCallerContextSchema = object({ baseUrl: string(), user: optional(nullable(mcpUserContextSchema)), homeConnectorId: optional(nullable(string())), + remoteConnectors: optional(nullable(array(remoteConnectorRefSchema))), storageContext: optional(nullable(mcpStorageContextSchema)), }) diff --git a/packages/shared/src/remote-connectors.ts b/packages/shared/src/remote-connectors.ts new file mode 100644 index 0000000000..c56772bfb3 --- /dev/null +++ b/packages/shared/src/remote-connectors.ts @@ -0,0 +1,43 @@ +import { type InferOutput } from 'remix/data-schema' +import { type mcpCallerContextSchema } from './chat.ts' + +type McpCallerContext = InferOutput + +export type RemoteConnectorRef = { + kind: string + instanceId: string +} + +function normalizeKind(kind: string): string { + return kind.trim().toLowerCase() +} + +function normalizeInstanceId(instanceId: string): string { + return instanceId.trim() +} + +/** + * Effective remote connectors for MCP execution, in order. + * When `remoteConnectors` is set (including empty), it wins. + * Otherwise `homeConnectorId` maps to `{ kind: "home", instanceId }`. + */ +export function normalizeRemoteConnectorRefs( + context: Pick, +): Array { + if ( + context.remoteConnectors !== undefined && + context.remoteConnectors !== null + ) { + return context.remoteConnectors + .map((ref) => ({ + kind: normalizeKind(ref.kind), + instanceId: normalizeInstanceId(ref.instanceId), + })) + .filter((ref) => ref.kind.length > 0 && ref.instanceId.length > 0) + } + const hid = context.homeConnectorId?.trim() + if (!hid) { + return [] + } + return [{ kind: 'home', instanceId: hid }] +} diff --git a/packages/worker/client/mcp-apps/kody-ui-utils.ts b/packages/worker/client/mcp-apps/kody-ui-utils.ts index 279501139b..493c272f29 100644 --- a/packages/worker/client/mcp-apps/kody-ui-utils.ts +++ b/packages/worker/client/mcp-apps/kody-ui-utils.ts @@ -833,8 +833,7 @@ async function executeGeneratedUiShellScript( return } const shouldAwaitLoad = - scriptDescriptor.executionMode === 'module' || - scriptDescriptor.src != null + scriptDescriptor.executionMode === 'module' || scriptDescriptor.src != null if (!scriptDescriptor.src) { script.textContent = scriptDescriptor.textContent if (!shouldAwaitLoad) { @@ -1147,7 +1146,8 @@ async function initializeShellHostDocument() { } } - const isStaleRender = (renderId: number) => renderId !== latestScheduledRenderId + const isStaleRender = (renderId: number) => + renderId !== latestScheduledRenderId async function renderEnvelope( envelope: RenderEnvelope | null, diff --git a/packages/worker/src/env-schema.ts b/packages/worker/src/env-schema.ts index 9293e45db8..b460642d89 100644 --- a/packages/worker/src/env-schema.ts +++ b/packages/worker/src/env-schema.ts @@ -74,6 +74,62 @@ const optionalAiModeSchema = createSchema< return fail(`Expected one of: ${aiModeValues.join(', ')}`, context.path) }) +const optionalRemoteConnectorSecretsSchema = createSchema< + unknown, + Record | undefined +>((value, context) => { + if (value === undefined) return { value: undefined } + if (typeof value !== 'string') return fail('Expected string', context.path) + + const trimmed = value.trim() + if (!trimmed) return { value: undefined } + + let parsed: unknown + try { + parsed = JSON.parse(trimmed) as unknown + } catch { + return fail( + 'REMOTE_CONNECTOR_SECRETS must be valid JSON when set.', + context.path, + ) + } + if (!parsed || typeof parsed !== 'object' || Array.isArray(parsed)) { + return fail( + 'REMOTE_CONNECTOR_SECRETS must be a JSON object mapping "kind:instanceId" keys to secret strings.', + context.path, + ) + } + + const out: Record = {} + for (const [rawKey, rawVal] of Object.entries(parsed)) { + const key = rawKey.trim() + const colon = key.indexOf(':') + if (colon <= 0 || colon === key.length - 1) { + return fail( + `REMOTE_CONNECTOR_SECRETS has invalid key "${rawKey}" (expected "kind:instanceId").`, + context.path, + ) + } + const kind = key.slice(0, colon).trim().toLowerCase() + const instanceId = key.slice(colon + 1).trim() + if (!kind || !instanceId) { + return fail( + `REMOTE_CONNECTOR_SECRETS has invalid key "${rawKey}" (kind and instanceId must be non-empty).`, + context.path, + ) + } + const canonicalKey = `${kind}:${instanceId}` + if (typeof rawVal !== 'string' || !rawVal.trim()) { + return fail( + `REMOTE_CONNECTOR_SECRETS value for "${canonicalKey}" must be a non-empty string.`, + context.path, + ) + } + out[canonicalKey] = rawVal.trim() + } + return { value: out } +}) + const optionalSentryTracesSampleRateSchema = createSchema< unknown, number | undefined @@ -125,6 +181,7 @@ export const EnvSchema = object({ CLOUDFLARE_API_BASE_URL: optionalUrlStringSchema, CAPABILITY_REINDEX_SECRET: optionalNonEmptyStringSchema, HOME_CONNECTOR_SHARED_SECRET: optionalNonEmptyStringSchema, + REMOTE_CONNECTOR_SECRETS: optionalRemoteConnectorSecretsSchema, }) export type AppEnv = InferOutput diff --git a/packages/worker/src/home/client-transport.ts b/packages/worker/src/home/client-transport.ts index 7143ae044f..603e90f276 100644 --- a/packages/worker/src/home/client-transport.ts +++ b/packages/worker/src/home/client-transport.ts @@ -3,6 +3,10 @@ import { type MessageExtraInfo, } from '@modelcontextprotocol/sdk/types.js' import { type Transport } from '@modelcontextprotocol/sdk/shared/transport.js' +import { + connectorIngressPath, + connectorSessionKey, +} from '#worker/remote-connector/connector-session-key.ts' import { type HomeConnectorJsonRpcResponse } from './types.ts' export class HomeConnectorClientTransport implements Transport { @@ -15,16 +19,21 @@ export class HomeConnectorClientTransport implements Transport { ) => void private readonly input: { - connectorId: string + kind: string + instanceId: string baseUrl: string } - constructor(input: { connectorId: string; baseUrl: string }) { - this.input = input + constructor(input: { kind?: string; instanceId: string; baseUrl: string }) { + this.input = { + kind: input.kind ?? 'home', + instanceId: input.instanceId, + baseUrl: input.baseUrl, + } } async start(): Promise { - this.sessionId = this.input.connectorId + this.sessionId = connectorSessionKey(this.input.kind, this.input.instanceId) } async send(message: JSONRPCMessage): Promise { @@ -47,16 +56,14 @@ export class HomeConnectorClientTransport implements Transport { private async forwardJsonRpc( message: JSONRPCMessage, ): Promise { - const response = await fetch( - `${this.input.baseUrl}/home/connectors/${this.input.connectorId}/rpc/jsonrpc`, - { - method: 'POST', - headers: { - 'Content-Type': 'application/json', - }, - body: JSON.stringify({ message }), + const path = `${connectorIngressPath(this.input.kind, this.input.instanceId)}/rpc/jsonrpc` + const response = await fetch(`${this.input.baseUrl}${path}`, { + method: 'POST', + headers: { + 'Content-Type': 'application/json', }, - ) + body: JSON.stringify({ message }), + }) if (!response.ok) { throw new Error( `Home connector bridge request failed with ${response.status}.`, diff --git a/packages/worker/src/home/client.ts b/packages/worker/src/home/client.ts index 47fe3fe5d9..423df638c5 100644 --- a/packages/worker/src/home/client.ts +++ b/packages/worker/src/home/client.ts @@ -1,4 +1,8 @@ import { type CallToolResult } from '@modelcontextprotocol/sdk/types.js' +import { + connectorIngressPath, + connectorSessionKey, +} from '#worker/remote-connector/connector-session-key.ts' import { type HomeConnectorSnapshot, type HomeToolDescriptor } from './types.ts' export type HomeMcpTool = HomeToolDescriptor @@ -12,13 +16,19 @@ export type HomeMcpClient = { getSnapshot(): Promise } -function createSessionUrl(connectorId: string, pathname: string): string { - return `https://home-connectors/${connectorId}${pathname}` +function createSessionUrl( + kind: string, + instanceId: string, + pathname: string, +): string { + const base = connectorIngressPath(kind, instanceId) + return `https://home-connectors${base}${pathname}` } -function getSessionStub(env: Env, connectorId: string) { +function getSessionStub(env: Env, kind: string, instanceId: string) { + const sessionKey = connectorSessionKey(kind, instanceId) return env.HOME_CONNECTOR_SESSION.get( - env.HOME_CONNECTOR_SESSION.idFromName(connectorId), + env.HOME_CONNECTOR_SESSION.idFromName(sessionKey), ) } @@ -29,16 +39,17 @@ async function parseJsonResponse(response: Response): Promise { return (await response.json()) as T } -export function createHomeMcpClient( +export function createRemoteConnectorMcpClient( env: Env, - connectorId: string, + kind: string, + instanceId: string, ): HomeMcpClient { - const stub = getSessionStub(env, connectorId) + const stub = getSessionStub(env, kind, instanceId) return { async listTools() { const response = await stub.fetch( - createSessionUrl(connectorId, '/rpc/tools-list'), + createSessionUrl(kind, instanceId, '/rpc/tools-list'), { method: 'POST', }, @@ -50,7 +61,7 @@ export function createHomeMcpClient( }, async callTool(name, args) { const response = await stub.fetch( - createSessionUrl(connectorId, '/rpc/tools-call'), + createSessionUrl(kind, instanceId, '/rpc/tools-call'), { method: 'POST', headers: { @@ -66,9 +77,16 @@ export function createHomeMcpClient( }, async getSnapshot() { const response = await stub.fetch( - createSessionUrl(connectorId, '/snapshot'), + createSessionUrl(kind, instanceId, '/snapshot'), ) return parseJsonResponse(response) }, } } + +export function createHomeMcpClient( + env: Env, + connectorId: string, +): HomeMcpClient { + return createRemoteConnectorMcpClient(env, 'home', connectorId) +} diff --git a/packages/worker/src/home/mcp.ts b/packages/worker/src/home/mcp.ts index 78f0d927cd..4e3ac5fe82 100644 --- a/packages/worker/src/home/mcp.ts +++ b/packages/worker/src/home/mcp.ts @@ -16,6 +16,7 @@ import { parseMcpCallerContext, type McpServerProps, } from '#worker/mcp/context.ts' +import { normalizeRemoteConnectorRefs } from '@kody-internal/shared/remote-connectors.ts' export type HomeMcpState = {} export type HomeMcpProps = McpServerProps @@ -42,15 +43,21 @@ type HomeMcpBridge = { getHomeClient(): Promise> } +function resolveHomeBridgeInstanceId( + callerContext: McpServerProps, +): string | null { + const refs = normalizeRemoteConnectorRefs(callerContext) + const home = refs.find((r) => r.kind === 'home') + return home?.instanceId ?? null +} + async function createHomeToolErrorResult( agent: HomeMcpBridge, error: unknown, ): Promise { const callerContext = agent.getCallerContext() - const status = await getHomeConnectorStatus( - agent.getEnv(), - callerContext.homeConnectorId ?? null, - ) + const homeInstanceId = resolveHomeBridgeInstanceId(callerContext) + const status = await getHomeConnectorStatus(agent.getEnv(), homeInstanceId) const fallbackMessage = error instanceof Error ? error.message : 'Unknown home connector error.' const message = @@ -177,14 +184,14 @@ class HomeMCPBase extends McpAgent { } async getHomeClient() { - const { homeConnectorId } = this.getCallerContext() - if (!homeConnectorId) { + const homeInstanceId = resolveHomeBridgeInstanceId(this.getCallerContext()) + if (!homeInstanceId) { throw new Error( 'No home connector is associated with this MCP caller context.', ) } - return createHomeMcpClient(this.env, homeConnectorId) + return createHomeMcpClient(this.env, homeInstanceId) } } diff --git a/packages/worker/src/home/session.ts b/packages/worker/src/home/session.ts index b65b5ae540..23142b7249 100644 --- a/packages/worker/src/home/session.ts +++ b/packages/worker/src/home/session.ts @@ -16,6 +16,8 @@ import { jsonResponse, stringifyHomeConnectorMessage, } from './utils.ts' +import { connectorSessionKey } from '#worker/remote-connector/connector-session-key.ts' +import { resolveRemoteConnectorSharedSecret } from '#worker/remote-connector/resolve-remote-connector-secret.ts' const connectorTag = 'connector' const stateStorageKey = 'home-connector-session-state' @@ -36,12 +38,15 @@ class HomeConnectorSessionBase extends DurableObject { private stateSnapshot: HomeConnectorSessionState = { persisted: { connectorId: null, + connectorKind: null, connectedAt: null, lastSeenAt: null, }, tools: [], } + private ingressSessionKeys = new WeakMap() + private pendingRequests = new Map() constructor(state: DurableObjectState, env: Env) { @@ -52,7 +57,10 @@ class HomeConnectorSessionBase extends DurableObject { async fetch(request: Request): Promise { const url = new URL(request.url) if (request.headers.get('Upgrade') === 'websocket') { - return this.handleWebSocketUpgrade(request) + const sessionKeyHeader = request.headers + .get('X-Kody-Connector-Session-Key') + ?.trim() + return this.handleWebSocketUpgrade(sessionKeyHeader || null) } if (request.method === 'GET' && url.pathname.endsWith('/snapshot')) { return jsonResponse(await this.getSnapshot()) @@ -138,10 +146,12 @@ class HomeConnectorSessionBase extends DurableObject { } async getSnapshot(): Promise { - const { connectorId, connectedAt, lastSeenAt } = + const { connectorId, connectorKind, connectedAt, lastSeenAt } = this.stateSnapshot.persisted if (!connectorId || !connectedAt || !lastSeenAt) return null + const kind = (connectorKind && connectorKind.trim()) || ('home' as const) return { + ...(kind !== 'home' ? { connectorKind: kind } : {}), connectorId, connectedAt, lastSeenAt, @@ -153,6 +163,9 @@ class HomeConnectorSessionBase extends DurableObject { const stored = await this.ctx.storage.get(stateStorageKey) if (!stored) return + if (stored.persisted.connectorKind === undefined) { + stored.persisted.connectorKind = null + } this.stateSnapshot = stored } @@ -160,7 +173,7 @@ class HomeConnectorSessionBase extends DurableObject { await this.ctx.storage.put(stateStorageKey, this.stateSnapshot) } - private async handleWebSocketUpgrade(_request: Request) { + private async handleWebSocketUpgrade(ingressSessionKey: string | null) { const pair = new WebSocketPair() const sockets = Object.values(pair) const client = sockets[0] @@ -169,6 +182,7 @@ class HomeConnectorSessionBase extends DurableObject { throw new Error('Failed to create WebSocket pair.') } this.ctx.acceptWebSocket(server, [connectorTag]) + this.stashIngressSessionKey(server, ingressSessionKey) server.send( stringifyHomeConnectorMessage({ type: 'server.ping', @@ -211,7 +225,48 @@ class HomeConnectorSessionBase extends DurableObject { } private async handleHello(ws: WebSocket, message: HomeConnectorHelloMessage) { - const expectedSecret = this.env.HOME_CONNECTOR_SHARED_SECRET?.trim() + const trimmedKind = (message.connectorKind ?? '').trim() + const declaredKind = ( + trimmedKind === '' ? 'home' : trimmedKind + ).toLowerCase() + const canonicalInstanceId = message.connectorId.trim() + const expectedSessionKey = connectorSessionKey( + declaredKind, + canonicalInstanceId, + ) + const ingressSessionKey = this.loadIngressSessionKey(ws) + if (ingressSessionKey && ingressSessionKey !== expectedSessionKey) { + Sentry.captureMessage( + 'Remote connector session rejected hello (session key mismatch).', + { + level: 'error', + tags: { + service: 'worker', + worker_component: 'home-connector-session', + }, + extra: { + connectorId: canonicalInstanceId, + declaredKind, + ingressSessionKey, + expectedSessionKey, + }, + }, + ) + ws.send( + stringifyHomeConnectorMessage({ + type: 'server.error', + message: 'Connector session key does not match this endpoint.', + }), + ) + ws.close(4003, 'session-mismatch') + return + } + + const expectedSecret = resolveRemoteConnectorSharedSecret( + declaredKind, + canonicalInstanceId, + this.env, + ) if (!expectedSecret || message.sharedSecret !== expectedSecret) { Sentry.captureMessage( 'Home connector session rejected websocket hello.', @@ -222,7 +277,8 @@ class HomeConnectorSessionBase extends DurableObject { worker_component: 'home-connector-session', }, extra: { - connectorId: message.connectorId, + connectorId: canonicalInstanceId, + declaredKind, hasExpectedSecret: Boolean(expectedSecret), }, }, @@ -239,7 +295,8 @@ class HomeConnectorSessionBase extends DurableObject { const now = new Date().toISOString() this.stateSnapshot.persisted = { - connectorId: message.connectorId, + connectorId: canonicalInstanceId, + connectorKind: declaredKind, connectedAt: this.stateSnapshot.persisted.connectedAt ?? now, lastSeenAt: now, } @@ -247,7 +304,7 @@ class HomeConnectorSessionBase extends DurableObject { ws.send( stringifyHomeConnectorMessage({ type: 'server.ack', - connectorId: message.connectorId, + connectorId: canonicalInstanceId, }), ) } @@ -326,6 +383,35 @@ class HomeConnectorSessionBase extends DurableObject { return response } + + private stashIngressSessionKey( + ws: WebSocket, + ingressSessionKey: string | null, + ) { + this.ingressSessionKeys.set(ws, ingressSessionKey) + try { + ws.serializeAttachment(ingressSessionKey ?? '') + } catch { + // No attachment support; keep in-memory map only. + } + } + + private loadIngressSessionKey(ws: WebSocket): string | null { + if (this.ingressSessionKeys.has(ws)) { + return this.ingressSessionKeys.get(ws) ?? null + } + let ingressSessionKey: string | null = null + try { + const attachment = ws.deserializeAttachment() + if (typeof attachment === 'string') { + ingressSessionKey = attachment || null + } + } catch { + // Ignore deserialization errors, we only enforce if we have a key. + } + this.ingressSessionKeys.set(ws, ingressSessionKey) + return ingressSessionKey + } } export const HomeConnectorSession = Sentry.instrumentDurableObjectWithSentry( diff --git a/packages/worker/src/home/status.ts b/packages/worker/src/home/status.ts index 5480630639..010248c3e3 100644 --- a/packages/worker/src/home/status.ts +++ b/packages/worker/src/home/status.ts @@ -1,8 +1,10 @@ -import { createHomeMcpClient } from './client.ts' +import { type RemoteConnectorRef } from '@kody-internal/shared/remote-connectors.ts' +import { createRemoteConnectorMcpClient } from './client.ts' import { type HomeConnectorSnapshot } from './types.ts' export type HomeConnectorStatus = { state: 'connected' | 'disconnected' | 'unavailable' | 'error' + connectorKind: string connectorId: string | null connected: boolean connectedAt: string | null @@ -12,12 +14,25 @@ export type HomeConnectorStatus = { error: string | null } +function connectorLabel(kind: string, connectorId: string) { + const k = kind.trim().toLowerCase() + if (k === 'home') { + return `home connector "${connectorId}"` + } + return `${k} connector "${connectorId}"` +} + function createConnectedStatus( snapshot: HomeConnectorSnapshot, + kind: string, ): HomeConnectorStatus { + const resolvedKind = + (snapshot.connectorKind ?? kind).trim().toLowerCase() || 'home' const toolCount = snapshot.tools.length + const label = connectorLabel(resolvedKind, snapshot.connectorId) return { state: 'connected', + connectorKind: resolvedKind, connectorId: snapshot.connectorId, connected: true, connectedAt: snapshot.connectedAt, @@ -25,21 +40,26 @@ function createConnectedStatus( toolCount, message: toolCount > 0 - ? `The home connector "${snapshot.connectorId}" is connected and exposing ${toolCount} tool${toolCount === 1 ? '' : 's'}.` - : `The home connector "${snapshot.connectorId}" is connected, but it has not exposed any tools yet.`, + ? `The ${label} is connected and exposing ${toolCount} tool${toolCount === 1 ? '' : 's'}.` + : `The ${label} is connected, but it has not exposed any tools yet.`, error: null, } } -function createDisconnectedStatus(connectorId: string): HomeConnectorStatus { +function createDisconnectedStatus( + ref: RemoteConnectorRef, +): HomeConnectorStatus { + const k = ref.kind.trim().toLowerCase() + const label = connectorLabel(k, ref.instanceId) return { state: 'disconnected', - connectorId, + connectorKind: k, + connectorId: ref.instanceId, connected: false, connectedAt: null, lastSeenAt: null, toolCount: 0, - message: `The home connector "${connectorId}" is not currently connected.`, + message: `The ${label} is not currently connected.`, error: null, } } @@ -47,6 +67,7 @@ function createDisconnectedStatus(connectorId: string): HomeConnectorStatus { function createUnavailableStatus(): HomeConnectorStatus { return { state: 'unavailable', + connectorKind: 'home', connectorId: null, connected: false, connectedAt: null, @@ -58,18 +79,21 @@ function createUnavailableStatus(): HomeConnectorStatus { } function createErrorStatus( - connectorId: string, + ref: RemoteConnectorRef, error: unknown, ): HomeConnectorStatus { const message = error instanceof Error ? error.message : String(error) + const k = ref.kind.trim().toLowerCase() + const label = connectorLabel(k, ref.instanceId) return { state: 'error', - connectorId, + connectorKind: k, + connectorId: ref.instanceId, connected: false, connectedAt: null, lastSeenAt: null, toolCount: 0, - message: `Kody could not determine the home connector status for "${connectorId}".`, + message: `Kody could not determine the status for the ${label}.`, error: message, } } @@ -77,16 +101,29 @@ function createErrorStatus( export function formatHomeConnectorUnavailableMessage( status: HomeConnectorStatus, ) { + return formatRemoteConnectorUnavailableMessage(status) +} + +export function formatRemoteConnectorUnavailableMessage( + status: HomeConnectorStatus, +) { + const isHome = status.connectorKind === 'home' switch (status.state) { case 'connected': if (status.toolCount > 0) { return status.message } - return `${status.message} Home capabilities cannot be searched or used until the connector exposes tools.` + return isHome + ? `${status.message} Home capabilities cannot be searched or used until the connector exposes tools.` + : `${status.message} Capabilities from this connector cannot be searched or used until it exposes tools.` case 'disconnected': - return `${status.message} Kody cannot search or use home capabilities until it reconnects. Ask the user to start or reconnect the home connector and then try again.` + return isHome + ? `${status.message} Kody cannot search or use home capabilities until it reconnects. Ask the user to start or reconnect the home connector and then try again.` + : `${status.message} Kody cannot use this connector until it reconnects. Ask the user to start or reconnect the connector and then try again.` case 'unavailable': - return `${status.message} Kody cannot search or use home capabilities from this session.` + return isHome + ? `${status.message} Kody cannot search or use home capabilities from this session.` + : `${status.message} Kody cannot use this connector from this session.` case 'error': return status.error ? `${status.message} Underlying error: ${status.error}` @@ -94,12 +131,28 @@ export function formatHomeConnectorUnavailableMessage( default: { const exhaustiveState: never = status.state throw new Error( - `Unhandled home connector status state: ${exhaustiveState}`, + `Unhandled remote connector status state: ${exhaustiveState}`, ) } } } +export async function getRemoteConnectorStatus( + env: Env, + ref: RemoteConnectorRef, +): Promise { + try { + const client = createRemoteConnectorMcpClient(env, ref.kind, ref.instanceId) + const snapshot = await client.getSnapshot() + if (!snapshot) { + return createDisconnectedStatus(ref) + } + return createConnectedStatus(snapshot, ref.kind) + } catch (error) { + return createErrorStatus(ref, error) + } +} + export async function getHomeConnectorStatus( env: Env, connectorId: string | null, @@ -108,13 +161,8 @@ export async function getHomeConnectorStatus( return createUnavailableStatus() } - try { - const snapshot = await createHomeMcpClient(env, connectorId).getSnapshot() - if (!snapshot) { - return createDisconnectedStatus(connectorId) - } - return createConnectedStatus(snapshot) - } catch (error) { - return createErrorStatus(connectorId, error) - } + return getRemoteConnectorStatus(env, { + kind: 'home', + instanceId: connectorId, + }) } diff --git a/packages/worker/src/home/types.ts b/packages/worker/src/home/types.ts index ed8209e862..1875ff0d3e 100644 --- a/packages/worker/src/home/types.ts +++ b/packages/worker/src/home/types.ts @@ -16,6 +16,8 @@ export type HomeToolDescriptor = { } export type HomeConnectorSnapshot = { + /** Logical connector kind (e.g. `home`). Defaults to `home` when omitted. */ + connectorKind?: string connectorId: string connectedAt: string lastSeenAt: string @@ -26,6 +28,8 @@ export type HomeConnectorHelloMessage = { type: 'connector.hello' connectorId: string sharedSecret: string + /** When omitted, treated as `home` for backward compatibility. */ + connectorKind?: string } export type HomeConnectorHeartbeatMessage = { @@ -63,6 +67,8 @@ export type HomeConnectorClientMessage = export type HomeConnectorPersistedState = { connectorId: string | null + /** Persisted connector kind; null means legacy sessions (treated as `home`). */ + connectorKind: string | null connectedAt: string | null lastSeenAt: string | null } diff --git a/packages/worker/src/home/utils.ts b/packages/worker/src/home/utils.ts index 6278e331f7..c47f50b6cb 100644 --- a/packages/worker/src/home/utils.ts +++ b/packages/worker/src/home/utils.ts @@ -42,6 +42,18 @@ export function parseHomeConnectorMessage( if (type === 'connector.hello') { const connectorId = (value as Record)['connectorId'] const sharedSecret = (value as Record)['sharedSecret'] + const record = value as Record + const hasConnectorKindKey = Object.hasOwn(record, 'connectorKind') + const connectorKindRaw = record['connectorKind'] + if (hasConnectorKindKey && typeof connectorKindRaw !== 'string') { + throw new Error( + 'Invalid connector hello: connectorKind must be a string.', + ) + } + const connectorKind = + typeof connectorKindRaw === 'string' && connectorKindRaw.trim() + ? connectorKindRaw.trim().toLowerCase() + : undefined if (typeof connectorId !== 'string' || typeof sharedSecret !== 'string') { throw new Error('Invalid connector hello payload.') } @@ -49,6 +61,7 @@ export function parseHomeConnectorMessage( type, connectorId, sharedSecret, + ...(connectorKind ? { connectorKind } : {}), } } if (type === 'connector.heartbeat') { diff --git a/packages/worker/src/index.ts b/packages/worker/src/index.ts index 1b94eb5ab1..cd779f64dc 100644 --- a/packages/worker/src/index.ts +++ b/packages/worker/src/index.ts @@ -33,6 +33,10 @@ import { handleMemoryReindexRequest } from './memory-maintenance.ts' import { handleSkillReindexRequest } from './skill-maintenance.ts' import { handleUiArtifactReindexRequest } from './ui-artifact-maintenance.ts' import { CodemodeFetchGateway } from '#mcp/fetch-gateway.ts' +import { + connectorSessionKey, + parseConnectorRoutePath, +} from './remote-connector/connector-session-key.ts' export { ChatAgent, CodemodeFetchGateway, HomeConnectorSession, HomeMCP, MCP } @@ -130,16 +134,20 @@ const appHandler = withCors({ return handleGeneratedUiApiRequest(request, env) } - if (url.pathname.startsWith('/home/connectors/')) { - const parts = url.pathname.split('/').filter(Boolean) - const connectorId = parts[2]?.trim() - if (!connectorId) { - return new Response('Connector ID is required.', { status: 400 }) - } + const connectorRoute = parseConnectorRoutePath(url.pathname) + if (connectorRoute) { + const sessionKey = connectorSessionKey( + connectorRoute.kind, + connectorRoute.instanceId, + ) const stub = env.HOME_CONNECTOR_SESSION.get( - env.HOME_CONNECTOR_SESSION.idFromName(connectorId), + env.HOME_CONNECTOR_SESSION.idFromName(sessionKey), ) - return stub.fetch(request) + const forwardUrl = new URL(request.url) + forwardUrl.pathname = connectorRoute.rest || '/' + const forwardRequest = new Request(forwardUrl.toString(), request) + forwardRequest.headers.set('X-Kody-Connector-Session-Key', sessionKey) + return stub.fetch(forwardRequest) } if (url.pathname.startsWith(`${chatAgentBasePath}/`)) { diff --git a/packages/worker/src/mcp-auth.workers.test.ts b/packages/worker/src/mcp-auth.workers.test.ts index b9cb230fd4..dd09d91dad 100644 --- a/packages/worker/src/mcp-auth.workers.test.ts +++ b/packages/worker/src/mcp-auth.workers.test.ts @@ -195,7 +195,7 @@ test('mcp request forwards when token is valid', async () => { }) expect(response.status).toBe(200) - expect(receivedProps).toEqual({ + expect(receivedProps).toMatchObject({ baseUrl: 'https://example.com', homeConnectorId: 'default', storageContext: null, diff --git a/packages/worker/src/mcp/capabilities/domain-metadata.ts b/packages/worker/src/mcp/capabilities/domain-metadata.ts index a5245ebc8d..21707d73dc 100644 --- a/packages/worker/src/mcp/capabilities/domain-metadata.ts +++ b/packages/worker/src/mcp/capabilities/domain-metadata.ts @@ -9,5 +9,8 @@ export const capabilityDomainNames = { values: 'values', } as const -export type CapabilityDomain = +export type BuiltinCapabilityDomain = (typeof capabilityDomainNames)[keyof typeof capabilityDomainNames] + +/** Built-in domain ids plus runtime remote-connector domains (e.g. `remote:home:default`). */ +export type CapabilityDomain = string diff --git a/packages/worker/src/mcp/capabilities/home/index.ts b/packages/worker/src/mcp/capabilities/home/index.ts index 9640e7ab78..7cbcb854ae 100644 --- a/packages/worker/src/mcp/capabilities/home/index.ts +++ b/packages/worker/src/mcp/capabilities/home/index.ts @@ -1,39 +1,56 @@ +import { type RemoteConnectorRef } from '@kody-internal/shared/remote-connectors.ts' import { defineCapability } from '#mcp/capabilities/define-capability.ts' import { defineDomain } from '#mcp/capabilities/define-domain.ts' -import { capabilityDomainNames } from '#mcp/capabilities/domain-metadata.ts' -import { createHomeMcpClient } from '#worker/home/client.ts' +import { type CapabilityDomain } from '#mcp/capabilities/domain-metadata.ts' +import { createRemoteConnectorMcpClient } from '#worker/home/client.ts' import { - formatHomeConnectorUnavailableMessage, - getHomeConnectorStatus, + formatRemoteConnectorUnavailableMessage, + getRemoteConnectorStatus, } from '#worker/home/status.ts' +import { + remoteConnectorCapabilityPrefix, + remoteConnectorDomainId, +} from '#worker/remote-connector/remote-domain-id.ts' import { type Capability, type DomainSpec } from '#mcp/capabilities/types.ts' import { type HomeConnectorSnapshot } from '#worker/home/types.ts' -type HomeCapabilityBinding = { +type RemoteToolCapabilityBinding = { capabilityName: string - connectorId: string + kind: string + instanceId: string mcpToolName: string } -export type SynthesizedHomeDomain = { +export type SynthesizedRemoteConnectorDomain = { domain: DomainSpec - bindings: Record + bindings: Record } -function createCapabilityName(toolName: string) { - return `home_${toolName.replaceAll(/[^\w]+/g, '_').replaceAll(/_+/g, '_')}` +function createCapabilityNameFromPrefix(prefix: string, toolName: string) { + const safeTool = toolName + .replaceAll(/[^\w]+/g, '_') + .replaceAll(/_+/g, '_') + .replace(/^_|_$/g, '') + return `${prefix}_${safeTool}` } function buildKeywords( snapshot: HomeConnectorSnapshot, tool: HomeConnectorSnapshot['tools'][number], + ref: RemoteConnectorRef, + extraRoots: ReadonlyArray, ) { + const kind = (snapshot.connectorKind ?? ref.kind).trim().toLowerCase() const words = [ - 'home', + ...extraRoots, + kind, + 'connector', + 'remote', tool.name, tool.title ?? '', tool.description ?? '', snapshot.connectorId, + ref.instanceId, ] return Array.from( new Set( @@ -46,25 +63,41 @@ function buildKeywords( ) } -function createCapabilityFromTool( - snapshot: HomeConnectorSnapshot, - tool: HomeConnectorSnapshot['tools'][number], -): { capability: Capability; binding: HomeCapabilityBinding } { - const capabilityName = createCapabilityName(tool.name) - const binding: HomeCapabilityBinding = { +function createCapabilityFromTool(input: { + snapshot: HomeConnectorSnapshot + tool: HomeConnectorSnapshot['tools'][number] + ref: RemoteConnectorRef + domainId: CapabilityDomain + capabilityPrefix: string + domainKeywordRoots: ReadonlyArray +}): { capability: Capability; binding: RemoteToolCapabilityBinding } { + const { + snapshot, + tool, + ref, + domainId, + capabilityPrefix, + domainKeywordRoots, + } = input + const capabilityName = createCapabilityNameFromPrefix( + capabilityPrefix, + tool.name, + ) + const binding: RemoteToolCapabilityBinding = { capabilityName, - connectorId: snapshot.connectorId, + kind: (snapshot.connectorKind ?? ref.kind).trim().toLowerCase(), + instanceId: ref.instanceId, mcpToolName: tool.name, } const capability = defineCapability({ name: capabilityName, - domain: capabilityDomainNames.home, + domain: domainId, description: tool.description?.trim() || tool.title?.trim() || - `Home automation action for ${tool.name}.`, - keywords: buildKeywords(snapshot, tool), + `Remote connector action (${ref.kind}) for ${tool.name}.`, + keywords: buildKeywords(snapshot, tool, ref, domainKeywordRoots), readOnly: Boolean( (tool.annotations as Record | undefined)?.[ 'readOnlyHint' @@ -83,23 +116,27 @@ function createCapabilityFromTool( inputSchema: tool.inputSchema ?? { type: 'object', properties: {} }, ...(tool.outputSchema ? { outputSchema: tool.outputSchema } : {}), async handler(args, ctx) { - const client = createHomeMcpClient(ctx.env, snapshot.connectorId) + const client = createRemoteConnectorMcpClient( + ctx.env, + binding.kind, + binding.instanceId, + ) let result: Awaited> try { result = await client.callTool(tool.name, args) } catch (error) { - const status = await getHomeConnectorStatus( - ctx.env, - snapshot.connectorId, - ) + const status = await getRemoteConnectorStatus(ctx.env, { + kind: binding.kind, + instanceId: binding.instanceId, + }) if (status.state !== 'connected' || status.toolCount === 0) { - throw new Error(formatHomeConnectorUnavailableMessage(status)) + throw new Error(formatRemoteConnectorUnavailableMessage(status)) } const message = - error instanceof Error - ? error.message - : 'Unknown home connector error.' - throw new Error(`Home capability "${tool.name}" failed: ${message}`) + error instanceof Error ? error.message : 'Unknown connector error.' + throw new Error( + `Remote capability "${binding.kind}:${binding.instanceId}:${tool.name}" failed: ${message}`, + ) } if ( result.structuredContent && @@ -117,36 +154,59 @@ function createCapabilityFromTool( return { capability, binding } } -export async function synthesizeHomeDomain( +export async function synthesizeRemoteToolDomain( env: Env, - input: { - connectorId: string | null - baseUrl: string - }, -): Promise { - if (!input.connectorId) { - return null - } - - const client = createHomeMcpClient(env, input.connectorId) + ref: RemoteConnectorRef, + allRefs: ReadonlyArray, +): Promise { + const client = createRemoteConnectorMcpClient(env, ref.kind, ref.instanceId) const snapshot = await client.getSnapshot() if (!snapshot || snapshot.tools.length === 0) return null + const domainId = remoteConnectorDomainId(ref) + const capabilityPrefix = remoteConnectorCapabilityPrefix(ref, allRefs) + const k = ref.kind.trim().toLowerCase() + const isOnlyBuiltinHomeDomain = + k === 'home' && + allRefs.length === 1 && + allRefs[0]?.kind === 'home' && + allRefs[0]?.instanceId.trim() === 'default' + + const domainIdForCapabilities: CapabilityDomain = isOnlyBuiltinHomeDomain + ? 'home' + : domainId + + const domainKeywordRoots = + k === 'home' + ? (['home', 'roku', 'lutron', 'automation', 'devices'] as const) + : [k, 'integration', 'connector'] + + const domainDescription = + k === 'home' + ? 'Home automation capabilities discovered from the connected home connector.' + : `Capabilities discovered from the connected "${ref.kind}" remote connector ("${ref.instanceId}").` + const capabilities: Array = [] - const bindings: Record = {} + const bindings: Record = {} for (const tool of snapshot.tools) { - const { capability, binding } = createCapabilityFromTool(snapshot, tool) + const { capability, binding } = createCapabilityFromTool({ + snapshot, + tool, + ref, + domainId: domainIdForCapabilities, + capabilityPrefix, + domainKeywordRoots, + }) capabilities.push(capability) bindings[binding.capabilityName] = binding } return { domain: defineDomain({ - name: capabilityDomainNames.home, - description: - 'Home automation capabilities discovered from the connected home connector.', - keywords: ['home', 'roku', 'lutron', 'automation', 'devices'], + name: domainIdForCapabilities, + description: domainDescription, + keywords: [...domainKeywordRoots], capabilities, }), bindings, diff --git a/packages/worker/src/mcp/capabilities/meta/domain.ts b/packages/worker/src/mcp/capabilities/meta/domain.ts index d8407060f7..6a65bf6778 100644 --- a/packages/worker/src/mcp/capabilities/meta/domain.ts +++ b/packages/worker/src/mcp/capabilities/meta/domain.ts @@ -7,6 +7,7 @@ import { metaMemorySearchCapability } from './meta-memory-search.ts' import { metaMemoryUpsertCapability } from './meta-memory-upsert.ts' import { metaMemoryVerifyCapability } from './meta-memory-verify.ts' import { metaGetHomeConnectorStatusCapability } from './meta-get-home-connector-status.ts' +import { metaListRemoteConnectorStatusCapability } from './meta-list-remote-connector-status.ts' import { metaGetMcpServerInstructionsCapability } from './meta-get-mcp-server-instructions.ts' import { metaGetSkillCapability } from './meta-get-skill.ts' import { metaListCapabilitiesCapability } from './meta-list-capabilities.ts' @@ -34,6 +35,7 @@ export const metaDomain = defineDomain({ metaGetMcpServerInstructionsCapability, metaSetMcpServerInstructionsCapability, metaGetHomeConnectorStatusCapability, + metaListRemoteConnectorStatusCapability, metaMemorySearchCapability, metaMemoryGetCapability, metaMemoryVerifyCapability, diff --git a/packages/worker/src/mcp/capabilities/meta/meta-get-home-connector-status.ts b/packages/worker/src/mcp/capabilities/meta/meta-get-home-connector-status.ts index 00185ec33b..b8d0e50b01 100644 --- a/packages/worker/src/mcp/capabilities/meta/meta-get-home-connector-status.ts +++ b/packages/worker/src/mcp/capabilities/meta/meta-get-home-connector-status.ts @@ -2,6 +2,7 @@ import { z } from 'zod' import { defineDomainCapability } from '#mcp/capabilities/define-domain-capability.ts' import { capabilityDomainNames } from '#mcp/capabilities/domain-metadata.ts' import { getHomeConnectorStatus } from '#worker/home/status.ts' +import { normalizeRemoteConnectorRefs } from '@kody-internal/shared/remote-connectors.ts' const outputSchema = z.object({ status: z.enum(['connected', 'disconnected', 'unavailable', 'error']), @@ -36,9 +37,11 @@ export const metaGetHomeConnectorStatusCapability = defineDomainCapability( inputSchema: z.object({}), outputSchema, async handler(_args, ctx) { + const refs = normalizeRemoteConnectorRefs(ctx.callerContext) + const homeRef = refs.find((r) => r.kind === 'home') const status = await getHomeConnectorStatus( ctx.env, - ctx.callerContext.homeConnectorId ?? null, + homeRef?.instanceId ?? null, ) return { status: status.state, diff --git a/packages/worker/src/mcp/capabilities/meta/meta-list-remote-connector-status.ts b/packages/worker/src/mcp/capabilities/meta/meta-list-remote-connector-status.ts new file mode 100644 index 0000000000..63a1a75b3e --- /dev/null +++ b/packages/worker/src/mcp/capabilities/meta/meta-list-remote-connector-status.ts @@ -0,0 +1,64 @@ +import { z } from 'zod' +import { defineDomainCapability } from '#mcp/capabilities/define-domain-capability.ts' +import { capabilityDomainNames } from '#mcp/capabilities/domain-metadata.ts' +import { getRemoteConnectorStatus } from '#worker/home/status.ts' +import { normalizeRemoteConnectorRefs } from '@kody-internal/shared/remote-connectors.ts' + +const connectorStatusSchema = z.object({ + connector_kind: z.string(), + connector_instance_id: z.string(), + status: z.enum(['connected', 'disconnected', 'unavailable', 'error']), + connected: z.boolean(), + connected_at: z.string().nullable(), + last_seen_at: z.string().nullable(), + tool_count: z.number().int().nonnegative(), + message: z.string(), + error: z.string().nullable(), +}) + +const outputSchema = z.object({ + connectors: z.array(connectorStatusSchema), +}) + +export const metaListRemoteConnectorStatusCapability = defineDomainCapability( + capabilityDomainNames.meta, + { + name: 'meta_list_remote_connector_status', + description: + 'Report connection status for each remote connector attached to this session (kind + instance id). Use when search results miss remote capabilities or a remote capability fails.', + keywords: [ + 'remote', + 'connector', + 'status', + 'connected', + 'disconnected', + 'home', + 'troubleshoot', + ], + readOnly: true, + idempotent: true, + destructive: false, + inputSchema: z.object({}), + outputSchema, + async handler(_args, ctx) { + const refs = normalizeRemoteConnectorRefs(ctx.callerContext) + const connectors = await Promise.all( + refs.map(async (ref) => { + const s = await getRemoteConnectorStatus(ctx.env, ref) + return { + connector_kind: s.connectorKind, + connector_instance_id: s.connectorId ?? ref.instanceId, + status: s.state, + connected: s.connected, + connected_at: s.connectedAt, + last_seen_at: s.lastSeenAt, + tool_count: s.toolCount, + message: s.message, + error: s.error, + } + }), + ) + return { connectors } + }, + }, +) diff --git a/packages/worker/src/mcp/capabilities/registry.ts b/packages/worker/src/mcp/capabilities/registry.ts index 3818020cfd..45dc87876a 100644 --- a/packages/worker/src/mcp/capabilities/registry.ts +++ b/packages/worker/src/mcp/capabilities/registry.ts @@ -3,8 +3,12 @@ import { type BuiltCapabilityRegistry, } from './build-capability-registry.ts' import { builtinDomains } from './builtin-domains.ts' -import { synthesizeHomeDomain } from './home/index.ts' +import { + type SynthesizedRemoteConnectorDomain, + synthesizeRemoteToolDomain, +} from './home/index.ts' import { type McpCallerContext } from '@kody-internal/shared/chat.ts' +import { normalizeRemoteConnectorRefs } from '@kody-internal/shared/remote-connectors.ts' const staticRegistry = buildCapabilityRegistry(builtinDomains) @@ -28,16 +32,31 @@ export async function getCapabilityRegistryForContext(input: { env: Env callerContext: McpCallerContext }): Promise { - const homeDomain = await synthesizeHomeDomain(input.env, { - connectorId: input.callerContext.homeConnectorId ?? null, - baseUrl: input.callerContext.baseUrl, - }) - if (!homeDomain) { + const refs = normalizeRemoteConnectorRefs(input.callerContext) + const synthesizedDomains: Array = + [] + const settled = await Promise.allSettled( + refs.map((ref) => synthesizeRemoteToolDomain(input.env, ref, refs)), + ) + for (const [index, outcome] of settled.entries()) { + if (outcome.status === 'fulfilled' && outcome.value) { + synthesizedDomains.push(outcome.value.domain) + continue + } + if (outcome.status === 'rejected') { + const ref = refs[index] + console.error( + `[getCapabilityRegistryForContext] synthesizeRemoteToolDomain failed for ${ref?.kind ?? '?'}:${ref?.instanceId ?? '?'}`, + outcome.reason, + ) + } + } + if (synthesizedDomains.length === 0) { return staticRegistry } const registry = buildCapabilityRegistry([ ...builtinDomains, - homeDomain.domain, + ...synthesizedDomains, ]) return registry } diff --git a/packages/worker/src/mcp/context.node.test.ts b/packages/worker/src/mcp/context.node.test.ts index a5df1b605c..0ffded03fa 100644 --- a/packages/worker/src/mcp/context.node.test.ts +++ b/packages/worker/src/mcp/context.node.test.ts @@ -9,26 +9,26 @@ test('createMcpCallerContext normalizes missing user to null', () => { ).toEqual({ baseUrl: 'https://example.com', homeConnectorId: null, + remoteConnectors: null, storageContext: null, user: null, }) }) test('parseMcpCallerContext validates caller context shape', () => { - expect( - parseMcpCallerContext({ - baseUrl: 'https://example.com', - user: { - userId: '123', - email: 'user@example.com', - displayName: 'user', - }, - storageContext: { - sessionId: 'session-123', - appId: 'app-123', - }, - }), - ).toEqual({ + const parsed = parseMcpCallerContext({ + baseUrl: 'https://example.com', + user: { + userId: '123', + email: 'user@example.com', + displayName: 'user', + }, + storageContext: { + sessionId: 'session-123', + appId: 'app-123', + }, + }) + expect(parsed).toMatchObject({ baseUrl: 'https://example.com', user: { userId: '123', @@ -40,4 +40,6 @@ test('parseMcpCallerContext validates caller context shape', () => { appId: 'app-123', }, }) + expect(parsed.homeConnectorId ?? null).toBeNull() + expect(parsed.remoteConnectors ?? null).toBeNull() }) diff --git a/packages/worker/src/mcp/context.ts b/packages/worker/src/mcp/context.ts index 5eae70f906..58aae180be 100644 --- a/packages/worker/src/mcp/context.ts +++ b/packages/worker/src/mcp/context.ts @@ -5,6 +5,7 @@ import { type McpStorageContext, type McpUserContext, } from '@kody-internal/shared/chat.ts' +import { type RemoteConnectorRef } from '@kody-internal/shared/remote-connectors.ts' export type McpServerProps = McpCallerContext @@ -12,12 +13,14 @@ export function createMcpCallerContext(input: { baseUrl: string user?: McpUserContext | null homeConnectorId?: string | null + remoteConnectors?: Array | null storageContext?: McpStorageContext | null }): McpCallerContext { return { baseUrl: input.baseUrl, user: input.user ?? null, homeConnectorId: input.homeConnectorId ?? null, + remoteConnectors: input.remoteConnectors ?? null, storageContext: input.storageContext ?? null, } } diff --git a/packages/worker/src/mcp/tools/search-format.ts b/packages/worker/src/mcp/tools/search-format.ts index 93f0014af2..fe6819111f 100644 --- a/packages/worker/src/mcp/tools/search-format.ts +++ b/packages/worker/src/mcp/tools/search-format.ts @@ -37,11 +37,19 @@ export type SearchResultStructuredContent = { retrievalQuery: string } homeConnectorStatus?: { + connectorKind: string connectorId: string state: string connected: boolean toolCount: number } + remoteConnectorStatuses?: Array<{ + connectorKind: string + connectorId: string + state: string + connected: boolean + toolCount: number + }> } export type SlimSearchMatch = diff --git a/packages/worker/src/mcp/tools/search.ts b/packages/worker/src/mcp/tools/search.ts index d95182c076..4b03f25b21 100644 --- a/packages/worker/src/mcp/tools/search.ts +++ b/packages/worker/src/mcp/tools/search.ts @@ -34,9 +34,11 @@ import { import { listValues } from '#mcp/values/service.ts' import { type ValueMetadata } from '#mcp/values/types.ts' import { - getHomeConnectorStatus, + getRemoteConnectorStatus, type HomeConnectorStatus, } from '#worker/home/status.ts' +import { type McpCallerContext } from '@kody-internal/shared/chat.ts' +import { normalizeRemoteConnectorRefs } from '@kody-internal/shared/remote-connectors.ts' import { buildSavedUiUrl } from '#worker/ui-artifact-urls.ts' import { callerContextFields, @@ -140,7 +142,7 @@ Persisted values use \`codemode.value_get\` / \`codemode.value_list\`. Connector use \`codemode.connector_get\` / \`codemode.connector_list\`. If results look incomplete: \`meta_list_capabilities\` (full registry) or -\`meta_get_home_connector_status\` (home connector). +\`meta_list_remote_connector_status\` / \`meta_get_home_connector_status\` (remote connectors). Domain hints for \`query\` / \`skill_collection\`: \`coding\`, \`meta\`, \`home\` (see server instructions). @@ -181,20 +183,19 @@ type SearchRowsAndRegistry = OptionalSearchRowsResult & { appSecretsByAppId: Awaited> } -function shouldIncludeHomeConnectorStatus(status: HomeConnectorStatus) { +function shouldIncludeRemoteConnectorStatus(status: HomeConnectorStatus) { return status.state !== 'connected' || status.toolCount === 0 } -function serializeHomeConnectorStatus(status: HomeConnectorStatus | null): - | { - connectorId: string - state: string - connected: boolean - toolCount: number - } - | undefined { - if (!status) return undefined +function serializeRemoteConnectorStatus(status: HomeConnectorStatus): { + connectorKind: string + connectorId: string + state: string + connected: boolean + toolCount: number +} { return { + connectorKind: status.connectorKind, connectorId: status.connectorId ?? 'unknown', state: status.state, connected: status.connected, @@ -202,15 +203,30 @@ function serializeHomeConnectorStatus(status: HomeConnectorStatus | null): } } +export async function loadDownRemoteConnectorStatuses(input: { + env: Env + callerContext: Pick +}): Promise> { + const refs = normalizeRemoteConnectorRefs(input.callerContext) + const statuses = await Promise.all( + refs.map((ref) => getRemoteConnectorStatus(input.env, ref)), + ) + return statuses.filter(shouldIncludeRemoteConnectorStatus) +} + +/** @deprecated Prefer loadDownRemoteConnectorStatuses with full caller context. */ export async function loadDownHomeConnectorStatus(input: { env: Env homeConnectorId: string | null }): Promise { - const status = await getHomeConnectorStatus(input.env, input.homeConnectorId) - if (!shouldIncludeHomeConnectorStatus(status)) { - return null - } - return status + const statuses = await loadDownRemoteConnectorStatuses({ + env: input.env, + callerContext: { + homeConnectorId: input.homeConnectorId, + remoteConnectors: null, + }, + }) + return statuses[0] ?? null } export async function loadOptionalSearchRows(input: { @@ -557,7 +573,7 @@ export async function registerSearchTool(agent: McpRegistrationAgent) { const limit = args.limit ?? defaultSearchLimit const maxResponseSize = args.maxResponseSize ?? defaultMaxResponseSize let warnings: Array = [] - let homeConnectorStatus: HomeConnectorStatus | null = null + let remoteConnectorDownStatuses: Array = [] const searchSpan = async () => { const searchRows = await loadSearchRowsAndRegistry({ @@ -566,9 +582,9 @@ export async function registerSearchTool(agent: McpRegistrationAgent) { userId, skillCollection: args.skill_collection, }) - homeConnectorStatus = await loadDownHomeConnectorStatus({ + remoteConnectorDownStatuses = await loadDownRemoteConnectorStatuses({ env: agent.getEnv(), - homeConnectorId: callerContext.homeConnectorId ?? null, + callerContext, }) warnings = searchRows.warnings @@ -689,8 +705,15 @@ export async function registerSearchTool(agent: McpRegistrationAgent) { } } + const normalizedRemoteConnectorStatuses = + remoteConnectorDownStatuses.length > 0 + ? remoteConnectorDownStatuses.map(serializeRemoteConnectorStatus) + : undefined const normalizedHomeConnectorStatus = - serializeHomeConnectorStatus(homeConnectorStatus) + remoteConnectorDownStatuses.length === 1 && + remoteConnectorDownStatuses[0]?.connectorKind === 'home' + ? serializeRemoteConnectorStatus(remoteConnectorDownStatuses[0]!) + : undefined const memoryToolContext = await loadRelevantMemoriesForTool({ env: agent.getEnv(), callerContext, @@ -714,11 +737,19 @@ export async function registerSearchTool(agent: McpRegistrationAgent) { warnings: Array memories?: SearchResultStructuredContent['memories'] homeConnectorStatus?: { + connectorKind: string connectorId: string state: string connected: boolean toolCount: number } + remoteConnectorStatuses?: Array<{ + connectorKind: string + connectorId: string + state: string + connected: boolean + toolCount: number + }> } = { matches: outcome.result.matches, offline: outcome.result.offline, @@ -728,6 +759,11 @@ export async function registerSearchTool(agent: McpRegistrationAgent) { memories: searchMemories, } : {}), + ...(normalizedRemoteConnectorStatuses + ? { + remoteConnectorStatuses: normalizedRemoteConnectorStatuses, + } + : {}), ...(normalizedHomeConnectorStatus ? { homeConnectorStatus: normalizedHomeConnectorStatus, @@ -787,6 +823,11 @@ export async function registerSearchTool(agent: McpRegistrationAgent) { ...(trimmedPayload.homeConnectorStatus ? { homeConnectorStatus: trimmedPayload.homeConnectorStatus } : {}), + ...(trimmedPayload.remoteConnectorStatuses + ? { + remoteConnectorStatuses: trimmedPayload.remoteConnectorStatuses, + } + : {}), matches: toSlimStructuredMatches({ matches: trimmedPayload.matches, baseUrl, diff --git a/packages/worker/src/remote-connector/connector-session-key.node.test.ts b/packages/worker/src/remote-connector/connector-session-key.node.test.ts new file mode 100644 index 0000000000..af9ed4371d --- /dev/null +++ b/packages/worker/src/remote-connector/connector-session-key.node.test.ts @@ -0,0 +1,51 @@ +import { expect, test } from 'vitest' +import { + connectorIngressPath, + connectorSessionKey, + parseConnectorRoutePath, +} from './connector-session-key.ts' + +test('connectorSessionKey preserves home instance id', () => { + expect(connectorSessionKey('home', 'default')).toBe('default') + expect(connectorSessionKey('HOME', 'living-room')).toBe('living-room') +}) + +test('connectorSessionKey prefixes home ids containing colons', () => { + expect(connectorSessionKey('home', 'other:default')).toBe( + 'home:other:default', + ) +}) + +test('connectorSessionKey prefixes non-home kinds', () => { + expect(connectorSessionKey('custom', 'alpha')).toBe('custom:alpha') +}) + +test('parseConnectorRoutePath handles generic and legacy paths', () => { + expect(parseConnectorRoutePath('/connectors/custom/my-id/snapshot')).toEqual({ + kind: 'custom', + instanceId: 'my-id', + rest: '/snapshot', + }) + expect( + parseConnectorRoutePath('/home/connectors/default/rpc/tools-list'), + ).toEqual({ + kind: 'home', + instanceId: 'default', + rest: '/rpc/tools-list', + }) + expect( + parseConnectorRoutePath('/connectors/home/default/rpc/tools-list'), + ).toEqual({ + kind: 'home', + instanceId: 'default', + rest: '/rpc/tools-list', + }) + expect(parseConnectorRoutePath('/home/connectors')).toBeNull() +}) + +test('connectorIngressPath prefers legacy home URL', () => { + expect(connectorIngressPath('home', 'default')).toBe( + '/home/connectors/default', + ) + expect(connectorIngressPath('custom', 'a b')).toBe('/connectors/custom/a%20b') +}) diff --git a/packages/worker/src/remote-connector/connector-session-key.ts b/packages/worker/src/remote-connector/connector-session-key.ts new file mode 100644 index 0000000000..e99d256cb6 --- /dev/null +++ b/packages/worker/src/remote-connector/connector-session-key.ts @@ -0,0 +1,67 @@ +/** + * Stable Durable Object id segment for a remote connector WebSocket session. + * For kind "home" and instanceId X, returns X unchanged so existing deployments + * keep the same DO id as before (idFromName(connectorId) only). Home instance + * ids containing ":" are prefixed to avoid collisions with non-home keys. + */ +export function connectorSessionKey(kind: string, instanceId: string): string { + const k = kind.trim().toLowerCase() + const id = instanceId.trim() + if (k === 'home') { + if (id.includes(':')) { + return `home:${id}` + } + return id + } + return `${k}:${id}` +} + +export function parseConnectorRoutePath(pathname: string): { + kind: string + instanceId: string + rest: string +} | null { + const parts = pathname.split('/').filter(Boolean) + const decodeSegment = (value: string) => { + try { + return decodeURIComponent(value) + } catch { + return null + } + } + // /connectors/:kind/:instanceId/... + if (parts.length >= 3 && parts[0] === 'connectors' && parts[1] && parts[2]) { + const decodedKind = decodeSegment(parts[1]!) + const decodedInstanceId = decodeSegment(parts[2]!) + if (!decodedKind || !decodedInstanceId) return null + const kind = decodedKind.trim() + const instanceId = decodedInstanceId.trim() + if (!kind || !instanceId) return null + const rest = parts.length > 3 ? `/${parts.slice(3).join('/')}` : '' + return { kind, instanceId, rest } + } + // /home/connectors/:instanceId/... + if ( + parts.length >= 3 && + parts[0] === 'home' && + parts[1] === 'connectors' && + parts[2] + ) { + const decodedInstanceId = decodeSegment(parts[2]!) + if (!decodedInstanceId) return null + const instanceId = decodedInstanceId.trim() + if (!instanceId) return null + const rest = parts.length > 3 ? `/${parts.slice(3).join('/')}` : '' + return { kind: 'home', instanceId, rest } + } + return null +} + +export function connectorIngressPath(kind: string, instanceId: string): string { + const k = kind.trim().toLowerCase() + const id = encodeURIComponent(instanceId.trim()) + if (k === 'home') { + return `/home/connectors/${id}` + } + return `/connectors/${encodeURIComponent(k)}/${id}` +} diff --git a/packages/worker/src/remote-connector/remote-connectors-shared.node.test.ts b/packages/worker/src/remote-connector/remote-connectors-shared.node.test.ts new file mode 100644 index 0000000000..403a1b186a --- /dev/null +++ b/packages/worker/src/remote-connector/remote-connectors-shared.node.test.ts @@ -0,0 +1,35 @@ +import { expect, test } from 'vitest' +import { normalizeRemoteConnectorRefs } from '@kody-internal/shared/remote-connectors.ts' + +test('normalizeRemoteConnectorRefs maps homeConnectorId when remoteConnectors unset', () => { + expect( + normalizeRemoteConnectorRefs({ + homeConnectorId: 'living-room', + remoteConnectors: undefined, + }), + ).toEqual([{ kind: 'home', instanceId: 'living-room' }]) +}) + +test('normalizeRemoteConnectorRefs uses remoteConnectors when provided', () => { + expect( + normalizeRemoteConnectorRefs({ + homeConnectorId: 'ignored', + remoteConnectors: [ + { kind: 'Home', instanceId: ' a ' }, + { kind: 'custom', instanceId: 'x' }, + ], + }), + ).toEqual([ + { kind: 'home', instanceId: 'a' }, + { kind: 'custom', instanceId: 'x' }, + ]) +}) + +test('normalizeRemoteConnectorRefs empty array does not fall back to homeConnectorId', () => { + expect( + normalizeRemoteConnectorRefs({ + homeConnectorId: 'living-room', + remoteConnectors: [], + }), + ).toEqual([]) +}) diff --git a/packages/worker/src/remote-connector/remote-domain-id.ts b/packages/worker/src/remote-connector/remote-domain-id.ts new file mode 100644 index 0000000000..d39a29b66c --- /dev/null +++ b/packages/worker/src/remote-connector/remote-domain-id.ts @@ -0,0 +1,43 @@ +import { type RemoteConnectorRef } from '@kody-internal/shared/remote-connectors.ts' + +export function remoteConnectorDomainId(ref: RemoteConnectorRef): string { + const k = ref.kind.trim().toLowerCase() + const id = + ref.instanceId + .trim() + .replaceAll(/[^\w-]+/g, '_') + .replaceAll(/_+/g, '_') + .replace(/^_|_$/g, '') || 'instance' + return `remote:${k}:${id}` +} + +/** + * Prefix for synthesized capability names. Keeps legacy `home_*` names when + * there is a single home connector with instance id `default`. + */ +export function remoteConnectorCapabilityPrefix( + ref: RemoteConnectorRef, + allRefs: ReadonlyArray, +): string { + const k = ref.kind.trim().toLowerCase() + const rawId = ref.instanceId.trim() + const slug = + rawId + .replaceAll(/[^\w]+/g, '_') + .replaceAll(/_+/g, '_') + .replace(/^_|_$/g, '') || 'instance' + + if (k === 'home') { + const isOnlyBuiltinHome = + rawId === 'default' && + allRefs.length === 1 && + allRefs[0]?.kind.trim().toLowerCase() === 'home' && + allRefs[0]?.instanceId.trim() === 'default' + if (isOnlyBuiltinHome) { + return 'home' + } + return `home_${slug}` + } + + return `${k}_${slug}` +} diff --git a/packages/worker/src/remote-connector/resolve-remote-connector-secret.node.test.ts b/packages/worker/src/remote-connector/resolve-remote-connector-secret.node.test.ts new file mode 100644 index 0000000000..7c807e5113 --- /dev/null +++ b/packages/worker/src/remote-connector/resolve-remote-connector-secret.node.test.ts @@ -0,0 +1,37 @@ +import { expect, test } from 'vitest' +import { resolveRemoteConnectorSharedSecret } from './resolve-remote-connector-secret.ts' + +test('falls back to HOME_CONNECTOR_SHARED_SECRET for home kind', () => { + expect( + resolveRemoteConnectorSharedSecret('home', 'any', { + HOME_CONNECTOR_SHARED_SECRET: 'legacy-secret', + } as Env), + ).toBe('legacy-secret') +}) + +test('REMOTE_CONNECTOR_SECRETS overrides per kind and instance', () => { + const env = { + HOME_CONNECTOR_SHARED_SECRET: 'legacy-secret', + REMOTE_CONNECTOR_SECRETS: { + 'custom:alpha': 'alpha-secret', + 'home:default': 'home-override', + }, + } as Env + expect(resolveRemoteConnectorSharedSecret('custom', 'alpha', env)).toBe( + 'alpha-secret', + ) + expect(resolveRemoteConnectorSharedSecret('home', 'default', env)).toBe( + 'home-override', + ) + expect(resolveRemoteConnectorSharedSecret('home', 'other', env)).toBe( + 'legacy-secret', + ) +}) + +test('non-home kind has no legacy fallback when map missing', () => { + expect( + resolveRemoteConnectorSharedSecret('custom', 'alpha', { + HOME_CONNECTOR_SHARED_SECRET: 'legacy-secret', + } as Env), + ).toBeUndefined() +}) diff --git a/packages/worker/src/remote-connector/resolve-remote-connector-secret.ts b/packages/worker/src/remote-connector/resolve-remote-connector-secret.ts new file mode 100644 index 0000000000..d47acefb29 --- /dev/null +++ b/packages/worker/src/remote-connector/resolve-remote-connector-secret.ts @@ -0,0 +1,47 @@ +/** + * Resolve the shared secret for a remote connector WebSocket hello. + * Precedence: REMOTE_CONNECTOR_SECRETS map key "kind:instanceId", then + * legacy HOME_CONNECTOR_SHARED_SECRET when kind is "home". + */ +function parseSecretsMapFromEnv(value: unknown): Record | null { + if (!value) return null + if (typeof value === 'object' && !Array.isArray(value)) { + return value as Record + } + if (typeof value !== 'string') return null + const trimmed = value.trim() + if (!trimmed) return null + try { + const parsed = JSON.parse(trimmed) as unknown + if (parsed && typeof parsed === 'object' && !Array.isArray(parsed)) { + return parsed as Record + } + } catch (error) { + const detail = error instanceof Error ? error.message : String(error) + console.error( + `[REMOTE_CONNECTOR_SECRETS] invalid JSON (ignored for map lookup): ${detail}`, + ) + } + return null +} + +export function resolveRemoteConnectorSharedSecret( + kind: string, + instanceId: string, + env: Env, +): string | undefined { + const k = kind.trim().toLowerCase() + const id = instanceId.trim() + const map = parseSecretsMapFromEnv(env.REMOTE_CONNECTOR_SECRETS as unknown) + if (map) { + const key = `${k}:${id}` + const fromMap = map[key] + if (typeof fromMap === 'string' && fromMap.trim()) { + return fromMap.trim() + } + } + if (k === 'home') { + return env.HOME_CONNECTOR_SHARED_SECRET?.trim() + } + return undefined +}