mirror of
https://github.com/stablyai/orca.git
synced 2026-09-30 16:02:56 +00:00
fix(mobile): re-ask a refused status probe, give every resume probe a fresh miss budget, persist the relay endpoint before booking a refresh
This commit is contained in:
@@ -52,9 +52,11 @@ export class MobileRelayCredentialRefresh {
|
||||
randomBytes: args.randomBytes
|
||||
})
|
||||
args.adoptBundle(result.bundle)
|
||||
// Why: a scheduled rotation can finish after the old credential enters the rejection gate.
|
||||
refreshed = true
|
||||
await args.persistResolvedRelay(result.relay)
|
||||
// Why after the persist: lifting the gate and dialing on a rotation whose endpoint
|
||||
// write failed would dial the pre-rotation cell. The credential itself is already
|
||||
// durable, so a failed write leaves a retry, not a lost credential.
|
||||
refreshed = true
|
||||
} catch {
|
||||
// Why: pending material remains durable; the next authenticated direct
|
||||
// opportunity must reconcile it before creating another install key.
|
||||
|
||||
@@ -96,14 +96,16 @@ export class RelayDialStageTracker implements RelayDialStageSource {
|
||||
|
||||
// Budget per stage once the cell holds the dial. awaiting-hello covers the cell's
|
||||
// assignment/reservation transactions (observed 14–16s under lock contention) plus its
|
||||
// 10s host-attach deadline; handshaking is two E2EE round trips; confirming is bounded
|
||||
// by the session's own 30s resume-confirmation request, with slack so that error wins.
|
||||
// 10s host-attach deadline; handshaking is two E2EE round trips. confirming is kept for
|
||||
// the bound's own completeness: the live session publishes 'connected' in the same turn
|
||||
// it enters the stage and bounds the resume confirm itself at 12 s, so this entry is
|
||||
// only reachable by a session that times the stage without publishing, never in production.
|
||||
const RELAY_DIAL_STAGE_BUDGET_MS: Record<Exclude<RelayDialStage, 'opening'>, number> = {
|
||||
'awaiting-hello': 30_000,
|
||||
handshaking: 12_000,
|
||||
confirming: 35_000
|
||||
}
|
||||
|
||||
export function relayDialStageBudgetMs(stage: Exclude<RelayDialStage, 'opening'>): number {
|
||||
export function relayDialStageBudgetMs(stage: keyof typeof RELAY_DIAL_STAGE_BUDGET_MS): number {
|
||||
return RELAY_DIAL_STAGE_BUDGET_MS[stage]
|
||||
}
|
||||
|
||||
@@ -256,6 +256,36 @@ describe('RpcSessionLivenessWatchdog', () => {
|
||||
expect(terminate).toHaveBeenCalledOnce()
|
||||
})
|
||||
|
||||
it('gives a second resume probe a fresh miss budget too', async () => {
|
||||
// Why: two app-resume nudges ~2 s apart on a cold radio are one observation each, not a
|
||||
// shared budget; otherwise the second inherits the first's miss and one more slow answer
|
||||
// kills a healthy socket.
|
||||
const terminate = vi.fn()
|
||||
const identity = {}
|
||||
const watchdog = new RpcSessionLivenessWatchdog({
|
||||
transport: 'relay',
|
||||
idleProbeMs: 20_000,
|
||||
probeTimeoutMs: 4_000,
|
||||
missedProbeLimit: 2,
|
||||
urgentProbeTimeoutMs: 2_000,
|
||||
urgentMissedProbeLimit: 2,
|
||||
shouldIdleProbe: () => true,
|
||||
sendProbe: () => true,
|
||||
terminate,
|
||||
now: Date.now
|
||||
})
|
||||
watchdog.start(identity)
|
||||
|
||||
watchdog.probeNow(identity, 'resume')
|
||||
await vi.advanceTimersByTimeAsync(2_000)
|
||||
expect(terminate).not.toHaveBeenCalled()
|
||||
watchdog.probeNow(identity, 'resume')
|
||||
await vi.advanceTimersByTimeAsync(2_000)
|
||||
expect(terminate).not.toHaveBeenCalled()
|
||||
await vi.advanceTimersByTimeAsync(2_000)
|
||||
expect(terminate).toHaveBeenCalledOnce()
|
||||
})
|
||||
|
||||
it('still reaches a verdict on a caller probe when the app backgrounds', async () => {
|
||||
// The gate covers the idle sweep only. A nudge or resume probe was asked for on
|
||||
// purpose, and abandoning it would leave a genuinely dead socket unreported.
|
||||
|
||||
@@ -125,6 +125,12 @@ export class RpcSessionLivenessWatchdog {
|
||||
return
|
||||
}
|
||||
this.lastVoluntaryProbeAt = now
|
||||
// Why: a second resume inside the urgent window (iOS active/inactive/active on a
|
||||
// control-centre swipe) is a new observation on the same cold radio; inheriting the
|
||||
// first resume probe's miss would spend the tolerated slow answer before it is sent.
|
||||
if (urgent) {
|
||||
this.missedProbes = 0
|
||||
}
|
||||
this.startProbe(identity, urgent ? this.urgentProfile : this.ordinaryProfile)
|
||||
}
|
||||
|
||||
|
||||
@@ -129,10 +129,11 @@ describe('startRuntimeCapabilityProbe', () => {
|
||||
cancel()
|
||||
})
|
||||
|
||||
// Was: an ok:false response was retried like a timeout. The probe now backs the gate that sits
|
||||
// above every /h/ route, so polling a host that already answered would run for the life of the
|
||||
// connection. A reply is an answer; only an unanswered request is retried.
|
||||
it('settles once on an ok:false response rather than polling the host', async () => {
|
||||
// Why the ceiling and not a stop: the probe backs the gate above every /h/ route, and
|
||||
// worktree.activate plus every capability wait on its verdict. One error reply, from a
|
||||
// transient host failure or an error frame the relay surfaced, must not withhold those
|
||||
// for the life of the connection; but a host that keeps saying no is not polled fast.
|
||||
it('re-asks an ok:false response at the slow ceiling instead of stopping for good', async () => {
|
||||
const failure: RpcResponse = {
|
||||
ok: false,
|
||||
id: '1',
|
||||
@@ -147,12 +148,15 @@ describe('startRuntimeCapabilityProbe', () => {
|
||||
onUnavailable: (isRetrying) => retrying.push(isRetrying)
|
||||
})
|
||||
await flushMicrotasks()
|
||||
expect(retrying).toEqual([false])
|
||||
expect(retrying).toEqual([true])
|
||||
expect(seen).toEqual([])
|
||||
|
||||
await vi.advanceTimersByTimeAsync(60_000)
|
||||
await vi.advanceTimersByTimeAsync(14_999)
|
||||
expect(calls()).toBe(1)
|
||||
expect(seen).toEqual([])
|
||||
await vi.advanceTimersByTimeAsync(1)
|
||||
await flushMicrotasks()
|
||||
expect(calls()).toBe(2)
|
||||
expect(seen).toEqual([['a.v1']])
|
||||
cancel()
|
||||
})
|
||||
|
||||
|
||||
@@ -11,11 +11,10 @@ const FAILURE_RETRY_MAX_DELAY_MS = 15_000
|
||||
|
||||
export type RuntimeStatusProbeHandlers = {
|
||||
onStatus: (status: Record<string, unknown>) => void
|
||||
// Fires once per attempt that produced no status. `retrying` is false when the host itself
|
||||
// answered with an error: that is a definitive reply, so the probe stops rather than polling a
|
||||
// host that has already said no. It is true when nothing reached us and a retry is armed, which
|
||||
// lets a caller that must not stay blocked fail open on the first miss and be upgraded later.
|
||||
onUnavailable?: (retrying: boolean) => void
|
||||
// Fires once per attempt that produced no status, with a retry always armed: an unanswered
|
||||
// request backs off from 1 s, an answered error re-asks at the 15 s ceiling. Either way a
|
||||
// caller that must not stay blocked can fail open on the first miss and be upgraded later.
|
||||
onUnavailable?: (retrying: true) => void
|
||||
}
|
||||
|
||||
// Single status.get producer for a connected client: one request, retried until it
|
||||
@@ -35,9 +34,11 @@ export function startRuntimeStatusProbe(
|
||||
return
|
||||
}
|
||||
if (!response.ok) {
|
||||
// Why not retry: the desktop replied. Re-asking every 15 s for the life of a connection
|
||||
// from a probe mounted above every /h/ route buys nothing a reconnect would not.
|
||||
handlers.onUnavailable?.(false)
|
||||
// Why retry: an error reply is a definitive answer for this attempt, not for the
|
||||
// connection. A transient host failure or an error frame surfaced by the relay must
|
||||
// not withhold worktree.activate and every capability until the socket is replaced.
|
||||
// Backed off to the 15 s ceiling so a host that keeps saying no is not hammered.
|
||||
scheduleRetry(false, true)
|
||||
return
|
||||
}
|
||||
const result = (response as RpcSuccess).result
|
||||
@@ -54,12 +55,16 @@ export function startRuntimeStatusProbe(
|
||||
)
|
||||
}
|
||||
|
||||
function scheduleRetry(cutover: boolean): void {
|
||||
function scheduleRetry(cutover: boolean, answered = false): void {
|
||||
// Why: cutover means the replacement transport is already authenticated —
|
||||
// re-ask promptly; other failures back off so a wedged host isn't hammered.
|
||||
// An answered error skips straight to the ceiling: the host is reachable and
|
||||
// said no, so only a slow re-ask is worth anything.
|
||||
const delay = cutover
|
||||
? CUTOVER_RETRY_DELAY_MS
|
||||
: Math.min(FAILURE_RETRY_BASE_DELAY_MS * 2 ** failureRetries++, FAILURE_RETRY_MAX_DELAY_MS)
|
||||
: answered
|
||||
? FAILURE_RETRY_MAX_DELAY_MS
|
||||
: Math.min(FAILURE_RETRY_BASE_DELAY_MS * 2 ** failureRetries++, FAILURE_RETRY_MAX_DELAY_MS)
|
||||
retryTimer = setTimeout(attempt, delay)
|
||||
handlers.onUnavailable?.(true)
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user