Remove HTTP queuing
This commit is contained in:
parent
2983357946
commit
fe66c0751d
1 changed files with 45 additions and 97 deletions
|
|
@ -10,74 +10,29 @@ import type {
|
||||||
RelativePath,
|
RelativePath,
|
||||||
VaultUpdateId,
|
VaultUpdateId,
|
||||||
} from "src/database/document-metadata";
|
} from "src/database/document-metadata";
|
||||||
import PQueue from "p-queue";
|
|
||||||
import { Logger } from "src/tracing/logger.js";
|
import { Logger } from "src/tracing/logger.js";
|
||||||
|
|
||||||
export interface RequestCountStatus {
|
|
||||||
waiting: number;
|
|
||||||
success: number;
|
|
||||||
failure: number;
|
|
||||||
}
|
|
||||||
|
|
||||||
export class SyncService {
|
export class SyncService {
|
||||||
private client: Client<paths>;
|
private client: Client<paths>;
|
||||||
|
|
||||||
private readonly promiseQueue: PQueue;
|
|
||||||
private readonly requestCountListeners: ((
|
|
||||||
status: RequestCountStatus
|
|
||||||
) => void)[] = [];
|
|
||||||
private readonly status: RequestCountStatus = {
|
|
||||||
waiting: 0,
|
|
||||||
success: 0,
|
|
||||||
failure: 0,
|
|
||||||
};
|
|
||||||
|
|
||||||
public constructor(private readonly database: Database) {
|
public constructor(private readonly database: Database) {
|
||||||
this.createClient(database.getSettings());
|
this.createClient(database.getSettings());
|
||||||
this.promiseQueue = new PQueue({
|
|
||||||
concurrency: database.getSettings().uploadConcurrency,
|
|
||||||
});
|
|
||||||
|
|
||||||
database.addOnSettingsChangeHandlers((s) => {
|
database.addOnSettingsChangeHandlers((s) => {
|
||||||
this.createClient(s);
|
this.createClient(s);
|
||||||
this.promiseQueue.concurrency = s.uploadConcurrency;
|
|
||||||
});
|
});
|
||||||
|
|
||||||
this.promiseQueue.on("active", () => {
|
|
||||||
this.status.waiting = this.promiseQueue.pending;
|
|
||||||
this.emitRequestCountChange();
|
|
||||||
});
|
|
||||||
|
|
||||||
this.promiseQueue.on("completed", () => {
|
|
||||||
this.status.success++;
|
|
||||||
this.emitRequestCountChange();
|
|
||||||
});
|
|
||||||
|
|
||||||
this.promiseQueue.on("error", () => {
|
|
||||||
this.status.failure++;
|
|
||||||
this.emitRequestCountChange();
|
|
||||||
});
|
|
||||||
}
|
|
||||||
|
|
||||||
public addRequestCountChangeListener(
|
|
||||||
listener: (status: RequestCountStatus) => void
|
|
||||||
): void {
|
|
||||||
this.requestCountListeners.push(listener);
|
|
||||||
listener({ ...this.status });
|
|
||||||
}
|
}
|
||||||
|
|
||||||
public async ping(): Promise<components["schemas"]["PingResponse"]> {
|
public async ping(): Promise<components["schemas"]["PingResponse"]> {
|
||||||
const response = await this.enqueue(async () =>
|
const response = await this.client.GET("/ping", {
|
||||||
this.client.GET("/ping", {
|
params: {
|
||||||
params: {
|
header: {
|
||||||
header: {
|
authorization: `Bearer ${
|
||||||
authorization: `Bearer ${
|
this.database.getSettings().token
|
||||||
this.database.getSettings().token
|
}`,
|
||||||
}`,
|
|
||||||
},
|
|
||||||
},
|
},
|
||||||
})
|
},
|
||||||
);
|
});
|
||||||
|
|
||||||
Logger.getInstance().debug(
|
Logger.getInstance().debug(
|
||||||
`Ping response: ${JSON.stringify(response.data)}`
|
`Ping response: ${JSON.stringify(response.data)}`
|
||||||
|
|
@ -99,8 +54,9 @@ export class SyncService {
|
||||||
contentBytes: Uint8Array;
|
contentBytes: Uint8Array;
|
||||||
createdDate: Date;
|
createdDate: Date;
|
||||||
}): Promise<components["schemas"]["DocumentVersion"]> {
|
}): Promise<components["schemas"]["DocumentVersion"]> {
|
||||||
const response = await this.enqueue(async () =>
|
const response = await this.client.POST(
|
||||||
this.client.POST("/vaults/{vault_id}/documents", {
|
"/vaults/{vault_id}/documents",
|
||||||
|
{
|
||||||
params: {
|
params: {
|
||||||
path: {
|
path: {
|
||||||
vault_id: this.database.getSettings().vaultName,
|
vault_id: this.database.getSettings().vaultName,
|
||||||
|
|
@ -116,7 +72,7 @@ export class SyncService {
|
||||||
createdDate: createdDate.toISOString(),
|
createdDate: createdDate.toISOString(),
|
||||||
relativePath,
|
relativePath,
|
||||||
},
|
},
|
||||||
})
|
}
|
||||||
);
|
);
|
||||||
|
|
||||||
if (!response.data) {
|
if (!response.data) {
|
||||||
|
|
@ -124,7 +80,9 @@ export class SyncService {
|
||||||
}
|
}
|
||||||
|
|
||||||
Logger.getInstance().debug(
|
Logger.getInstance().debug(
|
||||||
`Created document ${JSON.stringify(response.data)}`
|
`Created document ${JSON.stringify(
|
||||||
|
response.data.relativePath
|
||||||
|
)} with id ${response.data.documentId}`
|
||||||
);
|
);
|
||||||
|
|
||||||
return response.data;
|
return response.data;
|
||||||
|
|
@ -143,8 +101,9 @@ export class SyncService {
|
||||||
contentBytes: Uint8Array;
|
contentBytes: Uint8Array;
|
||||||
createdDate: Date;
|
createdDate: Date;
|
||||||
}): Promise<components["schemas"]["DocumentVersion"]> {
|
}): Promise<components["schemas"]["DocumentVersion"]> {
|
||||||
const response = await this.enqueue(async () =>
|
const response = await this.client.PUT(
|
||||||
this.client.PUT("/vaults/{vault_id}/documents/{document_id}", {
|
"/vaults/{vault_id}/documents/{document_id}",
|
||||||
|
{
|
||||||
params: {
|
params: {
|
||||||
path: {
|
path: {
|
||||||
vault_id: this.database.getSettings().vaultName,
|
vault_id: this.database.getSettings().vaultName,
|
||||||
|
|
@ -162,7 +121,7 @@ export class SyncService {
|
||||||
createdDate: createdDate.toISOString(),
|
createdDate: createdDate.toISOString(),
|
||||||
relativePath,
|
relativePath,
|
||||||
},
|
},
|
||||||
})
|
}
|
||||||
);
|
);
|
||||||
|
|
||||||
if (!response.data) {
|
if (!response.data) {
|
||||||
|
|
@ -170,7 +129,7 @@ export class SyncService {
|
||||||
}
|
}
|
||||||
|
|
||||||
Logger.getInstance().debug(
|
Logger.getInstance().debug(
|
||||||
`Updated document ${JSON.stringify(response.data)}`
|
`Updated document ${response.data.relativePath} with id ${response.data.documentId}`
|
||||||
);
|
);
|
||||||
|
|
||||||
return response.data;
|
return response.data;
|
||||||
|
|
@ -185,8 +144,9 @@ export class SyncService {
|
||||||
relativePath: RelativePath;
|
relativePath: RelativePath;
|
||||||
createdDate: Date;
|
createdDate: Date;
|
||||||
}): Promise<void> {
|
}): Promise<void> {
|
||||||
const response = await this.enqueue(async () =>
|
const response = await this.client.DELETE(
|
||||||
this.client.DELETE("/vaults/{vault_id}/documents/{document_id}", {
|
"/vaults/{vault_id}/documents/{document_id}",
|
||||||
|
{
|
||||||
params: {
|
params: {
|
||||||
path: {
|
path: {
|
||||||
vault_id: this.database.getSettings().vaultName,
|
vault_id: this.database.getSettings().vaultName,
|
||||||
|
|
@ -202,7 +162,7 @@ export class SyncService {
|
||||||
createdDate: createdDate.toISOString(),
|
createdDate: createdDate.toISOString(),
|
||||||
relativePath,
|
relativePath,
|
||||||
},
|
},
|
||||||
})
|
}
|
||||||
);
|
);
|
||||||
|
|
||||||
if (response.error) {
|
if (response.error) {
|
||||||
|
|
@ -210,7 +170,7 @@ export class SyncService {
|
||||||
}
|
}
|
||||||
|
|
||||||
Logger.getInstance().debug(
|
Logger.getInstance().debug(
|
||||||
`Updated document ${JSON.stringify(response.data)}`
|
`Deleted document ${relativePath} with id ${documentId}`
|
||||||
);
|
);
|
||||||
|
|
||||||
return response.data;
|
return response.data;
|
||||||
|
|
@ -221,8 +181,9 @@ export class SyncService {
|
||||||
}: {
|
}: {
|
||||||
documentId: DocumentId;
|
documentId: DocumentId;
|
||||||
}): Promise<components["schemas"]["DocumentVersion"]> {
|
}): Promise<components["schemas"]["DocumentVersion"]> {
|
||||||
const response = await this.enqueue(async () =>
|
const response = await this.client.GET(
|
||||||
this.client.GET("/vaults/{vault_id}/documents/{document_id}", {
|
"/vaults/{vault_id}/documents/{document_id}",
|
||||||
|
{
|
||||||
params: {
|
params: {
|
||||||
path: {
|
path: {
|
||||||
vault_id: this.database.getSettings().vaultName,
|
vault_id: this.database.getSettings().vaultName,
|
||||||
|
|
@ -234,7 +195,7 @@ export class SyncService {
|
||||||
}`,
|
}`,
|
||||||
},
|
},
|
||||||
},
|
},
|
||||||
})
|
}
|
||||||
);
|
);
|
||||||
|
|
||||||
if (!response.data) {
|
if (!response.data) {
|
||||||
|
|
@ -242,7 +203,7 @@ export class SyncService {
|
||||||
}
|
}
|
||||||
|
|
||||||
Logger.getInstance().debug(
|
Logger.getInstance().debug(
|
||||||
`Get document ${JSON.stringify(response.data)}`
|
`Get document ${response.data.relativePath} with id ${response.data.documentId}`
|
||||||
);
|
);
|
||||||
|
|
||||||
return response.data;
|
return response.data;
|
||||||
|
|
@ -251,23 +212,21 @@ export class SyncService {
|
||||||
public async getAll(
|
public async getAll(
|
||||||
since?: VaultUpdateId
|
since?: VaultUpdateId
|
||||||
): Promise<components["schemas"]["FetchLatestDocumentsResponse"]> {
|
): Promise<components["schemas"]["FetchLatestDocumentsResponse"]> {
|
||||||
const response = await this.enqueue(async () =>
|
const response = await this.client.GET("/vaults/{vault_id}/documents", {
|
||||||
this.client.GET("/vaults/{vault_id}/documents", {
|
params: {
|
||||||
params: {
|
path: {
|
||||||
path: {
|
vault_id: this.database.getSettings().vaultName,
|
||||||
vault_id: this.database.getSettings().vaultName,
|
|
||||||
},
|
|
||||||
header: {
|
|
||||||
authorization: `Bearer ${
|
|
||||||
this.database.getSettings().token
|
|
||||||
}`,
|
|
||||||
},
|
|
||||||
query: {
|
|
||||||
since_update_id: since,
|
|
||||||
},
|
|
||||||
},
|
},
|
||||||
})
|
header: {
|
||||||
);
|
authorization: `Bearer ${
|
||||||
|
this.database.getSettings().token
|
||||||
|
}`,
|
||||||
|
},
|
||||||
|
query: {
|
||||||
|
since_update_id: since,
|
||||||
|
},
|
||||||
|
},
|
||||||
|
});
|
||||||
|
|
||||||
const { error } = response;
|
const { error } = response;
|
||||||
if (error) {
|
if (error) {
|
||||||
|
|
@ -275,26 +234,15 @@ export class SyncService {
|
||||||
}
|
}
|
||||||
|
|
||||||
Logger.getInstance().debug(
|
Logger.getInstance().debug(
|
||||||
`Get document ${JSON.stringify(response.data)}`
|
`Got ${response.data.latestDocuments.length} document metadata`
|
||||||
);
|
);
|
||||||
|
|
||||||
return response.data;
|
return response.data;
|
||||||
}
|
}
|
||||||
|
|
||||||
private emitRequestCountChange(): void {
|
|
||||||
this.requestCountListeners.forEach((listener) => {
|
|
||||||
listener({ ...this.status });
|
|
||||||
});
|
|
||||||
}
|
|
||||||
|
|
||||||
private createClient(settings: SyncSettings): void {
|
private createClient(settings: SyncSettings): void {
|
||||||
this.client = createClient<paths>({
|
this.client = createClient<paths>({
|
||||||
baseUrl: settings.remoteUri,
|
baseUrl: settings.remoteUri,
|
||||||
});
|
});
|
||||||
}
|
}
|
||||||
|
|
||||||
private async enqueue<T>(fn: () => Promise<T>): Promise<T> {
|
|
||||||
// eslint-disable-next-line @typescript-eslint/no-unsafe-type-assertion
|
|
||||||
return this.promiseQueue.add(fn) as Promise<T>;
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue