Skip to content
Merged
Show file tree
Hide file tree
Changes from 7 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
7 changes: 7 additions & 0 deletions apps/ade-cli/src/cli.ts
Original file line number Diff line number Diff line change
Expand Up @@ -13978,6 +13978,7 @@ async function runServe(
{ createBrainProjectActionsSyncHandler },
{ buildRosterSnapshot, createForeignChatTranscriptResolver },
{ createSyncCloudRelayStore },
{ setSyncRuntimeRpcHandlerFactory },
] = await Promise.all([
import("./services/projects/machineLayout"),
import("./services/projects/projectRegistry"),
Expand All @@ -13989,6 +13990,7 @@ async function runServe(
import("./services/sync/brainProjectActionsSyncHandler"),
import("./services/sync/rosterBuilder"),
import("./services/sync/syncCloudRelayStore"),
import("./services/sync/syncPairedChannelService"),
]);

const layout = resolveMachineAdeLayout();
Expand Down Expand Up @@ -14377,6 +14379,7 @@ async function runServe(
},
});
const previousRole = process.env.ADE_DEFAULT_ROLE;
let clearSyncRuntimeRpcHandlerFactory: (() => void) | null = null;
process.env.ADE_DEFAULT_ROLE = options.role;
try {

Expand All @@ -14399,6 +14402,7 @@ async function runServe(
disposeScopesOnDispose: false,
onShutdown: finish,
});
clearSyncRuntimeRpcHandlerFactory = setSyncRuntimeRpcHandlerFactory(createHandler);
const startSyncHost = async () => {
if (preferredSyncProjectId) {
return await scopeRegistry.switchSyncHost(preferredSyncProjectId);
Expand Down Expand Up @@ -14435,6 +14439,8 @@ async function runServe(
if (sharedSyncListener) {
await sharedSyncListener.close().catch(() => {});
}
clearSyncRuntimeRpcHandlerFactory?.();
clearSyncRuntimeRpcHandlerFactory = null;
};

const listen = async (
Expand Down Expand Up @@ -14570,6 +14576,7 @@ async function runServe(
}
return null;
} finally {
clearSyncRuntimeRpcHandlerFactory?.();
removeRuntimeProcessErrorBoundary();
if (previousRole == null) delete process.env.ADE_DEFAULT_ROLE;
else process.env.ADE_DEFAULT_ROLE = previousRole;
Expand Down
71 changes: 58 additions & 13 deletions apps/ade-cli/src/services/sync/brainProjectActionsSyncHandler.ts
Original file line number Diff line number Diff line change
Expand Up @@ -32,8 +32,12 @@ import { nowIso } from "../../../../desktop/src/main/services/shared/utils";
import type { SharedSyncListenerConnectionHandler } from "./sharedSyncListener";
import { SYNC_HOST_BIND_LOOPBACK_ONLY } from "./sharedSyncListener";
import type { SyncCredentialStore } from "../credentials/credentialStore";
import { createSyncPairingStore } from "./syncPairingStore";
import { createSyncPairingStore, type SyncPairingRecord } from "./syncPairingStore";
import { createSyncDpopNonceCache, evaluatePairedHelloDpop } from "./syncDpop";
import {
createSyncPairedChannelService,
isPairedRuntimeEnvelopeType,
} from "./syncPairedChannelService";
import { createSyncPinStore } from "./syncPinStore";
import { createSyncSecurityStore } from "./syncSecurityStore";
import {
Expand All @@ -46,6 +50,7 @@ import {
import {
buildSyncHostHelloOkPayload,
buildSyncProjectCatalogMessages,
isRuntimeHostPairingRecord,
type SyncProjectCatalogProvider,
} from "./syncHostService";
import { resolveDeviceDisplayName } from "./deviceRegistryService";
Expand All @@ -71,9 +76,11 @@ type BrainProjectActionsSyncHandlerArgs = {
type BrainPeerState = {
ws: WebSocket;
authenticated: boolean;
authKind: "bootstrap" | "paired" | null;
authTimeout: ReturnType<typeof setTimeout> | null;
metadata: SyncPeerMetadata | null;
personalChatSubscriptions: Map<string, { transcriptPath: string; offset: number }>;
pairingRecord: SyncPairingRecord | null;
};

const WS_OPEN = 1;
Expand Down Expand Up @@ -222,22 +229,33 @@ function parsePairingRequestPayload(payload: unknown): SyncPairingRequestPayload
const code = optionalString(record.code);
const peer = normalizePeerMetadata(record.peer);
const dpopPublicKey = optionalString(record.dpopPublicKey);
return code && peer ? { code, peer, ...(dpopPublicKey ? { dpopPublicKey } : {}) } : null;
const runtimeHostGrant = optionalString(record.runtimeHostGrant);
return code && peer ? {
code,
peer,
...(dpopPublicKey ? { dpopPublicKey } : {}),
...(runtimeHostGrant ? { runtimeHostGrant } : {}),
} : null;
}

function send(
ws: WebSocket,
type: SyncEnvelope["type"],
payload: unknown,
requestId?: string | null,
): void {
if (ws.readyState !== WS_OPEN) return;
ws.send(encodeSyncEnvelope({
type,
requestId,
payload,
compressionThresholdBytes: DEFAULT_SYNC_COMPRESSION_THRESHOLD_BYTES,
}));
): boolean {
if (ws.readyState !== WS_OPEN) return false;
try {
ws.send(encodeSyncEnvelope({
type,
requestId,
payload,
compressionThresholdBytes: DEFAULT_SYNC_COMPRESSION_THRESHOLD_BYTES,
}));
return true;
} catch {
return false;
}
}

function sendProjectCatalog(
Expand Down Expand Up @@ -416,6 +434,11 @@ export function createBrainProjectActionsSyncHandler(
const pollIntervalMs = Math.max(100, Math.floor(args.pollIntervalMs ?? 1_500));
const authTimeoutMs = Math.max(1_000, Math.floor(args.authTimeoutMs ?? BRAIN_SYNC_AUTH_TIMEOUT_MS));
const pairFailures = createPairFailureTracker();
const pairedChannelService = createSyncPairedChannelService<WebSocket>({
logger: args.logger,
getBufferedAmount: (ws) => ws.bufferedAmount,
send: (ws, type, payload) => send(ws, type, payload),
});

const brainMetadata = (): SyncPeerMetadata => ({
deviceId: localDeviceId,
Expand Down Expand Up @@ -451,6 +474,17 @@ export function createBrainProjectActionsSyncHandler(
};

const handleAuthenticatedEnvelope = async (peer: BrainPeerState, envelope: ReturnType<typeof parseSyncEnvelope>): Promise<void> => {
if (isPairedRuntimeEnvelopeType(envelope.type)) {
await pairedChannelService.handleEnvelope(
peer.ws,
envelope.type,
envelope.payload,
peer.authKind === "paired",
// Runtime RPC channel + port-forward are desktop-runtime-host only.
peer.authKind === "paired" && isRuntimeHostPairingRecord(peer.pairingRecord),
);
return;
}
switch (envelope.type) {
case "project_catalog_request": {
sendProjectCatalog(peer.ws, await projectCatalog(args.projectCatalogProvider, args.logger), envelope.requestId);
Expand Down Expand Up @@ -724,9 +758,11 @@ export function createBrainProjectActionsSyncHandler(
const peer: BrainPeerState = {
ws,
authenticated: false,
authKind: null,
authTimeout: null,
metadata: null,
personalChatSubscriptions: new Map(),
pairingRecord: null,
};
let personalChatPumpRunning = false;
const personalChatPump = setInterval(() => {
Expand Down Expand Up @@ -809,6 +845,7 @@ export function createBrainProjectActionsSyncHandler(
try {
const paired = pairingStore.pairPeer(payload.peer, payload.code, {
dpopPublicKey: payload.dpopPublicKey ?? null,
runtimeHostGrant: payload.runtimeHostGrant ?? null,
});
pairFailures.clearAfterSuccess(remoteAddress ?? null);
send(ws, "pairing_result", { ok: true, deviceId: paired.deviceId, secret: paired.secret }, envelope.requestId);
Expand Down Expand Up @@ -862,14 +899,15 @@ export function createBrainProjectActionsSyncHandler(
return;
}
const auth = hello.auth;
let authenticatedPairingRecord: SyncPairingRecord | null = null;
const authFailed = (() => {
if (auth?.kind === "paired") {
if (auth.deviceId !== hello.peer.deviceId) return true;
if (!pairingStore.authenticate(auth.deviceId, auth.secret)) return true;
const record = pairingStore.getPairingRecord(auth.deviceId);
if (!record) return true;
authenticatedPairingRecord = pairingStore.getPairingRecord(auth.deviceId);
if (!authenticatedPairingRecord) return true;
const dpopFailure = evaluatePairedHelloDpop({
storedPublicKey: record.dpopPublicKey,
storedPublicKey: authenticatedPairingRecord.dpopPublicKey,
deviceId: auth.deviceId,
secret: auth.secret,
proof: auth.dpop ?? null,
Expand Down Expand Up @@ -920,8 +958,10 @@ export function createBrainProjectActionsSyncHandler(
return;
}
peer.authenticated = true;
peer.authKind = auth?.kind ?? null;
clearAuthTimeout();
peer.metadata = hello.peer;
peer.pairingRecord = auth?.kind === "paired" ? authenticatedPairingRecord : null;
const catalog = await projectCatalog(args.projectCatalogProvider, args.logger);
const brain = brainMetadata();
const personalDescriptors = personalChatCommandDescriptors(args.personalChatScope);
Expand All @@ -940,6 +980,10 @@ export function createBrainProjectActionsSyncHandler(
localCommandDescriptors: [],
compressionThresholdBytes: DEFAULT_SYNC_COMPRESSION_THRESHOLD_BYTES,
cloudRelayWssUrl: args.getCloudRelayWssUrl?.() ?? null,
// Advertise the runtime RPC channel + port-forward only to paired
// desktop runtime-hosts (phones/browsers stay on the allowlist).
runtimeChannelEnabled:
auth?.kind === "paired" && isRuntimeHostPairingRecord(authenticatedPairingRecord),
}), envelope.requestId);
return;
}
Expand All @@ -962,6 +1006,7 @@ export function createBrainProjectActionsSyncHandler(
clearAuthTimeout();
clearInterval(personalChatPump);
peer.personalChatSubscriptions.clear();
pairedChannelService.closePeer(ws, "Sync socket closed.", false);
});
};
}
Loading