[temp] networker drain
This commit is contained in:
parent
2adb7e3528
commit
d7d045818c
@ -9,6 +9,8 @@
|
|||||||
* https://github.com/zhukov/webogram/blob/master/LICENSE
|
* https://github.com/zhukov/webogram/blob/master/LICENSE
|
||||||
*/
|
*/
|
||||||
|
|
||||||
|
import type { DcId } from "../types";
|
||||||
|
|
||||||
const App = {
|
const App = {
|
||||||
id: 1025907,
|
id: 1025907,
|
||||||
hash: '452b0359b988148995f22ff0f4229750',
|
hash: '452b0359b988148995f22ff0f4229750',
|
||||||
@ -17,7 +19,7 @@ const App = {
|
|||||||
langPack: 'macos',
|
langPack: 'macos',
|
||||||
langPackCode: 'en',
|
langPackCode: 'en',
|
||||||
domains: [] as string[],
|
domains: [] as string[],
|
||||||
baseDcId: 2,
|
baseDcId: 2 as DcId,
|
||||||
isMainDomain: location.hostname === 'web.telegram.org',
|
isMainDomain: location.hostname === 'web.telegram.org',
|
||||||
suffix: 'K'
|
suffix: 'K'
|
||||||
};
|
};
|
||||||
|
@ -16,6 +16,7 @@ import { notifyAll, notifySomeone } from "../../helpers/context";
|
|||||||
import { getFileNameByLocation } from "../../helpers/fileName";
|
import { getFileNameByLocation } from "../../helpers/fileName";
|
||||||
import { nextRandomInt } from "../../helpers/random";
|
import { nextRandomInt } from "../../helpers/random";
|
||||||
import { InputFile, InputFileLocation, UploadFile } from "../../layer";
|
import { InputFile, InputFileLocation, UploadFile } from "../../layer";
|
||||||
|
import { DcId } from "../../types";
|
||||||
import CacheStorageController from "../cacheStorage";
|
import CacheStorageController from "../cacheStorage";
|
||||||
import cryptoWorker from "../crypto/cryptoworker";
|
import cryptoWorker from "../crypto/cryptoworker";
|
||||||
import FileManager from "../filemanager";
|
import FileManager from "../filemanager";
|
||||||
@ -30,7 +31,7 @@ type Delayed = {
|
|||||||
};
|
};
|
||||||
|
|
||||||
export type DownloadOptions = {
|
export type DownloadOptions = {
|
||||||
dcId: number,
|
dcId: DcId,
|
||||||
location: InputFileLocation,
|
location: InputFileLocation,
|
||||||
size?: number,
|
size?: number,
|
||||||
fileName?: string,
|
fileName?: string,
|
||||||
@ -155,7 +156,7 @@ export class ApiFileManager {
|
|||||||
return canceled;
|
return canceled;
|
||||||
}
|
}
|
||||||
|
|
||||||
public requestFilePart(dcId: number, location: InputFileLocation, offset: number, limit: number, id = 0, queueId = 0, checkCancel?: () => void) {
|
public requestFilePart(dcId: DcId, location: InputFileLocation, offset: number, limit: number, id = 0, queueId = 0, checkCancel?: () => void) {
|
||||||
return this.downloadRequest(dcId, id, async() => {
|
return this.downloadRequest(dcId, id, async() => {
|
||||||
checkCancel && checkCancel();
|
checkCancel && checkCancel();
|
||||||
|
|
||||||
|
@ -18,7 +18,7 @@ import networkerFactory from './networkerFactory';
|
|||||||
import authorizer from './authorizer';
|
import authorizer from './authorizer';
|
||||||
import dcConfigurator, { ConnectionType, TransportType } from './dcConfigurator';
|
import dcConfigurator, { ConnectionType, TransportType } from './dcConfigurator';
|
||||||
import { logger } from '../logger';
|
import { logger } from '../logger';
|
||||||
import type { InvokeApiOptions } from '../../types';
|
import type { DcId, InvokeApiOptions, TrueDcId } from '../../types';
|
||||||
import type { MethodDeclMap } from '../../layer';
|
import type { MethodDeclMap } from '../../layer';
|
||||||
import { CancellablePromise, deferredPromise } from '../../helpers/cancellablePromise';
|
import { CancellablePromise, deferredPromise } from '../../helpers/cancellablePromise';
|
||||||
import { bytesFromHex, bytesToHex } from '../../helpers/bytes';
|
import { bytesFromHex, bytesToHex } from '../../helpers/bytes';
|
||||||
@ -78,7 +78,7 @@ export class ApiManager {
|
|||||||
|
|
||||||
private cachedExportPromise: {[x: number]: Promise<unknown>} = {};
|
private cachedExportPromise: {[x: number]: Promise<unknown>} = {};
|
||||||
private gettingNetworkers: {[dcIdAndType: string]: Promise<MTPNetworker>} = {};
|
private gettingNetworkers: {[dcIdAndType: string]: Promise<MTPNetworker>} = {};
|
||||||
private baseDcId = 0;
|
private baseDcId: DcId = 0 as DcId;
|
||||||
|
|
||||||
//public telegramMeNotified = false;
|
//public telegramMeNotified = false;
|
||||||
|
|
||||||
@ -140,7 +140,7 @@ export class ApiManager {
|
|||||||
/// #endif
|
/// #endif
|
||||||
}
|
}
|
||||||
|
|
||||||
public setBaseDcId(dcId: number) {
|
public setBaseDcId(dcId: DcId) {
|
||||||
this.baseDcId = dcId;
|
this.baseDcId = dcId;
|
||||||
|
|
||||||
sessionStorage.set({
|
sessionStorage.set({
|
||||||
@ -163,14 +163,14 @@ export class ApiManager {
|
|||||||
const logoutPromises: Promise<any>[] = [];
|
const logoutPromises: Promise<any>[] = [];
|
||||||
for(let i = 0; i < storageResult.length; i++) {
|
for(let i = 0; i < storageResult.length; i++) {
|
||||||
if(storageResult[i]) {
|
if(storageResult[i]) {
|
||||||
logoutPromises.push(this.invokeApi('auth.logOut', {}, {dcId: i + 1, ignoreErrors: true}));
|
logoutPromises.push(this.invokeApi('auth.logOut', {}, {dcId: (i + 1) as DcId, ignoreErrors: true}));
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
const clear = () => {
|
const clear = () => {
|
||||||
//console.error('apiManager: logOut clear');
|
//console.error('apiManager: logOut clear');
|
||||||
|
|
||||||
this.baseDcId = 0;
|
this.baseDcId = undefined;
|
||||||
//this.telegramMeNotify(false);
|
//this.telegramMeNotify(false);
|
||||||
IDBStorage.closeDatabases();
|
IDBStorage.closeDatabases();
|
||||||
self.postMessage({type: 'clear'});
|
self.postMessage({type: 'clear'});
|
||||||
@ -189,7 +189,7 @@ export class ApiManager {
|
|||||||
}
|
}
|
||||||
|
|
||||||
// mtpGetNetworker
|
// mtpGetNetworker
|
||||||
public getNetworker(dcId: number, options: InvokeApiOptions = {}): Promise<MTPNetworker> {
|
public getNetworker(dcId: DcId, options: InvokeApiOptions = {}): Promise<MTPNetworker> {
|
||||||
const connectionType: ConnectionType = options.fileDownload ? 'download' : (options.fileUpload ? 'upload' : 'client');
|
const connectionType: ConnectionType = options.fileDownload ? 'download' : (options.fileUpload ? 'upload' : 'client');
|
||||||
//const connectionType: ConnectionType = 'client';
|
//const connectionType: ConnectionType = 'client';
|
||||||
|
|
||||||
@ -237,10 +237,10 @@ export class ApiManager {
|
|||||||
return this.gettingNetworkers[getKey];
|
return this.gettingNetworkers[getKey];
|
||||||
}
|
}
|
||||||
|
|
||||||
const ak = 'dc' + dcId + '_auth_key';
|
const ak: `dc${TrueDcId}_auth_key` = `dc${dcId}_auth_key` as any;
|
||||||
const ss = 'dc' + dcId + '_server_salt';
|
const ss: `dc${TrueDcId}_server_salt` = `dc${dcId}_server_salt` as any;
|
||||||
|
|
||||||
return this.gettingNetworkers[getKey] = Promise.all([ak, ss].map(key => sessionStorage.get(key as any)))
|
return this.gettingNetworkers[getKey] = Promise.all([ak, ss].map(key => sessionStorage.get(key)))
|
||||||
.then(async([authKeyHex, serverSaltHex]) => {
|
.then(async([authKeyHex, serverSaltHex]) => {
|
||||||
const transport = dcConfigurator.chooseServer(dcId, connectionType, transportType, false);
|
const transport = dcConfigurator.chooseServer(dcId, connectionType, transportType, false);
|
||||||
let networker: MTPNetworker;
|
let networker: MTPNetworker;
|
||||||
@ -273,6 +273,16 @@ export class ApiManager {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
if(networker.isFileNetworker) {
|
||||||
|
networker.onDrain = () => {
|
||||||
|
this.log('networker drain', networker);
|
||||||
|
|
||||||
|
networker.onDrain = undefined;
|
||||||
|
const idx = networkers.indexOf(networker);
|
||||||
|
networkers.splice(idx, 1);
|
||||||
|
};
|
||||||
|
}
|
||||||
|
|
||||||
/* networker.onConnectionStatusChange = (online) => {
|
/* networker.onConnectionStatusChange = (online) => {
|
||||||
console.log('status:', online);
|
console.log('status:', online);
|
||||||
}; */
|
}; */
|
||||||
@ -354,7 +364,7 @@ export class ApiManager {
|
|||||||
}
|
}
|
||||||
};
|
};
|
||||||
|
|
||||||
let dcId: number;
|
let dcId: DcId;
|
||||||
|
|
||||||
let cachedNetworker: MTPNetworker;
|
let cachedNetworker: MTPNetworker;
|
||||||
let stack = (new Error()).stack || 'empty stack';
|
let stack = (new Error()).stack || 'empty stack';
|
||||||
@ -400,7 +410,7 @@ export class ApiManager {
|
|||||||
this.invokeApi(method, params, options).then(deferred.resolve, rejectPromise);
|
this.invokeApi(method, params, options).then(deferred.resolve, rejectPromise);
|
||||||
}, rejectPromise);
|
}, rejectPromise);
|
||||||
} else if(error.code === 303) {
|
} else if(error.code === 303) {
|
||||||
const newDcId = +error.type.match(/^(PHONE_MIGRATE_|NETWORK_MIGRATE_|USER_MIGRATE_|FILE_MIGRATE_)(\d+)/)[2];
|
const newDcId = +error.type.match(/^(PHONE_MIGRATE_|NETWORK_MIGRATE_|USER_MIGRATE_|FILE_MIGRATE_)(\d+)/)[2] as DcId;
|
||||||
if(newDcId !== dcId) {
|
if(newDcId !== dcId) {
|
||||||
if(options.dcId) {
|
if(options.dcId) {
|
||||||
options.dcId = newDcId;
|
options.dcId = newDcId;
|
||||||
|
@ -74,6 +74,7 @@ export type MTMessage = InvokeApiOptions & MTMessageOptions & {
|
|||||||
};
|
};
|
||||||
|
|
||||||
const CONNECTION_TIMEOUT = 5000;
|
const CONNECTION_TIMEOUT = 5000;
|
||||||
|
const DRAIN_TIMEOUT = 10000;
|
||||||
let invokeAfterMsgConstructor: number;
|
let invokeAfterMsgConstructor: number;
|
||||||
|
|
||||||
export default class MTPNetworker {
|
export default class MTPNetworker {
|
||||||
@ -124,6 +125,11 @@ export default class MTPNetworker {
|
|||||||
|
|
||||||
private debug = DEBUG /* && false */ || Modes.debug;
|
private debug = DEBUG /* && false */ || Modes.debug;
|
||||||
|
|
||||||
|
public activeRequests = 0;
|
||||||
|
|
||||||
|
public onDrain: () => void;
|
||||||
|
private onDrainTimeout: number;
|
||||||
|
|
||||||
//private disconnectDelay: number;
|
//private disconnectDelay: number;
|
||||||
//private pingPromise: CancellablePromise<any>;
|
//private pingPromise: CancellablePromise<any>;
|
||||||
//public onConnectionStatusChange: (online: boolean) => void;
|
//public onConnectionStatusChange: (online: boolean) => void;
|
||||||
@ -687,7 +693,19 @@ export default class MTPNetworker {
|
|||||||
promise.finally(() => {
|
promise.finally(() => {
|
||||||
clearTimeout(timeout);
|
clearTimeout(timeout);
|
||||||
this.setConnectionStatus(true);
|
this.setConnectionStatus(true);
|
||||||
|
|
||||||
|
if(!--this.activeRequests && this.onDrain) {
|
||||||
|
this.onDrainTimeout = self.setTimeout(() => {
|
||||||
|
this.log('drain');
|
||||||
|
this.onDrain();
|
||||||
|
}, DRAIN_TIMEOUT);
|
||||||
|
}
|
||||||
});
|
});
|
||||||
|
|
||||||
|
++this.activeRequests;
|
||||||
|
if(this.onDrainTimeout !== undefined) {
|
||||||
|
clearTimeout(this.onDrainTimeout);
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
return promise;
|
return promise;
|
||||||
|
@ -6,11 +6,12 @@
|
|||||||
|
|
||||||
import type { AppInstance } from './mtproto/singleInstance';
|
import type { AppInstance } from './mtproto/singleInstance';
|
||||||
import type { UserAuth } from './mtproto/mtproto_config';
|
import type { UserAuth } from './mtproto/mtproto_config';
|
||||||
|
import type { DcId } from '../types';
|
||||||
import { MOUNT_CLASS_TO } from '../config/debug';
|
import { MOUNT_CLASS_TO } from '../config/debug';
|
||||||
import LocalStorageController from './localStorage';
|
import LocalStorageController from './localStorage';
|
||||||
|
|
||||||
const sessionStorage = new LocalStorageController<{
|
const sessionStorage = new LocalStorageController<{
|
||||||
dc: number,
|
dc: DcId,
|
||||||
user_auth: UserAuth,
|
user_auth: UserAuth,
|
||||||
dc1_auth_key: string,
|
dc1_auth_key: string,
|
||||||
dc2_auth_key: string,
|
dc2_auth_key: string,
|
||||||
|
@ -4,6 +4,7 @@
|
|||||||
* https://github.com/morethanwords/tweb/blob/master/LICENSE
|
* https://github.com/morethanwords/tweb/blob/master/LICENSE
|
||||||
*/
|
*/
|
||||||
|
|
||||||
|
import type { DcId } from '../types';
|
||||||
import apiManager from '../lib/mtproto/mtprotoworker';
|
import apiManager from '../lib/mtproto/mtprotoworker';
|
||||||
import Page from './page';
|
import Page from './page';
|
||||||
import serverTimeManager from '../lib/mtproto/serverTimeManager';
|
import serverTimeManager from '../lib/mtproto/serverTimeManager';
|
||||||
@ -65,7 +66,7 @@ let onFirstMount = async() => {
|
|||||||
cachedPromise = null;
|
cachedPromise = null;
|
||||||
}, true);
|
}, true);
|
||||||
|
|
||||||
let options: {dcId?: number, ignoreErrors: true} = {ignoreErrors: true};
|
let options: {dcId?: DcId, ignoreErrors: true} = {ignoreErrors: true};
|
||||||
let prevToken: Uint8Array | number[];
|
let prevToken: Uint8Array | number[];
|
||||||
|
|
||||||
const iterate = async(isLoop: boolean) => {
|
const iterate = async(isLoop: boolean) => {
|
||||||
@ -78,7 +79,7 @@ let onFirstMount = async() => {
|
|||||||
|
|
||||||
if(loginToken._ === 'auth.loginTokenMigrateTo') {
|
if(loginToken._ === 'auth.loginTokenMigrateTo') {
|
||||||
if(!options.dcId) {
|
if(!options.dcId) {
|
||||||
options.dcId = loginToken.dc_id;
|
options.dcId = loginToken.dc_id as DcId;
|
||||||
apiManager.setBaseDcId(loginToken.dc_id);
|
apiManager.setBaseDcId(loginToken.dc_id);
|
||||||
//continue;
|
//continue;
|
||||||
}
|
}
|
||||||
|
5
src/types.d.ts
vendored
5
src/types.d.ts
vendored
@ -1,8 +1,11 @@
|
|||||||
import { AuthSentCode } from "./layer";
|
import { AuthSentCode } from "./layer";
|
||||||
import type { ApiError } from "./lib/mtproto/apiManager";
|
import type { ApiError } from "./lib/mtproto/apiManager";
|
||||||
|
|
||||||
|
export type DcId = number;
|
||||||
|
export type TrueDcId = 1 | 2 | 3 | 4 | 5;
|
||||||
|
|
||||||
export type InvokeApiOptions = Partial<{
|
export type InvokeApiOptions = Partial<{
|
||||||
dcId: number,
|
dcId: DcId,
|
||||||
floodMaxTimeout: number,
|
floodMaxTimeout: number,
|
||||||
noErrorBox: true,
|
noErrorBox: true,
|
||||||
fileUpload: true,
|
fileUpload: true,
|
||||||
|
Loading…
x
Reference in New Issue
Block a user