Skip to content

Commit a2a38b0

Browse files
committed
feat(mothership): receive private Sim controls over outbound transport
1 parent 27287a8 commit a2a38b0

17 files changed

Lines changed: 537 additions & 14 deletions

File tree

‎apps/sim/instrumentation-node.ts‎

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -403,4 +403,7 @@ export async function register() {
403403

404404
const { startMemoryTelemetry } = await import('./lib/monitoring/memory-telemetry')
405405
startMemoryTelemetry()
406+
407+
const { startSimReceivers } = await import('./lib/mothership/transport/receiver')
408+
await startSimReceivers()
406409
}

‎apps/sim/lib/core/config/env.ts‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -557,6 +557,7 @@ export const env = createEnv({
557557
E2B_FUNCTION_TEMPLATE_ID: z.string().refine(isImmutableE2BTemplateRef, { message: `E2B_FUNCTION_TEMPLATE_ID ${IMMUTABLE_E2B_TEMPLATE_REF_ERROR}` }).optional(), // Immutable dedicated E2B build for Function JavaScript/Python/Shell and workspace sandbox layers; no Mothership fallback
558558
E2B_FUNCTION_TEMPLATE_GENERATION: z.string().refine(isValidSandboxReleaseGeneration, { message: `E2B_FUNCTION_TEMPLATE_GENERATION ${SANDBOX_RELEASE_GENERATION_ERROR}` }).optional(), // Monotonic release epoch printed by the Function E2B builder
559559
MOTHERSHIP_E2B_TEMPLATE_ID: z.string().optional(), // Mothership code-tool template; never a Function-base fallback
560+
MOTHERSHIP_SIM_TRANSPORT: z.enum(['direct', 'checkpoint']).optional(), // Server-side Sim delivery; hosted defaults to direct, self-hosted to outbound checkpoint delivery
560561
MOTHERSHIP_SANDBOX_CLI_ENDPOINT: z.string().optional(), // Sim API base the sandboxed sim CLI calls back to; defaults to NEXT_PUBLIC_APP_URL (set when the public URL is not reachable from the sandbox network)
561562
MOTHERSHIP_E2B_DOC_TEMPLATE_ID: z.string().optional(), // Dedicated E2B template with python-pptx/docx/openpyxl/reportlab for document generation; when set (and E2B enabled), docs compile via Python instead of the JS isolated-vm path
562563
E2B_PI_TEMPLATE_ID: z.string().optional(), // E2B template ID/alias with the Pi CLI + git baked in (Create PR, its Babysit continuation, and Review Code)

‎apps/sim/lib/mothership/agent-cli/sink.test.ts‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -40,6 +40,7 @@ describe('stdout sink publication', () => {
4040
expect(saved).toMatchObject({ exitCode: 0, stderr: result.stderr })
4141
expect(saved.stdout).toContain('Command succeeded')
4242
expect(saved.stdout).toContain('Do not repeat a mutation')
43+
expect(saved.sinkError).toContain('Command succeeded')
4344
expect(saved.stdout).toContain(result.stdout)
4445
expect(write).not.toHaveBeenCalled()
4546
})
@@ -48,6 +49,7 @@ describe('stdout sink publication', () => {
4849
write.mockResolvedValue({ outcome: 'error', detail: 'storage unavailable' })
4950
const saved = await applySink(sink, 'chat', result)
5051
expect(saved.exitCode).toBe(0)
52+
expect(saved.sinkError).toContain('could not be confirmed')
5153
expect(saved.stdout).toContain('Command succeeded')
5254
expect(saved.stdout).toContain(result.stdout)
5355
})

‎apps/sim/lib/mothership/agent-cli/sink.ts‎

Lines changed: 11 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -22,10 +22,18 @@ export async function applySink(
2222
signal?: AbortSignal,
2323
observeOutput?: (value: string) => Promise<SessionFileObserver>
2424
): Promise<AgentCliRawResult> {
25-
if (result.exitCode !== 0 || signal?.aborted) return result
25+
if (result.exitCode !== 0) return result
26+
if (signal?.aborted)
27+
return {
28+
...result,
29+
sinkError:
30+
'Command completed; output publication was cancelled. Inspect the file before using it. Do not repeat a mutation.',
31+
}
2632
if (!sessionKey) {
2733
return {
2834
...result,
35+
sinkError:
36+
'Command succeeded; outputFile was not written because no chat workbench exists. Do not repeat a mutation.',
2937
stdout: `${result.stdout}\n[outputFile not written: no chat-scoped machine — output returned inline instead]`,
3038
}
3139
}
@@ -49,6 +57,8 @@ export async function applySink(
4957
}
5058
return {
5159
...result,
60+
sinkError:
61+
'Command succeeded; output publication could not be confirmed. Inspect the destination before using it. Do not repeat a mutation.',
5262
stdout: `[Command succeeded; writing its output to ${sink.path} could not be confirmed. The output follows inline. Do not repeat a mutation to recover its output.]\n${result.stdout}`,
5363
}
5464
}

‎apps/sim/lib/mothership/generated/agent-cli.ts‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -57,6 +57,7 @@ export const AgentCliRequest = z.object({
5757
export type AgentCliRequest = z.infer<typeof AgentCliRequest>;
5858

5959
export const AgentCliRawResult = z.object({
60+
sinkError: z.string().optional(),
6061
exitCode: z.number().int(),
6162
stdout: z.string(),
6263
stderr: z.string().default(""),

‎apps/sim/lib/mothership/generated/protocol.ts‎

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -16,6 +16,7 @@
1616
*/
1717

1818
import { z } from "zod";
19+
import { SimConnection } from "./sim-transport";
1920

2021
export const PROTOCOL_VERSION = 1;
2122

@@ -73,6 +74,7 @@ const WorkspaceInventorySchema = z.object({
7374
});
7475

7576
export const ChatPayloadSchema = z.strictObject({
77+
simConnection: SimConnection.optional(),
7678
message: z.string().min(1),
7779
...ResponseReceiptSchema.shape,
7880
userId: z.string().min(1),
@@ -139,6 +141,7 @@ export interface StreamToolReplay {
139141

140142
/** POST /api/mothership — the chat request sim sends. */
141143
export interface ChatRequest extends StreamResponseReceipt {
144+
simConnection?: SimConnection | undefined;
142145
message: string;
143146
userId: string;
144147
/** Bump-gated (S43): senders include it; the worker 426s on mismatch. */
Lines changed: 49 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,49 @@
1+
// GENERATED — do not edit. Source of truth: mothership worker packages/contracts/src/sim-transport.ts
2+
// Regenerate with `bun run contracts:sync` in the worker.
3+
4+
import { z } from "zod";
5+
import { RunControlRequest } from "./run-control";
6+
import { TaskWakeRequest, WorkflowWatchRequest } from "./tasks";
7+
8+
export const SimConnection = z.discriminatedUnion("mode", [
9+
z.object({ mode: z.literal("direct") }),
10+
z.object({ mode: z.literal("checkpoint"), channelId: z.string().regex(/^[a-f0-9]{64}$/) }),
11+
]);
12+
export type SimConnection = z.infer<typeof SimConnection>;
13+
14+
export const SimScope = z.object({
15+
userId: z.string().min(1),
16+
workspaceId: z.uuid(),
17+
chatId: z.uuid(),
18+
});
19+
export type SimScope = z.infer<typeof SimScope>;
20+
21+
/** Only idempotent controls use the idle channel; tool effects use the run event log. */
22+
export const SimControlOperation = z.discriminatedUnion("kind", [
23+
z.object({ kind: z.literal("run_control"), input: RunControlRequest }),
24+
z.object({ kind: z.literal("workflow_status"), input: WorkflowWatchRequest }),
25+
z.object({ kind: z.literal("wake"), input: TaskWakeRequest }),
26+
]);
27+
export type SimControlOperation = z.infer<typeof SimControlOperation>;
28+
29+
export const SimControlRequest = z.object({
30+
id: z.uuid(),
31+
scope: SimScope,
32+
operation: SimControlOperation,
33+
expiresAt: z.number().int().positive(),
34+
});
35+
export type SimControlRequest = z.infer<typeof SimControlRequest>;
36+
37+
/** HTTP and outbound delivery preserve the same status and JSON response bytes. */
38+
export const SimControlResult = z.object({
39+
status: z.number().int().min(200).max(599),
40+
body: z.string().max(16_000_000),
41+
});
42+
export type SimControlResult = z.infer<typeof SimControlResult>;
43+
44+
export const SimChannelPoll = z.object({ channelId: z.string().regex(/^[a-f0-9]{64}$/) });
45+
export const SimChannelBatch = z.object({ requests: z.array(SimControlRequest).max(16) });
46+
export const SimChannelReply = SimChannelPoll.extend({
47+
id: z.uuid(),
48+
result: SimControlResult,
49+
});

‎apps/sim/lib/mothership/request/lifecycle/run.test.ts‎

Lines changed: 36 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -39,6 +39,7 @@ const {
3939
mockUpdateRunStatus: vi.fn(),
4040
mockCheckAttributedUsageLimits: vi.fn(),
4141
mockEnv: {
42+
INTERNAL_API_SECRET: 'transport-test-secret-000000000000000000',
4243
COPILOT_API_KEY: undefined as string | undefined,
4344
MSHIP_SYSPROMPT_OVERRIDE: undefined as string | undefined,
4445
},
@@ -547,7 +548,33 @@ describe('runCopilotLifecycle', () => {
547548
const sent = JSON.parse(capturedRequestBody)
548549
/** Receipt metadata leaves ordinary caller content unchanged. */
549550
expect(sent).not.toHaveProperty('byokApiKey')
550-
expect(sent).toEqual({ ...payload, receivedTextChars: 0 })
551+
expect(sent).toEqual({
552+
...payload,
553+
receivedTextChars: 0,
554+
simConnection: {
555+
mode: 'checkpoint',
556+
channelId: expect.stringMatching(/^[a-f0-9]{64}$/),
557+
},
558+
})
559+
})
560+
561+
it('stamps server-owned transport over caller claims without exposing the instance secret', async () => {
562+
let captured = ''
563+
mockRunStreamLoop.mockImplementationOnce(async (_url: string, request: RequestInit) => {
564+
captured = String(request.body)
565+
})
566+
await runCopilotLifecycle(
567+
{ message: 'hello', simConnection: { mode: 'direct' } },
568+
{
569+
userId: 'user-1',
570+
workspaceId: 'ws-1',
571+
}
572+
)
573+
expect(JSON.parse(captured).simConnection).toEqual({
574+
mode: 'checkpoint',
575+
channelId: expect.stringMatching(/^[a-f0-9]{64}$/),
576+
})
577+
expect(captured).not.toContain(mockEnv.INTERNAL_API_SECRET)
551578
})
552579

553580
it('attaches the resolved enterprise BYOK key to the outbound payload', async () => {
@@ -949,7 +976,14 @@ describe('runCopilotLifecycle', () => {
949976

950977
const sent = JSON.parse(capturedRequestBody)
951978
expect(sent).not.toHaveProperty('byokApiKey')
952-
expect(sent).toEqual({ ...payload, receivedTextChars: 0 })
979+
expect(sent).toEqual({
980+
...payload,
981+
receivedTextChars: 0,
982+
simConnection: {
983+
mode: 'checkpoint',
984+
channelId: expect.stringMatching(/^[a-f0-9]{64}$/),
985+
},
986+
})
953987
}
954988
)
955989

‎apps/sim/lib/mothership/request/lifecycle/run.ts‎

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -68,6 +68,7 @@ import type {
6868
import type { SecretMountPolicy } from '@/lib/mothership/secret-mount-policy'
6969
import { getMothershipBaseURL } from '@/lib/mothership/server/agent-url'
7070
import { prepareExecutionContext } from '@/lib/mothership/tools/handlers/context'
71+
import { getSimConnection } from '@/lib/mothership/transport/connection'
7172
import { filterModelSafeWorkspaceFileAttachments } from '@/lib/uploads/contexts/workspace/workspace-file-secret-provenance'
7273
import { refuseResolvedSecretProjection } from '@/executor/utils/resolved-secret-projection-refusal'
7374
import type { ResolvedSecretTraceRegistry } from '@/executor/utils/resolved-secret-trace-registry'
@@ -891,6 +892,10 @@ async function runCheckpointLoop(
891892
let retry: StreamRetryWindow | undefined
892893
const callerOnEvent = options.onEvent
893894
const mothershipBaseURL = await getMothershipBaseURL({ userId: options.userId })
895+
if (initialRoute === '/api/mothership' || initialRoute === '/api/copilot') {
896+
const simConnection = getSimConnection()
897+
payload = { ...payload, simConnection }
898+
}
894899
const lifecycleWorkspaceId = nonBlankString(options.workspaceId)
895900
const mothershipRequestId = nonBlankString(options.simRequestId) ?? generateId()
896901
if (!options.simRequestId) {

‎apps/sim/lib/mothership/tools/hosted-workbench.smoke.test.ts‎

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -102,6 +102,11 @@ describe.skipIf(!enabled)('configured hosted Mship workbench', () => {
102102
expect(written).toEqual({ outcome: 'written', path: '/home/user/input.bin' })
103103
expect(await listed(key)).toEqual([machine.sandboxId])
104104
expect(independent.sandboxId).not.toBe(machine.sandboxId)
105+
const privateBoundary = await machine.runCode(
106+
"import os\nassert not any(os.environ.get(k) for k in ['SIM_API_KEY', 'SIM_ENDPOINT', 'SIM_WORKSPACE'])\nprint('compute without Sim credentials')",
107+
{ timeoutMs: 20_000 }
108+
)
109+
expect(privateBoundary.stdout.trim()).toBe('compute without Sim credentials')
105110
const results = await Promise.all([
106111
machine.runCode(
107112
"from pathlib import Path\nb = Path('input.bin').read_bytes()\nPath('answer.bin').write_bytes(b[::-1])\nprint(sum(b))",

0 commit comments

Comments
 (0)