mirror of
https://github.com/stablyai/orca.git
synced 2026-09-28 16:02:45 +00:00
refactor(relay): extract the revoke-outbox flush into its own module
Behaviour-preserving move, no logic change: flushRevokeOutbox and flushRevoke leave DesktopRelayService as RelayRevokeOutboxFlusher.flushAll/flushItem. The two fixes that follow each add lines to this file, and it has no headroom under the 300-line cap; a max-lines disable is not an option. Why a new file rather than reusing something: nothing existing drains the outbox. RelayRevokeOutbox is the durable store — making it drain itself over a control would give the persistence layer a dependency on a broker. RelayControlRequests dedups in-flight reqIds on a single socket, not queued items across reconnects. RelayDemandLedger reads the outbox for demand and never mutates it. The drain only ever existed inline here, so this is an extraction, not a parallel implementation; the only new code is the options type and the class shell.
This commit is contained in:
@@ -12,11 +12,8 @@ import { RelayAuthCoordinator } from './relay-auth-coordinator'
|
||||
import { relayOfflineReasonMintFailureCode } from './relay-offline-reason'
|
||||
import { RelaySessionBroker, type RelayBrokerStatus } from './relay-session-broker'
|
||||
import type { PairingRelay } from '../../../shared/mobile-relay-pairing-offer'
|
||||
import type {
|
||||
RelayRevokeOutbox,
|
||||
RelayDeviceBinding,
|
||||
RelayRevokeOutboxItem
|
||||
} from './relay-revoke-outbox'
|
||||
import type { RelayDeviceBinding, RelayRevokeOutboxItem } from './relay-revoke-outbox'
|
||||
import { RelayRevokeOutboxFlusher } from './relay-revoke-outbox-flush'
|
||||
import { deriveRelayHostId } from './relay-http-client'
|
||||
import { RelayDemandLedger } from './relay-demand-ledger'
|
||||
import { createRelayRegionPreferenceReader } from './relay-region-preference'
|
||||
@@ -45,7 +42,7 @@ const RELAY_LIVENESS_INTERVAL_MS = 5 * 60_000
|
||||
|
||||
export class DesktopRelayService {
|
||||
private readonly coordinator: RelayAuthCoordinator
|
||||
private readonly revokeOutbox: RelayRevokeOutbox
|
||||
private readonly revokeFlusher: RelayRevokeOutboxFlusher
|
||||
private readonly runtimeRpc: OrcaRuntimeRpcServer
|
||||
private readonly demandLedger: RelayDemandLedger
|
||||
private readonly hostMobilePairingConnectionMode?: () => MobilePairingConnectionMode
|
||||
@@ -60,11 +57,15 @@ export class DesktopRelayService {
|
||||
throw new Error('mobile_runtime_not_ready')
|
||||
}
|
||||
this.runtimeRpc = options.runtimeRpc
|
||||
this.revokeOutbox = options.runtimeRpc.getRelayRevokeOutbox()
|
||||
const revokeOutbox = options.runtimeRpc.getRelayRevokeOutbox()
|
||||
this.revokeFlusher = new RelayRevokeOutboxFlusher({
|
||||
outbox: revokeOutbox,
|
||||
onDrained: () => this.refreshDemand()
|
||||
})
|
||||
this.hostMobilePairingConnectionMode = options.hostMobilePairingConnectionMode
|
||||
this.demandLedger = new RelayDemandLedger({
|
||||
deviceRegistry: options.runtimeRpc.getDeviceRegistry()!,
|
||||
revokeOutbox: this.revokeOutbox,
|
||||
revokeOutbox,
|
||||
relayHostId: deriveRelayHostId(keypair.publicKey),
|
||||
isRelayAllowedForDevice: (deviceId) => this.isRelayAllowedForDevice(deviceId)
|
||||
})
|
||||
@@ -89,7 +90,7 @@ export class DesktopRelayService {
|
||||
onAssignedCellActive: regionPreference.noteAssignedCell,
|
||||
onStatus: options.onStatus
|
||||
})
|
||||
void this.flushRevokeOutbox(broker)
|
||||
void this.revokeFlusher.flushAll(broker)
|
||||
return broker
|
||||
},
|
||||
onStatus: options.onStatus
|
||||
@@ -148,7 +149,7 @@ export class DesktopRelayService {
|
||||
broker.hostId === item.relayHostId &&
|
||||
broker.ownerIdentityKey === item.ownerIdentityKey
|
||||
) {
|
||||
void this.flushRevoke(broker, item)
|
||||
void this.revokeFlusher.flushItem(broker, item)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -233,12 +234,6 @@ export class DesktopRelayService {
|
||||
this.coordinator.stop()
|
||||
}
|
||||
|
||||
private async flushRevokeOutbox(broker: RelaySessionBroker): Promise<void> {
|
||||
for (const item of this.revokeOutbox.pendingFor(broker.ownerIdentityKey, broker.hostId)) {
|
||||
await this.flushRevoke(broker, item)
|
||||
}
|
||||
}
|
||||
|
||||
private requireMobileDevice(deviceId: string): void {
|
||||
if (this.runtimeRpc.getDeviceRegistry()?.getDevice(deviceId)?.scope !== 'mobile') {
|
||||
throw new Error('mobile_device_not_found')
|
||||
@@ -257,20 +252,6 @@ export class DesktopRelayService {
|
||||
}
|
||||
}
|
||||
|
||||
private async flushRevoke(
|
||||
broker: RelaySessionBroker,
|
||||
item: RelayRevokeOutboxItem
|
||||
): Promise<void> {
|
||||
try {
|
||||
await broker.revokeDevice(item.relayDeviceId, item.reqId)
|
||||
this.revokeOutbox.remove(item.reqId)
|
||||
this.refreshDemand()
|
||||
} catch {
|
||||
// Why: the durable item is the source of truth; reconnecting the same
|
||||
// account/control retries this stable reqId without delaying local revoke.
|
||||
}
|
||||
}
|
||||
|
||||
// Restrictive-only: the host setting withdraws Relay from an `automatic`
|
||||
// device, never grants it to a `local-only` one.
|
||||
private isRelayAllowedForDevice(deviceId: string): boolean {
|
||||
|
||||
@@ -0,0 +1,42 @@
|
||||
import type { RelayRevokeOutbox, RelayRevokeOutboxItem } from './relay-revoke-outbox'
|
||||
|
||||
type RevokingBroker = {
|
||||
hostId: string
|
||||
ownerIdentityKey: string
|
||||
revokeDevice(relayDeviceId: string, reqId: string): Promise<void>
|
||||
}
|
||||
|
||||
type RelayRevokeOutboxFlushOptions = {
|
||||
outbox: RelayRevokeOutbox
|
||||
/** Removing an item can retire the demand holding the broker open. */
|
||||
onDrained: () => void
|
||||
}
|
||||
|
||||
/**
|
||||
* Drains queued device revocations over a live relay control. Runs unawaited from
|
||||
* inside the coordinator's broker open, so it outlives its own caller.
|
||||
*/
|
||||
export class RelayRevokeOutboxFlusher {
|
||||
private readonly options: RelayRevokeOutboxFlushOptions
|
||||
|
||||
constructor(options: RelayRevokeOutboxFlushOptions) {
|
||||
this.options = options
|
||||
}
|
||||
|
||||
async flushAll(broker: RevokingBroker): Promise<void> {
|
||||
for (const item of this.options.outbox.pendingFor(broker.ownerIdentityKey, broker.hostId)) {
|
||||
await this.flushItem(broker, item)
|
||||
}
|
||||
}
|
||||
|
||||
async flushItem(broker: RevokingBroker, item: RelayRevokeOutboxItem): Promise<void> {
|
||||
try {
|
||||
await broker.revokeDevice(item.relayDeviceId, item.reqId)
|
||||
this.options.outbox.remove(item.reqId)
|
||||
this.options.onDrained()
|
||||
} catch {
|
||||
// Why: the durable item is the source of truth; reconnecting the same
|
||||
// account/control retries this stable reqId without delaying local revoke.
|
||||
}
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user