Add concurrency limit to service

This commit is contained in:
Andras Schmelczer 2024-12-18 20:37:54 +00:00
commit eb87de8e68
No known key found for this signature in database
GPG key ID: FC8F2C3D3D1A718C
9 changed files with 116 additions and 58 deletions

View file

@ -20,6 +20,7 @@
"obsidian": "1.7.2", "obsidian": "1.7.2",
"openapi-fetch": "0.13.3", "openapi-fetch": "0.13.3",
"openapi-typescript": "7.4.4", "openapi-typescript": "7.4.4",
"p-queue": "^8.0.1",
"tslib": "2.4.0", "tslib": "2.4.0",
"typescript": "5.7.2" "typescript": "5.7.2"
} }
@ -1543,6 +1544,13 @@
"node": ">=0.10.0" "node": ">=0.10.0"
} }
}, },
"node_modules/eventemitter3": {
"version": "5.0.1",
"resolved": "https://registry.npmjs.org/eventemitter3/-/eventemitter3-5.0.1.tgz",
"integrity": "sha512-GWkBvjiSZK87ELrYOSESUYeVIc9mvLLf/nXalMOS5dYrgZq9o5OVkbZAVM06CVxYsCwH9BDZFPlQTlPA1j4ahA==",
"dev": true,
"license": "MIT"
},
"node_modules/fast-deep-equal": { "node_modules/fast-deep-equal": {
"version": "3.1.3", "version": "3.1.3",
"resolved": "https://registry.npmjs.org/fast-deep-equal/-/fast-deep-equal-3.1.3.tgz", "resolved": "https://registry.npmjs.org/fast-deep-equal/-/fast-deep-equal-3.1.3.tgz",
@ -2118,6 +2126,36 @@
"url": "https://github.com/sponsors/sindresorhus" "url": "https://github.com/sponsors/sindresorhus"
} }
}, },
"node_modules/p-queue": {
"version": "8.0.1",
"resolved": "https://registry.npmjs.org/p-queue/-/p-queue-8.0.1.tgz",
"integrity": "sha512-NXzu9aQJTAzbBqOt2hwsR63ea7yvxJc0PwN/zobNAudYfb1B7R08SzB4TsLeSbUCuG467NhnoT0oO6w1qRO+BA==",
"dev": true,
"license": "MIT",
"dependencies": {
"eventemitter3": "^5.0.1",
"p-timeout": "^6.1.2"
},
"engines": {
"node": ">=18"
},
"funding": {
"url": "https://github.com/sponsors/sindresorhus"
}
},
"node_modules/p-timeout": {
"version": "6.1.3",
"resolved": "https://registry.npmjs.org/p-timeout/-/p-timeout-6.1.3.tgz",
"integrity": "sha512-UJUyfKbwvr/uZSV6btANfb+0t/mOhKV/KXcCUTp8FcQI+v/0d+wXqH4htrW0E4rR6WiEO/EPvUFiV9D5OI4vlw==",
"dev": true,
"license": "MIT",
"engines": {
"node": ">=14.16"
},
"funding": {
"url": "https://github.com/sponsors/sindresorhus"
}
},
"node_modules/parent-module": { "node_modules/parent-module": {
"version": "1.0.1", "version": "1.0.1",
"resolved": "https://registry.npmjs.org/parent-module/-/parent-module-1.0.1.tgz", "resolved": "https://registry.npmjs.org/parent-module/-/parent-module-1.0.1.tgz",

View file

@ -25,6 +25,7 @@
"openapi-fetch": "0.13.3", "openapi-fetch": "0.13.3",
"openapi-typescript": "7.4.4", "openapi-typescript": "7.4.4",
"tslib": "2.4.0", "tslib": "2.4.0",
"typescript": "5.7.2" "typescript": "5.7.2",
"p-queue": "^8.0.1"
} }
} }

View file

@ -3,6 +3,7 @@ export interface SyncSettings {
token: string; token: string;
vaultName: string; vaultName: string;
fetchChangesUpdateIntervalMs: number; fetchChangesUpdateIntervalMs: number;
uploadConcurrency: number;
isSyncEnabled: boolean; isSyncEnabled: boolean;
} }
@ -11,5 +12,6 @@ export const DEFAULT_SETTINGS: SyncSettings = {
token: "", token: "",
vaultName: "default", vaultName: "default",
fetchChangesUpdateIntervalMs: 1000, fetchChangesUpdateIntervalMs: 1000,
uploadConcurrency: 10,
isSyncEnabled: true, isSyncEnabled: true,
}; };

View file

@ -1,7 +1,7 @@
import { TAbstractFile, TFile } from "obsidian"; import { TAbstractFile, TFile } from "obsidian";
import { FileEventHandler } from "./file-event-handler"; import { FileEventHandler } from "./file-event-handler";
import { Logger } from "src/logger"; import { Logger } from "src/logger";
import { SyncServer } from "src/services/sync_service"; import { SyncService } from "src/services/sync_service";
import { Database } from "src/database/database"; import { Database } from "src/database/database";
import { syncLocallyDeletedFile } from "src/sync-operations/sync-locally-deleted-file"; import { syncLocallyDeletedFile } from "src/sync-operations/sync-locally-deleted-file";
import { syncLocallyUpdatedFile } from "src/sync-operations/sync-locally-updated-file"; import { syncLocallyUpdatedFile } from "src/sync-operations/sync-locally-updated-file";
@ -11,7 +11,7 @@ import { syncLocallyCreatedFile } from "src/sync-operations/sync-locally-created
export class SyncEventHandler implements FileEventHandler { export class SyncEventHandler implements FileEventHandler {
public constructor( public constructor(
private database: Database, private database: Database,
private syncServer: SyncServer, private syncServer: SyncService,
private operations: FileOperations private operations: FileOperations
) {} ) {}
@ -49,11 +49,11 @@ export class SyncEventHandler implements FileEventHandler {
return; return;
} }
await syncLocallyDeletedFile( await syncLocallyDeletedFile({
this.database, database: this.database,
this.syncServer, syncServer: this.syncServer,
file.path relativePath: file.path,
); });
} else { } else {
Logger.getInstance().info(`Folder deleted: ${file.path}, ignored`); Logger.getInstance().info(`Folder deleted: ${file.path}, ignored`);
} }

View file

@ -10,13 +10,22 @@ import {
RelativePath, RelativePath,
DocumentId, DocumentId,
} from "src/database/document-metadata.js"; } from "src/database/document-metadata.js";
import PQueue from "p-queue";
export class SyncServer { export class SyncService {
private promiseQueue: PQueue;
private client: Client<paths>; private client: Client<paths>;
public constructor(private database: Database) { public constructor(private database: Database) {
this.createClient(database.getSettings()); this.createClient(database.getSettings());
database.addOnSettingsChangeHandlers((s) => this.createClient(s)); this.promiseQueue = new PQueue({
concurrency: database.getSettings().uploadConcurrency,
});
database.addOnSettingsChangeHandlers((s) => {
this.createClient(s);
this.promiseQueue.concurrency = s.uploadConcurrency;
});
} }
private createClient(settings: SyncSettings) { private createClient(settings: SyncSettings) {
@ -25,15 +34,21 @@ export class SyncServer {
}); });
} }
private enqueue<T>(fn: () => Promise<T>): Promise<T> {
return this.promiseQueue.add(fn) as Promise<T>;
}
public async ping(): Promise<components["schemas"]["PingResponse"]> { public async ping(): Promise<components["schemas"]["PingResponse"]> {
const response = await this.client.GET("/ping", { const response = await this.enqueue(() =>
params: { this.client.GET("/ping", {
header: { params: {
authorization: header: {
"Bearer " + this.database.getSettings().token, authorization:
"Bearer " + this.database.getSettings().token,
},
}, },
}, })
}); );
Logger.getInstance().debug( Logger.getInstance().debug(
"Ping response: " + JSON.stringify(response.data) "Ping response: " + JSON.stringify(response.data)
@ -55,9 +70,8 @@ export class SyncServer {
contentBytes: Uint8Array; contentBytes: Uint8Array;
createdDate: Date; createdDate: Date;
}): Promise<components["schemas"]["DocumentVersion"]> { }): Promise<components["schemas"]["DocumentVersion"]> {
const response = await this.client.POST( const response = await this.enqueue(() =>
"/vaults/{vault_id}/documents", this.client.POST("/vaults/{vault_id}/documents", {
{
params: { params: {
path: { path: {
vault_id: this.database.getSettings().vaultName, vault_id: this.database.getSettings().vaultName,
@ -72,7 +86,7 @@ export class SyncServer {
createdDate: createdDate.toISOString(), createdDate: createdDate.toISOString(),
relativePath, relativePath,
}, },
} })
); );
if (!response.data) { if (!response.data) {
@ -99,9 +113,8 @@ export class SyncServer {
contentBytes: Uint8Array; contentBytes: Uint8Array;
createdDate: Date; createdDate: Date;
}): Promise<components["schemas"]["DocumentVersion"]> { }): Promise<components["schemas"]["DocumentVersion"]> {
const response = await this.client.PUT( const response = await this.enqueue(() =>
"/vaults/{vault_id}/documents/{document_id}", this.client.PUT("/vaults/{vault_id}/documents/{document_id}", {
{
params: { params: {
path: { path: {
vault_id: this.database.getSettings().vaultName, vault_id: this.database.getSettings().vaultName,
@ -118,7 +131,7 @@ export class SyncServer {
createdDate: createdDate.toISOString(), createdDate: createdDate.toISOString(),
relativePath, relativePath,
}, },
} })
); );
if (!response.data) { if (!response.data) {
@ -141,9 +154,8 @@ export class SyncServer {
relativePath: RelativePath; relativePath: RelativePath;
createdDate: Date; createdDate: Date;
}): Promise<void> { }): Promise<void> {
const response = await this.client.DELETE( const response = await this.enqueue(() =>
"/vaults/{vault_id}/documents/{document_id}", this.client.DELETE("/vaults/{vault_id}/documents/{document_id}", {
{
params: { params: {
path: { path: {
vault_id: this.database.getSettings().vaultName, vault_id: this.database.getSettings().vaultName,
@ -158,7 +170,7 @@ export class SyncServer {
createdDate: createdDate.toISOString(), createdDate: createdDate.toISOString(),
relativePath, relativePath,
}, },
} })
); );
if (response.error) { if (response.error) {
@ -177,9 +189,8 @@ export class SyncServer {
}: { }: {
documentId: DocumentId; documentId: DocumentId;
}): Promise<components["schemas"]["DocumentVersion"]> { }): Promise<components["schemas"]["DocumentVersion"]> {
const response = await this.client.GET( const response = await this.enqueue(() =>
"/vaults/{vault_id}/documents/{document_id}", this.client.GET("/vaults/{vault_id}/documents/{document_id}", {
{
params: { params: {
path: { path: {
vault_id: this.database.getSettings().vaultName, vault_id: this.database.getSettings().vaultName,
@ -190,7 +201,7 @@ export class SyncServer {
"Bearer " + this.database.getSettings().token, "Bearer " + this.database.getSettings().token,
}, },
}, },
} })
); );
if (!response.data) { if (!response.data) {
@ -207,20 +218,22 @@ export class SyncServer {
public async getAll( public async getAll(
since?: VaultUpdateId since?: VaultUpdateId
): Promise<components["schemas"]["FetchLatestDocumentsResponse"]> { ): Promise<components["schemas"]["FetchLatestDocumentsResponse"]> {
const response = await this.client.GET("/vaults/{vault_id}/documents", { const response = await this.enqueue(() =>
params: { this.client.GET("/vaults/{vault_id}/documents", {
path: { params: {
vault_id: this.database.getSettings().vaultName, path: {
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,
},
},
});
if (!response.data) { if (!response.data) {
throw new Error(`Failed to get documents: ${response.error}`); throw new Error(`Failed to get documents: ${response.error}`);

View file

@ -1,7 +1,7 @@
import * as lib from "../../../backend/sync_lib/pkg/sync_lib.js"; import * as lib from "../../../backend/sync_lib/pkg/sync_lib.js";
import { Database } from "src/database/database"; import { Database } from "src/database/database";
import { Logger } from "src/logger"; import { Logger } from "src/logger";
import { SyncServer } from "src/services/sync_service"; import { SyncService } from "src/services/sync_service";
import { hash } from "src/utils/hash"; import { hash } from "src/utils/hash";
import { unlockDocument, waitForDocumentLock } from "./locks.js"; import { unlockDocument, waitForDocumentLock } from "./locks.js";
import { FileOperations } from "src/file-operations/file-operations.js"; import { FileOperations } from "src/file-operations/file-operations.js";
@ -16,7 +16,7 @@ export async function syncLocallyCreatedFile({
filePath, filePath,
}: { }: {
database: Database; database: Database;
syncServer: SyncServer; syncServer: SyncService;
operations: FileOperations; operations: FileOperations;
updateTime: Date; updateTime: Date;
filePath: RelativePath; filePath: RelativePath;

View file

@ -1,14 +1,18 @@
import { Database } from "src/database/database"; import { Database } from "src/database/database";
import { RelativePath } from "src/database/document-metadata"; import { RelativePath } from "src/database/document-metadata";
import { Logger } from "src/logger"; import { Logger } from "src/logger";
import { SyncServer } from "src/services/sync_service"; import { SyncService } from "src/services/sync_service";
import { unlockDocument, waitForDocumentLock } from "./locks"; import { unlockDocument, waitForDocumentLock } from "./locks";
export async function syncLocallyDeletedFile( export async function syncLocallyDeletedFile({
database: Database, database,
syncServer: SyncServer, syncServer,
relativePath: RelativePath relativePath,
): Promise<void> { }: {
database: Database;
syncServer: SyncService;
relativePath: RelativePath;
}): Promise<void> {
await waitForDocumentLock(relativePath); await waitForDocumentLock(relativePath);
try { try {

View file

@ -1,7 +1,7 @@
import * as lib from "../../../backend/sync_lib/pkg/sync_lib.js"; import * as lib from "../../../backend/sync_lib/pkg/sync_lib.js";
import { Database } from "src/database/database"; import { Database } from "src/database/database";
import { Logger } from "src/logger"; import { Logger } from "src/logger";
import { SyncServer } from "src/services/sync_service"; import { SyncService } from "src/services/sync_service";
import { hash } from "src/utils/hash"; import { hash } from "src/utils/hash";
import { unlockDocument, waitForDocumentLock } from "./locks.js"; import { unlockDocument, waitForDocumentLock } from "./locks.js";
import { FileOperations } from "src/file-operations/file-operations.js"; import { FileOperations } from "src/file-operations/file-operations.js";
@ -17,7 +17,7 @@ export async function syncLocallyUpdatedFile({
oldPath, oldPath,
}: { }: {
database: Database; database: Database;
syncServer: SyncServer; syncServer: SyncService;
operations: FileOperations; operations: FileOperations;
updateTime: Date; updateTime: Date;
filePath: RelativePath; filePath: RelativePath;

View file

@ -1,6 +1,6 @@
import { Database } from "src/database/database"; import { Database } from "src/database/database";
import { unlockDocument, waitForDocumentLock } from "./locks"; import { unlockDocument, waitForDocumentLock } from "./locks";
import { SyncServer } from "src/services/sync_service"; import { SyncService } from "src/services/sync_service";
import * as lib from "../../../backend/sync_lib/pkg/sync_lib.js"; import * as lib from "../../../backend/sync_lib/pkg/sync_lib.js";
import { hash } from "src/utils/hash"; import { hash } from "src/utils/hash";
import { Logger } from "src/logger"; import { Logger } from "src/logger";
@ -14,7 +14,7 @@ export async function syncRemotelyUpdatedFile({
remoteVersion, remoteVersion,
}: { }: {
database: Database; database: Database;
syncServer: SyncServer; syncServer: SyncService;
operations: FileOperations; operations: FileOperations;
remoteVersion: components["schemas"]["DocumentVersionWithoutContent"]; remoteVersion: components["schemas"]["DocumentVersionWithoutContent"];
}): Promise<void> { }): Promise<void> {