-
Notifications
You must be signed in to change notification settings - Fork 367
Expand file tree
/
Copy pathmessage_queue.ts
More file actions
185 lines (158 loc) · 5.92 KB
/
Copy pathmessage_queue.ts
File metadata and controls
185 lines (158 loc) · 5.92 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
import EventEmitter from "eventemitter3";
import LoggerCore from "@App/app/logger/core";
import { type TMessage } from "./types";
import type { SystemConfigKey, SystemConfigValueType } from "@App/pkg/config/config";
export type TKeyValue<T extends SystemConfigKey> = {
key: T;
value: SystemConfigValueType<T> | undefined;
prev: SystemConfigValueType<T> | undefined;
};
// 中间件函数类型
type MiddlewareFunction<T = any> = (topic: string, message: T, next: () => void) => void | Promise<void>;
// 消息处理函数类型
type MessageHandler<T = any> = (message: T) => void;
export interface IMessageQueue {
// 创建子分组
group(name: string, middleware?: MiddlewareFunction): IMessageQueueExtended;
// 订阅消息
subscribe<T>(topic: string, handler: MessageHandler<T>): () => void;
// 发布消息
publish<T>(topic: string, message: NonNullable<T>): void;
// 只发布给当前环境
emit<T>(topic: string, message: NonNullable<T>): void;
}
export interface IMessageQueueExtra {
// 创建子分组
use(middleware: MiddlewareFunction): IMessageQueueExtended;
}
export interface IMessageQueueExtended extends IMessageQueue, IMessageQueueExtra {}
// 消息队列
export class MessageQueue implements IMessageQueue {
private EE = new EventEmitter<string, any>();
constructor() {
chrome.runtime.onMessage.addListener((msg: TMessage) => {
const lastError = chrome.runtime.lastError;
const topic = msg.msgQueue;
if (typeof topic !== "string") return;
if (lastError) {
console.error("chrome.runtime.lastError in chrome.runtime.onMessage:", lastError);
// 消息API发生错误因此不继续执行
return false;
}
this.handler(topic, msg.data);
});
}
handler(topic: string, { action, message }: { action: string; message: any }) {
LoggerCore.getInstance()
.logger({ service: "messageQueue" })
.trace("messageQueueHandler", { action, topic, message });
switch (action) {
case "message":
this.EE.emit(topic, message);
break;
default:
throw new Error("action not found");
}
}
subscribe<T>(topic: string, handler: (msg: T) => void) {
this.EE.on(topic, handler);
return this.EE.off.bind(this.EE, topic, handler) as () => void;
}
publish<T>(topic: string, message: NonNullable<T>) {
// chrome.runtime.sendMessage() 不带回调时返回 Promise。没有其它上下文在监听(例如尚未打开
// popup/options,或没有已注入的 content script)是完全正常的广播场景,不代表出错。
// Chrome 在这种情况下该 Promise 不会 reject;但 Firefox 会 reject 并抛出
// "Could not establish connection. Receiving end does not exist."——不接住就会变成
// 未处理的 Promise rejection。publish 本身是"广播给任何在监听的人",无人监听应静默忽略。
const messageQueueLogger = LoggerCore.getInstance().logger({ service: "messageQueue" });
chrome.runtime
.sendMessage({
msgQueue: topic,
data: { action: "message", message },
})
.catch((e) => {
const msg = JSON.stringify(e?.message || e);
if (msg.includes("Could not establish connection. Receiving end does not exist.")) {
messageQueueLogger.debug("No target audience for .publish", { msg });
} else {
messageQueueLogger.error("Unable to execute runtime.sendMessage for .publish", { msg });
}
});
this.EE.emit(topic, message);
//@ts-ignore
messageQueueLogger.trace("publish", { topic, message });
}
// 只发布给当前环境
emit<T>(topic: string, message: NonNullable<T>) {
this.EE.emit(topic, message);
}
// 创建分组
group(name: string, middleware?: MiddlewareFunction) {
return new MessageQueueGroup(this, name, middleware);
}
}
// 消息队列分组
export class MessageQueueGroup implements IMessageQueue {
private middlewares: MiddlewareFunction[] = [];
constructor(
private messageQueue: IMessageQueue,
private name: string,
middleware?: MiddlewareFunction
) {
if (!name.endsWith("/") && name.length > 0) {
this.name += "/";
}
if (middleware) {
this.middlewares.push(middleware);
}
}
// 创建子分组
group(name: string, middleware?: MiddlewareFunction) {
const newGroup = new MessageQueueGroup(this.messageQueue, `${this.name}${name}`, middleware);
// 继承父级的中间件
newGroup.middlewares = [...this.middlewares, ...newGroup.middlewares];
return newGroup;
}
// 添加中间件
use(middleware: MiddlewareFunction) {
this.middlewares.push(middleware);
return this;
}
// 订阅消息
subscribe<T>(topic: string, handler: MessageHandler<T>) {
const fullTopic = `${this.name}${topic}`;
if (this.middlewares.length === 0) {
// 没有中间件,直接订阅
return this.messageQueue.subscribe(fullTopic, handler);
} else {
// 有中间件,需要包装处理函数
const wrappedHandler = async (message: T) => {
let index = 0;
const next = async (): Promise<void> => {
if (index < this.middlewares.length) {
const middleware = this.middlewares[index++];
const result = middleware(fullTopic, message, next);
if (result instanceof Promise) {
await result;
}
} else {
// 所有中间件都执行完毕,执行最终的处理函数
handler(message);
}
};
await next();
};
return this.messageQueue.subscribe(fullTopic, wrappedHandler);
}
}
// 发布消息
publish<T>(topic: string, message: NonNullable<T>) {
const fullTopic = `${this.name}${topic}`;
this.messageQueue.publish(fullTopic, message);
}
// 只发布给当前环境
emit<T>(topic: string, message: NonNullable<T>) {
const fullTopic = `${this.name}${topic}`;
this.messageQueue.emit(fullTopic, message);
}
}