From f0234bbbed4b9ddab0f630cde5ee8966ebdf2363 Mon Sep 17 00:00:00 2001 From: Wenzhao Hu Date: Tue, 26 Nov 2024 19:13:43 +0800 Subject: [PATCH] fix: fix client may miss the INITIALIZE event (#4158) --- .../rpc/src/services/rpc/channel.service.ts | 3 +- packages/rpc/src/services/rpc/rpc.service.ts | 30 ++++++++++++++----- 2 files changed, 24 insertions(+), 9 deletions(-) diff --git a/packages/rpc/src/services/rpc/channel.service.ts b/packages/rpc/src/services/rpc/channel.service.ts index 5de377bc72..25c56163f6 100644 --- a/packages/rpc/src/services/rpc/channel.service.ts +++ b/packages/rpc/src/services/rpc/channel.service.ts @@ -15,9 +15,8 @@ */ import type { IDisposable } from '@univerjs/core'; -import { createIdentifier } from '@univerjs/core'; - import type { IChannel, IMessageProtocol } from './rpc.service'; +import { createIdentifier } from '@univerjs/core'; import { ChannelClient, ChannelServer } from './rpc.service'; export interface IRPCChannelService { diff --git a/packages/rpc/src/services/rpc/rpc.service.ts b/packages/rpc/src/services/rpc/rpc.service.ts index e03eb5b031..b53f3c4f67 100644 --- a/packages/rpc/src/services/rpc/rpc.service.ts +++ b/packages/rpc/src/services/rpc/rpc.service.ts @@ -122,14 +122,22 @@ export interface IChannelClient { getChannel(channelName: string): T; } -/** - * - */ export interface IChannelServer { registerChannel(channelName: string, channel: T): void; } enum RequestType { + /** + * In Univer, we cannot make sure that when IPCServer constructs, the process (or thread) + * where the corresponding IPCClient residents has bootstrapped and been ready to recieve messages. + * This may result in the IPCClient hanging there, waiting for the `INITIALIZE` message that it has + * already missed. So the client should send a REQUEST_INITIALIZATION in case of that. + * + * Later, we may want a more sophisticated RPC system where the server can serve more than + * one clients, and this event may be removed. + */ + REQUEST_INITIALIZATION = 50, + /** A simple remote calling wrapper in a Promise. */ CALL = 100, @@ -168,7 +176,8 @@ interface IRPCResponse { /** It should be the same as its corresponding requests' `seq`. */ seq: number; type: ResponseType; - data?: any; + + data?: any; // TODO: replace it with ISerializable. } interface IResponseHandler { @@ -187,8 +196,7 @@ export class ChannelClient extends RxDisposable implements IChannelClient { constructor(private readonly _protocol: IMessageProtocol) { super(); - // TODO: subscribe to the state of the protocol and see if it is connected \ - // and initialized. + this._protocol.send({ type: RequestType.REQUEST_INITIALIZATION }); this._protocol.onMessage.pipe(takeUntil(this.dispose$)).subscribe((message) => this._onMessage(message)); } @@ -334,7 +342,7 @@ export class ChannelServer extends RxDisposable implements IChannelServer { super(); this._protocol.onMessage.pipe(takeUntil(this.dispose$)).subscribe((message) => this._onRequest(message)); - this._sendResponse({ seq: -1, type: ResponseType.INITIALIZE }); + this._sendInitialize(); } override dispose(): void { @@ -350,6 +358,9 @@ export class ChannelServer extends RxDisposable implements IChannelServer { private _onRequest(request: IRPCRequest): void { switch (request.type) { + case RequestType.REQUEST_INITIALIZATION: + this._sendInitialize(); + break; case RequestType.CALL: this._onMethodCall(request); break; @@ -364,6 +375,10 @@ export class ChannelServer extends RxDisposable implements IChannelServer { } } + private _sendInitialize(): void { + this._sendResponse({ seq: -1, type: ResponseType.INITIALIZE }); + } + private _onMethodCall(request: IRPCRequest): void { const { channelName, method, args } = request; const channel = this._channels.get(channelName); @@ -436,3 +451,4 @@ export class ChannelServer extends RxDisposable implements IChannelServer { this._protocol.send(response); } } +