92 lines
2.2 KiB
TypeScript
92 lines
2.2 KiB
TypeScript
|
|
// Copyright (C) 2026 Whiterun LLC,
|
||
|
|
// This software is licensed under the GNU Lesser General Public License (LGPL), version 3.0 or later.
|
||
|
|
// A copy of the license can be found in the LICENSE file or at https://www.gnu.org/licenses/lgpl-3.0.html
|
||
|
|
|
||
|
|
import { ProtocolMessage } from "./protocols/hdwalletv1.js";
|
||
|
|
import { debug, Scope } from "./log.js";
|
||
|
|
|
||
|
|
interface QueuedMessage {
|
||
|
|
message: ProtocolMessage;
|
||
|
|
resolve: () => void;
|
||
|
|
reject: (_error: Error) => void;
|
||
|
|
}
|
||
|
|
|
||
|
|
export class MessageQueue {
|
||
|
|
private queue: QueuedMessage[] = [];
|
||
|
|
private isReady: boolean = false;
|
||
|
|
private logActivity: boolean = false;
|
||
|
|
|
||
|
|
constructor(options?: { logActivity?: boolean }) {
|
||
|
|
this.logActivity = options?.logActivity ?? false;
|
||
|
|
}
|
||
|
|
|
||
|
|
getReady(): boolean {
|
||
|
|
return this.isReady;
|
||
|
|
}
|
||
|
|
|
||
|
|
getQueueLength(): number {
|
||
|
|
return this.queue.length;
|
||
|
|
}
|
||
|
|
|
||
|
|
enqueue(message: ProtocolMessage): Promise<void> {
|
||
|
|
if (this.isReady) {
|
||
|
|
return Promise.resolve();
|
||
|
|
}
|
||
|
|
|
||
|
|
if (this.logActivity) {
|
||
|
|
debug(
|
||
|
|
Scope.Relay,
|
||
|
|
`net: Relays not ready, queuing message ${message.action}`,
|
||
|
|
);
|
||
|
|
}
|
||
|
|
|
||
|
|
return new Promise<void>((resolve, reject) => {
|
||
|
|
this.queue.push({ message, resolve, reject });
|
||
|
|
});
|
||
|
|
}
|
||
|
|
|
||
|
|
async setReady(
|
||
|
|
publishFn: (_message: ProtocolMessage) => Promise<void>,
|
||
|
|
): Promise<void> {
|
||
|
|
this.isReady = true;
|
||
|
|
|
||
|
|
const queuedMessages = [...this.queue];
|
||
|
|
this.queue = [];
|
||
|
|
|
||
|
|
if (queuedMessages.length > 0 && this.logActivity) {
|
||
|
|
debug(Scope.Relay, `Processing ${queuedMessages.length} queued messages`);
|
||
|
|
}
|
||
|
|
|
||
|
|
for (const queued of queuedMessages) {
|
||
|
|
try {
|
||
|
|
await publishFn(queued.message);
|
||
|
|
queued.resolve();
|
||
|
|
} catch (error) {
|
||
|
|
queued.reject(error as Error);
|
||
|
|
}
|
||
|
|
}
|
||
|
|
}
|
||
|
|
|
||
|
|
setNotReady(
|
||
|
|
errorMessage: string = "Connection closed before message could be sent",
|
||
|
|
): void {
|
||
|
|
this.isReady = false;
|
||
|
|
|
||
|
|
const queuedMessages = [...this.queue];
|
||
|
|
this.queue = [];
|
||
|
|
|
||
|
|
for (const queued of queuedMessages) {
|
||
|
|
queued.reject(new Error(errorMessage));
|
||
|
|
}
|
||
|
|
}
|
||
|
|
|
||
|
|
clear(): void {
|
||
|
|
const queuedMessages = [...this.queue];
|
||
|
|
this.queue = [];
|
||
|
|
|
||
|
|
for (const queued of queuedMessages) {
|
||
|
|
queued.reject(new Error("Queue cleared"));
|
||
|
|
}
|
||
|
|
}
|
||
|
|
}
|