Compare commits
No commits in common. "b274b82d3c66db0b0073e8c5b6a098daf01e1076" and "62058d31fd94887be74f275222e6721187f1d812" have entirely different histories.
b274b82d3c
...
62058d31fd
|
@ -40,13 +40,13 @@ export declare class Socket extends Client {
|
||||||
private _rawStream?;
|
private _rawStream?;
|
||||||
get rawStream(): Uint8Array;
|
get rawStream(): Uint8Array;
|
||||||
setup(): Promise<void>;
|
setup(): Promise<void>;
|
||||||
|
private _initSync;
|
||||||
on<T extends EventEmitter.EventNames<string | symbol>>(event: T, fn: EventEmitter.EventListener<string | symbol, T>, context?: any): this;
|
on<T extends EventEmitter.EventNames<string | symbol>>(event: T, fn: EventEmitter.EventListener<string | symbol, T>, context?: any): this;
|
||||||
off<T extends EventEmitter.EventNames<string | symbol>>(event: T, fn?: EventEmitter.EventListener<string | symbol, T>, context?: any, once?: boolean): this;
|
off<T extends EventEmitter.EventNames<string | symbol>>(event: T, fn?: EventEmitter.EventListener<string | symbol, T>, context?: any, once?: boolean): this;
|
||||||
write(message: string | Buffer): void;
|
write(message: string | Buffer): void;
|
||||||
end(): void;
|
end(): void;
|
||||||
private ensureEvent;
|
private ensureEvent;
|
||||||
private trackEvent;
|
private trackEvent;
|
||||||
syncProtomux(action: string, id: number): Promise<any>;
|
|
||||||
}
|
}
|
||||||
export declare const createClient: (...args: any) => SwarmClient;
|
export declare const createClient: (...args: any) => SwarmClient;
|
||||||
//# sourceMappingURL=index.d.ts.map
|
//# sourceMappingURL=index.d.ts.map
|
|
@ -1 +1 @@
|
||||||
{"version":3,"file":"index.d.ts","sourceRoot":"","sources":["../src/index.ts"],"names":[],"mappings":";AAAA,OAAO,EAAE,MAAM,EAAE,MAAM,QAAQ,CAAC;AAChC,OAAO,EAAE,MAAM,EAAW,MAAM,8BAA8B,CAAC;AAC/D,OAAO,EAAU,QAAQ,EAAY,MAAM,gBAAgB,CAAC;AAI5D,OAAO,KAAK,EAAE,YAAY,EAAE,MAAM,eAAe,CAAC;AASlD,qBAAa,WAAY,SAAQ,MAAM;IACrC,OAAO,CAAC,eAAe,CAAU;IACjC,OAAO,CAAC,EAAE,CAAa;IACvB,OAAO,CAAC,cAAc,CAAU;IAChC,OAAO,CAAC,eAAe,CAAM;IAE7B,OAAO,CAAC,MAAM,CAAC,CAAgB;IAC/B,OAAO,CAAC,mBAAmB,CAAC,CAG1B;IAEF,OAAO,CAAC,OAAO,CAA0C;IACzD,OAAO,CAAC,QAAQ,CAAkD;IAElE,IAAI,GAAG;;MAON;gBAEW,aAAa,UAAO,EAAE,aAAa,UAAQ;IAcvD,IAAI,KAAK,IAAI,MAAM,GAAG,SAAS,CAE9B;IAEY,OAAO,CAAC,MAAM,EAAE,MAAM,GAAG,UAAU,GAAG,OAAO,CAAC,MAAM,CAAC;IAiB5D,IAAI,IAAI,OAAO,CAAC,QAAQ,CAAC;IAIzB,KAAK,IAAI,OAAO,CAAC,IAAI,CAAC;IAkBtB,KAAK,IAAI,OAAO,CAAC,IAAI,CAAC;YAMd,OAAO;IA2BR,QAAQ,CAAC,MAAM,EAAE,MAAM,GAAG,OAAO,CAAC,IAAI,CAAC;IAIvC,WAAW,CAAC,MAAM,EAAE,MAAM,GAAG,OAAO,CAAC,IAAI,CAAC;IAI1C,WAAW,IAAI,OAAO,CAAC,IAAI,CAAC;IAI5B,SAAS,IAAI,OAAO,CAAC,MAAM,EAAE,CAAC;IAI9B,IAAI,CAAC,KAAK,EAAE,MAAM,GAAG,UAAU,GAAG,MAAM,GAAG,OAAO,CAAC,IAAI,CAAC;CAQtE;AAQD,qBAAa,MAAO,SAAQ,MAAM;IAChC,OAAO,CAAC,EAAE,CAAS;IACnB,OAAO,CAAC,YAAY,CAAqC;IAEzD,OAAO,CAAC,SAAS,CAAe;IAEhC,OAAO,CAAC,KAAK,CAAc;IAC3B,OAAO,CAAC,QAAQ,CAAC,CAAa;gBAElB,EAAE,EAAE,MAAM,EAAE,KAAK,EAAE,WAAW;IAM1C,OAAO,CAAC,gBAAgB,CAAC,CAAa;IAEtC,IAAI,eAAe,IAAI,UAAU,CAEhC;IAED,OAAO,CAAC,UAAU,CAAC,CAAa;IAEhC,IAAI,SAAS,IAAI,UAAU,CAE1B;IAEK,KAAK;IAQX,EAAE,CAAC,CAAC,SAAS,YAAY,CAAC,UAAU,CAAC,MAAM,GAAG,MAAM,CAAC,EACnD,KAAK,EAAE,CAAC,EACR,EAAE,EAAE,YAAY,CAAC,aAAa,CAAC,MAAM,GAAG,MAAM,EAAE,CAAC,CAAC,EAClD,OAAO,CAAC,EAAE,GAAG,GACZ,IAAI;IAiBP,GAAG,CAAC,CAAC,SAAS,YAAY,CAAC,UAAU,CAAC,MAAM,GAAG,MAAM,CAAC,EACpD,KAAK,EAAE,CAAC,EACR,EAAE,CAAC,EAAE,YAAY,CAAC,aAAa,CAAC,MAAM,GAAG,MAAM,EAAE,CAAC,CAAC,EACnD,OAAO,CAAC,EAAE,GAAG,EACb,IAAI,CAAC,EAAE,OAAO,GACb,IAAI;IASP,KAAK,CAAC,OAAO,EAAE,MAAM,GAAG,MAAM,GAAG,IAAI;IAIrC,GAAG,IAAI,IAAI;IAUX,OAAO,CAAC,WAAW;IAMnB,OAAO,CAAC,UAAU;IAKL,YAAY,CAAC,MAAM,EAAE,MAAM,EAAE,EAAE,EAAE,MAAM;CAGrD;AAID,eAAO,MAAM,YAAY,+BAA4C,CAAC"}
|
{"version":3,"file":"index.d.ts","sourceRoot":"","sources":["../src/index.ts"],"names":[],"mappings":";AAAA,OAAO,EAAE,MAAM,EAAE,MAAM,QAAQ,CAAC;AAChC,OAAO,EAAE,MAAM,EAAW,MAAM,8BAA8B,CAAC;AAC/D,OAAO,EAAU,QAAQ,EAAY,MAAM,gBAAgB,CAAC;AAI5D,OAAO,KAAK,EAAE,YAAY,EAAE,MAAM,eAAe,CAAC;AASlD,qBAAa,WAAY,SAAQ,MAAM;IACrC,OAAO,CAAC,eAAe,CAAU;IACjC,OAAO,CAAC,EAAE,CAAa;IACvB,OAAO,CAAC,cAAc,CAAU;IAChC,OAAO,CAAC,eAAe,CAAM;IAE7B,OAAO,CAAC,MAAM,CAAC,CAAgB;IAC/B,OAAO,CAAC,mBAAmB,CAAC,CAG1B;IAEF,OAAO,CAAC,OAAO,CAA0C;IACzD,OAAO,CAAC,QAAQ,CAAkD;IAElE,IAAI,GAAG;;MAON;gBAEW,aAAa,UAAO,EAAE,aAAa,UAAQ;IAcvD,IAAI,KAAK,IAAI,MAAM,GAAG,SAAS,CAE9B;IAEY,OAAO,CAAC,MAAM,EAAE,MAAM,GAAG,UAAU,GAAG,OAAO,CAAC,MAAM,CAAC;IAiB5D,IAAI,IAAI,OAAO,CAAC,QAAQ,CAAC;IAIzB,KAAK,IAAI,OAAO,CAAC,IAAI,CAAC;IAkBtB,KAAK,IAAI,OAAO,CAAC,IAAI,CAAC;YAMd,OAAO;IA2BR,QAAQ,CAAC,MAAM,EAAE,MAAM,GAAG,OAAO,CAAC,IAAI,CAAC;IAIvC,WAAW,CAAC,MAAM,EAAE,MAAM,GAAG,OAAO,CAAC,IAAI,CAAC;IAI1C,WAAW,IAAI,OAAO,CAAC,IAAI,CAAC;IAI5B,SAAS,IAAI,OAAO,CAAC,MAAM,EAAE,CAAC;IAI9B,IAAI,CAAC,KAAK,EAAE,MAAM,GAAG,UAAU,GAAG,MAAM,GAAG,OAAO,CAAC,IAAI,CAAC;CAQtE;AAQD,qBAAa,MAAO,SAAQ,MAAM;IAChC,OAAO,CAAC,EAAE,CAAS;IACnB,OAAO,CAAC,YAAY,CAAqC;IAEzD,OAAO,CAAC,SAAS,CAAe;IAEhC,OAAO,CAAC,KAAK,CAAc;IAC3B,OAAO,CAAC,QAAQ,CAAC,CAAa;gBAElB,EAAE,EAAE,MAAM,EAAE,KAAK,EAAE,WAAW;IAM1C,OAAO,CAAC,gBAAgB,CAAC,CAAa;IAEtC,IAAI,eAAe,IAAI,UAAU,CAEhC;IAED,OAAO,CAAC,UAAU,CAAC,CAAa;IAEhC,IAAI,SAAS,IAAI,UAAU,CAE1B;IAEK,KAAK;YASG,SAAS;IAgEvB,EAAE,CAAC,CAAC,SAAS,YAAY,CAAC,UAAU,CAAC,MAAM,GAAG,MAAM,CAAC,EACnD,KAAK,EAAE,CAAC,EACR,EAAE,EAAE,YAAY,CAAC,aAAa,CAAC,MAAM,GAAG,MAAM,EAAE,CAAC,CAAC,EAClD,OAAO,CAAC,EAAE,GAAG,GACZ,IAAI;IAiBP,GAAG,CAAC,CAAC,SAAS,YAAY,CAAC,UAAU,CAAC,MAAM,GAAG,MAAM,CAAC,EACpD,KAAK,EAAE,CAAC,EACR,EAAE,CAAC,EAAE,YAAY,CAAC,aAAa,CAAC,MAAM,GAAG,MAAM,EAAE,CAAC,CAAC,EACnD,OAAO,CAAC,EAAE,GAAG,EACb,IAAI,CAAC,EAAE,OAAO,GACb,IAAI;IASP,KAAK,CAAC,OAAO,EAAE,MAAM,GAAG,MAAM,GAAG,IAAI;IAIrC,GAAG,IAAI,IAAI;IAUX,OAAO,CAAC,WAAW;IAMnB,OAAO,CAAC,UAAU;CAInB;AAID,eAAO,MAAM,YAAY,+BAA4C,CAAC"}
|
|
@ -7,6 +7,7 @@ import Backoff from "backoff.js";
|
||||||
import { Mutex } from "async-mutex";
|
import { Mutex } from "async-mutex";
|
||||||
// @ts-ignore
|
// @ts-ignore
|
||||||
import Protomux from "protomux";
|
import Protomux from "protomux";
|
||||||
|
import defer from "p-defer";
|
||||||
export class SwarmClient extends Client {
|
export class SwarmClient extends Client {
|
||||||
useDefaultSwarm;
|
useDefaultSwarm;
|
||||||
id = 0;
|
id = 0;
|
||||||
|
@ -131,7 +132,57 @@ export class Socket extends Client {
|
||||||
let info = await this.callModuleReturn("socketGetInfo", { id: this.id });
|
let info = await this.callModuleReturn("socketGetInfo", { id: this.id });
|
||||||
this._remotePublicKey = info.remotePublicKey;
|
this._remotePublicKey = info.remotePublicKey;
|
||||||
this._rawStream = info.rawStream;
|
this._rawStream = info.rawStream;
|
||||||
Protomux.from(this, { slave: true });
|
this._initSync();
|
||||||
|
}
|
||||||
|
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) => {
|
||||||
|
if (data === true) {
|
||||||
|
updateSent.resolve();
|
||||||
|
updateSent = defer();
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
await this.syncMutex.acquire();
|
||||||
|
["remote", "local"].forEach((field) => {
|
||||||
|
const rField = `_${field}`;
|
||||||
|
data[field].forEach((item) => {
|
||||||
|
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) => {
|
||||||
|
if (mux._free[index] === null) {
|
||||||
|
mux._free[index] = undefined;
|
||||||
|
}
|
||||||
|
});
|
||||||
|
mux._free = mux._free.filter((item) => item !== undefined);
|
||||||
|
this.syncMutex.release();
|
||||||
|
setup.resolve();
|
||||||
|
updateReceived.resolve();
|
||||||
|
});
|
||||||
|
const send = async (mux) => {
|
||||||
|
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(event, fn, context) {
|
on(event, fn, context) {
|
||||||
const [update, promise] = this.connectModule("socketListenEvent", { id: this.id, event: event }, (data) => {
|
const [update, promise] = this.connectModule("socketListenEvent", { id: this.id, event: event }, (data) => {
|
||||||
|
@ -170,9 +221,6 @@ export class Socket extends Client {
|
||||||
this.ensureEvent(event);
|
this.ensureEvent(event);
|
||||||
this.eventUpdates[event].push(update);
|
this.eventUpdates[event].push(update);
|
||||||
}
|
}
|
||||||
async syncProtomux(action, id) {
|
|
||||||
return this.callModuleReturn("syncProtomux", { action, data: id });
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
const MODULE = "_AVKgzVYC8Sb_qiTA6kw5BDzQ4Ch-8D4sldQJl8dXF9oTw";
|
const MODULE = "_AVKgzVYC8Sb_qiTA6kw5BDzQ4Ch-8D4sldQJl8dXF9oTw";
|
||||||
export const createClient = factory(SwarmClient, MODULE);
|
export const createClient = factory(SwarmClient, MODULE);
|
||||||
|
|
71
src/index.ts
71
src/index.ts
|
@ -192,8 +192,73 @@ export class Socket extends Client {
|
||||||
this._remotePublicKey = info.remotePublicKey;
|
this._remotePublicKey = info.remotePublicKey;
|
||||||
this._rawStream = info.rawStream;
|
this._rawStream = info.rawStream;
|
||||||
|
|
||||||
Protomux.from(this, { slave: true });
|
this._initSync();
|
||||||
}
|
}
|
||||||
|
|
||||||
|
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>,
|
||||||
|
@ -253,10 +318,6 @@ 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";
|
||||||
|
|
Loading…
Reference in New Issue