diff --git a/apps/desktop/renderer-architecture.json b/apps/desktop/renderer-architecture.json index 36e9bb5369..fc1b4e558b 100644 --- a/apps/desktop/renderer-architecture.json +++ b/apps/desktop/renderer-architecture.json @@ -234,7 +234,6 @@ "src/renderer/use-new-task-choice.ts", "src/renderer/use-onboarding-snapshot.ts", "src/renderer/use-project-context.ts", - "src/renderer/use-session-collaboration-dialog.ts", "src/renderer/use-session-setting-intent.ts", "src/renderer/use-settings-modal.ts", "src/renderer/use-shell-appearance.ts", @@ -985,6 +984,7 @@ "./error-boundary": 1, "./features/goals": 1, "./features/module-hub": 1, + "./features/session-collaboration": 1, "./features/session-navigation": 1, "./features/task-entry": 1, "./features/workbar": 1, @@ -1020,7 +1020,6 @@ "./use-new-task-choice": 1, "./use-onboarding-snapshot": 1, "./use-project-context": 1, - "./use-session-collaboration-dialog": 1, "./use-session-setting-intent": 1, "./use-settings-modal": 1, "./use-shell-appearance": 1, @@ -1058,7 +1057,7 @@ "react": 1 }, "importSpecifiers": 186, - "nonTriviaTokens": 15825 + "nonTriviaTokens": 15800 }, "src/renderer/use-app-shell-composer-quotes.ts": { "importDeclarations": 3, @@ -4927,19 +4926,6 @@ "react": 1 } }, - "src/renderer/use-session-collaboration-dialog.ts": { - "bridgePaths": {}, - "environmentCapabilities": {}, - "hookCalls": { - "useState": 1 - }, - "lifecycleMethods": {}, - "unresolvedDependencies": 0, - "actionFactories": [], - "dependencyPaths": { - "react": 1 - } - }, "src/renderer/use-session-setting-intent.ts": { "bridgePaths": {}, "environmentCapabilities": {}, diff --git a/apps/desktop/src/main/__tests__/runtime-host-collaboration-ipc-main.test.ts b/apps/desktop/src/main/__tests__/runtime-host-collaboration-ipc-main.test.ts index 2863fec8e2..ce64f26182 100644 --- a/apps/desktop/src/main/__tests__/runtime-host-collaboration-ipc-main.test.ts +++ b/apps/desktop/src/main/__tests__/runtime-host-collaboration-ipc-main.test.ts @@ -29,7 +29,7 @@ import { registerRuntimeHostCollaborationIpc } from '../runtime-host-collaborati const ROOT_ID = 'a'.repeat(64); -test('requires Owner confirmation before issuing a plaintext collaboration invitation', async () => { +test('requires plaintext confirmation and reports the issued invitation routes', async () => { const handlers = new Map(); const ipcMain: ReconnectableReadIpcMain = { handle(channel, listener) { @@ -37,6 +37,7 @@ test('requires Owner confirmation before issuing a plaintext collaboration invit }, }; let prepareCalls = 0; + const queryCalls: Array = []; const client = { async prepareCollaborationInvitation(sessionId: string, grantKinds: readonly string[]) { prepareCalls += 1; @@ -56,6 +57,10 @@ test('requires Owner confirmation before issuing a plaintext collaboration invit async queryCollaborationAccess() { return { principals: [], grants: [] }; }, + async queryCollaborationTurnRequests(sessionId?: string) { + queryCalls.push(sessionId); + return { canRequestTurns: false, requests: [] }; + }, async revokeCollaborationPrincipal() { return { revoked: false }; }, @@ -92,9 +97,55 @@ test('requires Owner confirmation before issuing a plaintext collaboration invit assert.equal(prepareCalls, 1); assert.equal((result as { kind?: unknown }).kind, 'prepared'); const invitation = (result as { - invitation: { invitationCode: string }; + invitation: { invitationCode: string; connectivity: unknown }; }).invitation; + assert.deepEqual(invitation.connectivity, { kind: 'configured' }); const bundle = decodeDesktopCollaborationInvitation(invitation.invitationCode); assert.equal(decodeCollaborationInvitationCode(bundle.invitationCode).rootId, ROOT_ID); assert.equal(bundle.target.transport.kind, 'plaintext'); + + const query = handlers.get('session-collaboration:turn-request:query'); + assert.ok(query); + assert.deepEqual(await query({} as Parameters[0]), { + canRequestTurns: false, + requests: [], + }); + assert.deepEqual(await query({} as Parameters[0], 'session-1'), { + canRequestTurns: false, + requests: [], + }); + assert.deepEqual(queryCalls, [undefined, 'session-1']); + + const peerHandlers = new Map(); + registerRuntimeHostCollaborationIpc( + client as unknown as Parameters[0], + { + handle(channel, listener) { + peerHandlers.set(channel, listener); + }, + }, + async () => ({ + name: 'Peer Lab', + transport: { + kind: 'libp2p-direct', + peerId: '12D3KooWpeer', + routeHints: ['/ip4/192.0.2.1/udp/41000/quic-v1'], + coordinationRelays: [ + '/dns4/relay.example/udp/443/quic-v1/p2p/12D3KooWrelay', + ], + }, + }), + ); + const preparePeer = peerHandlers.get('session-collaboration:prepare'); + assert.ok(preparePeer); + const peerResult = await preparePeer( + {} as Parameters[0], + 'session-1', + 'observe', + false, + ); + assert.deepEqual( + (peerResult as { invitation: { connectivity: unknown } }).invitation.connectivity, + { kind: 'peer', coordinationRelayCount: 1 }, + ); }); diff --git a/apps/desktop/src/main/__tests__/runtime-host-guest-session-mounts.test.ts b/apps/desktop/src/main/__tests__/runtime-host-guest-session-mounts.test.ts index 3d8102b844..ba5ee594c5 100644 --- a/apps/desktop/src/main/__tests__/runtime-host-guest-session-mounts.test.ts +++ b/apps/desktop/src/main/__tests__/runtime-host-guest-session-mounts.test.ts @@ -19,12 +19,17 @@ import assert from 'node:assert/strict'; import test from 'node:test'; -import { encodeCollaborationInvitationCode } from '@maka/runtime-host/protocol'; +import type { ResolvedRuntimeHostProfile } from '@maka/runtime-host/client'; +import { + encodeCollaborationInvitationCode, + type HostPeerEndpoint, +} from '@maka/runtime-host/protocol'; import { encodeDesktopCollaborationInvitation } from '../runtime-host-collaboration-invitation.js'; import { createDesktopGuestSessionMountService, type GuestSessionMount, type GuestSessionMountStore, + registerDesktopGuestSessionMountIpc, } from '../runtime-host-guest-session-mounts.js'; import { RuntimeHostPairingFinalizationInterruptedError } from '../runtime-host-desktop-manager.js'; @@ -56,6 +61,49 @@ test('retains a successful Guest mount and rehydrates the same authority after r await second!.close(); }); +test('persists authenticated route rotation for reconnect and restart', async () => { + const store = memoryStore(); + let observePeerEndpoint!: (endpoint: HostPeerEndpoint) => void; + const first = service(store, { + mount: async (target, _signal, _onConnectionPhase, onPeerEndpoint) => { + assert.deepEqual( + target.profile.kind === 'remote' ? target.profile.transport : undefined, + { + kind: 'libp2p-direct', + peerId: '12D3KooWpeer', + routeHints: ['/ip4/192.0.2.1/udp/41000/quic-v1'], + coordinationRelays: ['/memory/stale-relay'], + }, + ); + assert.ok(onPeerEndpoint); + observePeerEndpoint = onPeerEndpoint; + }, + }); + + const imported = await first.importInvitation(peerInvitation('guest-routes'), false, 'routes'); + assert.equal(imported.kind, 'connected'); + observePeerEndpoint({ + peerId: '12D3KooWpeer', + routeHints: ['/ip4/198.51.100.2/udp/42000/quic-v1'], + coordinationRelays: ['/memory/fresh-relay'], + }); + await first.close(); + + let restarted!: ReturnType; + const restartedTarget = new Promise((resolve) => { + restarted = service(store, { mount: async (target) => resolve(target) }); + void restarted.start(); + }); + const target = await restartedTarget; + assert.deepEqual(target.profile.kind === 'remote' ? target.profile.transport : undefined, { + kind: 'libp2p-direct', + peerId: '12D3KooWpeer', + routeHints: ['/ip4/198.51.100.2/udp/42000/quic-v1'], + coordinationRelays: ['/memory/fresh-relay'], + }); + await restarted.close(); +}); + test('removes failed activation desire instead of creating recoverable profile state', async () => { const store = memoryStore(); const unmounted: string[] = []; @@ -273,6 +321,52 @@ test('cancels an in-flight import and removes its durable mount desire', async ( await mounts.close(); }); +test('reads an invitation from the clipboard only on explicit IPC invocation', async () => { + type IpcHandler = Parameters['handle']>[1]; + const handlers = new Map(); + let clipboardReads = 0; + const clipboardInvitation = invitation('guest-clipboard'); + let clipboardText = ` ${clipboardInvitation} `; + const mounts = service(memoryStore()); + const dispose = registerDesktopGuestSessionMountIpc( + { + handle(channel, handler) { + handlers.set(channel, handler); + }, + removeHandler(channel) { + handlers.delete(channel); + }, + }, + mounts, + () => { + clipboardReads += 1; + return clipboardText; + }, + ); + + assert.equal(clipboardReads, 0); + const handler = handlers.get('session-collaboration:invitation:read-clipboard'); + assert.ok(handler); + assert.equal(await handler({} as never), clipboardInvitation); + assert.equal(clipboardReads, 1); + + clipboardText = 'an unrelated clipboard secret'; + assert.throws(() => handler({} as never), /Invalid Desktop collaboration invitation/); + + clipboardText = ' '; + assert.equal(await handler({} as never), ''); + + clipboardText = 'x'.repeat(32 * 1024 + 1); + assert.throws( + () => handler({} as never), + /Clipboard content is too large to be a shared Session invitation/, + ); + + dispose(); + assert.equal(handlers.size, 0); + await mounts.close(); +}); + function service( store: GuestSessionMountStore, overrides: { @@ -314,6 +408,25 @@ function invitation(credential: string): string { }); } +function peerInvitation(credential: string): string { + return encodeDesktopCollaborationInvitation({ + invitationCode: encodeCollaborationInvitationCode({ + schemaVersion: 1, + rootId: ROOT_ID, + credential, + }), + target: { + name: 'Shared Host', + transport: { + kind: 'libp2p-direct', + peerId: '12D3KooWpeer', + routeHints: ['/ip4/192.0.2.1/udp/41000/quic-v1'], + coordinationRelays: ['/memory/stale-relay'], + }, + }, + }); +} + function retainedMount(mountId: string): GuestSessionMount { return { mountId, diff --git a/apps/desktop/src/main/__tests__/runtime-host-local-remote-access.test.ts b/apps/desktop/src/main/__tests__/runtime-host-local-remote-access.test.ts index 9e502095c6..c6f2e90813 100644 --- a/apps/desktop/src/main/__tests__/runtime-host-local-remote-access.test.ts +++ b/apps/desktop/src/main/__tests__/runtime-host-local-remote-access.test.ts @@ -44,17 +44,36 @@ test('enabling remote access hands the same root to one managed service before D const handlers = new Map[1]>(); let retired = false; let resumed = false; + const peer = { + peerId: '12D3KooWpeer', + routeHints: ['/ip4/192.0.2.1/udp/41000/quic-v1'], + coordinationRelays: [], + }; + const livePeer = { + ...peer, + coordinationRelays: ['/dns4/relay.example/udp/443/quic-v1/p2p/12D3KooWrelay'], + }; const manager = { async retireOwnedLocalHost() { retired = true; return { kind: 'retired' as const, resume: () => { resumed = true; } }; }, + async waitUntilReady(profileId: string) { + assert.equal(profileId, 'local'); + }, + current(profileId: string) { + assert.equal(profileId, 'local'); + return { + candidate: { + client: { + async status() { + return { peerEndpoint: livePeer }; + }, + }, + }, + }; + }, } as unknown as RuntimeHostDesktopManager; - const peer = { - peerId: '12D3KooWpeer', - routeHints: ['/ip4/192.0.2.1/udp/41000/quic-v1'], - coordinationRelays: [], - }; const deploymentId = '11111111-1111-4111-8111-111111111111'; const operator = { async runSetup(input: { @@ -123,7 +142,7 @@ test('enabling remote access hands the same root to one managed service before D assert.deepEqual(decodeRuntimeHostOwnerConnectionCode(result.connectionCode), { name: decodeRuntimeHostOwnerConnectionCode(result.connectionCode).name, rootId: 'a'.repeat(64), - transport: { kind: 'libp2p-direct', ...peer }, + transport: { kind: 'libp2p-direct', ...livePeer }, credential: 'pending-credential', }); const lifecycle = JSON.parse( @@ -134,6 +153,69 @@ test('enabling remote access hands the same root to one managed service before D assert.equal(lifecycle.deploymentId, deploymentId); }); +test('shares the running Local Host endpoint instead of its persisted startup routes', async (t) => { + const base = await mkdtemp(join(tmpdir(), 'maka-local-live-peer-share-')); + t.after(() => rm(base, { recursive: true, force: true })); + const clientDataRoot = join(base, 'client'); + const rootPath = join(clientDataRoot, 'workspaces', 'default'); + const rootId = 'a'.repeat(64); + await mkdir(rootPath, { recursive: true }); + await writeManagedLifecycle(clientDataRoot, rootPath, rootId); + const configuredPeer = { + peerId: '12D3KooWpeer', + routeHints: ['/ip4/192.0.2.1/udp/41000/quic-v1'], + coordinationRelays: [], + }; + const livePeer = { + ...configuredPeer, + coordinationRelays: ['/dns4/relay.example/udp/443/quic-v1/p2p/12D3KooWrelay'], + }; + const service = createDesktopLocalRuntimeHostRemoteAccess({ + ipcMain: { handle() {}, removeHandler() {} }, + clientDataRoot, + rootPath, + rootId, + directPeerAvailable: true, + manager: () => + ({ + current() { + return { + candidate: { + client: { + async status() { + return { peerEndpoint: livePeer }; + }, + }, + }, + }; + }, + }) as unknown as RuntimeHostDesktopManager, + resolveSetupPackage: async () => ({ kind: 'npm', specifier: 'maka-agent@0.2.0' }), + operator: { + async runPeer() { + return { + kind: 'result' as const, + action: 'status' as const, + status: { + state: 'enabled' as const, + serviceState: 'running', + rootId, + ...configuredPeer, + }, + }; + }, + async close() {}, + } as unknown as ReturnType, + }); + t.after(() => service.close()); + + const target = await service.createCollaborationConnectionTarget(); + assert.deepEqual(target, { + name: target.name, + transport: { kind: 'libp2p-direct', ...livePeer }, + }); +}); + test('revokes the one Local sharing authority without changing peer connectivity', async (t) => { const base = await mkdtemp(join(tmpdir(), 'maka-local-shared-access-revoke-')); t.after(() => rm(base, { recursive: true, force: true })); diff --git a/apps/desktop/src/main/__tests__/runtime-host-profile-service.test.ts b/apps/desktop/src/main/__tests__/runtime-host-profile-service.test.ts index d99b90a262..12a9886d42 100644 --- a/apps/desktop/src/main/__tests__/runtime-host-profile-service.test.ts +++ b/apps/desktop/src/main/__tests__/runtime-host-profile-service.test.ts @@ -716,12 +716,27 @@ test("keeps a managed Direct route on the SSH profile credential authority", asy await managedServices.save(MANAGED_PROFILE, MANAGED_SERVICE); const startup = await resolveDesktopRuntimeHostStartup(root, { catalog }); const activated: ResolvedRuntimeHostProfile[] = []; + const livePeer = { + peerId: "12D3KooWpeer", + routeHints: ["/ip4/192.0.2.9/udp/44002/quic-v1"], + coordinationRelays: ["/dns4/relay.example/udp/443/quic-v1/p2p/12D3KooWrelay"], + }; + let exposeReadyState = false; const service = createDesktopRuntimeHostProfileService({ clientDataRoot: root, startup, catalog, managedServices, - states: () => [connectingLocal()], + states: () => + exposeReadyState + ? [ + connectingLocal(), + readyWithPeerEndpoint( + { profile: MANAGED_PROFILE, credential: "owner-token" }, + livePeer, + ), + ] + : [connectingLocal()], enable: async (target) => { activated.push(target); }, @@ -755,10 +770,15 @@ test("keeps a managed Direct route on the SSH profile credential authority", asy routeHints: ["/ip4/192.0.2.8/udp/44001/quic-v1"], coordinationRelays: [], }); + exposeReadyState = true; assert.deepEqual( await service.resolveCollaborationConnectionTarget(MANAGED_PROFILE), - { name: MANAGED_PROFILE.name, transport: direct.profile.transport }, + { + name: MANAGED_PROFILE.name, + transport: { kind: "libp2p-direct", ...livePeer }, + }, ); + exposeReadyState = false; assert.equal((await catalog.resolve(MANAGED_PROFILE.id)).credential, "owner-token"); const beforeRejectedRemoval = { @@ -1602,6 +1622,27 @@ function ready(target: ResolvedRuntimeHostProfile): RuntimeHostDesktopTargetStat }; } +function readyWithPeerEndpoint( + target: ResolvedRuntimeHostProfile, + peerEndpoint: { + readonly peerId: string; + readonly routeHints: readonly string[]; + readonly coordinationRelays: readonly string[]; + }, +): RuntimeHostDesktopTargetState { + return { + epoch: `epoch-${target.profile.id}`, + target, + readiness: "ready", + candidate: { + client: { + hostId: target.profile.kind === "remote" ? target.profile.rootId : ROOT_ID, + status: async () => ({ peerEndpoint }), + }, + } as never, + }; +} + function unavailable( target: ResolvedRuntimeHostProfile, error: Error, diff --git a/apps/desktop/src/main/__tests__/runtime-host-session-catalog-preload.test.ts b/apps/desktop/src/main/__tests__/runtime-host-session-catalog-preload.test.ts index e44732fbe7..bed0ef61e3 100644 --- a/apps/desktop/src/main/__tests__/runtime-host-session-catalog-preload.test.ts +++ b/apps/desktop/src/main/__tests__/runtime-host-session-catalog-preload.test.ts @@ -21,8 +21,9 @@ import assert from 'node:assert/strict'; import test from 'node:test'; import type { DesktopSessionSummary } from '../../preload/bridge-contract.js'; import { - collectRuntimeHostSessionCatalogs, collectRuntimeHostSessionCatalogsWithCoverage, + reconcileRuntimeHostSessionCatalog, + resolveRuntimeHostSessionCatalog, } from '../../preload/runtime-host-session-catalog.js'; function session(id: string, activityAt: number): DesktopSessionSummary { @@ -30,13 +31,17 @@ function session(id: string, activityAt: number): DesktopSessionSummary { } test('keeps healthy Host catalogs when another Host rejects', async () => { - const sessions = await collectRuntimeHostSessionCatalogs([ - Promise.resolve([session('older', 1)]), - Promise.reject(new Error('remote unavailable')), - Promise.resolve([session('newer', 2)]), + const catalog = await collectRuntimeHostSessionCatalogsWithCoverage([ + { hostId: 'older', access: 'owner', sessions: Promise.resolve([session('older', 1)]) }, + { + hostId: 'unavailable', + access: 'owner', + sessions: Promise.reject(new Error('remote unavailable')), + }, + { hostId: 'newer', access: 'owner', sessions: Promise.resolve([session('newer', 2)]) }, ]); - assert.deepEqual(sessions.map(({ id }) => id), ['newer', 'older']); + assert.deepEqual(catalog.sessions.map(({ id }) => id), ['newer', 'older']); }); test('reports exactly which Host catalogs are complete', async () => { @@ -57,19 +62,85 @@ test('collapses overlapping Guest catalogs in favor of the Owner authority', asy const owner = session('shared-session', 2); const guest = { ...owner, shared: true as const }; - const sessions = await collectRuntimeHostSessionCatalogs([ - Promise.resolve([guest]), - Promise.resolve([owner]), + const catalog = await collectRuntimeHostSessionCatalogsWithCoverage([ + { hostId: 'shared', access: 'session_guest', sessions: Promise.resolve([guest]) }, + { hostId: 'shared', access: 'owner', sessions: Promise.resolve([owner]) }, ]); - assert.deepEqual(sessions, [owner]); + assert.deepEqual(catalog.sessions, [owner]); +}); + +test('retains a Guest catalog row only while its Runtime Host profile remains known', () => { + const shared = { + ...session('shared-session', 2), + runtimeHostId: 'host-guest', + profileId: 'guest-profile', + shared: true as const, + }; + + const reconnecting = reconcileRuntimeHostSessionCatalog([shared], { + sessions: [], + completeHostIds: [], + knownProfileIds: ['guest-profile'], + }); + assert.deepEqual(reconnecting, [shared]); + + assert.deepEqual( + reconcileRuntimeHostSessionCatalog(reconnecting, { + sessions: [{ ...shared, activityAt: 3 }], + completeHostIds: ['host-guest'], + knownProfileIds: ['guest-profile'], + }), + [{ ...shared, activityAt: 3 }], + ); + + assert.deepEqual( + reconcileRuntimeHostSessionCatalog(reconnecting, { + sessions: [], + completeHostIds: [], + knownProfileIds: [], + }), + [], + ); +}); + +test('keeps healthy Owner catalogs when Guest mount inventory is unavailable', async () => { + const owner = { + ...session('owner-session', 3), + runtimeHostId: 'owner-host', + profileId: 'owner-profile', + }; + const shared = { + ...session('shared-session', 2), + runtimeHostId: 'guest-host', + profileId: 'guest-profile', + shared: true as const, + }; + + assert.deepEqual( + await resolveRuntimeHostSessionCatalog( + [shared], + Promise.resolve({ sessions: [owner], completeHostIds: ['owner-host'] }), + () => ['owner-profile'], + Promise.reject(new Error('Guest mount store is unreadable')), + ), + [owner, shared], + ); }); test('fails when every Host catalog rejects', async () => { await assert.rejects( - collectRuntimeHostSessionCatalogs([ - Promise.reject(new Error('first unavailable')), - Promise.reject(new Error('second unavailable')), + collectRuntimeHostSessionCatalogsWithCoverage([ + { + hostId: 'first', + access: 'owner', + sessions: Promise.reject(new Error('first unavailable')), + }, + { + hostId: 'second', + access: 'owner', + sessions: Promise.reject(new Error('second unavailable')), + }, ]), /Every Runtime Host Session Catalog request failed/, ); diff --git a/apps/desktop/src/main/__tests__/runtime-host-turn-request-inbox-preload.test.ts b/apps/desktop/src/main/__tests__/runtime-host-turn-request-inbox-preload.test.ts new file mode 100644 index 0000000000..267504c866 --- /dev/null +++ b/apps/desktop/src/main/__tests__/runtime-host-turn-request-inbox-preload.test.ts @@ -0,0 +1,49 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ + +import assert from 'node:assert/strict'; +import test from 'node:test'; +import type { SessionTurnAccessRequest } from '@maka/runtime-host/protocol'; +import { collectAvailablePendingTurnRequests } from '../../preload/runtime-host-turn-request-inbox.js'; + +function request(requestId: string, createdAt: string): SessionTurnAccessRequest { + return { requestId, createdAt } as SessionTurnAccessRequest; +} + +test('keeps available collaboration inboxes when another Owner Host rejects', async () => { + const requests = await collectAvailablePendingTurnRequests([ + Promise.reject(new Error('Local Host does not expose collaboration authority')), + Promise.resolve([ + request('later', '2026-09-01T00:00:02.000Z'), + request('earlier', '2026-09-01T00:00:01.000Z'), + ]), + ]); + + assert.deepEqual(requests.map(({ requestId }) => requestId), ['earlier', 'later']); +}); + +test('retains the previous inbox projection when every Owner Host rejects', async () => { + await assert.rejects( + collectAvailablePendingTurnRequests([ + Promise.reject(new Error('first unavailable')), + Promise.reject(new Error('second unavailable')), + ]), + /Every Runtime Host collaboration inbox request failed/, + ); +}); diff --git a/apps/desktop/src/main/__tests__/session-turn-request-inbox-model.test.ts b/apps/desktop/src/main/__tests__/session-turn-request-inbox-model.test.ts new file mode 100644 index 0000000000..eba12f0ea8 --- /dev/null +++ b/apps/desktop/src/main/__tests__/session-turn-request-inbox-model.test.ts @@ -0,0 +1,64 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ + +import assert from 'node:assert/strict'; +import test from 'node:test'; +import type { SessionTurnAccessRequest } from '@maka/runtime-host/protocol'; +import { + groupPendingTurnRequests, + samePendingTurnRequests, + turnRequestPreview, + unseenTurnRequests, +} from '../../renderer/features/session-collaboration/testing.js'; + +test('Turn-request inbox keeps only actionable requests and detects new arrivals', () => { + const first = request('request-1', 'session-1', ' Review\nthis change '); + const second = request('request-2', 'session-1', 'Continue'); + const decided: SessionTurnAccessRequest = { + ...request('request-3', 'session-2', 'Already handled'), + state: { + kind: 'rejected', + decidedAt: '2026-09-01T00:00:03.000Z', + decidedBy: 'local_owner', + }, + }; + + assert.deepEqual([...groupPendingTurnRequests([first, decided, second])], [ + ['session-1', [first, second]], + ]); + assert.deepEqual(unseenTurnRequests([first, second], new Set(['request-1'])), [second]); + assert.equal(samePendingTurnRequests([first, second], [first, second]), true); + assert.equal(samePendingTurnRequests([first], [second]), false); + assert.equal(turnRequestPreview(first.intent.content.text), 'Review this change'); +}); + +function request( + requestId: string, + sessionId: string, + text: string, +): SessionTurnAccessRequest { + return { + requestId, + principalId: `session_guest:${requestId}`, + grantId: `grant-${requestId}`, + intent: { sessionId, turnId: `turn-${requestId}`, content: { text } }, + createdAt: '2026-09-01T00:00:00.000Z', + state: { kind: 'pending' }, + }; +} diff --git a/apps/desktop/src/main/runtime-host-boot.ts b/apps/desktop/src/main/runtime-host-boot.ts index dae1a3cbd6..e1f7ab5898 100644 --- a/apps/desktop/src/main/runtime-host-boot.ts +++ b/apps/desktop/src/main/runtime-host-boot.ts @@ -270,6 +270,9 @@ if (runtimeHostPeerConfiguration) { ...runtimeHostPeerConfiguration, dataRoot: join(userDataDir, 'peer-mesh'), endpointKind: 'client', + onBackgroundReconcileError: (error) => { + console.error('[runtime-host] Peer Mesh background synchronization failed:', error); + }, }); runtimeHostPeerClient = runtimeHostPeerOwner.client; runtimeHostPeerMesh = runtimeHostPeerOwner.mesh; @@ -532,7 +535,7 @@ const runtimeHostProfileService = createDesktopRuntimeHostProfileService({ }); const guestSessionMountService = createDesktopGuestSessionMountService({ store: createGuestSessionMountStore(runtimeHostCredentialStore), - mount: async (target, signal, onConnectionPhase) => { + mount: async (target, signal, onConnectionPhase, onPeerEndpoint) => { if (target.profile.kind !== 'remote' || !target.credential) { throw new Error('A shared Session requires a remote Guest target'); } @@ -541,6 +544,9 @@ const guestSessionMountService = createDesktopGuestSessionMountService({ { profile: target.profile, credential: target.credential }, signal, onConnectionPhase, + (status) => { + if (status.peerEndpoint) onPeerEndpoint?.(status.peerEndpoint); + }, ); }, finalizeAccess: async (mountId, signal) => { @@ -1603,7 +1609,11 @@ function registerPersistentClientIpc(): void { }); registerMarkdownSaveIpc({ ipcMain, mainWindowController }); registerDesktopRuntimeHostProfileIpc(ipcMain, runtimeHostProfileService); - registerDesktopGuestSessionMountIpc(ipcMain, guestSessionMountService); + registerDesktopGuestSessionMountIpc( + ipcMain, + guestSessionMountService, + () => clipboard.readText(), + ); registerClientSettingsIpc({ ipcMain, settingsStore, diff --git a/apps/desktop/src/main/runtime-host-client.ts b/apps/desktop/src/main/runtime-host-client.ts index 9bc0318393..29b98597c5 100644 --- a/apps/desktop/src/main/runtime-host-client.ts +++ b/apps/desktop/src/main/runtime-host-client.ts @@ -84,6 +84,7 @@ import { type MemoryQueryInput, type MemoryQueryResult, type GoalControlAction, + type HostStatusResult, type GoalProjection, type OperationInput, type OperationOutput, @@ -295,6 +296,10 @@ export class DesktopRuntimeHostClient { return this.#connectionClosed || this.#closeTask ? 'unavailable' : 'ready'; } + status(): Promise { + return this.connection.status(); + } + finalizeAccessCredential( timeoutMs?: number, ): Promise> { @@ -331,8 +336,11 @@ export class DesktopRuntimeHostClient { return this.request('collaboration.turn-request.create', { intent }); } - queryCollaborationTurnRequests(sessionId: string): Promise { - return this.request('collaboration.turn-request.query', { sessionId }); + queryCollaborationTurnRequests(sessionId?: string): Promise { + return this.request( + 'collaboration.turn-request.query', + sessionId === undefined ? {} : { sessionId }, + ); } acknowledgeCollaborationTurnRequest( diff --git a/apps/desktop/src/main/runtime-host-collaboration-invitation.ts b/apps/desktop/src/main/runtime-host-collaboration-invitation.ts index 7b9890c4fc..3137dc2d74 100644 --- a/apps/desktop/src/main/runtime-host-collaboration-invitation.ts +++ b/apps/desktop/src/main/runtime-host-collaboration-invitation.ts @@ -24,7 +24,7 @@ import { import { decodeCollaborationInvitationCode } from '@maka/runtime-host/protocol'; const SCHEMA_VERSION = 1; -const CODE_MAX_BYTES = 32 * 1024; +export const DESKTOP_COLLABORATION_INVITATION_CODE_MAX_BYTES = 32 * 1024; export interface DesktopCollaborationConnectionTarget { readonly name: string; @@ -54,7 +54,10 @@ export function encodeDesktopCollaborationInvitation( export function decodeDesktopCollaborationInvitation( code: string, ): DesktopCollaborationInvitation { - if (!code || Buffer.byteLength(code, 'utf8') > CODE_MAX_BYTES) { + if ( + !code || + Buffer.byteLength(code, 'utf8') > DESKTOP_COLLABORATION_INVITATION_CODE_MAX_BYTES + ) { throw new Error('Invalid Desktop collaboration invitation'); } let value: unknown; diff --git a/apps/desktop/src/main/runtime-host-collaboration-ipc-main.ts b/apps/desktop/src/main/runtime-host-collaboration-ipc-main.ts index 622cb75c7a..afe90c6f66 100644 --- a/apps/desktop/src/main/runtime-host-collaboration-ipc-main.ts +++ b/apps/desktop/src/main/runtime-host-collaboration-ipc-main.ts @@ -70,6 +70,13 @@ export function registerRuntimeHostCollaborationIpc( invitationCode: prepared.invitationCode, target, }), + connectivity: + target.transport.kind === 'libp2p-direct' + ? { + kind: 'peer' as const, + coordinationRelayCount: target.transport.coordinationRelays.length, + } + : { kind: 'configured' as const }, }, }; }, @@ -85,7 +92,9 @@ export function registerRuntimeHostCollaborationIpc( ipcMain, 'session-collaboration:turn-request:query', (_event, sessionId: unknown) => - client.queryCollaborationTurnRequests(requiredId(sessionId, 'Session')), + client.queryCollaborationTurnRequests( + sessionId === undefined ? undefined : requiredId(sessionId, 'Session'), + ), ); ipcMain.handle( 'session-collaboration:turn-request:acknowledge', diff --git a/apps/desktop/src/main/runtime-host-desktop-candidate.ts b/apps/desktop/src/main/runtime-host-desktop-candidate.ts index f621beb29a..8f2d460c44 100644 --- a/apps/desktop/src/main/runtime-host-desktop-candidate.ts +++ b/apps/desktop/src/main/runtime-host-desktop-candidate.ts @@ -52,6 +52,7 @@ import { import { INTERACTIVE_RUNTIME_HOST_COMPOSITION_ID, RUNTIME_HOST_PROTOCOL_VERSION, + type HostStatusResult, type WorkspaceTarget, } from "@maka/runtime-host/protocol"; import type { AttachmentApprovalRegistry } from "./attachment-approval.js"; @@ -204,6 +205,7 @@ export interface DesktopRuntimeHostCandidateStartInput readonly candidateLaunchBarrier?: RuntimeHostCandidateLaunchBarrier; readonly peerClient?: RuntimeHostPeerClient; readonly onConnectionPhase?: (phase: RuntimeHostConnectionPhase) => void; + readonly onHostStatus?: (status: HostStatusResult) => void; readonly profileTarget?: { readonly profile: PersistedRuntimeHostProfile; readonly credential?: string; @@ -440,6 +442,9 @@ async function startProfileDesktopRuntimeHostCandidate( ...(input.onConnectionPhase === undefined ? {} : { onConnectionPhase: input.onConnectionPhase }), + ...(input.onHostStatus === undefined + ? {} + : { onHostStatus: input.onHostStatus }), ...(profileTarget.sshInteraction === undefined ? {} : { sshInteraction: profileTarget.sshInteraction }), diff --git a/apps/desktop/src/main/runtime-host-desktop-manager.ts b/apps/desktop/src/main/runtime-host-desktop-manager.ts index 7953e0e1d8..de155eeb53 100644 --- a/apps/desktop/src/main/runtime-host-desktop-manager.ts +++ b/apps/desktop/src/main/runtime-host-desktop-manager.ts @@ -36,7 +36,7 @@ import { type RuntimeHostRetirementMode, type RuntimeHostSshInteraction, } from '@maka/runtime-host/client'; -import type { HostRegistration } from '@maka/runtime-host/protocol'; +import type { HostRegistration, HostStatusResult } from '@maka/runtime-host/protocol'; import type { DesktopTargetSessionRef } from '../shared/runtime-host-identity.js'; import { startDesktopRuntimeHostCandidate, @@ -72,6 +72,7 @@ export interface RuntimeHostDesktopManager { profileTarget: NonNullable, signal?: AbortSignal, onConnectionPhase?: (phase: RuntimeHostConnectionPhase) => void, + onHostStatus?: (status: HostStatusResult) => void, ): Promise; finalizeGuestAccess(mountId: string, signal?: AbortSignal): Promise; unmountGuest(mountId: string): Promise; @@ -503,12 +504,13 @@ class RuntimeHostDesktopManagerImpl implements RuntimeHostDesktopManager { profileTarget: NonNullable, signal?: AbortSignal, onConnectionPhase?: (phase: RuntimeHostConnectionPhase) => void, + onHostStatus?: (status: HostStatusResult) => void, ): Promise { if (!isSessionGuestProfile(profileTarget.profile)) { return Promise.reject(new Error('A Session Guest target is required')); } return this.#mutateTarget(profileTarget.profile.id, () => - this.#enable(profileTarget, true, signal, onConnectionPhase), + this.#enable(profileTarget, true, signal, onConnectionPhase, onHostStatus), ); } @@ -517,6 +519,7 @@ class RuntimeHostDesktopManagerImpl implements RuntimeHostDesktopManager { allowSameRoot: boolean, signal?: AbortSignal, onConnectionPhase?: (phase: RuntimeHostConnectionPhase) => void, + onHostStatus?: (status: HostStatusResult) => void, ): Promise { signal?.throwIfAborted(); if (this.#closed) throw new Error('Desktop Runtime Host manager is closed'); @@ -547,6 +550,7 @@ class RuntimeHostDesktopManagerImpl implements RuntimeHostDesktopManager { const target = this.#createTarget({ ...withRuntimeHostTarget(this.#baseInput, profileTarget), ...(onConnectionPhase ? { onConnectionPhase } : {}), + ...(onHostStatus ? { onHostStatus } : {}), }); this.#targets.set(profileId, target); this.#publishState(target, { diff --git a/apps/desktop/src/main/runtime-host-guest-session-mounts.ts b/apps/desktop/src/main/runtime-host-guest-session-mounts.ts index 0b890f4872..4b9e5a97e6 100644 --- a/apps/desktop/src/main/runtime-host-guest-session-mounts.ts +++ b/apps/desktop/src/main/runtime-host-guest-session-mounts.ts @@ -25,7 +25,10 @@ import { type RuntimeHostConnectionPhase, type RuntimeHostRemoteTransport, } from '@maka/runtime-host/client'; -import { decodeCollaborationInvitationCode } from '@maka/runtime-host/protocol'; +import { + decodeCollaborationInvitationCode, + type HostPeerEndpoint, +} from '@maka/runtime-host/protocol'; import type { CredentialStore } from '@maka/storage/credential-store'; import type { SessionCollaborationCancelResult, @@ -33,7 +36,10 @@ import type { SessionCollaborationImportResult, SessionCollaborationMountSummary, } from '../shared/session-collaboration.js'; -import { decodeDesktopCollaborationInvitation } from './runtime-host-collaboration-invitation.js'; +import { + decodeDesktopCollaborationInvitation, + DESKTOP_COLLABORATION_INVITATION_CODE_MAX_BYTES, +} from './runtime-host-collaboration-invitation.js'; import { RuntimeHostPairingFinalizationInterruptedError } from './runtime-host-desktop-manager.js'; const STORE_SCHEMA_VERSION = 1; @@ -127,6 +133,7 @@ export function createDesktopGuestSessionMountService(input: { target: ResolvedRuntimeHostProfile, signal: AbortSignal, onConnectionPhase?: (phase: RuntimeHostConnectionPhase) => void, + onPeerEndpoint?: (endpoint: HostPeerEndpoint) => void, ) => Promise; readonly finalizeAccess: (mountId: string, signal: AbortSignal) => Promise; readonly unmount: (mountId: string) => Promise; @@ -162,6 +169,38 @@ export function createDesktopGuestSessionMountService(input: { mounts = next; }; + const recordPeerEndpoint = (mount: GuestSessionMount, endpoint: HostPeerEndpoint): void => { + if ( + closed || + mount.transport.kind !== 'libp2p-direct' || + endpoint.peerId !== mount.transport.peerId || + (endpoint.routeHints.length === 0 && endpoint.coordinationRelays.length === 0) + ) return; + void mutate(async () => { + if (removingMounts.has(mount.mountId)) return; + const current = await load(); + const retained = current.get(mount.mountId); + if ( + retained?.transport.kind !== 'libp2p-direct' || + retained.transport.peerId !== endpoint.peerId || + ( + sameStrings(retained.transport.routeHints, endpoint.routeHints) && + sameStrings(retained.transport.coordinationRelays, endpoint.coordinationRelays) + ) + ) return; + const updated = decodeMount({ + ...retained, + transport: { + kind: 'libp2p-direct', + peerId: endpoint.peerId, + routeHints: endpoint.routeHints, + coordinationRelays: endpoint.coordinationRelays, + }, + }); + await persist(new Map(current).set(mount.mountId, updated)); + }).catch((error: unknown) => onError(asError(error), mount)); + }; + const activate = async ( activation: LiveGuestActivation, mount: GuestSessionMount, @@ -174,7 +213,10 @@ export function createDesktopGuestSessionMountService(input: { collaborationProgressForConnectionPhase(phase), ); } - }); + }, (endpoint) => recordPeerEndpoint(mount, endpoint)); + // waitForReady observes host.status before mount resolves. Commit that + // authenticated route snapshot before declaring the durable mount ready. + await mutationTail; activation.controller.signal.throwIfAborted(); if (removingMounts.has(mount.mountId)) { throw new Error('Shared Session mount was removed while connecting'); @@ -419,12 +461,14 @@ export function createDesktopGuestSessionMountService(input: { export function registerDesktopGuestSessionMountIpc( ipcMain: Pick, service: DesktopGuestSessionMountService, + readClipboardText: () => string, ): () => void { const channels = [ 'session-collaboration:import', 'session-collaboration:import:cancel', 'session-collaboration:mount:list', 'session-collaboration:mount:remove', + 'session-collaboration:invitation:read-clipboard', ] as const; ipcMain.handle( channels[0], @@ -441,6 +485,15 @@ export function registerDesktopGuestSessionMountIpc( service.cancelImport(requireOperationId(operationIdValue))); ipcMain.handle(channels[2], () => service.list()); ipcMain.handle(channels[3], (_event, mountId: string) => service.remove(mountId)); + ipcMain.handle(channels[4], () => { + const value = readClipboardText().trim(); + if (Buffer.byteLength(value, 'utf8') > DESKTOP_COLLABORATION_INVITATION_CODE_MAX_BYTES) { + throw new Error('Clipboard content is too large to be a shared Session invitation'); + } + if (!value) return value; + decodeDesktopCollaborationInvitation(value); + return value; + }); return () => { for (const channel of channels) ipcMain.removeHandler(channel); }; @@ -510,6 +563,10 @@ function isPeerPathUnavailable(error: unknown): boolean { return error.code === 'direct_path_unavailable' || error.code === 'transit_unavailable'; } +function sameStrings(left: readonly string[], right: readonly string[]): boolean { + return left.length === right.length && left.every((value, index) => value === right[index]); +} + function collaborationProgressForConnectionPhase( phase: RuntimeHostConnectionPhase, ): SessionCollaborationImportPhase { diff --git a/apps/desktop/src/main/runtime-host-local-remote-access.ts b/apps/desktop/src/main/runtime-host-local-remote-access.ts index 4ade5c5610..3623734975 100644 --- a/apps/desktop/src/main/runtime-host-local-remote-access.ts +++ b/apps/desktop/src/main/runtime-host-local-remote-access.ts @@ -339,11 +339,14 @@ export function createDesktopLocalRuntimeHostRemoteAccess(input: { }; const completed = await finishSetup(setup, 'request'); if (completed.kind === 'active_tasks') return completed; + const manager = requireManager(input.manager); + await manager.waitUntilReady('local'); + const livePeer = await readLivePeer(localClient(input.manager), completed.peer); return enabledResult( encodeRuntimeHostOwnerConnectionCode({ name: hostName(), rootId: completed.managed.rootId, - transport: { kind: 'libp2p-direct', ...completed.peer }, + transport: { kind: 'libp2p-direct', ...livePeer }, credential: completed.credential, }), ); @@ -515,9 +518,10 @@ export function createDesktopLocalRuntimeHostRemoteAccess(input: { ); const peer = await readPeer(input.operator, managed); if (!peer) throw new Error('Remote access is not enabled on this computer'); + const livePeer = await readLivePeer(localClient(input.manager), peer); return { name: hostName(), - transport: { kind: 'libp2p-direct' as const, ...peer }, + transport: { kind: 'libp2p-direct' as const, ...livePeer }, }; }); @@ -968,6 +972,7 @@ async function issueConnectionCode( peer: LocalPeerDescriptor, client: DesktopRuntimeHostClient, ): Promise { + const livePeer = await readLivePeer(client, peer); const prepared = await client.request('access.credential.prepare', { principalKind: 'remote_owner', principalId: LOCAL_REMOTE_ACCESS_PRINCIPAL_ID, @@ -984,11 +989,30 @@ async function issueConnectionCode( return encodeRuntimeHostOwnerConnectionCode({ name: hostName(), rootId, - transport: { kind: 'libp2p-direct', ...peer }, + transport: { kind: 'libp2p-direct', ...livePeer }, credential, }); } +async function readLivePeer( + client: DesktopRuntimeHostClient, + configured: LocalPeerDescriptor, +): Promise { + const endpoint = (await client.status()).peerEndpoint; + if (!endpoint) { + throw new Error('Runtime Host Direct peer is not available'); + } + if (endpoint.peerId !== configured.peerId) { + throw new Error('Runtime Host Direct peer identity changed'); + } + return requireEnabledPeer({ + state: 'enabled', + peerId: endpoint.peerId, + routeHints: endpoint.routeHints, + coordinationRelays: endpoint.coordinationRelays, + }); +} + async function hasSharedAccess( operator: DesktopRuntimeHostLocalOperator, target: LocalServiceTarget, diff --git a/apps/desktop/src/main/runtime-host-profile-service.ts b/apps/desktop/src/main/runtime-host-profile-service.ts index 33bb1380f7..a310141439 100644 --- a/apps/desktop/src/main/runtime-host-profile-service.ts +++ b/apps/desktop/src/main/runtime-host-profile-service.ts @@ -920,23 +920,49 @@ export function createDesktopRuntimeHostProfileService(input: { if (profile.kind !== 'remote') { throw new Error('This Runtime Host does not expose a shareable peer endpoint'); } - if (profile.transport.kind !== 'ssh') { + if (profile.transport.kind !== 'ssh' && profile.transport.kind !== 'libp2p-direct') { return { name: profile.name, transport: profile.transport }; } - const direct = (await catalog.read()).profiles.find( - (candidate) => candidate.id === managedDirectPeerProfileId(profile.id), - ); - if ( - !direct || - direct.kind !== 'remote' || - direct.rootId !== profile.rootId || - direct.transport.kind !== 'libp2p-direct' - ) { + let configuredPeerId: string | undefined; + if (profile.transport.kind === 'libp2p-direct') { + configuredPeerId = profile.transport.peerId; + } else { + const direct = (await catalog.read()).profiles.find( + (candidate) => candidate.id === managedDirectPeerProfileId(profile.id), + ); + if ( + direct?.kind === 'remote' && + direct.rootId === profile.rootId && + direct.transport.kind === 'libp2p-direct' + ) { + configuredPeerId = direct.transport.peerId; + } + } + if (!configuredPeerId) { throw new Error( 'Enable Direct peer access for this Runtime Host before sharing its Sessions', ); } - return { name: profile.name, transport: direct.transport }; + const active = input.states().find((state) => state.target.profile.id === profile.id); + if (!active || active.readiness !== 'ready') { + throw new Error('Connect this Runtime Host before sharing its Sessions'); + } + const endpoint = (await active.candidate.client.status()).peerEndpoint; + if (!endpoint) { + throw new Error('Runtime Host Direct peer is not available'); + } + if (configuredPeerId !== endpoint.peerId) { + throw new Error('Runtime Host Direct peer identity changed'); + } + return { + name: profile.name, + transport: { + kind: 'libp2p-direct' as const, + peerId: endpoint.peerId, + routeHints: endpoint.routeHints, + coordinationRelays: endpoint.coordinationRelays, + }, + }; }); }, resolveManagedAccess(profileId) { diff --git a/apps/desktop/src/preload/bridge-contract.d.ts b/apps/desktop/src/preload/bridge-contract.d.ts index 2bc4ce1aad..da3dc9b7fc 100644 --- a/apps/desktop/src/preload/bridge-contract.d.ts +++ b/apps/desktop/src/preload/bridge-contract.d.ts @@ -345,7 +345,14 @@ export type DesktopSessionCollaborationCancelResult = SessionCollaborationCancel export type DesktopSessionCollaborationPrepareResult = | { readonly kind: 'prepared'; - readonly invitation: CollaborationInvitationPrepareResult; + readonly invitation: CollaborationInvitationPrepareResult & { + readonly connectivity: + | { + readonly kind: 'peer'; + readonly coordinationRelayCount: number; + } + | { readonly kind: 'configured' }; + }; } | { readonly kind: 'insecure_confirmation_required' }; @@ -749,6 +756,8 @@ export interface MakaBridge { readonly operationId: string; }, onProgress?: (phase: DesktopSessionCollaborationImportPhase) => void): Promise; cancelImport(operationId: string): Promise; + /** Reads only after the user invokes the invitation paste action. */ + readInvitationClipboard(): Promise; listMounts(): Promise; removeMount(mountId: string): Promise; requestTurn( @@ -756,6 +765,8 @@ export interface MakaBridge { input: { readonly turnId: string; readonly text: string }, ): Promise; getTurnRequests(sessionId: string): Promise; + /** Pending Owner decisions across every connected Owner Runtime Host. */ + getPendingTurnRequests(): Promise; acknowledgeTurnRequest( sessionId: string, requestId: string, diff --git a/apps/desktop/src/preload/preload.ts b/apps/desktop/src/preload/preload.ts index c3a7726846..c462b5215b 100644 --- a/apps/desktop/src/preload/preload.ts +++ b/apps/desktop/src/preload/preload.ts @@ -121,9 +121,10 @@ import type { BotProvider } from '@maka/core/bot-chat-settings'; import type { BotOnboardingSnapshot, BotOnboardingStartInput } from '@maka/core/bot-onboarding'; import type { HealthSnapshot } from '@maka/core/health'; import { - collectRuntimeHostSessionCatalogs, collectRuntimeHostSessionCatalogsWithCoverage, + resolveRuntimeHostSessionCatalog, } from './runtime-host-session-catalog.js'; +import { collectAvailablePendingTurnRequests } from './runtime-host-turn-request-inbox.js'; import type { ExecutionBoundaryReadModel, SandboxBoundaryResponse } from '@maka/core/sandbox-boundary'; import type { ClientCapabilityResponse } from '@maka/core/client-capability-grant'; import type { @@ -222,6 +223,8 @@ import { type OperationInput, type OperationOutcome, type OperationOutput, + type CollaborationTurnRequestQueryResult, + type SessionTurnAccessRequest, } from '@maka/runtime-host/protocol'; import type { AgentGraphEpochDirectory } from '@maka/runtime-host/client'; import { @@ -260,6 +263,7 @@ const runtimeHostMetadata = new Map< } >(); const runtimeHostSessionScopes = new Map(); +let lastDesktopSessionCatalog: DesktopSessionSummary[] = []; const newTaskChangeListeners = new Set<() => void>(); let previousMainProcessInterruptionRead: Promise | undefined; @@ -894,19 +898,15 @@ async function listDesktopSessions( ) as SessionCatalogSummary[]; return sessions.map((session) => projectSessionSummary(parent.scope, session)); } - const scopes = await runtimeHostScopeList(); - const sessions = await collectRuntimeHostSessionCatalogs( - scopes.map(async (scope) => { - const sessions = await ipcRenderer.invoke( - 'sessions:list', - scope, - filter, - ) as SessionCatalogSummary[]; - return sessions.map((session) => projectSessionCatalogSummary(scope, session)); - }), + lastDesktopSessionCatalog = await resolveRuntimeHostSessionCatalog( + lastDesktopSessionCatalog, + listDesktopSessionsWithCoverage(), + () => [...runtimeHostMetadata.values()].map(({ profileId }) => profileId), + // Unknown Guest coverage cannot suppress healthy Owner catalogs or prove + // that a previously observed Guest mount was removed. + listKnownGuestMountProfileIds(), ); - recordSessionCatalogScopes(sessions); - return sessions; + return lastDesktopSessionCatalog; } async function listDesktopSessionsWithCoverage(): Promise<{ @@ -931,6 +931,24 @@ async function listDesktopSessionsWithCoverage(): Promise<{ return catalog; } +async function listKnownGuestMountProfileIds(): Promise { + const mounts: unknown = await ipcRenderer.invoke('session-collaboration:mount:list'); + if (!Array.isArray(mounts)) { + throw new Error('Desktop shared Session mounts are unavailable'); + } + return mounts.map((mount) => { + if ( + !mount || + typeof mount !== 'object' || + !('mountId' in mount) || + typeof mount.mountId !== 'string' + ) { + throw new Error('Desktop shared Session mount is invalid'); + } + return mount.mountId; + }); +} + async function createDesktopSessionOnScope( scope: DesktopTargetScope, input?: CreateSessionRequestInput, @@ -1305,6 +1323,9 @@ const makaBridge = { cancelImport(operationId) { return ipcRenderer.invoke('session-collaboration:import:cancel', operationId); }, + readInvitationClipboard() { + return ipcRenderer.invoke('session-collaboration:invitation:read-clipboard'); + }, listMounts() { return ipcRenderer.invoke('session-collaboration:mount:list'); }, @@ -1331,6 +1352,28 @@ const makaBridge = { session.sessionId, ); }, + async getPendingTurnRequests() { + const scopes = (await runtimeHostScopeList()).filter( + (scope) => runtimeHostMetadataFor(scope)?.profileAccess === 'owner', + ); + return collectAvailablePendingTurnRequests( + scopes.map(async (scope) => { + const result = await ipcRenderer.invoke( + 'session-collaboration:turn-request:query', + scope, + ) as CollaborationTurnRequestQueryResult; + return result.requests + .filter((request) => request.state.kind === 'pending') + .map((request): SessionTurnAccessRequest => ({ + ...request, + intent: { + ...request.intent, + sessionId: recordRuntimeHostSessionScope(scope, request.intent.sessionId), + }, + })); + }), + ); + }, async acknowledgeTurnRequest(sessionId, requestId) { const session = await runtimeHostSessionRef(sessionId); return ipcRenderer.invoke( diff --git a/apps/desktop/src/preload/runtime-host-session-catalog.ts b/apps/desktop/src/preload/runtime-host-session-catalog.ts index d4755a8949..d8c5955d36 100644 --- a/apps/desktop/src/preload/runtime-host-session-catalog.ts +++ b/apps/desktop/src/preload/runtime-host-session-catalog.ts @@ -30,6 +30,29 @@ export interface RuntimeHostSessionCatalogCoverage { readonly completeHostIds: string[]; } +export interface RuntimeHostSessionCatalogSnapshot extends RuntimeHostSessionCatalogCoverage { + /** Profiles still retained by Desktop, including unavailable Guest mounts. */ + readonly knownProfileIds: string[]; +} + +export async function resolveRuntimeHostSessionCatalog( + current: readonly DesktopSessionSummary[], + coverage: Promise, + knownRuntimeProfileIds: () => readonly string[], + guestMountProfileIds: Promise, +): Promise { + const [snapshot, knownGuestProfileIds] = await Promise.all([ + coverage, + guestMountProfileIds.catch(() => + current.flatMap((session) => session.shared === true ? [session.profileId] : []), + ), + ]); + return reconcileRuntimeHostSessionCatalog(current, { + ...snapshot, + knownProfileIds: [...knownRuntimeProfileIds(), ...knownGuestProfileIds], + }); +} + export async function collectRuntimeHostSessionCatalogsWithCoverage( requests: readonly RuntimeHostSessionCatalogRequest[], ): Promise { @@ -59,18 +82,25 @@ export async function collectRuntimeHostSessionCatalogsWithCoverage( }; } -export async function collectRuntimeHostSessionCatalogs( - requests: readonly Promise[], -): Promise { - const results = await Promise.allSettled(requests); - const groups = results.flatMap((result) => result.status === 'fulfilled' ? [result.value] : []); - if (requests.length > 0 && groups.length === 0) { - throw new AggregateError( - results.flatMap((result) => result.status === 'rejected' ? [result.reason] : []), - 'Every Runtime Host Session Catalog request failed', - ); - } - return sortSessionCatalogs(groups.flat()); +/** + * Commits complete Host catalogs authoritatively while retaining the last + * accepted rows for a Host that still exists but cannot answer this read. + * An explicitly removed profile is absent from knownProfileIds and therefore + * retires immediately; transport availability alone cannot change access. + */ +export function reconcileRuntimeHostSessionCatalog( + current: readonly DesktopSessionSummary[], + snapshot: RuntimeHostSessionCatalogSnapshot, +): DesktopSessionSummary[] { + const completeHostIds = new Set(snapshot.completeHostIds); + const knownProfileIds = new Set(snapshot.knownProfileIds); + return sortSessionCatalogs([ + ...snapshot.sessions, + ...current.filter( + (session) => + knownProfileIds.has(session.profileId) && !completeHostIds.has(session.runtimeHostId), + ), + ]); } function sortSessionCatalogs(sessions: DesktopSessionSummary[]): DesktopSessionSummary[] { diff --git a/apps/desktop/src/preload/runtime-host-turn-request-inbox.ts b/apps/desktop/src/preload/runtime-host-turn-request-inbox.ts new file mode 100644 index 0000000000..86199b9187 --- /dev/null +++ b/apps/desktop/src/preload/runtime-host-turn-request-inbox.ts @@ -0,0 +1,38 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ + +import type { SessionTurnAccessRequest } from '@maka/runtime-host/protocol'; + +export async function collectAvailablePendingTurnRequests( + queries: readonly Promise[], +): Promise { + const results = await Promise.allSettled(queries); + const available = results.flatMap( + (result) => result.status === 'fulfilled' ? [result.value] : [], + ); + if (queries.length > 0 && available.length === 0) { + throw new AggregateError( + results.flatMap((result) => result.status === 'rejected' ? [result.reason] : []), + 'Every Runtime Host collaboration inbox request failed', + ); + } + return available + .flat() + .sort((left, right) => left.createdAt.localeCompare(right.createdAt)); +} diff --git a/apps/desktop/src/renderer/app-shell.tsx b/apps/desktop/src/renderer/app-shell.tsx index 75a0b851ef..270e1d6b39 100644 --- a/apps/desktop/src/renderer/app-shell.tsx +++ b/apps/desktop/src/renderer/app-shell.tsx @@ -109,8 +109,8 @@ import { import { useNewTaskChoice } from './use-new-task-choice'; import { SessionCollaborationDialog } from './session-collaboration-dialog'; import { SessionTurnRequestComposer } from './session-turn-request-composer.js'; +import * as SessionCollaboration from './features/session-collaboration'; import { getSessionCollaborationCopy } from './locales/session-collaboration-copy'; -import { useSessionCollaborationDialog } from './use-session-collaboration-dialog'; import { NEW_TASK_PENDING_KEY } from './pending-items'; import { parseDesktopSlashCommand } from './desktop-slash-command'; import { @@ -242,6 +242,7 @@ type ComposerImportOwner = { */ const SETTLE_FALLBACK_GRACE_MS = 1000; const FIRST_SEND_OBSERVATION_TIMEOUT_MS = 30_000; +const { useSessionCollaborationDialog } = SessionCollaboration; type FirstSendObservationWaiter = { promise: Promise; resolve: () => void; @@ -646,6 +647,7 @@ function AppShellContent({ setUiLocalePreference, }); const shellCopy = getShellCopy(uiLocale).app; + const sessionCollaborationCopy = getSessionCollaborationCopy(uiLocale); const previousInterruptionCopy = getShellRemainingCopy(uiLocale).previousMainProcessInterruption; const desktopConversationCopy = getDesktopConversationCopy(uiLocale); @@ -830,13 +832,6 @@ function AppShellContent({ const activeMessageQueue = activeId ? messageQueueBySession[activeId] : undefined; const activeMessageSubmitting = transientMessages.length > 0; const activeDesktopSession = activeSession; - function openSessionSharing(session: DesktopSessionSummary): void { - sharedSessionDialog.open({ - sessionId: session.id, - sessionName: session.name, - requiresRemoteAccess: session.profileKind === 'local', - }); - } // The shell's reading of the active live turn: streaming/settled flags, the // in-flight tool signal, and the #646 turn-wait cues, all derived from the // semantic snapshot rather than the projection (#1985). @@ -1654,7 +1649,7 @@ function AppShellContent({ newSessionPermissionMode: newTaskPermissionMode, }; - const hasModalOpen = helpOpen || paletteOpen || searchModalOpen || sharedSessionDialog.target !== undefined; + const hasModalOpen = helpOpen || paletteOpen || searchModalOpen || sharedSessionDialog.isOpen; const shellObscured = hasModalOpen || settingsOpen; const contextCompactionPresentation = useMemo( () => @@ -2732,6 +2727,12 @@ function AppShellContent({ reportError={showSessionError} > +
void openSessionSharing(activeDesktopSession), + onClick: () => sharedSessionDialog.openSession(activeDesktopSession), } } onRenameSession={(name) => { @@ -2861,6 +2862,7 @@ function AppShellContent({ projects={localProjects} streamingSessionIds={streamingSessionIds} staleSessionIds={staleSessionIds} + SessionBadge={SessionCollaboration.SessionTurnRequestBadge} ports={sessionNavigationPorts} commandsRef={sessionNavigationCommandsRef} onExitWorkHub={exitWorkHub} @@ -2929,6 +2931,11 @@ function AppShellContent({ hidden={navSelection.section !== 'sessions'} composer={ <> + {!sharedSessionActive && activeId ? ( + + ) : null} {!sharedSessionActive && navSelection.section === 'sessions' && activeId && activeSessionForView && @@ -3274,21 +3281,11 @@ function AppShellContent({ - {sharedSessionDialog.target ? ( - { - const copy = getSessionCollaborationCopy(uiLocale); - sharedSessionDialog.close(); - toastApi.info(copy.enableRemoteAccessTitle, copy.enableRemoteAccessBody); - openSettingsSection('projects'); - }} - onClose={sharedSessionDialog.close} - /> - ) : null} + openSettingsSection('projects')} + onClose={sharedSessionDialog.close} + />
+
); diff --git a/apps/desktop/src/renderer/use-session-collaboration-dialog.ts b/apps/desktop/src/renderer/features/session-collaboration/controller/use-session-collaboration-dialog.ts similarity index 78% rename from apps/desktop/src/renderer/use-session-collaboration-dialog.ts rename to apps/desktop/src/renderer/features/session-collaboration/controller/use-session-collaboration-dialog.ts index 27f727f2c1..87de07e793 100644 --- a/apps/desktop/src/renderer/use-session-collaboration-dialog.ts +++ b/apps/desktop/src/renderer/features/session-collaboration/controller/use-session-collaboration-dialog.ts @@ -30,7 +30,19 @@ export function useSessionCollaborationDialog() { return { target, + isOpen: target !== undefined, open: setTarget, + openSession(session: { + readonly id: string; + readonly name: string; + readonly profileKind: string; + }) { + setTarget({ + sessionId: session.id, + sessionName: session.name, + requiresRemoteAccess: session.profileKind === 'local', + }); + }, close: () => setTarget(undefined), }; } diff --git a/apps/desktop/src/renderer/features/session-collaboration/controller/use-turn-request-inbox.ts b/apps/desktop/src/renderer/features/session-collaboration/controller/use-turn-request-inbox.ts new file mode 100644 index 0000000000..e7618ac395 --- /dev/null +++ b/apps/desktop/src/renderer/features/session-collaboration/controller/use-turn-request-inbox.ts @@ -0,0 +1,142 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ + +import { useCallback, useEffect, useMemo, useRef, useState } from 'react'; +import type { SessionTurnAccessRequest } from '@maka/runtime-host/protocol'; +import type { ToastApi } from '@maka/ui'; +import { + groupPendingTurnRequests, + samePendingTurnRequests, + turnRequestPreview, + unseenTurnRequests, +} from '../model/turn-request-inbox.js'; +import { useSessionCollaborationServices } from '../services-context.js'; + +const INBOX_REFRESH_INTERVAL_MS = 2_000; + +export function useSessionTurnRequestInbox(input: { + readonly sessions: readonly { readonly id: string; readonly name: string }[]; + readonly toast: ToastApi; + readonly onOpenSession: (sessionId: string) => void; + readonly copy: { + readonly sharedTask: string; + readonly newTurnRequestTitle: (count: number) => string; + readonly newTurnRequestSummary: (count: number) => string; + readonly reviewTurnRequest: string; + readonly turnRequests: string; + }; +}) { + const services = useSessionCollaborationServices(); + const [requests, setRequests] = useState([]); + const [workingRequestIds, setWorkingRequestIds] = useState>( + () => new Set(), + ); + const seenRequestIdsRef = useRef(new Set()); + const latestRef = useRef(input); + latestRef.current = input; + + const applyProjection = useCallback((next: readonly SessionTurnAccessRequest[]) => { + const unseen = unseenTurnRequests(next, seenRequestIdsRef.current); + for (const request of next) seenRequestIdsRef.current.add(request.requestId); + setRequests((current) => samePendingTurnRequests(current, next) ? current : next); + if (unseen.length === 0) return; + + const first = unseen[0]!; + const latest = latestRef.current; + const copy = latest.copy; + const sessionName = latest.sessions.find( + (session) => session.id === first.intent.sessionId, + )?.name ?? copy.sharedTask; + latest.toast.toast({ + variant: 'warning', + title: copy.newTurnRequestTitle(unseen.length), + description: unseen.length === 1 + ? `${sessionName} · ${turnRequestPreview(first.intent.content.text)}` + : copy.newTurnRequestSummary(unseen.length), + duration: 10_000, + action: { + label: copy.reviewTurnRequest, + onClick: () => latest.onOpenSession(first.intent.sessionId), + }, + }); + }, []); + + const refresh = useCallback(async () => { + const next = await services.getPendingTurnRequests(); + applyProjection(next); + }, [applyProjection, services]); + + useEffect(() => { + let disposed = false; + let timer: number | undefined; + const poll = async () => { + try { + const next = await services.getPendingTurnRequests(); + if (!disposed) applyProjection(next); + } catch { + // Keep the last durable projection while a Host reconnects and retry. + } finally { + if (!disposed) timer = window.setTimeout(() => void poll(), INBOX_REFRESH_INTERVAL_MS); + } + }; + void poll(); + return () => { + disposed = true; + if (timer !== undefined) window.clearTimeout(timer); + }; + }, [applyProjection, services]); + + const decide = useCallback(async ( + request: SessionTurnAccessRequest, + decision: 'approve' | 'reject', + ) => { + setWorkingRequestIds((current) => new Set(current).add(request.requestId)); + try { + await services.decideTurnRequest( + request.intent.sessionId, + request.requestId, + decision, + ); + setRequests((current) => current.filter( + (candidate) => candidate.requestId !== request.requestId, + )); + await refresh(); + } catch (error) { + const copy = latestRef.current.copy; + latestRef.current.toast.error(copy.turnRequests, errorMessage(error)); + } finally { + setWorkingRequestIds((current) => { + const next = new Set(current); + next.delete(request.requestId); + return next; + }); + } + }, [refresh]); + + return { + requests, + requestsBySession: useMemo(() => groupPendingTurnRequests(requests), [requests]), + workingRequestIds, + decide, + }; +} + +function errorMessage(error: unknown): string { + return error instanceof Error ? error.message : String(error); +} diff --git a/apps/desktop/src/renderer/features/session-collaboration/index.ts b/apps/desktop/src/renderer/features/session-collaboration/index.ts index d5c0973e3e..5ffb61409b 100644 --- a/apps/desktop/src/renderer/features/session-collaboration/index.ts +++ b/apps/desktop/src/renderer/features/session-collaboration/index.ts @@ -18,5 +18,9 @@ */ export { SessionCollaborationServicesProvider } from './services-context'; +export { useSessionCollaborationDialog } from './controller/use-session-collaboration-dialog'; export { SessionCollaborationJoinDialog } from './ui/session-collaboration-join-dialog'; +export { SessionTurnRequestApprovalForSession } from './ui/session-turn-request-approval'; +export { SessionTurnRequestBadge } from './ui/session-turn-request-badge'; +export { SessionTurnRequestInboxProvider } from './turn-request-inbox-context'; export type { SessionCollaborationServices } from './ports'; diff --git a/apps/desktop/src/renderer/features/session-collaboration/model/turn-request-inbox.ts b/apps/desktop/src/renderer/features/session-collaboration/model/turn-request-inbox.ts new file mode 100644 index 0000000000..387297f4ff --- /dev/null +++ b/apps/desktop/src/renderer/features/session-collaboration/model/turn-request-inbox.ts @@ -0,0 +1,62 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ + +import type { SessionTurnAccessRequest } from '@maka/runtime-host/protocol'; + +export function groupPendingTurnRequests( + requests: readonly SessionTurnAccessRequest[], +): ReadonlyMap { + const grouped = new Map(); + for (const request of requests) { + if (request.state.kind !== 'pending') continue; + const sessionId = request.intent.sessionId; + const sessionRequests = grouped.get(sessionId) ?? []; + sessionRequests.push(request); + grouped.set(sessionId, sessionRequests); + } + return grouped; +} + +export function unseenTurnRequests( + requests: readonly SessionTurnAccessRequest[], + seenRequestIds: ReadonlySet, +): readonly SessionTurnAccessRequest[] { + return requests.filter( + (request) => request.state.kind === 'pending' && !seenRequestIds.has(request.requestId), + ); +} + +export function samePendingTurnRequests( + left: readonly SessionTurnAccessRequest[], + right: readonly SessionTurnAccessRequest[], +): boolean { + return left.length === right.length && left.every((request, index) => { + const candidate = right[index]; + return candidate !== undefined && + request.requestId === candidate.requestId && + request.intent.sessionId === candidate.intent.sessionId && + request.intent.content.text === candidate.intent.content.text; + }); +} + +export function turnRequestPreview(text: string, maxLength = 120): string { + const collapsed = text.replace(/\s+/gu, ' ').trim(); + if (collapsed.length <= maxLength) return collapsed; + return `${collapsed.slice(0, Math.max(0, maxLength - 1)).trimEnd()}…`; +} diff --git a/apps/desktop/src/renderer/features/session-collaboration/ports.ts b/apps/desktop/src/renderer/features/session-collaboration/ports.ts index 6079577ad3..7e65b3a5ac 100644 --- a/apps/desktop/src/renderer/features/session-collaboration/ports.ts +++ b/apps/desktop/src/renderer/features/session-collaboration/ports.ts @@ -23,6 +23,10 @@ import type { SessionCollaborationImportResult, SessionCollaborationMountSummary, } from '../../../shared/session-collaboration.js'; +import type { + CollaborationTurnRequestDecideResult, + SessionTurnAccessRequest, +} from '@maka/runtime-host/protocol'; export type { SessionCollaborationCancelResult, @@ -38,7 +42,14 @@ export interface SessionCollaborationServices { readonly operationId: string; }, onProgress?: (phase: SessionCollaborationImportPhase) => void): Promise; cancelImport(operationId: string): Promise; + readInvitationClipboard(): Promise; listMounts(): Promise; removeMount(mountId: string): Promise; + getPendingTurnRequests(): Promise; + decideTurnRequest( + sessionId: string, + requestId: string, + decision: 'approve' | 'reject', + ): Promise; createOperationId(): string; } diff --git a/apps/desktop/src/renderer/features/session-collaboration/testing.ts b/apps/desktop/src/renderer/features/session-collaboration/testing.ts new file mode 100644 index 0000000000..a1d8da6866 --- /dev/null +++ b/apps/desktop/src/renderer/features/session-collaboration/testing.ts @@ -0,0 +1,25 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ + +export { + groupPendingTurnRequests, + samePendingTurnRequests, + turnRequestPreview, + unseenTurnRequests, +} from './model/turn-request-inbox.js'; diff --git a/apps/desktop/src/renderer/features/session-collaboration/turn-request-inbox-context.tsx b/apps/desktop/src/renderer/features/session-collaboration/turn-request-inbox-context.tsx new file mode 100644 index 0000000000..2b7af3eb61 --- /dev/null +++ b/apps/desktop/src/renderer/features/session-collaboration/turn-request-inbox-context.tsx @@ -0,0 +1,72 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ + +import { createContext, useContext, useMemo, type ReactNode } from 'react'; +import type { ToastApi } from '@maka/ui'; +import { useSessionTurnRequestInbox } from './controller/use-turn-request-inbox.js'; + +export interface SessionTurnRequestInboxCopy { + readonly sharedTask: string; + readonly newTurnRequestTitle: (count: number) => string; + readonly newTurnRequestSummary: (count: number) => string; + readonly reviewTurnRequest: string; + readonly turnRequests: string; + readonly ownerTurnRequestTitle: string; + readonly reject: string; + readonly approve: string; + readonly moreTurnRequests: (count: number) => string; + readonly pendingTurnRequestCount: (count: number) => string; +} + +type SessionTurnRequestInbox = ReturnType & { + readonly copy: SessionTurnRequestInboxCopy; +}; + +const SessionTurnRequestInboxContext = createContext(null); + +export function SessionTurnRequestInboxProvider(props: { + readonly sessions: readonly { readonly id: string; readonly name: string }[]; + readonly toast: ToastApi; + readonly onOpenSession: (sessionId: string) => void; + readonly copy: SessionTurnRequestInboxCopy; + readonly children?: ReactNode; +}) { + const inbox = useSessionTurnRequestInbox(props); + const value = useMemo( + () => ({ ...inbox, copy: props.copy }), + [ + inbox.decide, + inbox.requests, + inbox.requestsBySession, + inbox.workingRequestIds, + props.copy, + ], + ); + return ( + + {props.children} + + ); +} + +export function useSessionTurnRequestInboxContext(): SessionTurnRequestInbox { + const inbox = useContext(SessionTurnRequestInboxContext); + if (!inbox) throw new Error('SessionTurnRequestInboxProvider is missing'); + return inbox; +} diff --git a/apps/desktop/src/renderer/features/session-collaboration/ui/session-collaboration-join-dialog.tsx b/apps/desktop/src/renderer/features/session-collaboration/ui/session-collaboration-join-dialog.tsx index 414158f9ed..d047c09b2a 100644 --- a/apps/desktop/src/renderer/features/session-collaboration/ui/session-collaboration-join-dialog.tsx +++ b/apps/desktop/src/renderer/features/session-collaboration/ui/session-collaboration-join-dialog.tsx @@ -45,6 +45,9 @@ export interface SessionCollaborationJoinCopy { readonly invalidCode: string; readonly directPathUnavailable: string; readonly code: string; + readonly pasteInvitation: string; + readonly clipboardEmpty: string; + readonly clipboardUnavailable: string; readonly join: string; readonly validatingInvitation: string; readonly discoveringHost: string; @@ -66,6 +69,7 @@ export function SessionCollaborationJoinDialog(props: { const services = useSessionCollaborationServices(); const toast = useToast(); const [code, setCode] = useState(''); + const [readingClipboard, setReadingClipboard] = useState(false); const [mounts, setMounts] = useState([]); const [removingMountId, setRemovingMountId] = useState(); const activeOperationId = useRef(undefined); @@ -198,6 +202,27 @@ export function SessionCollaborationJoinDialog(props: { } } + async function pasteInvitation(): Promise { + if (working || readingClipboard) return; + setReadingClipboard(true); + try { + const invitation = await services.readInvitationClipboard(); + if (!open.current) return; + if (!invitation) { + toast.info(props.copy.joinTitle, props.copy.clipboardEmpty); + return; + } + setCode(invitation); + setJoinState({ kind: 'idle' }); + } catch { + if (open.current) { + toast.error(props.copy.joinTitle, props.copy.clipboardUnavailable); + } + } finally { + if (open.current) setReadingClipboard(false); + } + } + return ( { + setCode(value); + if (joinState.kind === 'failed') setJoinState({ kind: 'idle' }); + }} /> +
+
{mounts.length > 0 ? ( {mounts.map((mount) => ( @@ -272,7 +310,7 @@ export function SessionCollaborationJoinDialog(props: {
); diff --git a/apps/desktop/src/renderer/styles/settings/runtime-host.css b/apps/desktop/src/renderer/styles/settings/runtime-host.css index d1d28b5e58..cb19a2bb5b 100644 --- a/apps/desktop/src/renderer/styles/settings/runtime-host.css +++ b/apps/desktop/src/renderer/styles/settings/runtime-host.css @@ -835,6 +835,11 @@ gap: var(--space-4); } +.sessionCollaborationInvitationAction { + display: flex; + justify-content: flex-end; +} + .sessionCollaborationDisclosure { display: flex; flex-direction: column; diff --git a/apps/desktop/src/renderer/styles/shell-layout.css b/apps/desktop/src/renderer/styles/shell-layout.css index ca810d0622..edf355a89f 100644 --- a/apps/desktop/src/renderer/styles/shell-layout.css +++ b/apps/desktop/src/renderer/styles/shell-layout.css @@ -584,6 +584,54 @@ dialog:modal { background: var(--background); } +.sessionTurnRequestApproval { + display: grid; + gap: var(--space-1); + width: 100%; + max-width: 36rem; + margin-inline: auto; + padding: var(--space-2) var(--space-3) 0; + box-sizing: border-box; + background: var(--background); +} + +.sessionTurnRequestApprovalBanner { + min-width: 0; +} + +.sessionTurnRequestApprovalIntent { + display: block; + width: 100%; + max-width: 100%; + min-width: 0; + overflow: hidden; + text-overflow: ellipsis; + white-space: nowrap; + cursor: help; +} + +.sessionTurnRequestApprovalDetails { + width: min(32rem, calc(100vw - 2rem)); + max-height: min(16rem, 50vh); + overflow: auto; + white-space: pre-wrap; + overscroll-behavior: contain; + scrollbar-gutter: stable; + user-select: text; +} + +.sessionTurnRequestApprovalActions { + display: inline-flex; + align-items: center; + gap: var(--space-1); +} + +.sessionTurnRequestApprovalMore { + color: var(--foreground-secondary); + font: var(--typography-caption); + text-align: end; +} + .sessionTurnRequestComposer { padding: var(--space-3); } diff --git a/apps/desktop/src/renderer/styles/sidebar.css b/apps/desktop/src/renderer/styles/sidebar.css index e5a039a327..2e66caed60 100644 --- a/apps/desktop/src/renderer/styles/sidebar.css +++ b/apps/desktop/src/renderer/styles/sidebar.css @@ -357,6 +357,11 @@ max-width: 88px; } +.maka-session-row-attention-badge { + display: inline-flex; + flex: none; +} + .maka-session-row-host-badge > span { overflow: hidden; max-width: 100%; diff --git a/docs/astryx-surface-file-inventory.md b/docs/astryx-surface-file-inventory.md index 4cdf6cfa2b..2b832925c7 100644 --- a/docs/astryx-surface-file-inventory.md +++ b/docs/astryx-surface-file-inventory.md @@ -6,7 +6,7 @@ Generated against `@astryxdesign/core@0.5.2` (194 component exports). Wiki bar: Design Conventions · API Use-the-System · Theming · Container Padding. -**Totals:** 237 files — blocker 0, reimplementation 0, polish 1, aligned 236. +**Totals:** 240 files — blocker 0, reimplementation 0, polish 1, aligned 239. ## Exclusions (explicit) @@ -54,7 +54,10 @@ Wiki bar: Design Conventions · API Use-the-System · Theming · Container Paddi | `apps/desktop/src/renderer/features/runtime-host-management/ui/runtime-host-peer-mesh-dialog.tsx` | dialog-overlay | Badge, Banner, Button, Dialog, DialogHeader, HStack, Layout, LayoutContent, LayoutFooter, MoreMenu, SegmentedControl, SegmentedControlItem, Selector, Switch, Text, TextArea, TextInput, Tooltip | aligned — uses Astryx (Badge, Banner, Button, Dialog, DialogHeader, HStack, Layout, LayoutContent) | aligned | | `apps/desktop/src/renderer/features/runtime-host-management/ui/runtime-host-profile-pairing-actions.tsx` | other | Button, MoreMenu | aligned — uses Astryx (Button, MoreMenu) | aligned | | `apps/desktop/src/renderer/features/session-collaboration/services-context.tsx` | shell-chrome-or-panel | none | aligned — no raw controls; no Astryx JSX usage | aligned | +| `apps/desktop/src/renderer/features/session-collaboration/turn-request-inbox-context.tsx` | shell-chrome-or-panel | none | aligned — no raw controls; no Astryx JSX usage | aligned | | `apps/desktop/src/renderer/features/session-collaboration/ui/session-collaboration-join-dialog.tsx` | dialog-overlay | Banner, Button, Dialog, DialogHeader, FormLayout, Layout, LayoutContent, LayoutFooter, List, ListItem, TextArea | aligned — uses Astryx (Banner, Button, Dialog, DialogHeader, FormLayout, Layout, LayoutContent, LayoutFooter) | aligned | +| `apps/desktop/src/renderer/features/session-collaboration/ui/session-turn-request-approval.tsx` | shell-chrome-or-panel | Banner, Button, HoverCard | aligned — uses Astryx (Banner, Button, HoverCard) | aligned | +| `apps/desktop/src/renderer/features/session-collaboration/ui/session-turn-request-badge.tsx` | shell-chrome-or-panel | Badge | aligned — uses Astryx (Badge) | aligned | | `apps/desktop/src/renderer/features/session-navigation/services-context.tsx` | shell-chrome-or-panel | none | aligned — no raw controls; no Astryx JSX usage | aligned | | `apps/desktop/src/renderer/features/session-navigation/ui/session-navigation-provider.tsx` | shell-chrome-or-panel | none | aligned — no raw controls; no Astryx JSX usage | aligned | | `apps/desktop/src/renderer/features/task-entry/services-context.tsx` | other | none | aligned — no raw controls; no Astryx JSX usage | aligned | @@ -81,7 +84,7 @@ Wiki bar: Design Conventions · API Use-the-System · Theming · Container Paddi | `apps/desktop/src/renderer/plan-mode-panel.tsx` | shell-chrome-or-panel | Badge, Banner, Button, Collapsible | aligned — uses Astryx (Badge, Banner, Button, Collapsible) | aligned | | `apps/desktop/src/renderer/reference-shell.css` | styles | n/a (css) | aligned — no off-rhythm control heights flagged | aligned | | `apps/desktop/src/renderer/remote-project-directory-dialog.tsx` | dialog-overlay | Button, Dialog, DialogHeader, DropdownMenu, DropdownMenuItem, HStack, Layout, LayoutContent, LayoutFooter, Text | aligned — uses Astryx (Button, Dialog, DialogHeader, DropdownMenu, DropdownMenuItem, HStack, Layout, LayoutContent) | aligned | -| `apps/desktop/src/renderer/session-collaboration-dialog.tsx` | dialog-overlay | Button, Dialog, DialogHeader, FormLayout, Layout, LayoutContent, LayoutFooter, SegmentedControl, SegmentedControlItem, Text, TextArea | aligned — uses Astryx (Button, Dialog, DialogHeader, FormLayout, Layout, LayoutContent, LayoutFooter, SegmentedControl) | aligned | +| `apps/desktop/src/renderer/session-collaboration-dialog.tsx` | dialog-overlay | Banner, Button, Dialog, DialogHeader, FormLayout, Layout, LayoutContent, LayoutFooter, SegmentedControl, SegmentedControlItem, Text, TextArea | aligned — uses Astryx (Banner, Button, Dialog, DialogHeader, FormLayout, Layout, LayoutContent, LayoutFooter) | aligned | | `apps/desktop/src/renderer/session-turn-request-composer.tsx` | shell-chrome-or-panel | Button, Text, TextArea | aligned — uses Astryx (Button, Text, TextArea) | aligned | | `apps/desktop/src/renderer/settings/about-settings-page.tsx` | settings-page | Badge, Banner, Button, Kbd, Link, List, ListItem | aligned — uses Astryx (Badge, Banner, Button, Kbd, Link, List, ListItem) | aligned | | `apps/desktop/src/renderer/settings/appearance-settings-page.tsx` | settings-page | Button, Grid, HStack, NumberInput, SelectableCard, Switch, Text, VStack | aligned — uses Astryx (Button, Grid, HStack, NumberInput, SelectableCard, Switch, Text, VStack) | aligned | diff --git a/docs/astryx-surface-file-inventory.paths b/docs/astryx-surface-file-inventory.paths index d699489cd9..398cd565ec 100644 --- a/docs/astryx-surface-file-inventory.paths +++ b/docs/astryx-surface-file-inventory.paths @@ -25,7 +25,10 @@ apps/desktop/src/renderer/features/runtime-host-management/ui/peer-mesh-peer-id- apps/desktop/src/renderer/features/runtime-host-management/ui/runtime-host-peer-mesh-dialog.tsx apps/desktop/src/renderer/features/runtime-host-management/ui/runtime-host-profile-pairing-actions.tsx apps/desktop/src/renderer/features/session-collaboration/services-context.tsx +apps/desktop/src/renderer/features/session-collaboration/turn-request-inbox-context.tsx apps/desktop/src/renderer/features/session-collaboration/ui/session-collaboration-join-dialog.tsx +apps/desktop/src/renderer/features/session-collaboration/ui/session-turn-request-approval.tsx +apps/desktop/src/renderer/features/session-collaboration/ui/session-turn-request-badge.tsx apps/desktop/src/renderer/features/session-navigation/services-context.tsx apps/desktop/src/renderer/features/session-navigation/ui/session-navigation-provider.tsx apps/desktop/src/renderer/features/task-entry/services-context.tsx diff --git a/native/runtime-host-peer/src/engine.rs b/native/runtime-host-peer/src/engine.rs index 357901a468..7f650253f6 100644 --- a/native/runtime-host-peer/src/engine.rs +++ b/native/runtime-host-peer/src/engine.rs @@ -82,6 +82,7 @@ const IDLE_CONNECTION_TIMEOUT: Duration = Duration::from_secs(10); const TARGET_COORDINATION_RESERVATIONS: usize = 2; const MAX_AUTOMATIC_RELAY_CANDIDATES: usize = 8; const MAX_RELAY_ADDRESSES_PER_PEER: usize = 4; +const MAX_PUBLISHED_COORDINATION_RELAY_ADDRESSES: usize = 16; const TRANSIT_FALLBACK_DELAY: Duration = Duration::from_secs(3); const MAX_TRANSIT_RESERVATIONS: usize = 32; const MAX_TRANSIT_CIRCUITS: usize = 8; @@ -644,14 +645,12 @@ async fn run_endpoint_async( &mut swarm, &mut coordination_relays, &HashMap::new(), - false, Instant::now(), ); let startup_deadline = Instant::now() + Duration::from_secs(10); let mut address_quiet_deadline = None; let mut bound_addresses = HashSet::new(); - let mut startup_external_candidate_ready = false; loop { let deadline = address_quiet_deadline .unwrap_or(startup_deadline) @@ -668,13 +667,9 @@ async fn run_endpoint_async( address_quiet_deadline = Some(Instant::now() + LISTENER_ADDRESS_QUIET_PERIOD); } } - Ok(event) => handle_startup_event( - &mut swarm, - event, - &mut coordination_relays, - &mut startup_external_candidate_ready, - &mut transit, - ), + Ok(event) => { + handle_startup_event(&mut swarm, event, &mut coordination_relays, &mut transit) + } Err(_) if pending_listeners.is_empty() => break, Err(_) => { return Err(PeerError::new( @@ -698,7 +693,6 @@ async fn run_endpoint_async( let (stream_completed_tx, mut stream_completed_rx) = mpsc::channel::(MAX_ESTABLISHED_CONNECTIONS as usize); let mut direct = DirectConnectState::default(); - let mut external_candidate_ready = startup_external_candidate_ready; let mut deadline_tick = tokio::time::interval(Duration::from_millis(100)); deadline_tick.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip); let mut discovered_relays = options @@ -770,7 +764,6 @@ async fn run_endpoint_async( &mut direct, &coordination_relays, &stream_control, - external_candidate_ready, Instant::now(), ); start_pending_webrtc_upgrades( @@ -831,7 +824,6 @@ async fn run_endpoint_async( &mut direct, &coordination_relays, &stream_control, - external_candidate_ready, Instant::now(), ); let requests = direct.pending.keys().copied().collect::>(); @@ -1002,7 +994,6 @@ async fn run_endpoint_async( &mut swarm, &mut coordination_relays, &direct.active, - external_candidate_ready, Instant::now(), ); } @@ -1095,7 +1086,6 @@ async fn run_endpoint_async( &mut direct, &coordination_relays, &stream_control, - external_candidate_ready, Instant::now(), ); maybe_open_peer_stream( @@ -1190,7 +1180,6 @@ async fn run_endpoint_async( &mut direct, &coordination_relays, &stream_control, - external_candidate_ready, Instant::now(), ); maybe_open_peer_stream( @@ -1210,7 +1199,6 @@ async fn run_endpoint_async( event, &mut coordination_relays, &mut direct, - &mut external_candidate_ready, RouteRuntime { active_coordination_relays: &active_coordination_relays, transit: &mut transit, @@ -1230,7 +1218,6 @@ async fn run_endpoint_async( &mut swarm, &mut coordination_relays, &direct.active, - external_candidate_ready, Instant::now(), ); let requests = direct.pending.keys().copied().collect::>(); @@ -1268,7 +1255,6 @@ async fn run_endpoint_async( &mut swarm, &mut coordination_relays, &direct.active, - external_candidate_ready, now, ); retry_connect_routes( @@ -1276,7 +1262,6 @@ async fn run_endpoint_async( &mut direct, &coordination_relays, &stream_control, - external_candidate_ready, now, ); start_pending_webrtc_upgrades( @@ -1508,13 +1493,7 @@ fn start_connect( referenced.insert(relay_peer), )?; } - maintain_coordination_relays( - swarm, - coordination_relays, - active_streams, - false, - Instant::now(), - ); + maintain_coordination_relays(swarm, coordination_relays, active_streams, Instant::now()); Ok(StartedConnect { direct_routes: direct_targets, coordination_relay_peers: relay_peers, @@ -1793,7 +1772,6 @@ fn handle_swarm_event( event: SwarmEvent, coordination_relays: &mut HashMap, direct: &mut DirectConnectState, - external_candidate_ready: &mut bool, route_runtime: RouteRuntime<'_>, ) { match event { @@ -2016,13 +1994,7 @@ fn handle_swarm_event( if let Some(relay) = coordination_relays.get_mut(&peer_id) { relay.identify_received = true; } - request_coordination_reservation( - swarm, - coordination_relays, - peer_id, - *external_candidate_ready, - Instant::now(), - ); + request_coordination_reservation(swarm, coordination_relays, peer_id, Instant::now()); } SwarmEvent::Behaviour(BehaviourEvent::Identify(identify::Event::Sent { peer_id, .. @@ -2030,25 +2002,7 @@ fn handle_swarm_event( if let Some(relay) = coordination_relays.get_mut(&peer_id) { relay.identify_sent = true; } - request_coordination_reservation( - swarm, - coordination_relays, - peer_id, - *external_candidate_ready, - Instant::now(), - ); - } - SwarmEvent::NewExternalAddrCandidate { .. } => { - *external_candidate_ready = true; - for peer_id in coordination_relays.keys().copied().collect::>() { - request_coordination_reservation( - swarm, - coordination_relays, - peer_id, - true, - Instant::now(), - ); - } + request_coordination_reservation(swarm, coordination_relays, peer_id, Instant::now()); } SwarmEvent::ListenerClosed { listener_id, @@ -2093,7 +2047,6 @@ fn handle_startup_event( swarm: &mut Swarm, event: SwarmEvent, coordination_relays: &mut HashMap, - external_candidate_ready: &mut bool, transit: &mut TransitRuntime, ) { handle_swarm_event( @@ -2101,7 +2054,6 @@ fn handle_startup_event( event, coordination_relays, &mut DirectConnectState::default(), - external_candidate_ready, RouteRuntime { active_coordination_relays: &Arc::new(RwLock::new(Vec::new())), transit, @@ -2405,7 +2357,7 @@ fn reconcile_transit_reservations( } relay.direct_connection_addresses.clear(); } - maintain_coordination_relays(swarm, relays, active_streams, true, Instant::now()); + maintain_coordination_relays(swarm, relays, active_streams, Instant::now()); } fn retained_transit_candidates( @@ -2600,13 +2552,38 @@ fn publish_active_coordination_relays( relays: &HashMap, snapshot: &Arc>>, ) { - let mut addresses = relays - .values() - .filter(|relay| relay.reserve && relay.reservation_accepted) - .flat_map(|relay| relay.reservation_addresses.iter().cloned()) + let mut relay_routes = relays + .iter() + .filter(|(_, relay)| relay.reserve && relay.reservation_accepted) + .map(|(peer_id, relay)| { + let mut routes = relay.reservation_addresses.clone(); + routes.sort_unstable_by_key(ToString::to_string); + routes.dedup(); + (peer_id.to_string(), routes) + }) .collect::>(); - addresses.sort_unstable_by_key(ToString::to_string); - addresses.dedup(); + relay_routes.sort_unstable_by(|left, right| left.0.cmp(&right.0)); + + let mut addresses = Vec::new(); + let mut seen = HashSet::new(); + let route_count = relay_routes + .iter() + .map(|(_, routes)| routes.len()) + .max() + .unwrap_or_default(); + 'routes: for route_index in 0..route_count { + for (_, routes) in &relay_routes { + let Some(address) = routes.get(route_index) else { + continue; + }; + if seen.insert(address.clone()) { + addresses.push(address.clone()); + if addresses.len() == MAX_PUBLISHED_COORDINATION_RELAY_ADDRESSES { + break 'routes; + } + } + } + } if let Ok(mut current) = snapshot.write() { *current = addresses; } @@ -2817,7 +2794,6 @@ fn maintain_coordination_relays( swarm: &mut Swarm, relays: &mut HashMap, active_streams: &HashMap, - external_candidate_ready: bool, now: Instant, ) { for peer_id in relays.keys().copied().collect::>() { @@ -2870,7 +2846,7 @@ fn maintain_coordination_relays( } continue; } - request_coordination_reservation(swarm, relays, peer_id, external_candidate_ready, now); + request_coordination_reservation(swarm, relays, peer_id, now); } } @@ -2878,7 +2854,6 @@ fn request_coordination_reservation( swarm: &mut Swarm, relays: &mut HashMap, peer_id: PeerId, - external_candidate_ready: bool, now: Instant, ) { let Some(relay) = relays.get_mut(&peer_id) else { @@ -2898,7 +2873,6 @@ fn request_coordination_reservation( || relay.reservation_listener.is_some() || !identified || (transit_provider && !directly_connected) - || !external_candidate_ready || relay.next_reservation_attempt > now { return; @@ -2979,7 +2953,6 @@ fn retry_connect_routes( direct: &mut DirectConnectState, coordination_relays: &HashMap, stream_control: &application_stream::Control, - external_candidate_ready: bool, now: Instant, ) { for connect in direct.pending.values_mut() { @@ -3011,24 +2984,11 @@ fn retry_connect_routes( .any(|origin| *origin == DialOrigin::Coordination) && (connect.retry_coordination || !stream_control.has_relayed_connection(peer_id)) { - let mut targets = Vec::new(); - for relay in &connect.coordination_relays { - let relay_peer = coordination_relay_peer_id(relay) - .expect("coordination relay was validated before connecting"); - if !external_candidate_ready - || coordination_relays - .get(&relay_peer) - .is_none_or(|relay| !relay.identify_received || !relay.identify_sent) - { - continue; - } - targets.push( - relay - .clone() - .with(Protocol::P2pCircuit) - .with(Protocol::P2p(peer_id)), - ); - } + let targets = coordination_dial_targets( + peer_id, + &connect.coordination_relays, + coordination_relays, + ); if let Some(connection_id) = dial_direct_targets(swarm, peer_id, targets) { connect .dials @@ -3069,6 +3029,29 @@ fn retry_connect_routes( } } +fn coordination_dial_targets( + peer_id: PeerId, + addresses: &[Multiaddr], + relays: &HashMap, +) -> Vec { + addresses + .iter() + .filter(|address| { + let relay_peer = coordination_relay_peer_id(address) + .expect("coordination relay was validated before connecting"); + relays + .get(&relay_peer) + .is_some_and(|relay| relay.identify_received && relay.identify_sent) + }) + .map(|address| { + address + .clone() + .with(Protocol::P2pCircuit) + .with(Protocol::P2p(peer_id)) + }) + .collect() +} + fn dial_direct_targets( swarm: &mut Swarm, peer_id: PeerId, @@ -3121,6 +3104,36 @@ mod tests { ); } + #[test] + fn coordination_dial_only_requires_an_identified_relay() { + let relay_peer_id = PeerId::random(); + let target_peer_id = PeerId::random(); + let relay_address: Multiaddr = format!("/ip4/203.0.113.1/tcp/4001/p2p/{relay_peer_id}") + .parse() + .expect("valid relay address"); + let relays = HashMap::from([( + relay_peer_id, + CoordinationRelay { + identify_received: true, + identify_sent: true, + ..CoordinationRelay::default() + }, + )]); + + assert_eq!( + coordination_dial_targets( + target_peer_id, + std::slice::from_ref(&relay_address), + &relays, + ), + vec![ + relay_address + .with(Protocol::P2pCircuit) + .with(Protocol::P2p(target_peer_id)) + ], + ); + } + #[test] fn completion_at_the_immutable_deadline_cannot_commit() { let now = Instant::now(); @@ -3944,6 +3957,49 @@ mod tests { ); } + #[test] + fn active_coordination_routes_respect_the_shared_route_limit() { + let relays = (0..5) + .map(|relay_index| { + let peer_id = PeerId::random(); + let reservation_addresses = (0..MAX_RELAY_ADDRESSES_PER_PEER) + .map(|route_index| { + format!( + "/ip4/192.0.2.{}/tcp/{}/p2p/{peer_id}", + relay_index + 1, + 4001 + route_index, + ) + .parse() + .expect("valid relay address") + }) + .collect(); + ( + peer_id, + CoordinationRelay { + reserve: true, + reservation_accepted: true, + reservation_addresses, + ..CoordinationRelay::default() + }, + ) + }) + .collect::>(); + let snapshot = Arc::new(RwLock::new(Vec::new())); + + publish_active_coordination_relays(&relays, &snapshot); + + let routes = snapshot.read().expect("read snapshot"); + assert_eq!(routes.len(), MAX_PUBLISHED_COORDINATION_RELAY_ADDRESSES); + assert_eq!( + routes + .iter() + .map(|route| coordination_relay_peer_id(route).expect("published relay identity")) + .collect::>() + .len(), + relays.len(), + ); + } + #[test] fn coordination_relay_requires_one_terminal_peer_identity() { let relay = PeerId::random(); diff --git a/packages/runtime-host/src/__tests__/execution-service-listeners.test.ts b/packages/runtime-host/src/__tests__/execution-service-listeners.test.ts new file mode 100644 index 0000000000..8027a0aab7 --- /dev/null +++ b/packages/runtime-host/src/__tests__/execution-service-listeners.test.ts @@ -0,0 +1,62 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ + +import assert from 'node:assert/strict'; +import test from 'node:test'; +import type { RuntimeHostPeerMeshOwner } from '../peer-mesh/owner.js'; +import { attachPeerOwnerCleanup } from '../server/execution-service.js'; +import type { RuntimeHostListenerSet } from '../server/listener-set.js'; + +test('the execution service exposes relay reservations discovered after startup', async () => { + let coordinationRelays: readonly string[] = []; + let listenersClosed = false; + let ownerClosed = false; + const peerListener = { + peerId: '12D3KooWpeer', + listenAddresses: ['/ip4/192.0.2.1/udp/41000/quic-v1'], + get coordinationRelays() { + return coordinationRelays; + }, + }; + const listeners: RuntimeHostListenerSet = { + listeners: [], + localEndpoint: 'local', + websocketEndpoints: [], + peerListeners: [peerListener], + async closeAdmission() {}, + async cleanup() { + listenersClosed = true; + }, + }; + const owner = { + async close() { + ownerClosed = true; + }, + } as unknown as RuntimeHostPeerMeshOwner; + + const attached = attachPeerOwnerCleanup(listeners, owner); + assert.deepEqual(attached.peerListeners[0]?.coordinationRelays, []); + + coordinationRelays = ['/dns4/relay.example/udp/443/quic-v1/p2p/12D3KooWrelay']; + assert.deepEqual(attached.peerListeners[0]?.coordinationRelays, coordinationRelays); + + await attached.cleanup(); + assert.equal(listenersClosed, true); + assert.equal(ownerClosed, true); +}); diff --git a/packages/runtime-host/src/__tests__/peer-listener.test.ts b/packages/runtime-host/src/__tests__/peer-listener.test.ts index 2a5d900895..e0d40ff411 100644 --- a/packages/runtime-host/src/__tests__/peer-listener.test.ts +++ b/packages/runtime-host/src/__tests__/peer-listener.test.ts @@ -21,6 +21,7 @@ import assert from 'node:assert/strict'; import { setImmediate as waitForImmediate } from 'node:timers/promises'; import { test } from 'node:test'; import { createRuntimeHostPeerListener } from '../server/peer-listener.js'; +import { createRuntimeHostListenerSet } from '../server/listener-set.js'; import type { RuntimeHostPeerClient } from '../client/peer-client.js'; import type { RuntimeHostPeerNativeStream } from '../transport/peer-native.js'; @@ -122,6 +123,33 @@ test('bounds active application streams from one authenticated peer', async () = await listener.cleanup(); }); +test('projects newly accepted coordination relays from the running peer endpoint', async () => { + let coordinationRelays: readonly string[] = []; + const peer = { + ...peerWith([]), + identity: () => ({ + peerId: 'peer', + listenAddresses: ['/ip4/192.0.2.1/udp/41000/quic-v1'], + coordinationRelays, + }), + }; + const listener = createRuntimeHostPeerListener(peer, {} as never, () => {}); + const listeners = createRuntimeHostListenerSet( + { + kind: 'local_ipc', + endpoint: 'local', + closeAdmission: async () => undefined, + cleanup: async () => undefined, + }, + [listener], + ); + + assert.deepEqual(listeners.peerListeners[0]?.coordinationRelays, []); + coordinationRelays = ['/dns4/relay.example/udp/443/quic-v1/p2p/12D3KooWrelay']; + assert.deepEqual(listeners.peerListeners[0]?.coordinationRelays, coordinationRelays); + await listeners.cleanup(); +}); + function peerWith(streams: RuntimeHostPeerNativeStream[]): RuntimeHostPeerClient { return { identity: () => ({ peerId: 'peer', listenAddresses: [], coordinationRelays: [] }), @@ -140,6 +168,7 @@ function peerWith(streams: RuntimeHostPeerNativeStream[]): RuntimeHostPeerClient maxCircuitBytes: 256 * 1024 * 1024, }), configureTransit: async () => undefined, + observeAuthenticatedRoutes: () => undefined, connect: async () => { throw new Error('not used'); }, diff --git a/packages/runtime-host/src/__tests__/peer-mesh.test.ts b/packages/runtime-host/src/__tests__/peer-mesh.test.ts index bdeba658d1..be45b83810 100644 --- a/packages/runtime-host/src/__tests__/peer-mesh.test.ts +++ b/packages/runtime-host/src/__tests__/peer-mesh.test.ts @@ -352,6 +352,59 @@ test('reconciles changed routes, propagates removal, and recovers the verified c } }); +test('publishes a changed local route promptly and refreshes a live cached peer route', { + timeout: 10_000, +}, async () => { + const root = await mkdtemp(join(tmpdir(), 'maka-peer-mesh-route-liveness-')); + const network = new MemoryPeerNetwork(); + const authorityPeer = network.create('peer-a'); + const memberBPeer = network.create('peer-b'); + const memberCPeer = network.create('peer-c'); + const authority = await openPeerMeshNode({ + dataRoot: join(root, 'authority'), + peer: authorityPeer, + }); + const memberB = await openPeerMeshNode({ + dataRoot: join(root, 'member-b'), + peer: memberBPeer, + }); + const memberC = await openPeerMeshNode({ + dataRoot: join(root, 'member-c'), + peer: memberCPeer, + }); + // Keep the observer demand-driven so its explicit prepareRoutes() call, + // rather than a racing background reconciliation, owns the cache refresh. + const serving = [authority.serve(), memberB.serve()]; + try { + const mesh = await authority.create(); + await memberB.join(await authority.invite(mesh.roster.roster.meshId)); + await memberC.join(await authority.invite(mesh.roster.roster.meshId)); + await memberC.reconcile(); + assert.deepEqual(memberC.resolveRoutes('peer-b')?.routeHints, ['/memory/peer-b/p2p/peer-b']); + + const movedRoute = '/memory/peer-b-restarted/p2p/peer-b'; + const movedRelay = '/memory/relay/peer-b-restarted'; + memberBPeer.setRouteHints([movedRoute]); + memberBPeer.setCoordinationRelays([movedRelay]); + await waitForRoutes(authority, 'peer-b', [movedRoute], [movedRelay]); + + assert.deepEqual(memberC.resolveRoutes('peer-b')?.routeHints, ['/memory/peer-b/p2p/peer-b']); + assert.deepEqual(memberC.resolveRoutes('peer-b')?.coordinationRelays, ['/memory/relay/peer-b']); + await memberC.prepareRoutes('peer-b', AbortSignal.timeout(4_000)); + assert.deepEqual(memberC.resolveRoutes('peer-b')?.routeHints, [movedRoute]); + assert.deepEqual(memberC.resolveRoutes('peer-b')?.coordinationRelays, [movedRelay]); + } finally { + await Promise.allSettled([authority.close(), memberB.close(), memberC.close()]); + await Promise.allSettled([ + authorityPeer.close(), + memberBPeer.close(), + memberCPeer.close(), + ...serving, + ]); + await rm(root, { recursive: true, force: true }); + } +}); + test('reconciles one selected Mesh into signed transit routes and native policy', async () => { const root = await mkdtemp(join(tmpdir(), 'maka-peer-mesh-transit-')); const network = new MemoryPeerNetwork(); @@ -887,6 +940,7 @@ class MemoryPeerClient implements PeerMeshTransport { #responseDelayMs = 0; #reachable = true; #routeHints: readonly string[]; + #coordinationRelays: readonly string[]; #nextConnectionBarrier: | { readonly started: () => void; @@ -916,13 +970,14 @@ class MemoryPeerClient implements PeerMeshTransport { private readonly peers: ReadonlyMap, ) { this.#routeHints = [`/memory/${peerId}/p2p/${peerId}`]; + this.#coordinationRelays = [`/memory/relay/${peerId}`]; } identity() { return { peerId: this.peerId, listenAddresses: this.#routeHints, - coordinationRelays: [`/memory/relay/${this.peerId}`], + coordinationRelays: this.#coordinationRelays, } as const; } @@ -930,6 +985,10 @@ class MemoryPeerClient implements PeerMeshTransport { this.#routeHints = [...routeHints]; } + setCoordinationRelays(coordinationRelays: readonly string[]): void { + this.#coordinationRelays = [...coordinationRelays]; + } + setReachable(reachable: boolean): void { this.#reachable = reachable; } @@ -1128,6 +1187,29 @@ async function waitForAbortable(task: Promise, signal?: AbortSignal): Prom } } +async function waitForRoutes( + node: PeerMeshNode, + peerId: string, + expectedRouteHints: readonly string[], + expectedCoordinationRelays: readonly string[], +): Promise { + const deadline = Date.now() + 4_000; + while (Date.now() < deadline) { + const routes = node.resolveRoutes(peerId); + if ( + JSON.stringify(routes?.routeHints) === JSON.stringify(expectedRouteHints) && + JSON.stringify(routes?.coordinationRelays) === JSON.stringify(expectedCoordinationRelays) + ) + return; + await delay(25); + } + assert.deepEqual(node.resolveRoutes(peerId), { + routeHints: expectedRouteHints, + coordinationRelays: expectedCoordinationRelays, + transitRelayPeerIds: [], + }); +} + function memorySignature(peerId: string, payload: Buffer): Buffer { return createHash('sha256').update(peerId).update(payload).digest(); } diff --git a/packages/runtime-host/src/__tests__/peer-native.test.ts b/packages/runtime-host/src/__tests__/peer-native.test.ts index fc12167100..1e8c2d033b 100644 --- a/packages/runtime-host/src/__tests__/peer-native.test.ts +++ b/packages/runtime-host/src/__tests__/peer-native.test.ts @@ -51,6 +51,7 @@ let finishMeshAccept; const pending = new Map(); const stats = { starts: 0, closes: 0, requests: [], cancellations: [] }; let missFirstCancellation = true; +let selfContainedRoutesPrepared = false; const stream = { read: async () => null, write: async () => {}, close: async () => {}, abort: () => {} }; module.exports = { stats, @@ -72,12 +73,15 @@ module.exports = { connect: ({ requestId, peerId, routeHints, coordinationRelays, transitRelayPeerIds }) => { stats.requests.push({ requestId, peerId, routeHints, coordinationRelays, transitRelayPeerIds }); if (peerId === 'unreachable') return Promise.reject(Object.assign(new Error('transit_unavailable: no approved route'), { code: 'GenericFailure' })); - if (peerId === 'ready' || peerId === 'fallback') return Promise.resolve(stream); + if (peerId === 'ready' || peerId === 'fallback' || peerId === 'observed' || (peerId === 'self-contained' && selfContainedRoutesPrepared)) return Promise.resolve(stream); return new Promise((resolve, reject) => pending.set(requestId, { resolve, reject })); }, connectMeshControl: ({ requestId, peerId, routeHints, coordinationRelays, transitRelayPeerIds }) => { stats.requests.push({ requestId, peerId, routeHints, coordinationRelays, transitRelayPeerIds }); - if (peerId === 'ready') return Promise.resolve(stream); + if (peerId === 'ready' || peerId === 'self-contained') { + if (peerId === 'self-contained') selfContainedRoutesPrepared = true; + return Promise.resolve(stream); + } return new Promise((resolve, reject) => pending.set(requestId, { resolve, reject })); }, configureTransit: async () => {}, @@ -100,16 +104,22 @@ module.exports = { `, ); let routesPrepared = false; - const client = createRuntimeHostPeerClient({ + const preparedPeerIds: string[] = []; + let client!: ReturnType; + client = createRuntimeHostPeerClient({ nativePath, keyPath: join(directory, 'peer.key'), routeResolver: { prepareRoutes: async (peerId) => { + preparedPeerIds.push(peerId); routesPrepared = true; if (peerId === 'fallback') throw new Error('Mesh refresh failed'); + if (peerId === 'self-contained') { + await client.connectMeshControl(peerConnectInput(peerId)); + } }, - resolveRoutes: () => - routesPrepared + resolveRoutes: (peerId) => + routesPrepared && peerId !== 'observed' ? { routeHints: ['/memory/discovered'], coordinationRelays: ['/memory/relay'], @@ -151,6 +161,22 @@ module.exports = { await assert.rejects(client.connect(peerConnectInput('unreachable')), (failure: unknown) => { return failure instanceof RuntimeHostPeerError && failure.code === 'transit_unavailable'; }); + const selfContained = client.connect({ + ...peerConnectInput('self-contained'), + coordinationRelays: ['/memory/explicit-relay'], + }); + await selfContained; + assert.equal(preparedPeerIds.includes('self-contained'), true); + client.observeAuthenticatedRoutes({ + peerId: 'observed', + routeHints: ['/memory/fresh'], + coordinationRelays: ['/memory/fresh-relay'], + }); + await client.connect({ + ...peerConnectInput('observed'), + routeHints: ['/memory/stale'], + coordinationRelays: ['/memory/stale-relay'], + }); assert.deepEqual(native.default.stats, { starts: 1, closes: 0, @@ -197,6 +223,27 @@ module.exports = { coordinationRelays: ['/memory/relay'], transitRelayPeerIds: ['transit-peer'], }, + { + requestId: 7, + peerId: 'self-contained', + routeHints: ['/memory/1'], + coordinationRelays: [], + transitRelayPeerIds: [], + }, + { + requestId: 8, + peerId: 'self-contained', + routeHints: ['/memory/discovered', '/memory/1'], + coordinationRelays: ['/memory/relay', '/memory/explicit-relay'], + transitRelayPeerIds: ['transit-peer'], + }, + { + requestId: 9, + peerId: 'observed', + routeHints: ['/memory/fresh', '/memory/stale'], + coordinationRelays: ['/memory/fresh-relay', '/memory/stale-relay'], + transitRelayPeerIds: [], + }, ], cancellations: [1, 1], }); diff --git a/packages/runtime-host/src/__tests__/protocol.test.ts b/packages/runtime-host/src/__tests__/protocol.test.ts index db9f2d70f6..3bbe698872 100644 --- a/packages/runtime-host/src/__tests__/protocol.test.ts +++ b/packages/runtime-host/src/__tests__/protocol.test.ts @@ -1924,6 +1924,36 @@ describe('Runtime Host bootstrap protocol', () => { ); }); + test('publishes a bounded live Direct peer endpoint through Host status', () => { + const status = { + hostEpoch: 'epoch-1', + compositionId: 'maka.interactive', + compositionRevision: '1', + state: 'ready', + connections: 1, + activeOperations: 0, + activeResidencies: 0, + peerEndpoint: { + peerId: '12D3KooWhost', + routeHints: ['/ip4/192.0.2.1/udp/41000/quic-v1'], + coordinationRelays: ['/dns4/relay.example/udp/443/quic-v1/p2p/12D3KooWrelay'], + }, + }; + assert.deepEqual(HOST_BOOTSTRAP_OPERATION_SPECS['host.status'].decodeOutput(status), status); + assert.throws(() => + HOST_BOOTSTRAP_OPERATION_SPECS['host.status'].decodeOutput({ + ...status, + peerEndpoint: { + ...status.peerEndpoint, + coordinationRelays: [ + status.peerEndpoint.coordinationRelays[0], + status.peerEndpoint.coordinationRelays[0], + ], + }, + }), + ); + }); + test('keeps Runtime Host logs within the diagnostics operation contract', () => { for (let index = 0; index < 257; index += 1) { runtimeHostLogBuffer.append('info', `entry ${index}`); diff --git a/packages/runtime-host/src/__tests__/session-collaboration-authority.test.ts b/packages/runtime-host/src/__tests__/session-collaboration-authority.test.ts index c94f198afb..13a3185638 100644 --- a/packages/runtime-host/src/__tests__/session-collaboration-authority.test.ts +++ b/packages/runtime-host/src/__tests__/session-collaboration-authority.test.ts @@ -156,6 +156,52 @@ test('Turn requests reject execution input that the Owner cannot review', () => ); }); +test('Owner can query one durable Turn-request inbox without exposing another Guest', async () => { + const directory = await mkdtemp(join(tmpdir(), 'maka-session-turn-request-inbox-')); + const authority = await openRuntimeHostAccessAuthority(directory); + try { + assert.deepEqual(HOST_OPERATION_SPECS['collaboration.turn-request.query'].decodeInput({}), {}); + const firstGuest = await activateTurnGuest(authority, 'session-1'); + const secondGuest = await activateTurnGuest(authority, 'session-2'); + await authority.createTurnAccessRequest(firstGuest, { + intent: { + sessionId: 'session-1', + turnId: 'turn-1', + content: { text: 'First request' }, + }, + }); + await authority.createTurnAccessRequest(secondGuest, { + intent: { + sessionId: 'session-2', + turnId: 'turn-2', + content: { text: 'Second request' }, + }, + }); + + assert.deepEqual( + authority + .queryTurnAccessRequests(LOCAL_OWNER, {}) + .requests.map((request) => request.intent.sessionId), + ['session-1', 'session-2'], + ); + assert.deepEqual( + authority + .queryTurnAccessRequests(sessionGuest(firstGuest), {}) + .requests.map((request) => request.intent.sessionId), + ['session-1'], + ); + assert.deepEqual( + authority + .queryTurnAccessRequests(LOCAL_OWNER, { sessionId: 'session-2' }) + .requests.map((request) => request.intent.sessionId), + ['session-2'], + ); + } finally { + await authority.close(); + await rm(directory, { recursive: true, force: true }); + } +}); + test('an approved exact Turn request survives restart and is admitted once', async () => { const directory = await mkdtemp(join(tmpdir(), 'maka-session-turn-request-')); let authority = await openRuntimeHostAccessAuthority(directory); @@ -592,9 +638,12 @@ test('drain does not terminalize an in-flight admission failure', async () => { } }); -async function activateTurnGuest(authority: RuntimeHostAccessAuthority): Promise { +async function activateTurnGuest( + authority: RuntimeHostAccessAuthority, + sessionId = 'session-1', +): Promise { const prepared = await authority.prepareCollaborationInvitation('root-1', { - sessionId: 'session-1', + sessionId, grantKinds: ['session_turn_request'], }); const invitation = decodeCollaborationInvitationCode(prepared.invitationCode); diff --git a/packages/runtime-host/src/__tests__/session-subscription-client.test.ts b/packages/runtime-host/src/__tests__/session-subscription-client.test.ts index 6791d4d51f..7f5ba135a3 100644 --- a/packages/runtime-host/src/__tests__/session-subscription-client.test.ts +++ b/packages/runtime-host/src/__tests__/session-subscription-client.test.ts @@ -55,6 +55,7 @@ import { type SessionTranscriptFragment, type SessionTranscriptPage, type HostFrame, + type HostStatusResult, type RequestFrame, type SubscriptionFrame, } from '../protocol/index.js'; @@ -1275,9 +1276,48 @@ test('close stops transcript pagination after the in-flight page', async () => { assert.equal(pageRequests, 1); }); +test('probes an otherwise idle accepted Runtime Host connection', { timeout: 2_000 }, async () => { + const probed = deferred(); + const observed = deferred(); + await withProtocolPeer( + async (transport, hostEpoch, rootId) => { + const hello = decodeClientFrame(await transport.read(1_000)); + assert.ok('kind' in hello && hello.kind === 'hello'); + await writeProtocolFrame(transport, { + kind: 'accepted', + rootId, + hostEpoch, + connectionId: 'connection-idle-liveness', + selectedProtocol: RUNTIME_HOST_PROTOCOL_VERSION, + compatibilityEpoch: RUNTIME_HOST_COMPATIBILITY_EPOCH, + compositionId: 'maka.interactive', + compositionRevision: '1', + state: 'ready', + }); + await answerStatus(transport, hostEpoch); + }, + async () => { + const status = await observed.promise; + assert.equal(status.state, 'ready'); + assert.equal(status.compositionId, 'maka.interactive'); + await probed.promise; + }, + { + livenessIntervalMs: 20, + onLivenessProbe: probed.resolve, + onHostStatus: observed.resolve, + }, + ); +}); + async function withProtocolPeer( serve: (transport: FramedTransport, hostEpoch: string, rootId: string) => Promise, run: (connection: RuntimeHostConnection) => Promise, + connectionOptions: { + readonly livenessIntervalMs?: number; + readonly onLivenessProbe?: () => void; + readonly onHostStatus?: (status: HostStatusResult) => void; + } = {}, ): Promise { const base = await mkdtemp(join(tmpdir(), 'maka-runtime-host-subscription-')); const capability = await resolveStorageRoot({ @@ -1318,6 +1358,7 @@ async function withProtocolPeer( const connected = await connectRuntimeHost({ rootPath: join(base, 'root'), protocol: PROTOCOL, + ...connectionOptions, }); assert.equal(connected.kind, 'connected'); if (connected.kind !== 'connected') return; diff --git a/packages/runtime-host/src/client/connection.ts b/packages/runtime-host/src/client/connection.ts index 55a1749f3f..9788c79f06 100644 --- a/packages/runtime-host/src/client/connection.ts +++ b/packages/runtime-host/src/client/connection.ts @@ -98,9 +98,9 @@ export interface ConnectRuntimeHostInput { connectTimeoutMs?: number; handshakeTimeoutMs?: number; /** - * Interval between liveness probes while a domain request is outstanding. - * Injectable so tests exercise requests that outlive a probe cycle without - * waiting the real cadence; defaults to DEFAULT_LIVENESS_INTERVAL_MS (2s). + * Maximum quiet interval before an end-to-end Host liveness probe. + * Any valid inbound frame restarts the interval. Injectable so tests can + * exercise the cadence without waiting the real default (2s). */ livenessIntervalMs?: number; /** @@ -111,6 +111,8 @@ export interface ConnectRuntimeHostInput { * connection health. */ onLivenessProbe?: () => void; + /** Receives each identity-validated Host status observation. */ + onHostStatus?: (status: HostStatusResult) => void; } export type RuntimeHostUnavailableReason = @@ -164,6 +166,7 @@ export interface ConnectRemoteRuntimeHostInput { readonly handshakeTimeoutMs?: number; readonly livenessIntervalMs?: number; readonly onLivenessProbe?: () => void; + readonly onHostStatus?: (status: HostStatusResult) => void; readonly connectionResource?: RuntimeHostConnectionResource; } @@ -176,6 +179,7 @@ export interface ConnectRuntimeHostMessageTransportInput { readonly handshakeTimeoutMs?: number; readonly livenessIntervalMs?: number; readonly onLivenessProbe?: () => void; + readonly onHostStatus?: (status: HostStatusResult) => void; readonly connectionResource?: RuntimeHostConnectionResource; readonly peerPath?: RuntimeHostPeerConnectionPath; } @@ -197,6 +201,7 @@ export type ConnectRemoteRuntimeHostResult = | 'unreachable' | 'connect_failed' | 'handshake_failed' + | 'handshake_timed_out' | 'root_mismatch' | 'composition_mismatch'; }; @@ -354,6 +359,7 @@ class RuntimeHostConnectionImpl implements RuntimeHostConnection { #terminalError: Error | undefined; readonly #livenessIntervalMs: number; readonly #onLivenessProbe: (() => void) | undefined; + readonly #onHostStatus: ((status: HostStatusResult) => void) | undefined; constructor( transport: RuntimeHostMessageTransport, @@ -370,12 +376,14 @@ class RuntimeHostConnectionImpl implements RuntimeHostConnection { options?: { livenessIntervalMs?: number; onLivenessProbe?: () => void; + onHostStatus?: (status: HostStatusResult) => void; connectionResource?: RuntimeHostConnectionResource; peerPath?: RuntimeHostPeerConnectionPath; }, ) { this.#livenessIntervalMs = options?.livenessIntervalMs ?? DEFAULT_LIVENESS_INTERVAL_MS; this.#onLivenessProbe = options?.onLivenessProbe; + this.#onHostStatus = options?.onHostStatus; this.#transport = transport; this.rootId = accepted.rootId; this.hostEpoch = accepted.hostEpoch; @@ -421,6 +429,7 @@ class RuntimeHostConnectionImpl implements RuntimeHostConnection { onFailure: (error) => this.#fail(error), }); void this.#readResponses(); + this.#scheduleLivenessCheck(); } request( @@ -494,7 +503,6 @@ class RuntimeHostConnectionImpl implements RuntimeHostConnection { ...(isDomainRequest ? { domainState: 'queued' as const } : {}), timer, }); - this.#scheduleLivenessCheck(); }); const frame = { requestId, @@ -549,6 +557,11 @@ class RuntimeHostConnectionImpl implements RuntimeHostConnection { this.#fail(error); throw error; } + try { + this.#onHostStatus?.(status); + } catch { + // Observation cannot control the authenticated connection it watches. + } return status; } @@ -704,7 +717,6 @@ class RuntimeHostConnectionImpl implements RuntimeHostConnection { if (retired?.operation === frame.operation) { this.#retiredRequests.delete(frame.requestId); this.#releaseDomainSlot(retired); - this.#scheduleLivenessCheck(); return; } this.#fail(new Error('Runtime Host returned an unmatched operation response')); @@ -716,7 +728,6 @@ class RuntimeHostConnectionImpl implements RuntimeHostConnection { } this.#pendingRequests.delete(frame.requestId); if (pending.timer) clearTimeout(pending.timer); - this.#scheduleLivenessCheck(); if (frame.ok) { try { const accepted = pending.accept(frame.result); @@ -795,7 +806,6 @@ class RuntimeHostConnectionImpl implements RuntimeHostConnection { pending.reject( interruptedRequestError(pending.operation, 'not_dispatched', 'timeout', error), ); - this.#scheduleLivenessCheck(); return; } this.#retiredRequests.set(requestId, { @@ -803,7 +813,6 @@ class RuntimeHostConnectionImpl implements RuntimeHostConnection { ...(pending.domainState === 'in_flight' ? { domainState: pending.domainState } : {}), }); pending.reject(interruptedRequestError(pending.operation, 'dispatched', 'timeout', error)); - this.#scheduleLivenessCheck(); } #releaseDomainSlot(request: PendingRequest | RetiredRequest): void { @@ -814,22 +823,16 @@ class RuntimeHostConnectionImpl implements RuntimeHostConnection { } #resetLivenessCheck(): void { + // A status observer needs periodic route observations even while other + // Session traffic keeps the connection active. + if (this.#onHostStatus) return; if (this.#livenessTimer) clearTimeout(this.#livenessTimer); this.#livenessTimer = undefined; this.#scheduleLivenessCheck(); } #scheduleLivenessCheck(): void { - if ( - this.#terminalError || - this.#livenessTimer || - this.#livenessProbePending || - !this.#hasOutstandingDomainRequest() - ) { - if (!this.#hasOutstandingDomainRequest() && this.#livenessTimer) { - clearTimeout(this.#livenessTimer); - this.#livenessTimer = undefined; - } + if (this.#terminalError || this.#livenessTimer || this.#livenessProbePending) { return; } this.#livenessTimer = setTimeout(() => { @@ -838,18 +841,8 @@ class RuntimeHostConnectionImpl implements RuntimeHostConnection { }, this.#livenessIntervalMs); } - #hasOutstandingDomainRequest(): boolean { - if (this.#retiredRequests.size > 0) return true; - for (const pending of this.#pendingRequests.values()) { - if (pending.operation !== 'host.status') return true; - } - return false; - } - #startLivenessProbe(): void { - if (this.#terminalError || this.#livenessProbePending || !this.#hasOutstandingDomainRequest()) { - return; - } + if (this.#terminalError || this.#livenessProbePending) return; this.#livenessProbePending = true; void this.#requestOperation( 'host.status', @@ -1043,6 +1036,7 @@ export async function connectRemoteRuntimeHost( handshakeTimeoutMs: normalized.handshakeTimeoutMs, livenessIntervalMs: normalized.livenessIntervalMs, onLivenessProbe: input.onLivenessProbe, + onHostStatus: input.onHostStatus, connectionResource: input.connectionResource, }); } catch (error) { @@ -1056,11 +1050,13 @@ export async function connectRuntimeHostMessageTransport( ): Promise { let resourceTransferred = false; let timer: ReturnType | undefined; + let handshakeTimedOut = false; try { const normalized = normalizeConnectRuntimeHostInput(input); const compositionId = requireHostCompositionId(input.compositionId); const expectedRootId = requireHostRootId(input.expectedRootId); timer = setTimeout(() => { + handshakeTimedOut = true; input.transport.abort(new Error('Timed out handshaking with Runtime Host')); }, normalized.handshakeTimeoutMs); const result = await exchangeRuntimeHostHandshake({ @@ -1071,6 +1067,7 @@ export async function connectRuntimeHostMessageTransport( expectedRootId, livenessIntervalMs: normalized.livenessIntervalMs, onLivenessProbe: input.onLivenessProbe, + onHostStatus: input.onHostStatus, connectionResource: input.connectionResource, ...(input.peerPath ? { peerPath: input.peerPath } : {}), }); @@ -1088,7 +1085,10 @@ export async function connectRuntimeHostMessageTransport( if (error instanceof RuntimeHostCompositionMismatchError) { return { kind: 'unavailable', reason: 'composition_mismatch' }; } - return { kind: 'unavailable', reason: 'handshake_failed' }; + return { + kind: 'unavailable', + reason: handshakeTimedOut ? 'handshake_timed_out' : 'handshake_failed', + }; } finally { if (timer) clearTimeout(timer); if (!resourceTransferred) await input.connectionResource?.close().catch(() => undefined); @@ -1279,6 +1279,7 @@ export async function connectResolvedRuntimeHost( hostProtocol: { min: registration.protocolMin, max: registration.protocolMax }, livenessIntervalMs, onLivenessProbe: input.onLivenessProbe, + onHostStatus: input.onHostStatus, }); if (result.kind === 'connected') { if ( @@ -1376,6 +1377,7 @@ interface ExchangeRuntimeHostHandshakeInput { readonly expectedCompositionRevision?: string; readonly livenessIntervalMs?: number; readonly onLivenessProbe?: () => void; + readonly onHostStatus?: (status: HostStatusResult) => void; readonly connectionResource?: RuntimeHostConnectionResource; readonly peerPath?: RuntimeHostPeerConnectionPath; } @@ -1456,6 +1458,7 @@ async function exchangeRuntimeHostHandshake( connection: new RuntimeHostConnectionImpl(input.transport, handshake, { livenessIntervalMs: input.livenessIntervalMs, onLivenessProbe: input.onLivenessProbe, + onHostStatus: input.onHostStatus, connectionResource: input.connectionResource, ...(input.peerPath ? { peerPath: input.peerPath } : {}), }), diff --git a/packages/runtime-host/src/client/host-profile.ts b/packages/runtime-host/src/client/host-profile.ts index ea0a72ac9c..1e77cb64bb 100644 --- a/packages/runtime-host/src/client/host-profile.ts +++ b/packages/runtime-host/src/client/host-profile.ts @@ -27,6 +27,7 @@ import { INTERACTIVE_RUNTIME_HOST_COMPOSITION_ID, isCanonicalRuntimeHostWebSocketPath, RUNTIME_HOST_PROTOCOL_VERSION, + type HostStatusResult, requireClientInstanceId, requireHostRootId, } from '../protocol/index.js'; @@ -71,6 +72,7 @@ const PROFILE_ID_PATTERN = /^[A-Za-z0-9][A-Za-z0-9._-]{0,63}$/; const PEER_ID_MAX_BYTES = 160; const PEER_ADDRESS_MAX_BYTES = 2 * 1024; const PEER_ROUTE_MAX = 16; +const DEFAULT_PEER_HANDSHAKE_TIMEOUT_MS = 5_000; const PROFILE_CREDENTIAL_RECORD_PREFIX = 'maka-runtime-host-profile-credential-v1:'; const PROFILE_INCARNATION_ID_MAX_BYTES = 128; export const RUNTIME_HOST_ACCESS_CREDENTIAL_MAX_BYTES = 8 * 1024; @@ -382,6 +384,7 @@ export async function connectRuntimeHostProfile( readonly sshInteraction?: RuntimeHostSshInteraction; readonly peerClient?: RuntimeHostPeerClient; readonly onConnectionPhase?: (phase: RuntimeHostConnectionPhase) => void; + readonly onHostStatus?: (status: HostStatusResult) => void; }, overrides: { connect?: typeof connectRemoteRuntimeHost; @@ -438,6 +441,7 @@ export async function connectRemoteRuntimeHostProfile( readonly sshInteraction?: RuntimeHostSshInteraction; readonly peerClient?: RuntimeHostPeerClient; readonly onConnectionPhase?: (phase: RuntimeHostConnectionPhase) => void; + readonly onHostStatus?: (status: HostStatusResult) => void; }, overrides: { connect?: typeof connectRemoteRuntimeHost; @@ -466,6 +470,7 @@ export async function connectRemoteRuntimeHostProfile( ...(input.onConnectionPhase === undefined ? {} : { onConnectionPhase: input.onConnectionPhase }), + ...(input.onHostStatus === undefined ? {} : { onHostStatus: input.onHostStatus }), }); } else { notifyConnectionPhase(input.onConnectionPhase, 'connecting'); @@ -511,6 +516,7 @@ export async function connectRemoteRuntimeHostProfile( ...(input.handshakeTimeoutMs === undefined ? {} : { handshakeTimeoutMs: input.handshakeTimeoutMs }), + ...(input.onHostStatus === undefined ? {} : { onHostStatus: input.onHostStatus }), }); try { input.signal?.throwIfAborted(); @@ -586,8 +592,10 @@ export async function connectPeerRuntimeHost(input: { readonly connectTimeoutMs?: number; readonly handshakeTimeoutMs?: number; readonly onConnectionPhase?: (phase: RuntimeHostConnectionPhase) => void; + readonly onHostStatus?: (status: HostStatusResult) => void; }): Promise { input.signal?.throwIfAborted(); + const handshakeTimeoutMs = input.handshakeTimeoutMs ?? DEFAULT_PEER_HANDSHAKE_TIMEOUT_MS; const stream = await input.peerClient.connect( { peerId: input.transport.peerId, @@ -608,7 +616,7 @@ export async function connectPeerRuntimeHost(input: { await writeRuntimeHostPeerAuthentication(stream, input.credential); const authentication = await readRuntimeHostPeerAuthenticationResult( stream, - input.handshakeTimeoutMs, + handshakeTimeoutMs, ); if (!authentication.accepted) { throw new RuntimeHostProfileConnectionError( @@ -628,9 +636,14 @@ export async function connectPeerRuntimeHost(input: { max: RUNTIME_HOST_PROTOCOL_VERSION, }, clientInstanceId: input.clientInstanceId, - ...(input.handshakeTimeoutMs === undefined - ? {} - : { handshakeTimeoutMs: input.handshakeTimeoutMs }), + handshakeTimeoutMs, + onHostStatus: (status) => { + const endpoint = status.peerEndpoint; + if (endpoint?.peerId === input.transport.peerId) { + input.peerClient.observeAuthenticatedRoutes(endpoint); + } + input.onHostStatus?.(status); + }, ...(stream.path ? { peerPath: stream.path } : {}), }); input.signal?.throwIfAborted(); @@ -700,6 +713,9 @@ export function remoteRuntimeHostUnavailableError( case 'unreachable': message = `${subject} could not reach its endpoint`; break; + case 'handshake_timed_out': + message = `${subject} timed out while establishing its protocol session`; + break; default: message = `${subject} is unavailable (${reason})`; } diff --git a/packages/runtime-host/src/client/peer-client.ts b/packages/runtime-host/src/client/peer-client.ts index 543efdb2f0..705ba0f9fb 100644 --- a/packages/runtime-host/src/client/peer-client.ts +++ b/packages/runtime-host/src/client/peer-client.ts @@ -66,6 +66,11 @@ export interface RuntimeHostPeerClient { readonly approvedRelayPeerIds: readonly string[]; readonly relayCandidates: readonly RuntimeHostPeerTransitRelayCandidate[]; }): Promise; + observeAuthenticatedRoutes(input: { + readonly peerId: string; + readonly routeHints: readonly string[]; + readonly coordinationRelays: readonly string[]; + }): void; connect( input: RuntimeHostPeerConnectInput, signal?: AbortSignal, @@ -129,6 +134,13 @@ class RuntimeHostPeerClientImpl implements RuntimeHostPeerClient { readonly #automaticRelayDiscovery: boolean; readonly #webRtcStunUrls: readonly string[] | undefined; readonly #routeResolver: RuntimeHostPeerRouteResolver | undefined; + readonly #authenticatedRoutes = new Map< + string, + Readonly<{ + routeHints: readonly string[]; + coordinationRelays: readonly string[]; + }> + >(); #endpoint: RuntimeHostPeerNativeEndpoint | undefined; #draining: Promise | undefined; #meshDraining: Promise | undefined; @@ -210,6 +222,21 @@ class RuntimeHostPeerClientImpl implements RuntimeHostPeerClient { }); } + observeAuthenticatedRoutes(input: { + readonly peerId: string; + readonly routeHints: readonly string[]; + readonly coordinationRelays: readonly string[]; + }): void { + if (input.routeHints.length === 0 && input.coordinationRelays.length === 0) return; + this.#authenticatedRoutes.set( + input.peerId, + Object.freeze({ + routeHints: Object.freeze([...input.routeHints]), + coordinationRelays: Object.freeze([...input.coordinationRelays]), + }), + ); + } + async connect( input: RuntimeHostPeerConnectInput, signal?: AbortSignal, @@ -337,13 +364,18 @@ class RuntimeHostPeerClientImpl implements RuntimeHostPeerClient { const requestId = this.#allocateRequestId(); const discovered = kind === 'application' ? this.#routeResolver?.resolveRoutes(input.peerId) : undefined; + const authenticated = + kind === 'application' ? this.#authenticatedRoutes.get(input.peerId) : undefined; const connection = endpoint[kind === 'application' ? 'connect' : 'connectMeshControl']({ ...input, - routeHints: mergeAddresses(discovered?.routeHints ?? [], input.routeHints), - coordinationRelays: mergeAddresses( - discovered?.coordinationRelays ?? [], - input.coordinationRelays, - ), + routeHints: mergeAddresses(discovered?.routeHints ?? [], [ + ...(authenticated?.routeHints ?? []), + ...input.routeHints, + ]), + coordinationRelays: mergeAddresses(discovered?.coordinationRelays ?? [], [ + ...(authenticated?.coordinationRelays ?? []), + ...(input.coordinationRelays ?? []), + ]), transitRelayPeerIds: mergeValues( discovered?.transitRelayPeerIds ?? [], input.transitRelayPeerIds, diff --git a/packages/runtime-host/src/peer-mesh/node.ts b/packages/runtime-host/src/peer-mesh/node.ts index d49068723e..be4c224b46 100644 --- a/packages/runtime-host/src/peer-mesh/node.ts +++ b/packages/runtime-host/src/peer-mesh/node.ts @@ -76,6 +76,7 @@ const ROUTE_PAGE_SIZE = 8; const RECONCILE_CONCURRENCY = 4; const RECONCILE_DEADLINE_MS = 60 * 1_000; const RECONCILE_INTERVAL_MS = 30 * 1_000; +const LOCAL_ROUTE_OBSERVATION_INTERVAL_MS = 1_000; interface RedeemInvitationRequest { readonly kind: 'redeem-invitation'; @@ -224,6 +225,7 @@ export async function openPeerMeshNode(input: { readonly peer: PeerMeshTransport; readonly endpointKind?: 'client' | 'host'; readonly now?: () => number; + readonly onBackgroundReconcileError?: (error: unknown) => void; }): Promise { const store = await openPeerMeshStateStore(input.dataRoot, input.peer.identity().peerId); const node = new PeerMeshNodeImpl({ ...input, store }); @@ -241,6 +243,7 @@ class PeerMeshNodeImpl implements PeerMeshNode { readonly #peer: PeerMeshTransport; readonly #endpointKind: 'client' | 'host' | undefined; readonly #now: () => number; + readonly #onBackgroundReconcileError: ((error: unknown) => void) | undefined; readonly #activeControlStreams = new Set(); readonly #lifetime = new AbortController(); #admissionTail = Promise.resolve(); @@ -256,11 +259,13 @@ class PeerMeshNodeImpl implements PeerMeshNode { readonly peer: PeerMeshTransport; readonly endpointKind?: 'client' | 'host'; readonly now?: () => number; + readonly onBackgroundReconcileError?: (error: unknown) => void; }) { this.#store = input.store; this.#peer = input.peer; this.#endpointKind = input.endpointKind; this.#now = input.now ?? Date.now; + this.#onBackgroundReconcileError = input.onBackgroundReconcileError; } async initialize(): Promise { @@ -848,11 +853,11 @@ class PeerMeshNodeImpl implements PeerMeshNode { isActiveMembership(mesh, localPeerId) && mesh.roster.roster.members.includes(peerId), ); if (!visible) return; - if ( - stored.routes.some(({ route }) => route.peerId === peerId && route.expiresAt > this.#now()) - ) { - return; - } + // A signed route can remain within its TTL after a peer restarted or + // rotated Relay reservations. Every connection establishment therefore + // asks the Mesh control plane for its newest record. Callers with a + // self-contained invitation run this reconciliation in parallel with the + // first dial; callers without usable routes wait for it. await this.reconcile(signal); } @@ -910,10 +915,45 @@ class PeerMeshNodeImpl implements PeerMeshNode { } async #runReconciliation(signal: AbortSignal): Promise { + let failureReported = false; while (!signal.aborted) { - await this.reconcile(signal).catch(() => undefined); - await delay(RECONCILE_INTERVAL_MS, undefined, { signal }).catch(() => undefined); + try { + await this.reconcile(signal); + failureReported = false; + } catch (error) { + if (!signal.aborted && !failureReported) { + failureReported = true; + try { + this.#onBackgroundReconcileError?.(error); + } catch { + // Diagnostics cannot control Peer Mesh reconciliation. + } + } + } + await this.#waitForReconciliationTrigger(signal).catch(() => undefined); + } + } + + async #waitForReconciliationTrigger(signal: AbortSignal): Promise { + let remainingMs = RECONCILE_INTERVAL_MS; + while (!signal.aborted && remainingMs > 0) { + const waitMs = Math.min(LOCAL_ROUTE_OBSERVATION_INTERVAL_MS, remainingMs); + await delay(waitMs, undefined, { signal }); + if (this.#localRouteRequiresRefresh()) return; + remainingMs -= waitMs; + } + } + + #localRouteRequiresRefresh(): boolean { + const identity = this.#peer.identity(); + const current = this.#store.read(); + if (!current.meshes.some((state) => isActiveMembership(state, identity.peerId))) { + return false; } + const existing = current.routes + .filter(({ route }) => route.peerId === identity.peerId) + .sort((left, right) => right.route.sequence - left.route.sequence)[0]; + return !isCurrentLocalRoute(existing, identity, current, this.#endpointKind, this.#now()); } async #reconcile(signal?: AbortSignal): Promise { @@ -1157,15 +1197,7 @@ class PeerMeshNodeImpl implements PeerMeshNode { const existing = current.routes .filter(({ route }) => route.peerId === identity.peerId) .sort((left, right) => right.route.sequence - left.route.sequence)[0]; - if ( - existing && - existing.route.expiresAt > now + ROUTE_REFRESH_LEAD_MS && - sameAddresses(existing.route.routeHints, identity.listenAddresses) && - sameAddresses(existing.route.coordinationRelays, identity.coordinationRelays) && - existing.route.endpointKind === this.#endpointKind && - existing.route.displayName === (current.displayName ?? undefined) && - existing.route.transitMeshId === current.transitMeshId - ) { + if (isCurrentLocalRoute(existing, identity, current, this.#endpointKind, now)) { return { state: current, result: existing }; } const route = await this.#signLocalRoute(current); @@ -2124,6 +2156,24 @@ function sameAddresses(left: readonly string[], right: readonly string[]): boole return left.length === right.length && left.every((address, index) => address === right[index]); } +function isCurrentLocalRoute( + existing: SignedPeerMeshRouteRecordV1 | undefined, + identity: ReturnType, + current: Pick, + endpointKind: 'client' | 'host' | undefined, + now: number, +): existing is SignedPeerMeshRouteRecordV1 { + return Boolean( + existing && + existing.route.expiresAt > now + ROUTE_REFRESH_LEAD_MS && + sameAddresses(existing.route.routeHints, identity.listenAddresses) && + sameAddresses(existing.route.coordinationRelays, identity.coordinationRelays) && + existing.route.endpointKind === endpointKind && + existing.route.displayName === (current.displayName ?? undefined) && + existing.route.transitMeshId === current.transitMeshId, + ); +} + function mergeAddresses( primary: readonly string[], fallback: readonly string[], diff --git a/packages/runtime-host/src/peer-mesh/owner.ts b/packages/runtime-host/src/peer-mesh/owner.ts index 5bde0a8d10..035ab2a46b 100644 --- a/packages/runtime-host/src/peer-mesh/owner.ts +++ b/packages/runtime-host/src/peer-mesh/owner.ts @@ -44,6 +44,7 @@ export async function openRuntimeHostPeerMeshOwner(input: { readonly coordinationRelays?: readonly string[]; readonly automaticRelayDiscovery?: boolean; readonly webRtcStunUrls?: readonly string[]; + readonly onBackgroundReconcileError?: (error: unknown) => void; }): Promise { let mesh: PeerMeshNode | undefined; let resolverMesh: PeerMeshNode | undefined; @@ -77,6 +78,9 @@ export async function openRuntimeHostPeerMeshOwner(input: { dataRoot: join(input.dataRoot, client.identity().peerId), peer: client, endpointKind: input.endpointKind, + ...(input.onBackgroundReconcileError + ? { onBackgroundReconcileError: input.onBackgroundReconcileError } + : {}), }); resolverMesh = mesh; } catch (error) { diff --git a/packages/runtime-host/src/protocol/host-status.ts b/packages/runtime-host/src/protocol/host-status.ts index 4aac3a9603..a28ef97438 100644 --- a/packages/runtime-host/src/protocol/host-status.ts +++ b/packages/runtime-host/src/protocol/host-status.ts @@ -51,6 +51,9 @@ export type HostUpgradePrepareResult = export const HOST_DIAGNOSTICS_RESULT_MAX_BYTES = 72 * 1024; export const HOST_DIAGNOSTIC_LOG_MAX_ENTRIES = 256; export const HOST_DIAGNOSTIC_LOG_MAX_ENTRY_BYTES = 10 * 1024; +const HOST_PEER_ID_MAX_BYTES = 160; +const HOST_PEER_ADDRESS_MAX_BYTES = 2 * 1024; +const HOST_PEER_ROUTE_MAX = 16; export interface HostStatusResult { hostEpoch: string; @@ -60,6 +63,13 @@ export interface HostStatusResult { connections: number; activeOperations: number; activeResidencies: number; + peerEndpoint?: HostPeerEndpoint; +} + +export interface HostPeerEndpoint { + readonly peerId: string; + readonly routeHints: readonly string[]; + readonly coordinationRelays: readonly string[]; } export interface HostDiagnosticsResult extends HostStatusResult { @@ -100,12 +110,26 @@ export const HOST_BOOTSTRAP_OPERATION_SPECS = { }), } as const; +function decodePeerAddresses(value: unknown, label: string): readonly string[] { + if (!Array.isArray(value) || value.length > HOST_PEER_ROUTE_MAX) { + throw invalidProtocolFrame(`Invalid ${label}`); + } + const addresses = value.map((address) => + requireString(address, label, HOST_PEER_ADDRESS_MAX_BYTES), + ); + if (new Set(addresses).size !== addresses.length) { + throw invalidProtocolFrame(`Duplicate ${label}`); + } + return Object.freeze(addresses); +} + function decodeEmptyHostInput(value: unknown, label: string): HostStatusInput { requireExactRecord(value, label, []); return {}; } function decodeHostStatusResult(value: unknown): HostStatusResult { + const valueRecord = requireRecord(value, 'host.status result'); const record = requireExactRecord(value, 'host.status result', [ 'hostEpoch', 'compositionId', @@ -114,6 +138,7 @@ function decodeHostStatusResult(value: unknown): HostStatusResult { 'connections', 'activeOperations', 'activeResidencies', + ...(valueRecord.peerEndpoint === undefined ? [] : ['peerEndpoint']), ]); return decodeHostStatusFields(record); } @@ -124,6 +149,7 @@ function decodeHostDiagnosticsResult(value: unknown): HostDiagnosticsResult { 'host.diagnostics.query result', HOST_DIAGNOSTICS_RESULT_MAX_BYTES, ); + const valueRecord = requireRecord(value, 'host.diagnostics.query result'); const record = requireExactRecord(value, 'host.diagnostics.query result', [ 'hostEpoch', 'compositionId', @@ -132,6 +158,7 @@ function decodeHostDiagnosticsResult(value: unknown): HostDiagnosticsResult { 'connections', 'activeOperations', 'activeResidencies', + ...(valueRecord.peerEndpoint === undefined ? [] : ['peerEndpoint']), 'compositionModules', 'residencies', 'protocolVersion', @@ -262,6 +289,25 @@ function decodeHostStatusFields(record: Record): HostStatusResu connections: requireCount(record.connections, 'connections'), activeOperations: requireCount(record.activeOperations, 'activeOperations'), activeResidencies: requireCount(record.activeResidencies, 'activeResidencies'), + ...(record.peerEndpoint === undefined + ? {} + : { peerEndpoint: decodeHostPeerEndpoint(record.peerEndpoint) }), + }; +} + +function decodeHostPeerEndpoint(value: unknown): HostPeerEndpoint { + const record = requireExactRecord(value, 'Runtime Host peer endpoint', [ + 'peerId', + 'routeHints', + 'coordinationRelays', + ]); + return { + peerId: requireString(record.peerId, 'Runtime Host peer id', HOST_PEER_ID_MAX_BYTES), + routeHints: decodePeerAddresses(record.routeHints, 'Runtime Host peer route hints'), + coordinationRelays: decodePeerAddresses( + record.coordinationRelays, + 'Runtime Host peer coordination relays', + ), }; } diff --git a/packages/runtime-host/src/protocol/index.ts b/packages/runtime-host/src/protocol/index.ts index a173fdbc4a..7b7c8a2cf8 100644 --- a/packages/runtime-host/src/protocol/index.ts +++ b/packages/runtime-host/src/protocol/index.ts @@ -100,7 +100,10 @@ export const RUNTIME_HOST_REGISTRATION_SCHEMA_VERSION = 1 as const; export const RUNTIME_HOST_PROTOCOL_VERSION = 0 as const; // Increment when the same protocol version no longer guarantees safe Client-Host // interoperability. Mismatches are rejected before domain commands are admitted. -export const RUNTIME_HOST_COMPATIBILITY_EPOCH = 90 as const; +export const RUNTIME_HOST_COMPATIBILITY_EPOCH = 92 as const; +// 92: Owners can query their complete pending Session Turn-request inbox. +// 91: Host status publishes the live Direct peer endpoint so newly issued +// connection invitations do not preserve stale startup routes. // 90: `session.create.mode` accepts the Bot session mode. A Host that predates // it rejects the value as an invalid Session start mode. // 89: The Host refreshes its models.dev catalog at startup and announces the diff --git a/packages/runtime-host/src/protocol/operations.ts b/packages/runtime-host/src/protocol/operations.ts index 1f6f900805..f87902936e 100644 --- a/packages/runtime-host/src/protocol/operations.ts +++ b/packages/runtime-host/src/protocol/operations.ts @@ -71,6 +71,7 @@ export type { HostDiagnosticsResult, HostActivitySnapshot, HostLifecycleState, + HostPeerEndpoint, HostStatusInput, HostStatusResult, HostUpgradePrepareInput, diff --git a/packages/runtime-host/src/protocol/session-collaboration.ts b/packages/runtime-host/src/protocol/session-collaboration.ts index cdfc55df09..81eac55fdb 100644 --- a/packages/runtime-host/src/protocol/session-collaboration.ts +++ b/packages/runtime-host/src/protocol/session-collaboration.ts @@ -136,7 +136,8 @@ export interface CollaborationTurnRequestCreateInput { } export interface CollaborationTurnRequestQueryInput { - readonly sessionId: string; + /** Omit to query every request visible to the authenticated principal. */ + readonly sessionId?: string; } export interface CollaborationTurnRequestQueryResult { @@ -424,8 +425,10 @@ function decodeCollaborationTurnRequestCreateInput( function decodeCollaborationTurnRequestQueryInput( value: unknown, ): CollaborationTurnRequestQueryInput { - const record = requireExactRecord(value, 'collaboration Turn request query', ['sessionId']); - return { sessionId: requireEntityId(record.sessionId, 'sessionId') }; + const record = requireShapedRecord(value, 'collaboration Turn request query', [], ['sessionId']); + return record.sessionId === undefined + ? {} + : { sessionId: requireEntityId(record.sessionId, 'sessionId') }; } function decodeCollaborationTurnRequestQueryResult( diff --git a/packages/runtime-host/src/server/access-authority.ts b/packages/runtime-host/src/server/access-authority.ts index 2d7176ce28..f25fbc0e0a 100644 --- a/packages/runtime-host/src/server/access-authority.ts +++ b/packages/runtime-host/src/server/access-authority.ts @@ -439,12 +439,13 @@ class FileRuntimeHostAccessAuthority implements RuntimeHostAccessAuthority { return { canRequestTurns: guest && + input.sessionId !== undefined && this.activeSessionGrant(principal.principalId, input.sessionId, 'session_turn_request') !== undefined, requests: this.#file.turnAccessRequests.filter( (request) => (!guest || request.principalId === principal.principalId) && - request.intent.sessionId === input.sessionId, + (input.sessionId === undefined || request.intent.sessionId === input.sessionId), ), }; } diff --git a/packages/runtime-host/src/server/execution-service.ts b/packages/runtime-host/src/server/execution-service.ts index 1ae258c3ff..b2ee32967e 100644 --- a/packages/runtime-host/src/server/execution-service.ts +++ b/packages/runtime-host/src/server/execution-service.ts @@ -99,6 +99,9 @@ export async function startExecutionRuntimeHostService( ...options.peer, dataRoot: options.peer.meshDataRoot, endpointKind: 'host', + onBackgroundReconcileError: (error) => { + console.error('[runtime-host] Peer Mesh background synchronization failed:', error); + }, }); } catch (error) { console.error( @@ -160,7 +163,7 @@ export async function startExecutionRuntimeHostService( } } -function attachPeerOwnerCleanup( +export function attachPeerOwnerCleanup( listeners: RuntimeHostListenerSet, owner: RuntimeHostPeerMeshOwner, ): RuntimeHostListenerSet { diff --git a/packages/runtime-host/src/server/host-kernel.ts b/packages/runtime-host/src/server/host-kernel.ts index e86705ce28..491e663212 100644 --- a/packages/runtime-host/src/server/host-kernel.ts +++ b/packages/runtime-host/src/server/host-kernel.ts @@ -825,6 +825,7 @@ export class RuntimeHostKernel { } #statusSnapshot(): HostStatusResult { + const peer = this.peerListeners[0]; return { hostEpoch: this.hostEpoch, compositionId: this.compositionDescriptor.id, @@ -833,6 +834,15 @@ export class RuntimeHostKernel { connections: this.#acceptedTransports.size, activeOperations: this.#activeOperations, activeResidencies: this.#residencies.activeCount, + ...(peer + ? { + peerEndpoint: { + peerId: peer.peerId, + routeHints: peer.listenAddresses, + coordinationRelays: peer.coordinationRelays, + }, + } + : {}), }; } diff --git a/packages/runtime-host/src/server/listener-set.ts b/packages/runtime-host/src/server/listener-set.ts index c26508e581..87f8353b11 100644 --- a/packages/runtime-host/src/server/listener-set.ts +++ b/packages/runtime-host/src/server/listener-set.ts @@ -46,11 +46,13 @@ export interface RuntimeHostPeerListener extends RuntimeHostListener { readonly kind: 'libp2p_direct'; readonly peerId: string; readonly listenAddresses: readonly string[]; + readonly coordinationRelays: readonly string[]; } export interface RuntimeHostPeerListenerDescriptor { readonly peerId: string; readonly listenAddresses: readonly string[]; + readonly coordinationRelays: readonly string[]; } export type RuntimeHostListenerKind = 'local_ipc' | 'websocket' | 'libp2p_direct'; @@ -127,6 +129,17 @@ export function createRuntimeHostListenerSet( additional: readonly RuntimeHostListener[] = [], ): RuntimeHostListenerSet { const listeners = Object.freeze([local, ...additional]); + const peerListeners = Object.freeze( + additional.filter(isRuntimeHostPeerListener).map((listener) => + Object.freeze({ + peerId: listener.peerId, + listenAddresses: Object.freeze([...listener.listenAddresses]), + get coordinationRelays() { + return Object.freeze([...listener.coordinationRelays]); + }, + }), + ), + ); return { listeners, localEndpoint: local.endpoint, @@ -135,14 +148,7 @@ export function createRuntimeHostListenerSet( .filter((listener) => listener.kind === 'websocket') .map((listener) => listener.endpoint), ), - peerListeners: Object.freeze( - additional.filter(isRuntimeHostPeerListener).map((listener) => - Object.freeze({ - peerId: listener.peerId, - listenAddresses: Object.freeze([...listener.listenAddresses]), - }), - ), - ), + peerListeners, closeAdmission: () => settleListeners(listeners, (listener) => listener.closeAdmission()), cleanup: () => settleListeners([...listeners].reverse(), (listener) => listener.cleanup()), }; diff --git a/packages/runtime-host/src/server/peer-listener.ts b/packages/runtime-host/src/server/peer-listener.ts index 750f40bcce..3907669793 100644 --- a/packages/runtime-host/src/server/peer-listener.ts +++ b/packages/runtime-host/src/server/peer-listener.ts @@ -120,6 +120,10 @@ class RuntimeHostPeerListener implements RuntimeHostPeerListenerContract { .catch(captureFailure); } + get coordinationRelays(): readonly string[] { + return this.#client.identity().coordinationRelays; + } + closeAdmission(): Promise { this.#closeAdmissionTask ??= (async () => { this.#admitting = false; diff --git a/packages/ui/src/session-history-list.tsx b/packages/ui/src/session-history-list.tsx index a455f9e74e..51a136574e 100644 --- a/packages/ui/src/session-history-list.tsx +++ b/packages/ui/src/session-history-list.tsx @@ -63,6 +63,7 @@ import { dotForStatus } from './status-vocabulary.js'; import { SessionRenameDialog, type SessionRenameTarget } from './session-rename-dialog.js'; import { CheckboxInput } from '@astryxdesign/core/CheckboxInput'; import { + type SessionRailData, useSessionRailData, useSessionRailRowSelection, useSessionRailSelection, @@ -286,6 +287,7 @@ function SessionListGroups(props: { stale={rail.staleSessionIds?.has(session.id) ?? false} worktree={rail.worktreeSessionIds?.has(session.id) ?? false} meta={rail.sessionMeta?.(session)} + sessionBadge={rail.sessionBadge} onSelectSession={rail.onSelectSession} actions={(session as SessionSummary & { readonly shared?: true }).shared ? undefined @@ -490,6 +492,7 @@ const SessionNavRow = memo(function SessionNavRow(props: { stale: boolean; worktree: boolean; meta?: string; + sessionBadge?: SessionRailData['sessionBadge']; onSelectSession(sessionId: string): void; actions?: SessionRowActions; onStartRename(target: SessionRenameTarget, opener: HTMLElement | null): void; @@ -604,6 +607,11 @@ const SessionNavRow = memo(function SessionNavRow(props: { // the two on hover or keyboard focus. The span is rendered even with // no timestamp so the column exists on every row. + {props.sessionBadge ? ( + + {props.sessionBadge(props.session)} + + ) : null} {props.meta ? ( diff --git a/packages/ui/src/session-rail-context.tsx b/packages/ui/src/session-rail-context.tsx index ca747c830b..17bdcc3f85 100644 --- a/packages/ui/src/session-rail-context.tsx +++ b/packages/ui/src/session-rail-context.tsx @@ -52,6 +52,7 @@ export interface SessionRailData { groups?: ReadonlyArray; groupVariant: SessionViewMode; sessionMeta?(session: SessionSummary): string | undefined; + sessionBadge?(session: SessionSummary): ReactNode; onSelectSession(sessionId: string): void; rowActions?: SessionRowActions; projectActions?: ProjectRowActions;