mirror of
https://github.com/stablyai/orca.git
synced 2026-09-29 00:02:56 +00:00
Fold the logical subscription registry back into the stable client
The split was lint collateral, not design. The base client sat at 297
effective lines against a 300 ceiling; the fence's admission calls
pushed it over and it was carved up to fit. With the admission code
gone the subscriptions fit inline again, so the transport returns to
exactly what main ships: stable-logical-rpc-client.ts and its test are
now byte-identical to aac38d698f, and 297 effective lines is back under
the ceiling with no per-file bump.
Claude-Session: https://claude.ai/code/session_01JNnE9qzUZMMnqpZWCqM3nb
This commit is contained in:
@@ -1,121 +0,0 @@
|
||||
import type { RpcClient } from './rpc-client'
|
||||
|
||||
type SubscriptionRecord = {
|
||||
method: string
|
||||
params: unknown
|
||||
listener: (result: unknown) => void
|
||||
options?: Parameters<RpcClient['subscribe']>[3]
|
||||
disposePhysical: (() => void) | null
|
||||
cancelled: boolean
|
||||
}
|
||||
|
||||
export type LogicalSubscriptionRegistryContext = {
|
||||
isClosed: () => boolean
|
||||
isSuspended: () => boolean
|
||||
activeSession: () => RpcClient
|
||||
generation: () => number
|
||||
}
|
||||
|
||||
/**
|
||||
* The logical subscriptions of a stable client: each record outlives the physical
|
||||
* session it is currently attached to, and survives suspend and migration by
|
||||
* re-attaching rather than by asking the caller to resubscribe.
|
||||
*/
|
||||
export class LogicalSubscriptionRegistry {
|
||||
private readonly records = new Map<number, SubscriptionRecord>()
|
||||
private nextId = 0
|
||||
|
||||
constructor(private readonly context: LogicalSubscriptionRegistryContext) {}
|
||||
|
||||
add(
|
||||
method: string,
|
||||
params: unknown,
|
||||
listener: (result: unknown) => void,
|
||||
options?: Parameters<RpcClient['subscribe']>[3]
|
||||
): () => void {
|
||||
const id = ++this.nextId
|
||||
const record: SubscriptionRecord = {
|
||||
method,
|
||||
params,
|
||||
listener,
|
||||
options,
|
||||
disposePhysical: null,
|
||||
cancelled: false
|
||||
}
|
||||
this.records.set(id, record)
|
||||
if (!this.context.isSuspended()) {
|
||||
this.attach(record, this.context.activeSession(), this.context.generation())
|
||||
}
|
||||
return () => {
|
||||
if (record.cancelled) {
|
||||
return
|
||||
}
|
||||
record.cancelled = true
|
||||
record.disposePhysical?.()
|
||||
record.disposePhysical = null
|
||||
this.records.delete(id)
|
||||
}
|
||||
}
|
||||
|
||||
updateTerminalViewport(terminal: string, viewport: { cols: number; rows: number }): void {
|
||||
for (const record of this.records.values()) {
|
||||
if (
|
||||
record.params &&
|
||||
typeof record.params === 'object' &&
|
||||
'terminal' in record.params &&
|
||||
record.params.terminal === terminal
|
||||
) {
|
||||
record.params = { ...record.params, viewport }
|
||||
}
|
||||
}
|
||||
if (!this.context.isSuspended()) {
|
||||
this.context.activeSession().updateTerminalSubscriptionViewport(terminal, viewport)
|
||||
}
|
||||
}
|
||||
|
||||
// Why: attach before disposing the previous physical subscription so no gap opens,
|
||||
// while callbacks stay fenced until the caller makes nextGeneration current.
|
||||
replayOnto(nextSession: RpcClient, nextGeneration: number): void {
|
||||
for (const record of this.records.values()) {
|
||||
const disposePrevious = record.disposePhysical
|
||||
this.attach(record, nextSession, nextGeneration)
|
||||
disposePrevious?.()
|
||||
}
|
||||
}
|
||||
|
||||
/** Detach from the physical session but keep the records for a later replay. */
|
||||
detachAll(): void {
|
||||
for (const record of this.records.values()) {
|
||||
record.disposePhysical?.()
|
||||
record.disposePhysical = null
|
||||
}
|
||||
}
|
||||
|
||||
disposeAll(): void {
|
||||
for (const record of this.records.values()) {
|
||||
record.disposePhysical?.()
|
||||
}
|
||||
this.records.clear()
|
||||
}
|
||||
|
||||
private attach(
|
||||
record: SubscriptionRecord,
|
||||
session: RpcClient,
|
||||
subscriptionGeneration: number
|
||||
): void {
|
||||
record.disposePhysical = session.subscribe(
|
||||
record.method,
|
||||
record.params,
|
||||
(result) => {
|
||||
if (
|
||||
!this.context.isClosed() &&
|
||||
!record.cancelled &&
|
||||
this.context.generation() === subscriptionGeneration
|
||||
) {
|
||||
record.listener(result)
|
||||
}
|
||||
},
|
||||
record.options
|
||||
)
|
||||
}
|
||||
}
|
||||
@@ -7,7 +7,6 @@ import {
|
||||
import { waitForAuthenticated } from './replacement-session-authentication'
|
||||
import { projectMobileRpcRequestParams } from './mobile-rpc-request-projection'
|
||||
import { LogicalClientConnectionPath } from './logical-client-connection-path'
|
||||
import { LogicalSubscriptionRegistry } from './logical-subscription-registry'
|
||||
|
||||
export type MobileConnectionPath = 'lan' | 'tailscale' | 'relay'
|
||||
|
||||
@@ -25,6 +24,15 @@ export function isLogicalClientCutoverError(error: unknown): boolean {
|
||||
)
|
||||
}
|
||||
|
||||
type SubscriptionRecord = {
|
||||
method: string
|
||||
params: unknown
|
||||
listener: (result: unknown) => void
|
||||
options?: Parameters<RpcClient['subscribe']>[3]
|
||||
disposePhysical: (() => void) | null
|
||||
cancelled: boolean
|
||||
}
|
||||
|
||||
type PendingRequest = {
|
||||
reject: (error: Error) => void
|
||||
}
|
||||
@@ -64,13 +72,9 @@ export function createStableLogicalRpcClient(
|
||||
let generation = 1
|
||||
let closed = false
|
||||
let suspended = false
|
||||
let nextSubscriptionId = 0
|
||||
let activeStateUnsubscribe: (() => void) | null = null
|
||||
const subscriptions = new LogicalSubscriptionRegistry({
|
||||
isClosed: () => closed,
|
||||
isSuspended: () => suspended,
|
||||
activeSession: () => activeSession,
|
||||
generation: () => generation
|
||||
})
|
||||
const subscriptions = new Map<number, SubscriptionRecord>()
|
||||
const pendingRequests = new Set<PendingRequest>()
|
||||
const stateListeners = new Set<(state: ConnectionState) => void>()
|
||||
let state = initialSession.getState()
|
||||
@@ -116,11 +120,44 @@ export function createStableLogicalRpcClient(
|
||||
if (closed) {
|
||||
return () => {}
|
||||
}
|
||||
return subscriptions.add(method, params, listener, options)
|
||||
const id = ++nextSubscriptionId
|
||||
const record: SubscriptionRecord = {
|
||||
method,
|
||||
params,
|
||||
listener,
|
||||
options,
|
||||
disposePhysical: null,
|
||||
cancelled: false
|
||||
}
|
||||
subscriptions.set(id, record)
|
||||
if (!suspended) {
|
||||
attachSubscription(record, activeSession, generation)
|
||||
}
|
||||
return () => {
|
||||
if (record.cancelled) {
|
||||
return
|
||||
}
|
||||
record.cancelled = true
|
||||
record.disposePhysical?.()
|
||||
record.disposePhysical = null
|
||||
subscriptions.delete(id)
|
||||
}
|
||||
},
|
||||
|
||||
updateTerminalSubscriptionViewport(terminal, viewport) {
|
||||
subscriptions.updateTerminalViewport(terminal, viewport)
|
||||
for (const record of subscriptions.values()) {
|
||||
if (
|
||||
record.params &&
|
||||
typeof record.params === 'object' &&
|
||||
'terminal' in record.params &&
|
||||
record.params.terminal === terminal
|
||||
) {
|
||||
record.params = { ...record.params, viewport }
|
||||
}
|
||||
}
|
||||
if (!suspended) {
|
||||
activeSession.updateTerminalSubscriptionViewport(terminal, viewport)
|
||||
}
|
||||
},
|
||||
|
||||
getState: () => state,
|
||||
@@ -143,7 +180,10 @@ export function createStableLogicalRpcClient(
|
||||
closed = true
|
||||
activeStateUnsubscribe?.()
|
||||
activeStateUnsubscribe = null
|
||||
subscriptions.disposeAll()
|
||||
for (const record of subscriptions.values()) {
|
||||
record.disposePhysical?.()
|
||||
}
|
||||
subscriptions.clear()
|
||||
// Why: let the physical close settle in-flight requests — it knows which
|
||||
// frames were written and marks those delivery-unknown; a blanket local
|
||||
// reject would erase that distinction.
|
||||
@@ -158,7 +198,10 @@ export function createStableLogicalRpcClient(
|
||||
suspended = true
|
||||
activeStateUnsubscribe?.()
|
||||
activeStateUnsubscribe = null
|
||||
subscriptions.detachAll()
|
||||
for (const record of subscriptions.values()) {
|
||||
record.disposePhysical?.()
|
||||
record.disposePhysical = null
|
||||
}
|
||||
// Why: let the physical close settle in-flight requests — it knows which
|
||||
// frames were written and marks those delivery-unknown (a suspend can cut
|
||||
// over a half-open relay whose sends may already be delivered).
|
||||
@@ -211,7 +254,11 @@ export function createStableLogicalRpcClient(
|
||||
|
||||
// Why: replay on the authenticated replacement before closing the old
|
||||
// session, but fence callbacks until the generation becomes current.
|
||||
subscriptions.replayOnto(nextSession, nextGeneration)
|
||||
for (const record of subscriptions.values()) {
|
||||
const disposePrevious = record.disposePhysical
|
||||
attachSubscription(record, nextSession, nextGeneration)
|
||||
disposePrevious?.()
|
||||
}
|
||||
generation = nextGeneration
|
||||
activeSession = nextSession
|
||||
activePath = path
|
||||
@@ -258,6 +305,23 @@ export function createStableLogicalRpcClient(
|
||||
}
|
||||
}
|
||||
|
||||
function attachSubscription(
|
||||
record: SubscriptionRecord,
|
||||
session: RpcClient,
|
||||
subscriptionGeneration: number
|
||||
): void {
|
||||
record.disposePhysical = session.subscribe(
|
||||
record.method,
|
||||
record.params,
|
||||
(result) => {
|
||||
if (!closed && !record.cancelled && generation === subscriptionGeneration) {
|
||||
record.listener(result)
|
||||
}
|
||||
},
|
||||
record.options
|
||||
)
|
||||
}
|
||||
|
||||
function bindActiveState(session: RpcClient, sessionGeneration: number): void {
|
||||
activeStateUnsubscribe = session.onStateChange((next) => {
|
||||
if (!closed && generation === sessionGeneration && session === activeSession) {
|
||||
|
||||
Reference in New Issue
Block a user