mirror of
https://github.com/stablyai/orca.git
synced 2026-09-29 16:02:50 +00:00
fix(mobile): let a slow capability handshake still reach connected
The mobile capability update is an advisory whose result is discarded, yet an unanswered one was fatal while an explicit rejection was tolerated. A 5s timeout on the direct client force-closed the socket, and on the relay path it failed `confirmResume` before `connected` was ever published, so a consistently slow link redialled forever. Both paths now share one helper that settles every ambiguous outcome (timeout, mid-flight drop) like a rejection and rejects only when the frame never reached the wire — the one case nothing else recovers from, since the socket's own desync force-close is gated on already being connected. The generation guard still keeps a replaced session from connecting. Retained structured-session operation ids were capped at 128 with oldest-first eviction, but every retained id belongs to a send whose outcome is unknown, so eviction turned a user's retry into a second message on the host. Bound the map by expiry against the id's own embedded timestamp instead, mirroring the host's operation ledger, so no id is released while the host would still honour it. Also give the mobile CI install the root install's lockfile drift guard (mobile's lockfile carries patchedDependencies a silent rewrite would drop), gate mobile_dependencies on should_run, and key the pnpm store cache on both lockfiles.
This commit is contained in:
@@ -39,6 +39,9 @@ runs:
|
||||
with:
|
||||
install: false
|
||||
|
||||
# Why both lockfiles: setup-node keys the pnpm store on the root lockfile alone, so
|
||||
# jobs that also install mobile restored a store with none of the React Native tree
|
||||
# in it and re-downloaded the lot on every run.
|
||||
- name: Setup Node.js
|
||||
id: default-node
|
||||
if: inputs.node-version == ''
|
||||
@@ -46,6 +49,9 @@ runs:
|
||||
with:
|
||||
node-version-file: package.json
|
||||
cache: pnpm
|
||||
cache-dependency-path: |
|
||||
pnpm-lock.yaml
|
||||
mobile/pnpm-lock.yaml
|
||||
|
||||
- name: Setup requested Node.js
|
||||
id: requested-node
|
||||
@@ -54,6 +60,9 @@ runs:
|
||||
with:
|
||||
node-version: ${{ inputs.node-version }}
|
||||
cache: pnpm
|
||||
cache-dependency-path: |
|
||||
pnpm-lock.yaml
|
||||
mobile/pnpm-lock.yaml
|
||||
|
||||
- name: Validate native runtime
|
||||
shell: bash
|
||||
|
||||
@@ -100,10 +100,20 @@ jobs:
|
||||
# resolves types from mobile/node_modules. Mobile is a separate pnpm project,
|
||||
# so the root install above leaves it empty and every mobile type degrades to
|
||||
# an `error` type — reported as phantom findings against the changed lines.
|
||||
# Why no --ignore-scripts, unlike the root install: mobile's postinstall generates
|
||||
# the gitignored terminal/mermaid webview engine modules that tracked source imports,
|
||||
# and skipping it degrades those very types the step exists to resolve. The drift
|
||||
# guard mirrors the root install so a stale mobile lockfile fails by name — mobile's
|
||||
# lockfile carries patchedDependencies that a silent rewrite would drop.
|
||||
- name: Install mobile dependencies
|
||||
if: needs.code_paths.outputs.mobile_dependencies == 'true'
|
||||
working-directory: mobile
|
||||
run: pnpm install --frozen-lockfile
|
||||
run: |
|
||||
pnpm install --frozen-lockfile
|
||||
if [ "$(git -C "$GITHUB_WORKSPACE" rev-parse --is-inside-work-tree 2>/dev/null)" = true ]; then
|
||||
git -C "$GITHUB_WORKSPACE" diff --exit-code -- \
|
||||
mobile/package.json mobile/pnpm-lock.yaml mobile/pnpm-workspace.yaml
|
||||
fi
|
||||
|
||||
- name: Enforce changed-code quality
|
||||
run: pnpm run check:code-quality:changed -- "${{ github.event.pull_request.base.sha }}"
|
||||
|
||||
@@ -282,7 +282,7 @@ export function classifyPrJobs(changedFiles) {
|
||||
return {
|
||||
should_run: shouldRun,
|
||||
native_cache_changed: shouldRun && (emptyDiff || changedFiles.some(isNativeCacheInputPath)),
|
||||
mobile_dependencies: needsMobileDependencies(changedFiles),
|
||||
mobile_dependencies: shouldRun && needsMobileDependencies(changedFiles),
|
||||
...jobs
|
||||
}
|
||||
}
|
||||
|
||||
@@ -316,7 +316,11 @@ describe('per-job path classification', () => {
|
||||
expect(
|
||||
classifyPrJobs(['src/main/index.ts', 'mobile/src/session/a.test.ts']).mobile_dependencies
|
||||
).toBe(true)
|
||||
expect(classifyPrJobs(['mobile/package.json']).mobile_dependencies).toBe(true)
|
||||
// Why false: a mobile-only diff skips every desktop job, so the install step's own
|
||||
// job never runs and claiming the install is needed contradicts should_run.
|
||||
expect(classifyPrJobs(['mobile/package.json']).mobile_dependencies).toBe(false)
|
||||
expect(classifyPrJobs(['mobile/package.json']).should_run).toBe(false)
|
||||
expect(classifyPrJobs(['README.md', 'mobile/src/a.ts']).mobile_dependencies).toBe(false)
|
||||
})
|
||||
|
||||
it('keeps unit-test-only diffs out of packaging', () => {
|
||||
|
||||
@@ -1,3 +1,7 @@
|
||||
import {
|
||||
AGENT_SESSION_MAX_NEW_OPERATION_AGE_MS,
|
||||
parseAgentSessionOperationTimestamp
|
||||
} from '../../../src/shared/agent-session-host-authority'
|
||||
import type { AgentSessionMutationResult } from '../../../src/shared/agent-session-wire'
|
||||
import {
|
||||
createStructuredAgentSessionOperationId,
|
||||
@@ -57,21 +61,28 @@ export function structuredSessionOperationId(): string {
|
||||
return createStructuredAgentSessionOperationId(randomUuid)
|
||||
}
|
||||
|
||||
const MAX_RETAINED_OPERATION_IDS = 128
|
||||
|
||||
/**
|
||||
* Bounded by expiry, never by count: every retained id belongs to a send whose outcome is still
|
||||
* unknown, so dropping one turns the user's retry into a second message on the host. Only an id
|
||||
* the host would already refuse — unparseable, or past the window in which it can be admitted —
|
||||
* is safe to release, which matches the host's own tombstone retention.
|
||||
*/
|
||||
export function retainStructuredSessionOperationId(
|
||||
operationIds: Map<string, string>,
|
||||
key: string,
|
||||
operationId = structuredSessionOperationId()
|
||||
operationId = structuredSessionOperationId(),
|
||||
now: number = Date.now()
|
||||
): string {
|
||||
operationIds.delete(key)
|
||||
operationIds.set(key, operationId)
|
||||
while (operationIds.size > MAX_RETAINED_OPERATION_IDS) {
|
||||
const oldest = operationIds.keys().next().value
|
||||
if (oldest === undefined) {
|
||||
break
|
||||
for (const [retainedKey, retainedId] of operationIds) {
|
||||
if (retainedKey === key) {
|
||||
continue
|
||||
}
|
||||
const timestamp = parseAgentSessionOperationTimestamp(retainedId)
|
||||
if (timestamp === null || now - timestamp > AGENT_SESSION_MAX_NEW_OPERATION_AGE_MS) {
|
||||
operationIds.delete(retainedKey)
|
||||
}
|
||||
operationIds.delete(oldest)
|
||||
}
|
||||
return operationId
|
||||
}
|
||||
|
||||
@@ -1,15 +1,58 @@
|
||||
import { describe, expect, it } from 'vitest'
|
||||
import { AGENT_SESSION_MAX_NEW_OPERATION_AGE_MS } from '../../../src/shared/agent-session-host-authority'
|
||||
import { retainStructuredSessionOperationId } from './mobile-structured-agent-session-rpc'
|
||||
|
||||
const NOW = 1_900_000_000_000
|
||||
|
||||
function operationIdAt(timestamp: number, entropy: string): string {
|
||||
return `${timestamp}-${entropy.repeat(32).slice(0, 32)}`
|
||||
}
|
||||
|
||||
describe('structured session operation retention', () => {
|
||||
it('evicts the oldest ambiguous operation id at the retention bound', () => {
|
||||
it('keeps every unconfirmed operation id past the old 128-entry cap', () => {
|
||||
const operationIds = new Map<string, string>()
|
||||
for (let index = 0; index < 129; index += 1) {
|
||||
retainStructuredSessionOperationId(operationIds, `request-${index}`, `operation-${index}`)
|
||||
for (let index = 0; index < 400; index += 1) {
|
||||
retainStructuredSessionOperationId(
|
||||
operationIds,
|
||||
`request-${index}`,
|
||||
operationIdAt(NOW, 'a'),
|
||||
NOW
|
||||
)
|
||||
}
|
||||
|
||||
expect(operationIds).toHaveLength(128)
|
||||
expect(operationIds.has('request-0')).toBe(false)
|
||||
expect(operationIds.get('request-128')).toBe('operation-128')
|
||||
expect(operationIds.size).toBe(400)
|
||||
// Why: the first send is exactly the one a retry would duplicate if it were evicted.
|
||||
expect(operationIds.get('request-0')).toBe(operationIdAt(NOW, 'a'))
|
||||
})
|
||||
|
||||
it('releases only ids the host would already refuse as expired', () => {
|
||||
const operationIds = new Map<string, string>()
|
||||
const expired = operationIdAt(NOW - AGENT_SESSION_MAX_NEW_OPERATION_AGE_MS - 1, 'b')
|
||||
const admissible = operationIdAt(NOW - AGENT_SESSION_MAX_NEW_OPERATION_AGE_MS, 'c')
|
||||
retainStructuredSessionOperationId(operationIds, 'stale', expired, NOW)
|
||||
retainStructuredSessionOperationId(operationIds, 'live', admissible, NOW)
|
||||
|
||||
retainStructuredSessionOperationId(operationIds, 'fresh', operationIdAt(NOW, 'd'), NOW)
|
||||
|
||||
expect(operationIds.has('stale')).toBe(false)
|
||||
expect(operationIds.get('live')).toBe(admissible)
|
||||
expect(operationIds.get('fresh')).toBe(operationIdAt(NOW, 'd'))
|
||||
})
|
||||
|
||||
it('drops ids the host could never admit and re-keys a repeated send', () => {
|
||||
const operationIds = new Map<string, string>()
|
||||
retainStructuredSessionOperationId(operationIds, 'unparseable', 'not-an-operation-id', NOW)
|
||||
const reused = retainStructuredSessionOperationId(
|
||||
operationIds,
|
||||
'send',
|
||||
operationIdAt(NOW, 'e'),
|
||||
NOW
|
||||
)
|
||||
|
||||
// A retry of the same send reuses the retained id rather than minting a duplicate.
|
||||
expect(
|
||||
retainStructuredSessionOperationId(operationIds, 'send', operationIds.get('send'), NOW)
|
||||
).toBe(reused)
|
||||
expect(operationIds.has('unparseable')).toBe(false)
|
||||
})
|
||||
})
|
||||
|
||||
@@ -54,7 +54,7 @@ function openSession() {
|
||||
})
|
||||
}
|
||||
|
||||
async function authenticateSession(capabilitySupported = true) {
|
||||
async function confirmResume() {
|
||||
const session = openSession()
|
||||
fakes.linkOptions!.onHello({
|
||||
type: 'relay-hello',
|
||||
@@ -99,6 +99,11 @@ async function authenticateSession(capabilitySupported = true) {
|
||||
deviceToken: string
|
||||
params: { clientCapabilities?: string[] }
|
||||
}
|
||||
return { session, confirmationRequest: request, capabilityRequest }
|
||||
}
|
||||
|
||||
async function authenticateSession(capabilitySupported = true) {
|
||||
const { session, confirmationRequest, capabilityRequest } = await confirmResume()
|
||||
expect(session.getState()).toBe('handshaking')
|
||||
fakes.linkOptions!.onText(
|
||||
JSON.stringify(
|
||||
@@ -119,7 +124,7 @@ async function authenticateSession(capabilitySupported = true) {
|
||||
)
|
||||
await vi.waitFor(() => expect(session.getState()).toBe('connected'))
|
||||
fakes.sendText.mockClear()
|
||||
return { session, confirmationRequest: request, capabilityRequest }
|
||||
return { session, confirmationRequest, capabilityRequest }
|
||||
}
|
||||
|
||||
describe('mobile relay RPC session', () => {
|
||||
@@ -162,6 +167,15 @@ describe('mobile relay RPC session', () => {
|
||||
expect(session.getFailure()).toBeNull()
|
||||
})
|
||||
|
||||
it('connects when the relay never answers capability negotiation', async () => {
|
||||
const { session } = await confirmResume()
|
||||
|
||||
// Why: the advisory's own deadline used to fail confirmResume, so a link too slow to
|
||||
// answer within the request timeout never published 'connected' — it just redialled.
|
||||
await vi.waitFor(() => expect(session.getState()).toBe('connected'), { timeout: 5_000 })
|
||||
expect(session.getFailure()).toBeNull()
|
||||
})
|
||||
|
||||
it('rejects a mismatched outer credential version and closes the physical link', () => {
|
||||
const session = openSession()
|
||||
fakes.linkOptions!.onHello({
|
||||
|
||||
@@ -10,7 +10,7 @@ import { markRpcDeliveryUnknown } from './rpc-delivery-ambiguity'
|
||||
import { openRpcRequestBudget, resolvePostConnectRequestTimeout } from './rpc-request-budget'
|
||||
import { isRpcResponse } from './rpc-response-shape'
|
||||
import { RpcSessionLivenessWatchdog } from './rpc-session-liveness-watchdog'
|
||||
import { requestMobileRuntimeCapabilities } from './mobile-runtime-capability-negotiation'
|
||||
import { settleMobileRuntimeCapabilities } from './mobile-runtime-capability-negotiation'
|
||||
import type { RpcClient } from './rpc-client'
|
||||
import type { ConnectionLogSink, ConnectionState, RpcResponse } from './types'
|
||||
|
||||
@@ -183,7 +183,8 @@ export function connectMobileRelayRpcSession(args: {
|
||||
resumeConfirmation = result.resumeConfirmation
|
||||
resumeExpiresAt = result.resumeConfirmation.resumeExpiresAt
|
||||
lastConnectedAt = Date.now()
|
||||
await requestMobileRuntimeCapabilities((method, params) =>
|
||||
// Why: an unanswered advisory must not keep a slow relay from ever reaching connected.
|
||||
await settleMobileRuntimeCapabilities((method, params) =>
|
||||
sendRpc(method, params, requestTimeoutMs, true)
|
||||
)
|
||||
livenessWatchdog.start(livenessIdentity)
|
||||
|
||||
@@ -0,0 +1,69 @@
|
||||
import { describe, expect, it, vi } from 'vitest'
|
||||
import { negotiateMobileRuntimeCapabilities } from './mobile-runtime-capability-negotiation'
|
||||
import { markRpcDeliveryUnknown } from './rpc-delivery-ambiguity'
|
||||
import type { RpcResponse } from './types'
|
||||
|
||||
function negotiate(args: { reject: unknown; current?: boolean }): {
|
||||
onReady: ReturnType<typeof vi.fn>
|
||||
onFailure: ReturnType<typeof vi.fn>
|
||||
} {
|
||||
const onReady = vi.fn()
|
||||
const onFailure = vi.fn()
|
||||
negotiateMobileRuntimeCapabilities({
|
||||
sendRequest: () => Promise.reject(args.reject),
|
||||
current: () => args.current ?? true,
|
||||
onReady,
|
||||
onFailure
|
||||
})
|
||||
return { onReady, onFailure }
|
||||
}
|
||||
|
||||
describe('mobile runtime capability negotiation', () => {
|
||||
it('proceeds when the host never answers, so a slow link still reaches connected', async () => {
|
||||
const timedOut = markRpcDeliveryUnknown(
|
||||
new Error('Request timed out: runtime.clientCapabilities.update')
|
||||
)
|
||||
const { onReady, onFailure } = negotiate({ reject: timedOut })
|
||||
|
||||
await vi.waitFor(() => expect(onReady).toHaveBeenCalledTimes(1))
|
||||
expect(onFailure).not.toHaveBeenCalled()
|
||||
})
|
||||
|
||||
it('proceeds when the socket drops the request mid-flight', async () => {
|
||||
const interrupted = markRpcDeliveryUnknown(new Error('Connection interrupted'))
|
||||
const { onReady, onFailure } = negotiate({ reject: interrupted })
|
||||
|
||||
await vi.waitFor(() => expect(onReady).toHaveBeenCalledTimes(1))
|
||||
expect(onFailure).not.toHaveBeenCalled()
|
||||
})
|
||||
|
||||
it('fails a socket that could not put the advisory on the wire', async () => {
|
||||
const { onReady, onFailure } = negotiate({ reject: new Error('Connection interrupted') })
|
||||
|
||||
await vi.waitFor(() => expect(onFailure).toHaveBeenCalledTimes(1))
|
||||
expect(onReady).not.toHaveBeenCalled()
|
||||
})
|
||||
|
||||
it('leaves a replaced session alone on an unanswered request', async () => {
|
||||
const timedOut = markRpcDeliveryUnknown(new Error('Request timed out'))
|
||||
const { onReady, onFailure } = negotiate({ reject: timedOut, current: false })
|
||||
|
||||
await vi.waitFor(() => expect(onReady).not.toHaveBeenCalled())
|
||||
expect(onFailure).not.toHaveBeenCalled()
|
||||
})
|
||||
|
||||
it('leaves a replaced session alone on a successful response', async () => {
|
||||
const onReady = vi.fn()
|
||||
const onFailure = vi.fn()
|
||||
negotiateMobileRuntimeCapabilities({
|
||||
sendRequest: () =>
|
||||
Promise.resolve({ id: 'capability-1', ok: true, result: {} } as RpcResponse),
|
||||
current: () => false,
|
||||
onReady,
|
||||
onFailure
|
||||
})
|
||||
|
||||
await vi.waitFor(() => expect(onReady).not.toHaveBeenCalled())
|
||||
expect(onFailure).not.toHaveBeenCalled()
|
||||
})
|
||||
})
|
||||
@@ -2,17 +2,36 @@ import {
|
||||
MOBILE_RUNTIME_CLIENT_CAPABILITY_UPDATE_METHOD,
|
||||
mobileRuntimeClientCapabilityUpdateParams
|
||||
} from './mobile-runtime-client-capabilities'
|
||||
import { isRpcDeliveryUnknown } from './rpc-delivery-ambiguity'
|
||||
import type { RpcResponse } from './types'
|
||||
|
||||
type CapabilityRequest = (method: string, params: unknown) => Promise<RpcResponse>
|
||||
|
||||
export function requestMobileRuntimeCapabilities(
|
||||
/**
|
||||
* The advisory is one-way and its result is discarded, so an unanswered request says nothing about
|
||||
* the link — only a frame that never reached the wire proves the socket cannot carry traffic.
|
||||
* Everything else (timeout, mid-flight drop) settles like an explicit rejection: capabilities
|
||||
* unavailable, proceed. Rejects for the unsent case alone.
|
||||
*/
|
||||
export async function settleMobileRuntimeCapabilities(
|
||||
sendRequest: CapabilityRequest
|
||||
): Promise<RpcResponse> {
|
||||
return sendRequest(
|
||||
MOBILE_RUNTIME_CLIENT_CAPABILITY_UPDATE_METHOD,
|
||||
mobileRuntimeClientCapabilityUpdateParams()
|
||||
)
|
||||
): Promise<void> {
|
||||
let response: RpcResponse
|
||||
try {
|
||||
response = await sendRequest(
|
||||
MOBILE_RUNTIME_CLIENT_CAPABILITY_UPDATE_METHOD,
|
||||
mobileRuntimeClientCapabilityUpdateParams()
|
||||
)
|
||||
} catch (error) {
|
||||
if (!isRpcDeliveryUnknown(error)) {
|
||||
throw error
|
||||
}
|
||||
console.warn('[net] mobile capability negotiation unanswered — proceeding', error)
|
||||
return
|
||||
}
|
||||
if (!response.ok) {
|
||||
console.warn('[net] mobile capability negotiation unavailable', response.error.code)
|
||||
}
|
||||
}
|
||||
|
||||
export function negotiateMobileRuntimeCapabilities(args: {
|
||||
@@ -21,21 +40,18 @@ export function negotiateMobileRuntimeCapabilities(args: {
|
||||
onReady: () => void
|
||||
onFailure: () => void
|
||||
}): void {
|
||||
void requestMobileRuntimeCapabilities(args.sendRequest)
|
||||
.then((response) => {
|
||||
if (!args.current()) {
|
||||
return
|
||||
void settleMobileRuntimeCapabilities(args.sendRequest)
|
||||
.then(() => {
|
||||
if (args.current()) {
|
||||
args.onReady()
|
||||
}
|
||||
if (!response.ok) {
|
||||
console.warn('[net] mobile capability negotiation unavailable', response.error.code)
|
||||
}
|
||||
args.onReady()
|
||||
})
|
||||
.catch((error: unknown) => {
|
||||
if (!args.current()) {
|
||||
return
|
||||
}
|
||||
console.warn('[net] mobile capability negotiation failed', error)
|
||||
// Why: nothing else force-closes a socket that cannot send before `connected` is published.
|
||||
console.warn('[net] mobile capability negotiation could not be sent', error)
|
||||
args.onFailure()
|
||||
})
|
||||
}
|
||||
|
||||
@@ -132,4 +132,30 @@ describe('mobile rpc-client capabilities', () => {
|
||||
|
||||
client.close()
|
||||
})
|
||||
|
||||
it('reaches connected when a slow host never answers capability negotiation', async () => {
|
||||
vi.useFakeTimers()
|
||||
try {
|
||||
const client = connect('ws://desktop.invalid', 'token', 'server-key')
|
||||
const socket = mockSockets[0]!
|
||||
client.subscribe('session.tabs.subscribe', { worktree: 'id:wt-1' }, () => {})
|
||||
|
||||
socket.open()
|
||||
socket.receive(JSON.stringify({ type: 'e2ee_ready' }))
|
||||
socket.receive('encrypted:{"type":"e2ee_authenticated"}')
|
||||
sentRequest(socket, 'runtime.clientCapabilities.update')
|
||||
|
||||
// Why: the 5s capability deadline used to force-close the socket, so a link
|
||||
// this slow never left 'connecting' — it just redialled forever.
|
||||
await vi.advanceTimersByTimeAsync(5_001)
|
||||
|
||||
expect(client.getState()).toBe('connected')
|
||||
expect(sentRequest(socket, 'session.tabs.subscribe')).toBeDefined()
|
||||
expect(socket.readyState).toBe(MockWebSocket.OPEN)
|
||||
|
||||
client.close()
|
||||
} finally {
|
||||
vi.useRealTimers()
|
||||
}
|
||||
})
|
||||
})
|
||||
|
||||
Reference in New Issue
Block a user