Files
orca/src/main/runtime/remote-runtime-request-connection.integration.test.ts
T
Brennan Benson bbb3e7e5ee fix(native-chat): mirror multi-line launch drafts into the chat composer (#11253)
* fix(native-chat): mirror multi-line launch drafts into the chat composer

seedNativeChatLaunchDraftForAgentTab rejected any text containing a newline,
so every Linear launch ("Linked Linear issue: X\n<url>") and any GitHub launch
with a typed note was invisible in chat. The rejection existed because the send
path pre-cleared the TUI with a single Ctrl+U, which cannot clear a buffer with
embedded newlines.

Orca injects the draft itself, so when the composer still holds exactly what was
injected the buffer already IS the message: the send becomes the submit key
alone — no clear, no paste, nothing that can concatenate, and multi-line submits
as one turn for free. Only the edited case needs real buffer replacement, and
that now clears every line and verifies against the agent's rendered input line
instead of firing blind.

Measured on real PTYs against Claude Code and codex (both agree exactly):
clearing N logical lines costs 2N-1 Ctrl+U. See src/shared/agent-tui-input-clear.ts
for the law, the sequences that do NOT work, and why an upper bound is safe.

* fix(native-chat): send the mobile clear burst as its own write

Live QA caught the bundled form failing: a multi-line burst prefixed onto the
body in the SAME terminal.send reached the agent as LITERAL Ctrl+U characters,
so the parked draft survived and the message arrived as
draft + 21x \x15 + body. Sending the burst as its own non-submitting write —
the shape the image paste has always used — clears as intended.

The body write's own single-Ctrl+U prefix is dropped once that dedicated clear
ran, for the same reason: a Ctrl+U immediately followed by body text in one
write lands as a literal control character and headed the received message.

Re-verified live end to end: received prompt is exactly the draft, one turn,
zero control characters.

* test(native-chat): invert the multi-line Linear launch-draft mirror expectation

The Linear work-item launch seeds `Linked Linear issue: ENG-42\n<url>\n`.
This test pinned the pre-relaxation rule (multi-line drafts withheld), which
the send path no longer needs now that it submits the TUI buffer in place or
clears every line first — so it asserted the exact behavior the fix removes.

Assert the seeded payload instead of absence, so the test fails if the mirror
regresses to single-line-only.

* fix(native-chat): preserve launch draft send contents

* fix(native-chat): preserve confirmed send queue ordering

* fix(native-chat): preserve send pacing after renderer stalls

* test(native-chat): align activation with multiline draft mirroring

* fix(native-chat): clear launch drafts from any cursor

* fix(native-chat): retire mobile-consumed launch drafts

* test(mobile): stabilize QR capacity boundary fixture
2026-07-30 11:08:56 -07:00

820 lines
28 KiB
TypeScript

import { mkdtempSync, rmSync } from 'node:fs'
import { tmpdir } from 'node:os'
import { join } from 'node:path'
import { describe, expect, it, vi } from 'vitest'
import { getDefaultRepoHookSettings } from '../../shared/constants'
import type { Repo } from '../../shared/types'
import { parsePairingCode } from '../../shared/pairing'
import { RemoteRuntimeRequestConnection } from '../../shared/remote-runtime-request-connection'
import { RemoteRuntimeSharedControlConnection } from '../../shared/remote-runtime-shared-control-connection'
import { subscribeRemoteRuntimeRequest } from '../../shared/remote-runtime-client'
import type {
RuntimeClientEvent,
RuntimeClientEventStreamMessage
} from '../../shared/runtime-client-events'
import type { OrcaRuntimeService } from './orca-runtime'
import { OrcaRuntimeRpcServer } from './runtime-rpc'
import { REMOTE_RUNTIME_SHARED_CONTROL_CAPABILITY } from '../../shared/protocol-version'
const REMOTE_RUNTIME_TEST_TIMEOUT_MS = 15_000
const REMOTE_RUNTIME_REQUEST_TIMEOUT_MS = 5_000
// worktree.create routes through the runtime's clientMutationId idempotency
// wrapper; these stubs run the create straight through (no dedupe).
const passthroughDedupe = <T>(_repo: string, _id: string | undefined, run: () => Promise<T>) =>
run()
describe('remote runtime request connection integration', () => {
it(
'binds encrypted close-intent capability to the real runtime RPC context',
{ timeout: REMOTE_RUNTIME_TEST_TIMEOUT_MS },
async () => {
const userDataPath = mkdtempSync(join(tmpdir(), 'orca-runtime-close-intent-'))
const refuseUnattributedMobileSessionTabClose = vi.fn().mockResolvedValue({
closed: true,
refused: true,
refusalReason: 'missing-intent',
snapshotRepublished: true
})
const closeMobileSessionTab = vi.fn()
const runtime = {
getRuntimeId: () => 'close-intent-runtime-test',
getStartedAt: () => 1,
cleanupSubscriptionsForConnection: () => {},
cancelMobileDictationForConnection: () => {},
onClientDisconnected: () => {},
refuseUnattributedMobileSessionTabClose,
closeMobileSessionTab
} as unknown as OrcaRuntimeService
const server = new OrcaRuntimeRpcServer({
runtime,
userDataPath,
enableWebSocket: true,
wsPort: 0
})
await server.start()
try {
const offer = server.createPairingOffer({ name: 'integration', scope: 'runtime' })
if (!offer.available) {
throw new Error('pairing unavailable')
}
const pairing = parsePairingCode(offer.pairingUrl)
if (!pairing) {
throw new Error('invalid pairing')
}
const connection = new RemoteRuntimeRequestConnection(pairing)
try {
await expect(
connection.request(
'session.tabs.close',
{ worktree: 'id:wt-1', tabId: 'tab-1' },
REMOTE_RUNTIME_REQUEST_TIMEOUT_MS
)
).resolves.toMatchObject({
ok: true,
result: {
refused: true,
refusalReason: 'missing-intent',
snapshotRepublished: true
}
})
expect(refuseUnattributedMobileSessionTabClose).toHaveBeenCalledWith('id:wt-1', 'tab-1')
expect(closeMobileSessionTab).not.toHaveBeenCalled()
} finally {
connection.close()
}
} finally {
await server.stop()
rmSync(userDataPath, { recursive: true, force: true })
}
}
)
it(
'fetches repos through the real E2EE WebSocket runtime',
{ timeout: REMOTE_RUNTIME_TEST_TIMEOUT_MS },
async () => {
const userDataPath = mkdtempSync(join(tmpdir(), 'orca-runtime-request-'))
const repoPath = join(userDataPath, 'repo')
const repos: Repo[] = [
{
id: 'repo-1',
path: repoPath,
displayName: 'repo',
badgeColor: 'blue',
addedAt: 1,
hookSettings: getDefaultRepoHookSettings(),
worktreeBaseRef: 'main',
kind: 'git'
}
]
const runtime = {
getRuntimeId: () => 'fetch-runtime-test',
getStartedAt: () => 1,
cleanupSubscriptionsForConnection: () => {},
cancelMobileDictationForConnection: () => {},
onClientDisconnected: () => {},
listRepos: () => repos
} as unknown as OrcaRuntimeService
const server = new OrcaRuntimeRpcServer({
runtime,
userDataPath,
enableWebSocket: true,
wsPort: 0
})
await server.start()
try {
const offer = server.createPairingOffer({ name: 'integration', scope: 'runtime' })
if (!offer.available) {
throw new Error('pairing unavailable')
}
const pairing = parsePairingCode(offer.pairingUrl)
if (!pairing) {
throw new Error('invalid pairing')
}
const connection = new RemoteRuntimeRequestConnection(pairing)
try {
await expect(
connection.request('repo.list', undefined, REMOTE_RUNTIME_REQUEST_TIMEOUT_MS)
).resolves.toMatchObject({
ok: true,
result: { repos }
})
} finally {
connection.close()
}
} finally {
await server.stop()
rmSync(userDataPath, { recursive: true, force: true })
}
}
)
it(
'streams server worktree changes to another remote client',
{ timeout: REMOTE_RUNTIME_TEST_TIMEOUT_MS },
async () => {
const userDataPath = mkdtempSync(join(tmpdir(), 'orca-runtime-request-events-'))
const repoPath = join(userDataPath, 'repo')
const repo: Repo = {
id: 'repo-1',
path: repoPath,
displayName: 'repo',
badgeColor: 'blue',
addedAt: 1,
hookSettings: getDefaultRepoHookSettings(),
worktreeBaseRef: 'main',
kind: 'git'
}
const worktrees: unknown[] = [
{
id: 'repo-1::main',
repoId: repo.id,
path: repoPath,
branch: 'main',
displayName: 'repo',
isMainWorktree: true
}
]
const clientEventListeners = new Set<(event: RuntimeClientEvent) => void>()
const subscriptionCleanups = new Map<string, () => void>()
const runtime = {
getRuntimeId: () => 'events-runtime-test',
getStartedAt: () => 1,
cleanupSubscriptionsForConnection: (connectionId: string) => {
for (const [id, cleanup] of subscriptionCleanups) {
if (id.includes(connectionId)) {
cleanup()
subscriptionCleanups.delete(id)
}
}
},
registerSubscriptionCleanup: (id: string, cleanup: () => void) => {
subscriptionCleanups.set(id, cleanup)
},
cleanupSubscription: (id: string) => {
subscriptionCleanups.get(id)?.()
subscriptionCleanups.delete(id)
},
cancelMobileDictationForConnection: () => {},
onClientDisconnected: () => {},
showRepo: (selector: string) => {
if (selector !== repo.id && selector !== `id:${repo.id}`) {
throw new Error('repo_not_found')
}
return repo
},
onClientEvent: (listener: (event: RuntimeClientEvent) => void) => {
clientEventListeners.add(listener)
return () => clientEventListeners.delete(listener)
},
listDetectedManagedWorktrees: () => ({
repoId: repo.id,
authoritative: true,
source: 'git',
worktrees
}),
dedupeWorktreeCreate: passthroughDedupe,
createManagedWorktree: ({ name }: { name?: string }) => {
const worktree = {
id: `repo-1::${name || 'created'}`,
repoId: repo.id,
path: join(userDataPath, name || 'created'),
branch: name || 'created',
displayName: name || 'created',
isMainWorktree: false
}
worktrees.push(worktree)
for (const listener of clientEventListeners) {
listener({ type: 'worktreesChanged', repoId: repo.id })
}
return { worktree }
}
} as unknown as OrcaRuntimeService
const server = new OrcaRuntimeRpcServer({
runtime,
userDataPath,
enableWebSocket: true,
wsPort: 0
})
await server.start()
try {
const offer = server.createPairingOffer({ name: 'integration', scope: 'runtime' })
if (!offer.available) {
throw new Error('pairing unavailable')
}
const pairing = parsePairingCode(offer.pairingUrl)
if (!pairing) {
throw new Error('invalid pairing')
}
const events: RuntimeClientEventStreamMessage[] = []
const subscription = await subscribeRemoteRuntimeRequest<RuntimeClientEventStreamMessage>(
pairing,
'runtime.clientEvents.subscribe',
undefined,
REMOTE_RUNTIME_REQUEST_TIMEOUT_MS,
{
onResponse: (response) => {
if (response.ok) {
events.push(response.result)
}
},
onError: (error) => {
throw error
}
}
)
const desktop = new RemoteRuntimeRequestConnection(pairing)
const mobile = new RemoteRuntimeRequestConnection(pairing)
try {
await waitFor(() => events.some((event) => event.type === 'ready'))
await expect(
desktop.request<{ worktrees: unknown[] }>(
'worktree.detectedList',
{ repo: repo.id },
REMOTE_RUNTIME_REQUEST_TIMEOUT_MS
)
).resolves.toMatchObject({
ok: true,
result: { worktrees: [{ id: 'repo-1::main' }] }
})
await expect(
mobile.request(
'worktree.create',
{ repo: repo.id, name: 'mobile-created' },
REMOTE_RUNTIME_REQUEST_TIMEOUT_MS
)
).resolves.toMatchObject({
ok: true,
result: { worktree: { id: 'repo-1::mobile-created' } }
})
await waitFor(() =>
events.some((event) => event.type === 'worktreesChanged' && event.repoId === repo.id)
)
await expect(
desktop.request<{ worktrees: unknown[] }>(
'worktree.detectedList',
{ repo: repo.id },
REMOTE_RUNTIME_REQUEST_TIMEOUT_MS
)
).resolves.toMatchObject({
ok: true,
result: {
worktrees: [{ id: 'repo-1::main' }, { id: 'repo-1::mobile-created' }]
}
})
} finally {
subscription.close()
desktop.close()
mobile.close()
}
} finally {
await server.stop()
rmSync(userDataPath, { recursive: true, force: true })
}
}
)
it(
'delivers one host-authoritative sleep and ordered disposition to two remote clients',
{ timeout: REMOTE_RUNTIME_TEST_TIMEOUT_MS },
async () => {
const userDataPath = mkdtempSync(join(tmpdir(), 'orca-runtime-remote-sleep-'))
const clientEventListeners = new Set<(event: RuntimeClientEvent) => void>()
const subscriptionCleanups = new Map<string, () => void>()
const worktreeId = 'repo-1::C:\\repo\\feature'
const ptyId = `${worktreeId}@@pty-1`
let sleepSnapshot: RuntimeClientEvent[] = []
const launchDraftResolutionSnapshot: RuntimeClientEvent[] = [
{
type: 'nativeChatLaunchDraftResolved',
tabId: 'tab-1',
text: 'seed',
createdAt: 7
}
]
const emit = (event: RuntimeClientEvent): void => {
for (const listener of clientEventListeners) {
listener(event)
}
}
const runtime = {
getRuntimeId: () => 'remote-sleep-runtime-test',
getStartedAt: () => 1,
cleanupSubscriptionsForConnection: (connectionId: string) => {
for (const [id, cleanup] of subscriptionCleanups) {
if (id.includes(connectionId)) {
cleanup()
subscriptionCleanups.delete(id)
}
}
},
registerSubscriptionCleanup: (id: string, cleanup: () => void) => {
subscriptionCleanups.set(id, cleanup)
},
cleanupSubscription: (id: string) => {
subscriptionCleanups.get(id)?.()
subscriptionCleanups.delete(id)
},
cancelMobileDictationForConnection: () => {},
onClientDisconnected: () => {},
onClientEvent: (listener: (event: RuntimeClientEvent) => void) => {
clientEventListeners.add(listener)
return () => clientEventListeners.delete(listener)
},
getTerminalSleepClientEventSnapshot: () => sleepSnapshot,
getNativeChatLaunchDraftResolutionClientEventSnapshot: () => launchDraftResolutionSnapshot,
sleepTerminalsForWorktree: async () => {
emit({
type: 'worktreeTerminalSleepState',
worktreeId,
generation: 1,
phase: 'started',
ptyIds: [ptyId],
terminalHandles: ['terminal-handle-1']
})
emit({
type: 'worktreeTerminalSleepState',
worktreeId,
generation: 1,
phase: 'committed',
ptyIds: [ptyId],
terminalHandles: ['terminal-handle-1']
})
sleepSnapshot = [
{
type: 'worktreeTerminalSleepState',
worktreeId,
generation: 1,
phase: 'committed',
ptyIds: [ptyId],
terminalHandles: ['terminal-handle-1']
}
]
return {
stopped: 1,
stoppedPtyIds: [ptyId],
livePtyIds: [ptyId],
postStopVerified: true
}
}
} as unknown as OrcaRuntimeService
const server = new OrcaRuntimeRpcServer({
runtime,
userDataPath,
enableWebSocket: true,
wsPort: 0
})
await server.start()
try {
const offer = server.createPairingOffer({ name: 'remote-sleep', scope: 'runtime' })
if (!offer.available) {
throw new Error('pairing unavailable')
}
const pairing = parsePairingCode(offer.pairingUrl)
if (!pairing) {
throw new Error('invalid pairing')
}
const clientEvents: RuntimeClientEventStreamMessage[][] = [[], []]
const subscriptions = await Promise.all(
clientEvents.map((events) =>
subscribeRemoteRuntimeRequest<RuntimeClientEventStreamMessage>(
pairing,
'runtime.clientEvents.subscribe',
undefined,
REMOTE_RUNTIME_REQUEST_TIMEOUT_MS,
{
onResponse: (response) => {
if (response.ok) {
events.push(response.result)
}
},
onError: (error) => {
throw error
}
}
)
)
)
const requester = new RemoteRuntimeRequestConnection(pairing)
try {
await waitFor(() =>
clientEvents.every((events) => events.some((e) => e.type === 'ready'))
)
for (const events of clientEvents) {
expect(events).toContainEqual(launchDraftResolutionSnapshot[0])
}
await expect(
requester.request(
'terminal.sleep',
{ worktree: `id:${worktreeId}` },
REMOTE_RUNTIME_REQUEST_TIMEOUT_MS
)
).resolves.toMatchObject({
ok: true,
result: { stopped: 1, postStopVerified: true }
})
await waitFor(() =>
clientEvents.every(
(events) =>
events.filter((event) => event.type === 'worktreeTerminalSleepState').length === 2
)
)
for (const events of clientEvents) {
expect(
events
.filter((event) => event.type === 'worktreeTerminalSleepState')
.map((event) => event.phase)
).toEqual(['started', 'committed'])
}
const reconnectedEvents: RuntimeClientEventStreamMessage[] = []
const reconnected = await subscribeRemoteRuntimeRequest<RuntimeClientEventStreamMessage>(
pairing,
'runtime.clientEvents.subscribe',
undefined,
REMOTE_RUNTIME_REQUEST_TIMEOUT_MS,
{
onResponse: (response) => {
if (response.ok) {
reconnectedEvents.push(response.result)
}
},
onError: (error) => {
throw error
}
}
)
try {
await waitFor(() => reconnectedEvents.some((event) => event.type === 'ready'))
expect(
reconnectedEvents
.filter((event) => event.type === 'worktreeTerminalSleepState')
.map((event) => event.phase)
).toEqual(['committed'])
expect(reconnectedEvents).toContainEqual(launchDraftResolutionSnapshot[0])
sleepSnapshot = []
emit({
type: 'worktreeTerminalSleepState',
worktreeId,
generation: 1,
phase: 'woken',
ptyIds: [ptyId],
terminalHandles: ['terminal-handle-1']
})
await waitFor(() =>
reconnectedEvents.some(
(event) => event.type === 'worktreeTerminalSleepState' && event.phase === 'woken'
)
)
} finally {
reconnected.close()
}
} finally {
requester.close()
for (const subscription of subscriptions) {
subscription.close()
}
}
} finally {
await server.stop()
rmSync(userDataPath, { recursive: true, force: true })
}
}
)
it(
'multiplexes shared-control calls and passive subscriptions through the real runtime',
{ timeout: REMOTE_RUNTIME_TEST_TIMEOUT_MS },
async () => {
const userDataPath = mkdtempSync(join(tmpdir(), 'orca-runtime-shared-control-'))
const repoPath = join(userDataPath, 'repo')
const repo: Repo = {
id: 'repo-1',
path: repoPath,
displayName: 'repo',
badgeColor: 'blue',
addedAt: 1,
hookSettings: getDefaultRepoHookSettings(),
worktreeBaseRef: 'main',
kind: 'git'
}
const worktrees: unknown[] = [
{
id: 'repo-1::main',
repoId: repo.id,
path: repoPath,
branch: 'main',
displayName: 'repo',
isMainWorktree: true
}
]
const clientEventListeners = new Set<(event: RuntimeClientEvent) => void>()
const accountsListeners = new Set<(snapshot: unknown) => void>()
const notificationListeners = new Set<(event: unknown) => void>()
const sessionTabListeners = new Set<(snapshot: unknown) => void>()
const subscriptionCleanups = new Map<string, () => void>()
const sessionTabSnapshot = {
worktree: 'wt-1',
publicationEpoch: 'epoch-1',
snapshotVersion: 1,
activeGroupId: null,
activeTabId: null,
activeTabType: null,
tabs: []
}
const runtime = {
getRuntimeId: () => 'shared-runtime-test',
getStartedAt: () => 1,
getStatus: () => ({
runtimeId: 'shared-runtime-test',
startedAt: 1,
version: '1.0.0',
protocolVersion: 1,
minCompatibleDesktopVersion: '1.0.0',
minCompatibleMobileVersion: '1.0.0',
capabilities: [REMOTE_RUNTIME_SHARED_CONTROL_CAPABILITY]
}),
cleanupSubscriptionsForConnection: (connectionId: string) => {
for (const [id, cleanup] of Array.from(subscriptionCleanups)) {
if (id.includes(connectionId)) {
cleanup()
subscriptionCleanups.delete(id)
}
}
},
registerSubscriptionCleanup: (id: string, cleanup: () => void) => {
subscriptionCleanups.set(id, cleanup)
},
cleanupSubscription: (id: string) => {
subscriptionCleanups.get(id)?.()
subscriptionCleanups.delete(id)
},
cleanupSubscriptionsByPrefix: (prefix: string) => {
for (const [id, cleanup] of Array.from(subscriptionCleanups)) {
if (id.startsWith(prefix)) {
cleanup()
subscriptionCleanups.delete(id)
}
}
},
cancelMobileDictationForConnection: () => {},
onClientDisconnected: () => {},
onClientEvent: (listener: (event: RuntimeClientEvent) => void) => {
clientEventListeners.add(listener)
return () => clientEventListeners.delete(listener)
},
getAccountsSnapshot: () => ({ claude: null, codex: null }),
refreshAccountsForMobile: async () => {
for (const listener of accountsListeners) {
listener({ claude: null, codex: null })
}
},
refreshAccountsForMobileSubscriber: async () => {
for (const listener of accountsListeners) {
listener({ claude: null, codex: null })
}
},
onAccountsChanged: (listener: (snapshot: unknown) => void) => {
accountsListeners.add(listener)
return () => accountsListeners.delete(listener)
},
onNotificationDispatched: (listener: (event: unknown) => void) => {
notificationListeners.add(listener)
return () => notificationListeners.delete(listener)
},
listMobileSessionTabs: () => sessionTabSnapshot,
listAllMobileSessionTabs: () => [sessionTabSnapshot],
onMobileSessionTabsChanged: (listener: (snapshot: unknown) => void) => {
sessionTabListeners.add(listener)
return () => sessionTabListeners.delete(listener)
},
watchFileExplorer: async () => () => {},
listRepos: () => [repo],
showRepo: (selector: string) => {
if (selector !== repo.id && selector !== `id:${repo.id}`) {
throw new Error('repo_not_found')
}
return repo
},
listDetectedManagedWorktrees: () => ({
repoId: repo.id,
authoritative: true,
source: 'git',
worktrees
}),
dedupeWorktreeCreate: passthroughDedupe,
createManagedWorktree: ({ name }: { name?: string }) => {
const worktree = {
id: `repo-1::${name || 'created'}`,
repoId: repo.id,
path: join(userDataPath, name || 'created'),
branch: name || 'created',
displayName: name || 'created',
isMainWorktree: false
}
worktrees.push(worktree)
for (const listener of clientEventListeners) {
listener({ type: 'worktreesChanged', repoId: repo.id })
}
return { worktree }
}
} as unknown as OrcaRuntimeService
const server = new OrcaRuntimeRpcServer({
runtime,
userDataPath,
enableWebSocket: true,
wsPort: 0
})
await server.start()
try {
const offer = server.createPairingOffer({ name: 'integration', scope: 'runtime' })
if (!offer.available) {
throw new Error('pairing unavailable')
}
const pairing = parsePairingCode(offer.pairingUrl)
if (!pairing) {
throw new Error('invalid pairing')
}
const events: RuntimeClientEventStreamMessage[] = []
const shared = new RemoteRuntimeSharedControlConnection(pairing)
const subscription = await shared.subscribe<RuntimeClientEventStreamMessage>(
'runtime.clientEvents.subscribe',
undefined,
REMOTE_RUNTIME_REQUEST_TIMEOUT_MS,
{
onResponse: (response) => {
if (response.ok) {
events.push(response.result)
}
},
onError: (error) => {
throw error
}
}
)
try {
await waitFor(() => events.some((event) => event.type === 'ready'))
await expect(
shared.request('status.get', undefined, REMOTE_RUNTIME_REQUEST_TIMEOUT_MS)
).resolves.toMatchObject({
ok: true,
result: { capabilities: [REMOTE_RUNTIME_SHARED_CONTROL_CAPABILITY] }
})
await expect(
shared.request('repo.list', undefined, REMOTE_RUNTIME_REQUEST_TIMEOUT_MS)
).resolves.toMatchObject({
ok: true,
result: { repos: [repo] }
})
await expect(
shared.request(
'worktree.create',
{ repo: repo.id, name: 'shared-created' },
REMOTE_RUNTIME_REQUEST_TIMEOUT_MS
)
).resolves.toMatchObject({
ok: true,
result: { worktree: { id: 'repo-1::shared-created' } }
})
await waitFor(() =>
events.some((event) => event.type === 'worktreesChanged' && event.repoId === repo.id)
)
const mixedEvents: unknown[] = []
const mixedMethods = [
['runtime.clientEvents.subscribe', undefined],
['session.tabs.subscribe', { worktree: 'id:wt-1' }],
['accounts.subscribe', undefined],
['notifications.subscribe', undefined],
['files.watch', { worktree: 'id:wt-1' }]
] as const
const mixedSubscriptions = await Promise.all(
Array.from({ length: 30 }, (_value, index) => {
const [method, params] = mixedMethods[index % mixedMethods.length]!
return shared.subscribe(method, params, REMOTE_RUNTIME_REQUEST_TIMEOUT_MS, {
onResponse: (response) => {
if (response.ok) {
mixedEvents.push(response.result)
}
},
onError: (error) => {
throw error
}
})
})
)
await waitFor(
() => subscriptionCleanups.size >= mixedSubscriptions.length + 1,
5000,
() => `cleanup count ${subscriptionCleanups.size}, event count ${mixedEvents.length}`
)
await waitFor(
() => mixedEvents.length > 0,
REMOTE_RUNTIME_REQUEST_TIMEOUT_MS,
() => `cleanup count ${subscriptionCleanups.size}, event count ${mixedEvents.length}`
)
expect(server.getMobileSocketWiring()?.connectionCount).toBe(1)
for (const mixed of mixedSubscriptions) {
mixed.close()
}
const extraSubscriptions = await Promise.all(
Array.from({ length: 30 }, () =>
shared.subscribe<RuntimeClientEventStreamMessage>(
'runtime.clientEvents.subscribe',
undefined,
REMOTE_RUNTIME_REQUEST_TIMEOUT_MS,
{
onResponse: () => {},
onError: (error) => {
throw error
}
}
)
)
)
expect(server.getMobileSocketWiring()?.connectionCount).toBe(1)
for (const extra of extraSubscriptions) {
extra.close()
}
} finally {
subscription.close()
shared.close()
}
} finally {
await server.stop()
rmSync(userDataPath, { recursive: true, force: true })
}
}
)
})
async function waitFor(
predicate: () => boolean,
timeoutMs = 1000,
describeTimeout?: () => string
): Promise<void> {
const start = Date.now()
while (Date.now() - start < timeoutMs) {
if (predicate()) {
return
}
await new Promise((resolve) => setTimeout(resolve, 10))
}
throw new Error(
`Timed out waiting for condition${describeTimeout ? `: ${describeTimeout()}` : ''}`
)
}