feat: msg push

This commit is contained in:
手瓜一十雪
2024-11-14 20:18:19 +08:00
parent a0a50755d3
commit 4487db4e0a
15 changed files with 345 additions and 245 deletions

View File

@@ -6,6 +6,7 @@ export type OB11EmitEventContent = OB11BaseEvent | OB11Message;
export interface IOB11NetworkAdapter {
actions?: ActionMap;
name: string;
onEvent<T extends OB11EmitEventContent>(event: T): void;
@@ -15,19 +16,34 @@ export interface IOB11NetworkAdapter {
}
export class OB11NetworkManager {
adapters: IOB11NetworkAdapter[] = [];
adapters: Map<string, IOB11NetworkAdapter> = new Map();
async openAllAdapters() {
return Promise.all(this.adapters.map(adapter => adapter.open()));
return Promise.all(Array.from(this.adapters.values()).map(adapter => adapter.open()));
}
async emitEvent(event: OB11EmitEventContent) {
//console.log('adapters', this.adapters.length);
return Promise.all(this.adapters.map(adapter => adapter.onEvent(event)));
return Promise.all(Array.from(this.adapters.values()).map(adapter => adapter.onEvent(event)));
}
async emitEventByName(names: string[], event: OB11EmitEventContent) {
return Promise.all(names.map(name => {
const adapter = this.adapters.get(name);
if (adapter) {
return adapter.onEvent(event);
}
}));
}
async emitEventByNames(map:Map<string,OB11EmitEventContent>){
return Promise.all(Array.from(map.entries()).map(([name, event]) => {
const adapter = this.adapters.get(name);
if (adapter) {
return adapter.onEvent(event);
}
}));
}
registerAdapter(adapter: IOB11NetworkAdapter) {
this.adapters.push(adapter);
this.adapters.set(adapter.name, adapter);
}
async registerAdapterAndOpen(adapter: IOB11NetworkAdapter) {
@@ -36,24 +52,28 @@ export class OB11NetworkManager {
}
async closeSomeAdapters(adaptersToClose: IOB11NetworkAdapter[]) {
this.adapters = this.adapters.filter(adapter => !adaptersToClose.includes(adapter));
await Promise.all(adaptersToClose.map(adapter => adapter.close()));
for (const adapter of adaptersToClose) {
this.adapters.delete(adapter.name);
await adapter.close();
}
}
findSomeAdapter(name: string) {
return this.adapters.get(name);
}
/**
* Close all adapters that satisfy the predicate.
*/
async closeAdapterByPredicate(closeFilter: (adapter: IOB11NetworkAdapter) => boolean) {
await this.closeSomeAdapters(this.adapters.filter(closeFilter));
const adaptersToClose = Array.from(this.adapters.values()).filter(closeFilter);
await this.closeSomeAdapters(adaptersToClose);
}
async closeAllAdapters() {
await Promise.all(this.adapters.map(adapter => adapter.close()));
this.adapters = [];
await Promise.all(Array.from(this.adapters.values()).map(adapter => adapter.close()));
this.adapters.clear();
}
}
export * from './active-http';
export * from './active-websocket';
export * from './passive-http';
export * from './passive-websocket';
export * from './passive-websocket';