mirror of
https://github.com/stablyai/orca.git
synced 2026-10-09 00:02:39 +00:00
refactor(mobile): drop shell ledgers and stream caps the generic lane made dead
Codex review of the shell lane: the workspace authority keeps only the worktree handle map; the single-entry subscription wrapper, duplicate byte checks and arbitrary 128-stream caps go; ended streams are removed before reconnect replay; a retired snapshot pager cannot restore a stale workspace handle after client replacement. Claude-Session: https://claude.ai/code/session_01JNnE9qzUZMMnqpZWCqM3nb
This commit is contained in:
@@ -12,7 +12,7 @@ import {
|
||||
} from './mobile-web-broker-error'
|
||||
import { MobileWebOperationRateLimiter } from './mobile-web-operation-rate-limiter'
|
||||
import { MobileWebCommitMessageGeneration } from './mobile-web-commit-message-generation'
|
||||
import { MobileWebCapabilitySubscriptions } from './mobile-web-capability-subscriptions'
|
||||
import { MobileWebHostSubscriptions } from './mobile-web-host-subscriptions'
|
||||
import { MOBILE_WEB_PRODUCTION_GRANT_INDEX } from './mobile-web-production-grants'
|
||||
import { MobileWebTerminalStreams } from './mobile-web-terminal-streams'
|
||||
import { MOBILE_WEB_TERMINAL_CLIENT_CLOSURE } from './mobile-web-terminal-stream-retirement'
|
||||
@@ -41,7 +41,7 @@ type PendingRequest = { operationKey: string; subscriptionId?: string; cancelled
|
||||
export class MobileWebCapabilityBroker {
|
||||
private readonly pending = new Map<string, PendingRequest>()
|
||||
private readonly replay = new MobileWebBrokerReplayWindow()
|
||||
private readonly subscriptions: MobileWebCapabilitySubscriptions
|
||||
private readonly subscriptions: MobileWebHostSubscriptions
|
||||
private readonly terminalStreams: MobileWebTerminalStreams
|
||||
private readonly speechAuthority: MobileWebSpeechAuthority
|
||||
private readonly rateLimiter: MobileWebOperationRateLimiter
|
||||
@@ -60,7 +60,7 @@ export class MobileWebCapabilityBroker {
|
||||
postMessage: options.postMessage
|
||||
})
|
||||
const posts = this.messages.subscriptionPosts()
|
||||
this.subscriptions = new MobileWebCapabilitySubscriptions({
|
||||
this.subscriptions = new MobileWebHostSubscriptions({
|
||||
...posts,
|
||||
workspaceAuthority: this.authorities.workspace
|
||||
})
|
||||
@@ -151,6 +151,8 @@ export class MobileWebCapabilityBroker {
|
||||
}
|
||||
|
||||
const isHostRequest = mobileWebIsHostRequest(request)
|
||||
const isHostForward =
|
||||
isHostRequest || (request.capability === 'workspace' && request.operation === 'hostSubscribe')
|
||||
const grant = MOBILE_WEB_PRODUCTION_GRANT_INDEX.get(mobileWebOperationKey(request))
|
||||
const expectsSubscription = mobileWebRequestExpectsSubscription(request)
|
||||
if (!grant || (request.mode === 'subscription') !== expectsSubscription) {
|
||||
@@ -161,7 +163,11 @@ export class MobileWebCapabilityBroker {
|
||||
await this.messages.error(request.requestId, 'invalid_request', false)
|
||||
return
|
||||
}
|
||||
if (mobileWebEncodedByteLength(request.payload) > grant.limits.maxRequestBytes) {
|
||||
// Generic forwarding checks the envelope after rewriting the workspace handle.
|
||||
if (
|
||||
!isHostForward &&
|
||||
mobileWebEncodedByteLength(request.payload) > grant.limits.maxRequestBytes
|
||||
) {
|
||||
await this.messages.error(request.requestId, 'too_large', false)
|
||||
return
|
||||
}
|
||||
@@ -201,7 +207,7 @@ export class MobileWebCapabilityBroker {
|
||||
return
|
||||
}
|
||||
this.pending.delete(request.requestId)
|
||||
await (mobileWebEncodedByteLength(payload) > grant.limits.maxResponseBytes
|
||||
await (!isHostForward && mobileWebEncodedByteLength(payload) > grant.limits.maxResponseBytes
|
||||
? this.messages.error(request.requestId, 'unavailable', false)
|
||||
: this.messages.success(request.requestId, payload))
|
||||
} catch (error) {
|
||||
@@ -232,7 +238,7 @@ export class MobileWebCapabilityBroker {
|
||||
terminalClientId: this.options.terminalClientId,
|
||||
nativeAuthority: this.options.nativeAuthority,
|
||||
speechAuthority: this.speechAuthority,
|
||||
hostSubscriptions: this.subscriptions.host,
|
||||
hostSubscriptions: this.subscriptions,
|
||||
terminalStreams: this.terminalStreams,
|
||||
commitMessageGeneration: this.commitMessageGeneration,
|
||||
nativeChatAuthority: this.authorities.nativeChat,
|
||||
|
||||
@@ -38,7 +38,15 @@ describe('mobile web capability schema corpus', () => {
|
||||
it('applies every production request-byte limit before host or native access', async () => {
|
||||
for (const [index, grant] of MOBILE_WEB_PRODUCTION_GRANTS.entries()) {
|
||||
const harness = createHarness()
|
||||
const payload = { value: 'x'.repeat(grant.limits.maxRequestBytes + 1) }
|
||||
const oversized = { value: 'x'.repeat(grant.limits.maxRequestBytes + 1) }
|
||||
const payload =
|
||||
grant.capability === 'workspace' &&
|
||||
(grant.operation === 'hostRequest' || grant.operation === 'hostSubscribe')
|
||||
? {
|
||||
method: grant.operation === 'hostSubscribe' ? 'future.subscribe' : 'future.read',
|
||||
params: oversized
|
||||
}
|
||||
: oversized
|
||||
|
||||
await harness.broker.handle(requestFor(grant, payload, index))
|
||||
|
||||
|
||||
@@ -1,66 +0,0 @@
|
||||
import { MobileWebHostSubscriptions } from './mobile-web-host-subscriptions'
|
||||
import type { MobileWebSubscriptionClosure } from './mobile-web-subscription-closure'
|
||||
import type {
|
||||
MobileWebSubscriptionLedgerConfig,
|
||||
MobileWebSubscriptionLedgerHandle
|
||||
} from './mobile-web-subscription-ledger'
|
||||
import type { MobileWebWorkspaceAuthority } from './mobile-web-workspace-authority'
|
||||
|
||||
export class MobileWebCapabilitySubscriptions {
|
||||
readonly host: MobileWebHostSubscriptions
|
||||
private readonly ledgers: MobileWebSubscriptionLedgerHandle[]
|
||||
|
||||
constructor(
|
||||
args: MobileWebSubscriptionLedgerConfig<unknown> & {
|
||||
workspaceAuthority: MobileWebWorkspaceAuthority
|
||||
}
|
||||
) {
|
||||
const shared = {
|
||||
isActive: args.isActive,
|
||||
postEvent: args.postEvent,
|
||||
postClosed: args.postClosed
|
||||
}
|
||||
this.host = new MobileWebHostSubscriptions({
|
||||
...shared,
|
||||
workspaceAuthority: args.workspaceAuthority
|
||||
})
|
||||
this.ledgers = [this.host]
|
||||
}
|
||||
|
||||
countForOperation(operationKey: string): number {
|
||||
let count = 0
|
||||
for (const ledger of this.ledgers) {
|
||||
count += ledger.countForOperation(operationKey)
|
||||
}
|
||||
return count
|
||||
}
|
||||
|
||||
cancel(subscriptionId: string): string | null {
|
||||
for (const ledger of this.ledgers) {
|
||||
const requestId = ledger.cancel(subscriptionId)
|
||||
if (requestId !== null) {
|
||||
return requestId
|
||||
}
|
||||
}
|
||||
return null
|
||||
}
|
||||
|
||||
cancelByRequest(requestId: string): void {
|
||||
for (const ledger of this.ledgers) {
|
||||
ledger.cancelByRequest(requestId)
|
||||
}
|
||||
}
|
||||
|
||||
/** Used when the page survives but its host feed does not, so every live entry learns it is over. */
|
||||
closeAll(closure: MobileWebSubscriptionClosure): void {
|
||||
for (const ledger of this.ledgers) {
|
||||
ledger.closeAll(closure)
|
||||
}
|
||||
}
|
||||
|
||||
dispose(): void {
|
||||
for (const ledger of this.ledgers) {
|
||||
ledger.dispose()
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -65,15 +65,6 @@ describe('mobile web host navigation route', () => {
|
||||
truncated: false
|
||||
},
|
||||
'unavailable'
|
||||
],
|
||||
[
|
||||
'malformed target',
|
||||
{
|
||||
worktrees: [{ worktreeId: HOST_WORKSPACE_ID }],
|
||||
totalCount: 1,
|
||||
truncated: false
|
||||
},
|
||||
'unavailable'
|
||||
]
|
||||
])('rejects a %s Desktop snapshot', async (_label, result, code) => {
|
||||
await expect(
|
||||
@@ -85,6 +76,19 @@ describe('mobile web host navigation route', () => {
|
||||
).rejects.toMatchObject({ code })
|
||||
})
|
||||
|
||||
it('registers a workspace without repository metadata', async () => {
|
||||
const route = await resolveMobileWebHostNavigationRoute(
|
||||
HOST_WORKSPACE_ID,
|
||||
hostClient({ worktrees: [{ worktreeId: HOST_WORKSPACE_ID }] }),
|
||||
workspaceAuthority()
|
||||
)
|
||||
expect(route).toEqual({
|
||||
kind: 'session',
|
||||
workspaceId: `workspace_0_${'07'.repeat(16)}`,
|
||||
workspaceName: 'Workspace'
|
||||
})
|
||||
})
|
||||
|
||||
it('rejects an oversized Desktop snapshot before registering its target', async () => {
|
||||
const result = {
|
||||
worktrees: [
|
||||
|
||||
@@ -53,11 +53,7 @@ export async function resolveMobileWebHostNavigationRoute(
|
||||
}
|
||||
|
||||
const match = matches[0]
|
||||
const hostRepoId = boundedRequiredText(match.repoId, 512)
|
||||
if (!hostRepoId) {
|
||||
throw new MobileWebBrokerError('unavailable')
|
||||
}
|
||||
const pageWorkspaceId = authority.registerWorkspace(hostWorkspaceId, hostRepoId)
|
||||
const pageWorkspaceId = authority.registerWorkspace(hostWorkspaceId)
|
||||
return {
|
||||
kind: 'session',
|
||||
workspaceId: pageWorkspaceId,
|
||||
|
||||
@@ -73,13 +73,14 @@ export class MobileWebHostSubscriptions extends MobileWebSubscriptionLedger<
|
||||
) {
|
||||
return
|
||||
}
|
||||
if (mobileWebHostPayloadByteLength(event) === undefined) {
|
||||
const bytes = mobileWebHostPayloadByteLength(event)
|
||||
if (bytes === undefined) {
|
||||
this.cancel(subscriptionId, { code: 'too_large', retryable: false })
|
||||
return
|
||||
}
|
||||
const type =
|
||||
typeof event === 'object' && event !== null && 'type' in event ? event.type : undefined
|
||||
record.closing = type === 'end' || type === 'error'
|
||||
this.enqueue(subscriptionId, record, event, record.closing)
|
||||
this.enqueue(subscriptionId, record, event, record.closing, bytes)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -114,7 +114,7 @@ describe('mobile web native capability operations', () => {
|
||||
it('resolves opaque workspace authority before reading or writing a shell draft', async () => {
|
||||
const harness = createHarness()
|
||||
const workspaceAuthority = new MobileWebWorkspaceAuthority((length) => new Uint8Array(length))
|
||||
const workspaceId = workspaceAuthority.registerWorkspace('host-workspace', 'host-repo')
|
||||
const workspaceId = workspaceAuthority.registerWorkspace('host-workspace')
|
||||
const sessionChatDraftRead = vi.fn().mockResolvedValue('saved draft')
|
||||
const sessionChatDraftWrite = vi.fn().mockResolvedValue(undefined)
|
||||
harness.authority.sessionChatDraftRead = sessionChatDraftRead
|
||||
|
||||
@@ -60,6 +60,10 @@ export async function runMobileWebPackageRefresh(args: {
|
||||
downloaded.commit.buildId,
|
||||
MOBILE_WEB_BRIDGE_PROTOCOL_VERSION
|
||||
)
|
||||
if (!args.isCurrent()) {
|
||||
await ExpoMobileWebShell.closeSession(session.sessionId).catch(() => {})
|
||||
return { kind: 'stale' }
|
||||
}
|
||||
if (!(await args.publish(session, startedAt))) {
|
||||
return { kind: 'stale' }
|
||||
}
|
||||
|
||||
@@ -44,8 +44,7 @@ export function capabilityGrants<TCapability extends MobileWebBridgeCapability>(
|
||||
}))
|
||||
}
|
||||
|
||||
/** Grants indexed by `capability.operation`. The broker resolves one per request, including every
|
||||
* terminal keystroke, so a linear scan of all 226 entries is not an option. */
|
||||
// Indexed for every request, including terminal keystrokes.
|
||||
export function indexGrants(
|
||||
grants: readonly MobileWebOperationGrant[]
|
||||
): ReadonlyMap<string, MobileWebOperationGrant> {
|
||||
|
||||
@@ -13,13 +13,7 @@ export function mobileWebRequestExpectsSubscription(request: {
|
||||
}): boolean {
|
||||
return (
|
||||
(request.capability === 'workspace' && request.operation === 'hostSubscribe') ||
|
||||
((request.capability === 'workspace' ||
|
||||
request.capability === 'account' ||
|
||||
request.capability === 'session' ||
|
||||
request.capability === 'sourceControl' ||
|
||||
request.capability === 'terminal' ||
|
||||
request.capability === 'nativeChat' ||
|
||||
request.capability === 'speech') &&
|
||||
((request.capability === 'terminal' || request.capability === 'speech') &&
|
||||
request.operation === 'subscribe')
|
||||
)
|
||||
}
|
||||
|
||||
@@ -3,7 +3,6 @@ import type { MobileWebBridgeShellMessage } from '../../../src/shared/mobile-web
|
||||
import type { MobileWebSubscriptionClosure } from './mobile-web-subscription-closure'
|
||||
import type { RpcClient } from '../transport/rpc-client'
|
||||
import { MobileWebHostSubscriptions } from './mobile-web-host-subscriptions'
|
||||
import { MobileWebCapabilitySubscriptions } from './mobile-web-capability-subscriptions'
|
||||
import { MobileWebBrokerMessageSender } from './mobile-web-broker-message-sender'
|
||||
import { MobileWebWorkspaceAuthority } from './mobile-web-workspace-authority'
|
||||
import {
|
||||
@@ -80,7 +79,7 @@ describe('subscription ledger teardown', () => {
|
||||
expect(ledger.countForOperation('workspace.hostSubscribe')).toBe(0)
|
||||
})
|
||||
|
||||
it('fans closeAll out across every capability ledger', () => {
|
||||
it('posts host-feed closures through the broker message sender', () => {
|
||||
const messages: MobileWebBridgeShellMessage[] = []
|
||||
const sender = new MobileWebBrokerMessageSender({
|
||||
context: MOBILE_WEB_BRIDGE_ROUNDTRIP_CONTEXT,
|
||||
@@ -91,12 +90,12 @@ describe('subscription ledger teardown', () => {
|
||||
})
|
||||
const workspaceAuthority = new MobileWebWorkspaceAuthority(randomBytes)
|
||||
workspaceAuthority.synchronize(['host-workspace'])
|
||||
const subscriptions = new MobileWebCapabilitySubscriptions({
|
||||
const subscriptions = new MobileWebHostSubscriptions({
|
||||
...sender.subscriptionPosts(),
|
||||
workspaceAuthority
|
||||
})
|
||||
const client = stubClient(() => {})
|
||||
subscriptions.host.start({
|
||||
subscriptions.start({
|
||||
requestId: 'r2',
|
||||
subscriptionId: 'host-1',
|
||||
payload: { method: 'mobileWeb.workspace.subscribe', params: {} },
|
||||
|
||||
@@ -21,15 +21,6 @@ export type MobileWebSubscriptionLedgerOptions<TEvent> = {
|
||||
postClosed: MobileWebPostSubscriptionClosed
|
||||
}
|
||||
|
||||
/** The ledger surface callers share across every event and record type. */
|
||||
export type MobileWebSubscriptionLedgerHandle = {
|
||||
cancel: (subscriptionId: string, closure?: MobileWebSubscriptionClosure) => string | null
|
||||
cancelByRequest: (requestId: string) => void
|
||||
countForOperation: (operationKey: string) => number
|
||||
closeAll: (closure: MobileWebSubscriptionClosure) => void
|
||||
dispose: () => void
|
||||
}
|
||||
|
||||
/** What a concrete ledger supplies; the operation key is its own identity, not the caller's. */
|
||||
export type MobileWebSubscriptionLedgerConfig<TEvent> = Omit<
|
||||
MobileWebSubscriptionLedgerOptions<TEvent>,
|
||||
@@ -55,7 +46,6 @@ export class MobileWebSubscriptionLedger<
|
||||
}
|
||||
record.active = false
|
||||
this.records.delete(subscriptionId)
|
||||
this.retire(record)
|
||||
try {
|
||||
record.unsubscribe()
|
||||
} catch {
|
||||
@@ -138,7 +128,8 @@ export class MobileWebSubscriptionLedger<
|
||||
subscriptionId: string,
|
||||
record: TRecord,
|
||||
event: TEvent,
|
||||
retireAfterDelivery = false
|
||||
retireAfterDelivery = false,
|
||||
retainedBytes = mobileWebEncodedByteLength(event)
|
||||
): void {
|
||||
const sequence = record.sequence
|
||||
record.sequence += 1
|
||||
@@ -151,7 +142,7 @@ export class MobileWebSubscriptionLedger<
|
||||
this.cancel(subscriptionId)
|
||||
}
|
||||
},
|
||||
mobileWebEncodedByteLength(event)
|
||||
retainedBytes
|
||||
)
|
||||
}
|
||||
|
||||
@@ -186,7 +177,4 @@ export class MobileWebSubscriptionLedger<
|
||||
protected canDeliver(_subscriptionId: string, _record: TRecord): boolean {
|
||||
return true
|
||||
}
|
||||
|
||||
/** Runs while the record is being removed, before the host handle is released. */
|
||||
protected retire(_record: TRecord): void {}
|
||||
}
|
||||
|
||||
@@ -25,19 +25,6 @@ describe('mobile web workspace authority', () => {
|
||||
expect(authority.pageWorkspaceId('repo::C:\\private\\second')).toBe(second)
|
||||
})
|
||||
|
||||
it('revokes a repository handle only after workspace and catalog authority both omit it', () => {
|
||||
const authority = new MobileWebWorkspaceAuthority((length) => new Uint8Array(length).fill(8))
|
||||
authority.synchronize(['folder-workspace'])
|
||||
authority.synchronizeRepositories(['folder-repo'])
|
||||
const repo = authority.pageRepoId('folder-repo')
|
||||
|
||||
authority.synchronize([])
|
||||
expect(authority.pageRepoId('folder-repo')).toBe(repo)
|
||||
|
||||
authority.synchronizeRepositories([])
|
||||
expect(() => authority.hostRepoId(repo)).toThrow('not_found')
|
||||
})
|
||||
|
||||
it('revokes every mapping when the shell session is cleared', () => {
|
||||
const authority = new MobileWebWorkspaceAuthority((length) => new Uint8Array(length))
|
||||
authority.synchronize(['host-workspace'])
|
||||
@@ -46,6 +33,6 @@ describe('mobile web workspace authority', () => {
|
||||
authority.clear()
|
||||
|
||||
expect(() => authority.hostWorkspaceId(pageWorkspaceId)).toThrow('not_found')
|
||||
expect(() => authority.pageRepoId('host-repo')).toThrow('not_found')
|
||||
expect(() => authority.pageWorkspaceId('host-workspace')).toThrow('not_found')
|
||||
})
|
||||
})
|
||||
|
||||
@@ -1,19 +1,10 @@
|
||||
import { MobileWebBrokerError } from './mobile-web-broker-error'
|
||||
import {
|
||||
getRepoExecutionHostId,
|
||||
parseExecutionHostId,
|
||||
type ExecutionHostId
|
||||
} from '../../../src/shared/execution-host'
|
||||
import { getProjectIdentityKey } from '../../../src/shared/project-host-setup-projection'
|
||||
import type { NewWorkspaceRepository } from '../worktree/host-workspace-creation-operations'
|
||||
|
||||
declare const hostWorkspaceIdBrand: unique symbol
|
||||
declare const hostRepoIdBrand: unique symbol
|
||||
|
||||
/** Host-side identifiers the authority mints from an opaque page handle. Branding keeps a page
|
||||
* handle from being passed back in as if it were already resolved. */
|
||||
export type MobileWebHostWorkspaceId = string & { readonly [hostWorkspaceIdBrand]: true }
|
||||
export type MobileWebHostRepoId = string & { readonly [hostRepoIdBrand]: true }
|
||||
|
||||
/** A value the host itself reported as a host identifier, so it never came from the page. Every
|
||||
* use is a boundary annotation, not a conversion of a page handle. */
|
||||
@@ -21,20 +12,9 @@ export function mobileWebHostWorkspaceIdFromHost(value: string): MobileWebHostWo
|
||||
return value as MobileWebHostWorkspaceId
|
||||
}
|
||||
|
||||
export function mobileWebHostRepoIdFromHost(value: string): MobileWebHostRepoId {
|
||||
return value as MobileWebHostRepoId
|
||||
}
|
||||
|
||||
export class MobileWebWorkspaceAuthority {
|
||||
private readonly pageWorkspaceIdByHostId = new Map<string, string>()
|
||||
private readonly hostWorkspaceIdByPageId = new Map<string, string>()
|
||||
private readonly hostRepoIdByHostWorkspaceId = new Map<string, string>()
|
||||
private readonly pageRepoIdByHostId = new Map<string, string>()
|
||||
private readonly hostRepoIdByPageId = new Map<string, string>()
|
||||
private catalogRepoIds = new Set<string>()
|
||||
private readonly hostConnectionIdByPageRepoId = new Map<string, string>()
|
||||
private readonly pageProjectIdByHostId = new Map<string, string>()
|
||||
private readonly pageExecutionHostIdByHostId = new Map<string, ExecutionHostId>()
|
||||
private nextHandle = 0
|
||||
|
||||
constructor(private readonly randomBytes: (length: number) => Uint8Array) {}
|
||||
@@ -50,11 +30,9 @@ export class MobileWebWorkspaceAuthority {
|
||||
}
|
||||
}
|
||||
}
|
||||
this.hostRepoIdByHostWorkspaceId.clear()
|
||||
for (const hostWorkspaceId of hostWorkspaceIds) {
|
||||
this.rememberWorkspace(hostWorkspaceId)
|
||||
}
|
||||
this.revokeUnreferencedRepos()
|
||||
}
|
||||
|
||||
pageWorkspaceId(hostWorkspaceId: string): string {
|
||||
@@ -65,70 +43,6 @@ export class MobileWebWorkspaceAuthority {
|
||||
return pageWorkspaceId
|
||||
}
|
||||
|
||||
pageRepoId(hostRepoId: string): string {
|
||||
const pageRepoId = this.pageRepoIdByHostId.get(hostRepoId)
|
||||
if (!pageRepoId) {
|
||||
throw new MobileWebBrokerError('not_found')
|
||||
}
|
||||
return pageRepoId
|
||||
}
|
||||
|
||||
hostRepoId(pageRepoId: string): MobileWebHostRepoId {
|
||||
const hostRepoId = this.hostRepoIdByPageId.get(pageRepoId)
|
||||
if (!hostRepoId) {
|
||||
throw new MobileWebBrokerError('not_found')
|
||||
}
|
||||
return hostRepoId as MobileWebHostRepoId
|
||||
}
|
||||
|
||||
assertHostRepoBinding(pageRepoId: string, expectedHostRepoId: MobileWebHostRepoId): void {
|
||||
if (this.hostRepoId(pageRepoId) !== expectedHostRepoId) {
|
||||
throw new MobileWebBrokerError('conflict')
|
||||
}
|
||||
}
|
||||
|
||||
synchronizeRepositories(hostRepoIds: readonly string[]): void {
|
||||
this.catalogRepoIds = new Set(hostRepoIds)
|
||||
hostRepoIds.forEach((hostRepoId) => this.rememberRepo(hostRepoId))
|
||||
this.revokeUnreferencedRepos()
|
||||
}
|
||||
|
||||
synchronizeCreationRepositories(repositories: readonly NewWorkspaceRepository[]): void {
|
||||
this.synchronizeRepositories(repositories.map((repo) => repo.id))
|
||||
this.hostConnectionIdByPageRepoId.clear()
|
||||
for (const repo of repositories) {
|
||||
if (repo.connectionId) {
|
||||
this.hostConnectionIdByPageRepoId.set(this.pageRepoId(repo.id), repo.connectionId)
|
||||
}
|
||||
this.rememberProject(getProjectIdentityKey(repo))
|
||||
this.rememberExecutionHost(getRepoExecutionHostId(repo))
|
||||
}
|
||||
}
|
||||
|
||||
pageProjectId(hostProjectId: string): string {
|
||||
const pageProjectId = this.pageProjectIdByHostId.get(hostProjectId)
|
||||
if (!pageProjectId) {
|
||||
throw new MobileWebBrokerError('not_found')
|
||||
}
|
||||
return pageProjectId
|
||||
}
|
||||
|
||||
pageExecutionHostId(hostExecutionHostId: ExecutionHostId): ExecutionHostId {
|
||||
const pageExecutionHostId = this.pageExecutionHostIdByHostId.get(hostExecutionHostId)
|
||||
if (!pageExecutionHostId) {
|
||||
throw new MobileWebBrokerError('not_found')
|
||||
}
|
||||
return pageExecutionHostId
|
||||
}
|
||||
|
||||
hostConnectionId(pageRepoId: string): string {
|
||||
const connectionId = this.hostConnectionIdByPageRepoId.get(pageRepoId)
|
||||
if (!connectionId) {
|
||||
throw new MobileWebBrokerError('not_found')
|
||||
}
|
||||
return connectionId
|
||||
}
|
||||
|
||||
hostWorkspaceId(pageWorkspaceId: string): MobileWebHostWorkspaceId {
|
||||
const hostWorkspaceId = this.hostWorkspaceIdByPageId.get(pageWorkspaceId)
|
||||
if (!hostWorkspaceId) {
|
||||
@@ -146,85 +60,33 @@ export class MobileWebWorkspaceAuthority {
|
||||
}
|
||||
}
|
||||
|
||||
registerWorkspace(hostWorkspaceId: string, hostRepoId?: string): string {
|
||||
registerWorkspace(hostWorkspaceId: string): string {
|
||||
this.rememberWorkspace(hostWorkspaceId)
|
||||
if (hostRepoId !== undefined) {
|
||||
this.hostRepoIdByHostWorkspaceId.set(hostWorkspaceId, hostRepoId)
|
||||
this.rememberRepo(hostRepoId)
|
||||
}
|
||||
return this.pageWorkspaceId(hostWorkspaceId)
|
||||
}
|
||||
|
||||
clear(): void {
|
||||
this.pageWorkspaceIdByHostId.clear()
|
||||
this.hostWorkspaceIdByPageId.clear()
|
||||
this.hostRepoIdByHostWorkspaceId.clear()
|
||||
this.pageRepoIdByHostId.clear()
|
||||
this.hostRepoIdByPageId.clear()
|
||||
this.catalogRepoIds.clear()
|
||||
this.hostConnectionIdByPageRepoId.clear()
|
||||
this.pageProjectIdByHostId.clear()
|
||||
this.pageExecutionHostIdByHostId.clear()
|
||||
}
|
||||
|
||||
private rememberWorkspace(hostWorkspaceId: string): void {
|
||||
if (this.pageWorkspaceIdByHostId.has(hostWorkspaceId)) {
|
||||
return
|
||||
}
|
||||
const pageWorkspaceId = this.createHandle('workspace')
|
||||
const pageWorkspaceId = this.createHandle()
|
||||
this.pageWorkspaceIdByHostId.set(hostWorkspaceId, pageWorkspaceId)
|
||||
this.hostWorkspaceIdByPageId.set(pageWorkspaceId, hostWorkspaceId)
|
||||
}
|
||||
|
||||
private rememberRepo(hostRepoId: string): void {
|
||||
if (!this.pageRepoIdByHostId.has(hostRepoId)) {
|
||||
const pageRepoId = this.createHandle('repo')
|
||||
this.pageRepoIdByHostId.set(hostRepoId, pageRepoId)
|
||||
this.hostRepoIdByPageId.set(pageRepoId, hostRepoId)
|
||||
}
|
||||
}
|
||||
|
||||
private revokeUnreferencedRepos(): void {
|
||||
const workspaceRepoIds = new Set(this.hostRepoIdByHostWorkspaceId.values())
|
||||
for (const [hostRepoId, pageRepoId] of this.pageRepoIdByHostId) {
|
||||
if (workspaceRepoIds.has(hostRepoId) || this.catalogRepoIds.has(hostRepoId)) {
|
||||
continue
|
||||
}
|
||||
this.pageRepoIdByHostId.delete(hostRepoId)
|
||||
this.hostRepoIdByPageId.delete(pageRepoId)
|
||||
this.hostConnectionIdByPageRepoId.delete(pageRepoId)
|
||||
}
|
||||
}
|
||||
|
||||
private rememberProject(hostProjectId: string): void {
|
||||
if (!this.pageProjectIdByHostId.has(hostProjectId)) {
|
||||
this.pageProjectIdByHostId.set(hostProjectId, this.createHandle('project'))
|
||||
}
|
||||
}
|
||||
|
||||
private rememberExecutionHost(hostExecutionHostId: ExecutionHostId): void {
|
||||
if (this.pageExecutionHostIdByHostId.has(hostExecutionHostId)) {
|
||||
return
|
||||
}
|
||||
const host = parseExecutionHostId(hostExecutionHostId)
|
||||
if (!host || host.kind === 'local') {
|
||||
this.pageExecutionHostIdByHostId.set(hostExecutionHostId, 'local')
|
||||
return
|
||||
}
|
||||
this.pageExecutionHostIdByHostId.set(
|
||||
hostExecutionHostId,
|
||||
`${host.kind}:${this.createHandle('executionHost')}`
|
||||
)
|
||||
}
|
||||
|
||||
private createHandle(prefix: 'workspace' | 'repo' | 'project' | 'executionHost'): string {
|
||||
private createHandle(): string {
|
||||
const bytes = this.randomBytes(16)
|
||||
if (bytes.byteLength !== 16) {
|
||||
throw new MobileWebBrokerError('internal')
|
||||
}
|
||||
const counter = this.nextHandle.toString(36)
|
||||
this.nextHandle += 1
|
||||
return `${prefix}_${counter}_${Array.from(bytes, byteToHex).join('')}`
|
||||
return `workspace_${counter}_${Array.from(bytes, byteToHex).join('')}`
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -29,6 +29,27 @@ describe('mobile web workspace snapshot pager', () => {
|
||||
).rejects.toMatchObject({ code: 'invalid_request' })
|
||||
})
|
||||
|
||||
it('retires an in-flight snapshot on cleanup without blocking the replacement client', async () => {
|
||||
const authority = new MobileWebWorkspaceAuthority((length) => new Uint8Array(length))
|
||||
const pager = new MobileWebWorkspaceSnapshotPager((length) => new Uint8Array(length))
|
||||
const oldResponse = Promise.withResolvers<Awaited<ReturnType<RpcClient['sendRequest']>>>()
|
||||
const oldClient = { sendRequest: vi.fn(() => oldResponse.promise) } as unknown as RpcClient
|
||||
const pending = pager.snapshot({ limit: 1 }, oldClient, authority)
|
||||
const rejection = expect(pending).rejects.toMatchObject({ code: 'cancelled' })
|
||||
|
||||
pager.clear()
|
||||
authority.clear()
|
||||
const current = await pager.snapshot({ limit: 1 }, workspaceClient(2), authority)
|
||||
oldResponse.resolve({ ok: true, result: { worktrees: [{ worktreeId: 'retired-workspace' }] } })
|
||||
await rejection
|
||||
|
||||
expect(() => authority.pageWorkspaceId('retired-workspace')).toThrow('not_found')
|
||||
expect(authority.hostWorkspaceId(current.workspaces[0]!.id)).toBe('host-workspace-0')
|
||||
await expect(
|
||||
pager.snapshot({ limit: 1, cursor: current.nextCursor! }, workspaceClient(0), authority)
|
||||
).resolves.toMatchObject({ workspaces: [{ name: 'Workspace 1' }], nextCursor: null })
|
||||
})
|
||||
|
||||
it('revokes continuations on lifecycle cleanup and rejects oversized host lists', async () => {
|
||||
const authority = new MobileWebWorkspaceAuthority((length) => new Uint8Array(length))
|
||||
const pager = new MobileWebWorkspaceSnapshotPager((length) => new Uint8Array(length))
|
||||
|
||||
@@ -20,7 +20,7 @@ type Continuation = {
|
||||
|
||||
export class MobileWebWorkspaceSnapshotPager {
|
||||
private continuation: Continuation | null = null
|
||||
private active = false
|
||||
private active: object | null = null
|
||||
private nextCursorNumber = 0
|
||||
|
||||
constructor(private readonly randomBytes: (length: number) => Uint8Array) {}
|
||||
@@ -33,12 +33,16 @@ export class MobileWebWorkspaceSnapshotPager {
|
||||
if (this.active) {
|
||||
throw new MobileWebBrokerError('rate_limited')
|
||||
}
|
||||
this.active = true
|
||||
const request = {}
|
||||
this.active = request
|
||||
try {
|
||||
const payload = MobileWebWorkspaceSnapshotPayloadSchema.parse(payloadValue)
|
||||
const continuation = payload.cursor
|
||||
? this.consumeContinuation(payload.cursor)
|
||||
: await this.begin(client)
|
||||
if (this.active !== request) {
|
||||
throw new MobileWebBrokerError('cancelled')
|
||||
}
|
||||
const page = mobileWebWorkspaceSnapshotPage(
|
||||
continuation.hostResult,
|
||||
payload.limit,
|
||||
@@ -52,16 +56,19 @@ export class MobileWebWorkspaceSnapshotPager {
|
||||
page.nextOffset === null ? null : this.retain(continuation, page.nextOffset)
|
||||
return MobileWebWorkspaceSnapshotResultSchema.parse({ ...page.snapshot, nextCursor })
|
||||
} finally {
|
||||
this.active = false
|
||||
if (this.active === request) {
|
||||
this.active = null
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
clear(): void {
|
||||
this.active = null
|
||||
this.continuation = null
|
||||
}
|
||||
|
||||
private async begin(client: RpcClient): Promise<Continuation> {
|
||||
this.clear()
|
||||
this.continuation = null
|
||||
const response = await client.sendRequest('worktree.ps', {
|
||||
limit: MOBILE_WEB_WORKSPACE_LIST_LIMIT + 1
|
||||
})
|
||||
@@ -79,7 +86,7 @@ export class MobileWebWorkspaceSnapshotPager {
|
||||
|
||||
private consumeContinuation(cursor: string): Continuation {
|
||||
const continuation = this.continuation
|
||||
this.clear()
|
||||
this.continuation = null
|
||||
if (!continuation || continuation.cursor !== cursor) {
|
||||
throw new MobileWebBrokerError('invalid_request')
|
||||
}
|
||||
|
||||
@@ -307,6 +307,53 @@ describe('useMobileWebPackageSession', () => {
|
||||
})
|
||||
})
|
||||
|
||||
it('allows recovery again after the desktop upgrades its build', async () => {
|
||||
native.openSession.mockResolvedValue(SESSION_A)
|
||||
downloadPackage.mockResolvedValue({ commit: { buildId: SESSION_A.buildId } })
|
||||
await mount('connected')
|
||||
await act(async () => {
|
||||
packageSession?.handleLoadFailure('mobile_web_generation_invalid')
|
||||
await flushPromises()
|
||||
})
|
||||
|
||||
await update('disconnected')
|
||||
native.openSession.mockResolvedValue(SESSION_B)
|
||||
downloadPackage.mockResolvedValue({ commit: { buildId: SESSION_B.buildId } })
|
||||
await update('connected')
|
||||
expect(packageSession?.session).toEqual(SESSION_B)
|
||||
removeHostCache.mockClear()
|
||||
await act(async () => {
|
||||
packageSession?.handleLoadFailure('mobile_web_generation_invalid')
|
||||
await flushPromises()
|
||||
})
|
||||
|
||||
expect(removeHostCache).toHaveBeenCalledWith(HOST.publicKeyB64)
|
||||
})
|
||||
|
||||
it('does not reopen a new host when an old host cache removal finishes', async () => {
|
||||
native.openSession.mockResolvedValue(SESSION_A)
|
||||
await mount('disconnected')
|
||||
const removal = deferred<void>()
|
||||
removeHostCache.mockReturnValue(removal.promise)
|
||||
await act(async () => {
|
||||
packageSession?.handleLoadFailure('mobile_web_generation_invalid')
|
||||
await flushPromises()
|
||||
})
|
||||
native.openSession.mockResolvedValue(SESSION_B)
|
||||
await update('disconnected', HOST_B)
|
||||
native.openSession.mockClear()
|
||||
native.closeSession.mockClear()
|
||||
|
||||
await act(async () => {
|
||||
removal.resolve()
|
||||
await flushPromises()
|
||||
})
|
||||
|
||||
expect(packageSession?.session).toEqual(SESSION_B)
|
||||
expect(native.openSession).not.toHaveBeenCalled()
|
||||
expect(native.closeSession).not.toHaveBeenCalled()
|
||||
})
|
||||
|
||||
it('allows another drop after the host is selected again', async () => {
|
||||
native.openSession.mockResolvedValue(SESSION_A)
|
||||
downloadPackage.mockResolvedValue({ commit: { buildId: SESSION_A.buildId } })
|
||||
@@ -376,6 +423,26 @@ describe('useMobileWebPackageSession', () => {
|
||||
expect(native.closeSession).toHaveBeenCalledWith(SESSION_A.sessionId)
|
||||
})
|
||||
|
||||
it('closes a refreshed session opened after its connection was abandoned', async () => {
|
||||
const refreshed = deferred<typeof SESSION_B>()
|
||||
native.openSession.mockImplementation((_host: string, buildId: string | null) =>
|
||||
buildId ? refreshed.promise : Promise.resolve(SESSION_A)
|
||||
)
|
||||
downloadPackage.mockResolvedValue({ commit: { buildId: SESSION_B.buildId } })
|
||||
await mount('connected')
|
||||
expect(packageSession?.session).toEqual(SESSION_A)
|
||||
await update('disconnected')
|
||||
|
||||
await act(async () => {
|
||||
refreshed.resolve(SESSION_B)
|
||||
await flushPromises()
|
||||
})
|
||||
|
||||
expect(packageSession?.session).toEqual(SESSION_A)
|
||||
expect(native.closeSession).toHaveBeenCalledWith(SESSION_B.sessionId)
|
||||
expect(native.closeSession).not.toHaveBeenCalledWith(SESSION_A.sessionId)
|
||||
})
|
||||
|
||||
it('closes the previous host session when the selected host changes', async () => {
|
||||
native.openSession.mockResolvedValue(SESSION_A)
|
||||
await mount('disconnected')
|
||||
|
||||
@@ -37,8 +37,8 @@ export function useMobileWebPackageSession({
|
||||
const hostEpochRef = useRef(0)
|
||||
const ownedSessionRef = useRef<MobileWebShellSession | null>(null)
|
||||
const cachedBuildRef = useRef<Promise<string | null>>(Promise.resolve(null))
|
||||
// One re-download per host selection: a page that cannot load will not loop over the bundle.
|
||||
const droppedGenerationRef = useRef(false)
|
||||
// One re-download per build: a desktop upgrade gets its own recovery attempt.
|
||||
const droppedGenerationRef = useRef<string | null>(null)
|
||||
const droppedHostRef = useRef<string | undefined>(undefined)
|
||||
const retryRef = useRef({ hostId: '', loadEpoch: -1, attempts: 0 })
|
||||
const connectionId = currentConnectionId(client)
|
||||
@@ -105,7 +105,7 @@ export function useMobileWebPackageSession({
|
||||
hostEpochRef.current = hostEpoch
|
||||
if (droppedHostRef.current !== host?.id) {
|
||||
droppedHostRef.current = host?.id
|
||||
droppedGenerationRef.current = false
|
||||
droppedGenerationRef.current = null
|
||||
}
|
||||
dispatch({ type: 'reopening', hasHost: Boolean(host) })
|
||||
const closing = ownedSessionRef.current
|
||||
@@ -242,7 +242,7 @@ export function useMobileWebPackageSession({
|
||||
return
|
||||
}
|
||||
mobileWebDiagnosticsStore.warning(hostId, reason ?? 'mobile_web_document_unavailable')
|
||||
if (droppedGenerationRef.current || !owned) {
|
||||
if (!owned || droppedGenerationRef.current === owned.buildId) {
|
||||
dispatch({
|
||||
type: 'warning',
|
||||
warning: { message: 'Couldn’t open Orca.', code: reason }
|
||||
@@ -251,8 +251,8 @@ export function useMobileWebPackageSession({
|
||||
}
|
||||
// The cached generation cannot render, so it is deleted and downloaded again; an unreachable
|
||||
// desktop leaves the shell in its offline state until the connection returns.
|
||||
droppedGenerationRef.current = true
|
||||
hostEpochRef.current += 1
|
||||
droppedGenerationRef.current = owned.buildId
|
||||
const dropEpoch = ++hostEpochRef.current
|
||||
ownedSessionRef.current = null
|
||||
dispatch({
|
||||
type: 'generation-dropping',
|
||||
@@ -261,7 +261,9 @@ export function useMobileWebPackageSession({
|
||||
void (async () => {
|
||||
await ExpoMobileWebShell.closeSession(owned.sessionId).catch(() => {})
|
||||
await removeMobileWebHostCache(host.publicKeyB64).catch(() => {})
|
||||
dispatch({ type: 'reload' })
|
||||
if (hostEpochRef.current === dropEpoch) {
|
||||
dispatch({ type: 'reload' })
|
||||
}
|
||||
})()
|
||||
},
|
||||
[host?.id, host?.publicKeyB64, host?.name]
|
||||
|
||||
@@ -8,10 +8,7 @@ import {
|
||||
buildTerminalUnsubscribeParams,
|
||||
updateTerminalSubscriptionViewport
|
||||
} from './rpc-client-terminal-subscription'
|
||||
import {
|
||||
buildReadyStreamUnsubscribe,
|
||||
buildServerSubscriptionUnsubscribe
|
||||
} from './rpc-client-server-subscription'
|
||||
import { buildServerSubscriptionUnsubscribe } from './rpc-client-server-subscription'
|
||||
import type { RpcClient } from './rpc-client'
|
||||
import { routeTerminalMultiplexFrame } from './rpc-client-terminal-multiplex'
|
||||
import type { RpcResponse, RpcSuccess } from './types'
|
||||
@@ -76,13 +73,6 @@ export class MobileRelayRpcStreams {
|
||||
listener: (result: unknown) => void,
|
||||
subscribeOptions?: Parameters<RpcClient['subscribe']>[3]
|
||||
): () => void {
|
||||
if (
|
||||
subscribeOptions?.serverUnsubscribeMethod &&
|
||||
this.streams.size + this.cancelledSubscriptions.size >= 128
|
||||
) {
|
||||
listener({ type: 'error', message: 'Stream capacity reached' })
|
||||
return () => {}
|
||||
}
|
||||
const id = this.options.nextId()
|
||||
const stream: StreamRecord = {
|
||||
method,
|
||||
@@ -131,7 +121,7 @@ export class MobileRelayRpcStreams {
|
||||
this.options.sendFrame({ id: this.options.nextId(), ...cancelled.unsubscribe })
|
||||
} else if (typeof result.subscriptionId === 'string') {
|
||||
this.cancelledSubscriptions.delete(response.id)
|
||||
const unsubscribe = buildReadyStreamUnsubscribe(
|
||||
const unsubscribe = buildServerSubscriptionUnsubscribe(
|
||||
cancelled.method,
|
||||
result.subscriptionId,
|
||||
cancelled.serverUnsubscribeMethod
|
||||
@@ -229,7 +219,7 @@ export class MobileRelayRpcStreams {
|
||||
}
|
||||
} else {
|
||||
const unsubscribe = stream.subscriptionId
|
||||
? buildReadyStreamUnsubscribe(
|
||||
? buildServerSubscriptionUnsubscribe(
|
||||
stream.method,
|
||||
stream.subscriptionId,
|
||||
stream.serverUnsubscribeMethod
|
||||
|
||||
@@ -68,18 +68,27 @@ describe('Desktop-advertised stream cleanup', () => {
|
||||
)
|
||||
})
|
||||
|
||||
it('bounds cancelled streams that have not supplied their cleanup token yet', () => {
|
||||
const { registry } = fixture()
|
||||
for (let index = 0; index < 128; index++) {
|
||||
registry.subscribe('future.watch', {}, () => {}, {
|
||||
it.each([true, false])(
|
||||
'removes ended streams before reconnect replay: streaming=%s',
|
||||
(streaming) => {
|
||||
const { registry, send, ready } = fixture()
|
||||
const listener = vi.fn()
|
||||
registry.subscribe('future.watch', {}, listener, {
|
||||
serverUnsubscribeMethod: 'future.release'
|
||||
})()
|
||||
})
|
||||
ready('lease')
|
||||
registry.handleResponse({
|
||||
id: '1',
|
||||
ok: true,
|
||||
streaming,
|
||||
result: { type: 'end' },
|
||||
_meta: { runtimeId: 'host' }
|
||||
})
|
||||
expect(listener).toHaveBeenLastCalledWith({ type: 'end' })
|
||||
expect(registry.size()).toBe(0)
|
||||
registry.markForReplay()
|
||||
registry.replayAfterAuthentication()
|
||||
expect(send).toHaveBeenCalledOnce()
|
||||
}
|
||||
const listener = vi.fn()
|
||||
registry.subscribe('future.watch', {}, listener, { serverUnsubscribeMethod: 'future.release' })
|
||||
expect(registry.size()).toBe(128)
|
||||
expect(listener).toHaveBeenCalledWith(
|
||||
expect.objectContaining({ type: 'error', error: { code: 'rate_limited' } })
|
||||
)
|
||||
})
|
||||
)
|
||||
})
|
||||
|
||||
@@ -13,11 +13,3 @@ export function buildServerSubscriptionUnsubscribe(
|
||||
}[method]
|
||||
return unsubscribeMethod ? { method: unsubscribeMethod, params: { subscriptionId } } : null
|
||||
}
|
||||
|
||||
export function buildReadyStreamUnsubscribe(
|
||||
method: string,
|
||||
subscriptionId: string,
|
||||
cleanupMethod?: string
|
||||
): { method: string; params: { subscriptionId: string } } | null {
|
||||
return buildServerSubscriptionUnsubscribe(method, subscriptionId, cleanupMethod)
|
||||
}
|
||||
|
||||
@@ -7,7 +7,7 @@ import {
|
||||
buildTerminalUnsubscribeParams,
|
||||
updateTerminalSubscriptionViewport
|
||||
} from './rpc-client-terminal-subscription'
|
||||
import { buildReadyStreamUnsubscribe } from './rpc-client-server-subscription'
|
||||
import { buildServerSubscriptionUnsubscribe } from './rpc-client-server-subscription'
|
||||
import {
|
||||
isStreamingSubscriptionReadyResult,
|
||||
isTerminalSubscribedResult
|
||||
@@ -58,14 +58,6 @@ export class RpcClientStreamRegistry {
|
||||
listener: RpcStreamingListener,
|
||||
subscribeOptions?: RpcStreamSubscribeOptions
|
||||
): () => void {
|
||||
if (subscribeOptions?.serverUnsubscribeMethod && this.streams.size >= 128) {
|
||||
listener({
|
||||
type: 'error',
|
||||
error: { code: 'rate_limited' },
|
||||
message: 'Too many pending streams'
|
||||
})
|
||||
return () => {}
|
||||
}
|
||||
const id = this.options.nextId()
|
||||
const stream: StreamRequest = {
|
||||
method,
|
||||
@@ -127,10 +119,6 @@ export class RpcClientStreamRegistry {
|
||||
}
|
||||
|
||||
handleResponse(response: RpcResponse): boolean {
|
||||
if (response.ok && response.streaming === true) {
|
||||
this.handleStreamingResponse(response)
|
||||
return true
|
||||
}
|
||||
const stream = this.streams.get(response.id)
|
||||
if (response.ok) {
|
||||
const result = (response as RpcSuccess).result as Record<string, unknown> | null
|
||||
@@ -141,6 +129,10 @@ export class RpcClientStreamRegistry {
|
||||
this.remove(response.id)
|
||||
return true
|
||||
}
|
||||
if (response.streaming === true) {
|
||||
this.handleStreamingResponse(response)
|
||||
return true
|
||||
}
|
||||
if (stream && result?.type === 'scrollback') {
|
||||
stream.listener(result)
|
||||
return true
|
||||
@@ -256,12 +248,11 @@ export class RpcClientStreamRegistry {
|
||||
if (!stream.subscriptionId) {
|
||||
return
|
||||
}
|
||||
const unsubscribe = stream.serverUnsubscribeMethod
|
||||
? {
|
||||
method: stream.serverUnsubscribeMethod,
|
||||
params: { subscriptionId: stream.subscriptionId }
|
||||
}
|
||||
: buildReadyStreamUnsubscribe(stream.method, stream.subscriptionId)
|
||||
const unsubscribe = buildServerSubscriptionUnsubscribe(
|
||||
stream.method,
|
||||
stream.subscriptionId,
|
||||
stream.serverUnsubscribeMethod
|
||||
)
|
||||
if (unsubscribe) {
|
||||
this.sendRpc(unsubscribe.method, unsubscribe.params)
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user