*Switch to simpler protomux sync approach

This commit is contained in:
Derrick Hammer 2023-04-06 16:36:13 -04:00
parent 62058d31fd
commit 9ae392e1ab
Signed by: pcfreak30
GPG Key ID: C997C339BE476FF2
1 changed files with 5 additions and 66 deletions

View File

@ -192,73 +192,8 @@ export class Socket extends Client {
this._remotePublicKey = info.remotePublicKey; this._remotePublicKey = info.remotePublicKey;
this._rawStream = info.rawStream; this._rawStream = info.rawStream;
this._initSync(); Protomux.from(this, { slave: true });
} }
private async _initSync() {
this.userData = null;
const mux = Protomux.from(this, { slave: true });
let updateSent = defer();
let updateReceived = defer();
const setup = defer();
const [update] = this.connectModule(
"syncProtomux",
{ id: this.id },
async (data: any) => {
if (data === true) {
updateSent.resolve();
updateSent = defer();
return;
}
await this.syncMutex.acquire();
["remote", "local"].forEach((field) => {
const rField = `_${field}`;
data[field].forEach((item: any) => {
if (!mux[rField][item]) {
while (item > mux[rField].length) {
mux[rField].push(null);
}
}
if (!mux[rField][item]) {
mux[rField][item] = null;
}
});
});
data.free.forEach((index: number) => {
if (mux._free[index] === null) {
mux._free[index] = undefined;
}
});
mux._free = mux._free.filter((item: any) => item !== undefined);
this.syncMutex.release();
setup.resolve();
updateReceived.resolve();
}
);
const send = async (mux: any) => {
update({
remote: Object.keys(mux._remote),
local: Object.keys(mux._local),
free: mux._free,
});
updateReceived = defer();
await updateSent.promise;
await updateReceived.promise;
};
mux.syncState = send.bind(undefined, mux);
await setup.promise;
this.swarm.emit("setup", this);
}
on<T extends EventEmitter.EventNames<string | symbol>>( on<T extends EventEmitter.EventNames<string | symbol>>(
event: T, event: T,
fn: EventEmitter.EventListener<string | symbol, T>, fn: EventEmitter.EventListener<string | symbol, T>,
@ -318,6 +253,10 @@ export class Socket extends Client {
this.ensureEvent(event as string); this.ensureEvent(event as string);
this.eventUpdates[event].push(update); this.eventUpdates[event].push(update);
} }
public async syncProtomux(action: string, id: number) {
return this.callModuleReturn("syncProtomux", { action, data: id });
}
} }
const MODULE = "_AVKgzVYC8Sb_qiTA6kw5BDzQ4Ch-8D4sldQJl8dXF9oTw"; const MODULE = "_AVKgzVYC8Sb_qiTA6kw5BDzQ4Ch-8D4sldQJl8dXF9oTw";