Refactor and improve Syncer API
This commit is contained in:
parent
b0192aae23
commit
8b07507090
2 changed files with 152 additions and 152 deletions
|
|
@ -1,64 +0,0 @@
|
||||||
import type { Database } from "../persistence/database";
|
|
||||||
import type { SyncService } from "src/services/sync-service";
|
|
||||||
import { Logger } from "src/tracing/logger";
|
|
||||||
import type { Syncer } from "./syncer";
|
|
||||||
import type { Settings } from "src/persistence/settings";
|
|
||||||
|
|
||||||
let isRunning = false;
|
|
||||||
|
|
||||||
export async function applyRemoteChangesLocally({
|
|
||||||
settings,
|
|
||||||
database,
|
|
||||||
syncService,
|
|
||||||
syncer
|
|
||||||
}: {
|
|
||||||
settings: Settings;
|
|
||||||
database: Database;
|
|
||||||
syncService: SyncService;
|
|
||||||
syncer: Syncer;
|
|
||||||
}): Promise<void> {
|
|
||||||
if (!settings.getSettings().isSyncEnabled) {
|
|
||||||
Logger.getInstance().debug(
|
|
||||||
`Syncing is disabled, not fetching remote changes`
|
|
||||||
);
|
|
||||||
return;
|
|
||||||
} else if (isRunning) {
|
|
||||||
Logger.getInstance().debug(
|
|
||||||
"Applying remote changes locally is already in progress, skipping invocation"
|
|
||||||
);
|
|
||||||
return;
|
|
||||||
}
|
|
||||||
|
|
||||||
isRunning = true;
|
|
||||||
|
|
||||||
try {
|
|
||||||
const remote = await syncService.getAll(database.getLastSeenUpdateId());
|
|
||||||
|
|
||||||
if (remote.latestDocuments.length === 0) {
|
|
||||||
Logger.getInstance().debug("No remote changes to apply");
|
|
||||||
return;
|
|
||||||
}
|
|
||||||
|
|
||||||
Logger.getInstance().info("Applying remote changes locally");
|
|
||||||
|
|
||||||
await Promise.all(
|
|
||||||
remote.latestDocuments.map(async (remoteDocument) =>
|
|
||||||
syncer.syncRemotelyUpdatedFile(remoteDocument)
|
|
||||||
)
|
|
||||||
);
|
|
||||||
|
|
||||||
const lastSeenUpdateId = database.getLastSeenUpdateId();
|
|
||||||
if (
|
|
||||||
lastSeenUpdateId === undefined ||
|
|
||||||
remote.lastUpdateId > lastSeenUpdateId
|
|
||||||
) {
|
|
||||||
await database.setLastSeenUpdateId(remote.lastUpdateId);
|
|
||||||
}
|
|
||||||
} catch (e) {
|
|
||||||
Logger.getInstance().error(
|
|
||||||
`Failed to apply remote changes locally: ${e}`
|
|
||||||
);
|
|
||||||
} finally {
|
|
||||||
isRunning = false;
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
@ -4,7 +4,6 @@ import type {
|
||||||
RelativePath
|
RelativePath
|
||||||
} from "../persistence/database";
|
} from "../persistence/database";
|
||||||
|
|
||||||
import type { FileOperations } from "src/file-operations";
|
|
||||||
import type { SyncService } from "src/services/sync-service";
|
import type { SyncService } from "src/services/sync-service";
|
||||||
import { Logger } from "src/tracing/logger";
|
import { Logger } from "src/tracing/logger";
|
||||||
import type { SyncHistory } from "src/tracing/sync-history";
|
import type { SyncHistory } from "src/tracing/sync-history";
|
||||||
|
|
@ -15,6 +14,7 @@ import { EMPTY_HASH, hash } from "src/utils/hash";
|
||||||
import type { components } from "src/services/types";
|
import type { components } from "src/services/types";
|
||||||
import { deserialize } from "src/utils/deserialize";
|
import { deserialize } from "src/utils/deserialize";
|
||||||
import type { Settings } from "src/persistence/settings";
|
import type { Settings } from "src/persistence/settings";
|
||||||
|
import { FileOperations } from "src/file-operations/file-operations";
|
||||||
|
|
||||||
export class Syncer {
|
export class Syncer {
|
||||||
private readonly remainingOperationsListeners: ((
|
private readonly remainingOperationsListeners: ((
|
||||||
|
|
@ -23,7 +23,10 @@ export class Syncer {
|
||||||
|
|
||||||
private readonly syncQueue: PQueue;
|
private readonly syncQueue: PQueue;
|
||||||
|
|
||||||
private isRunningOfflineSync = false;
|
private runningScheduleSyncForOfflineChanges: Promise<void> | undefined =
|
||||||
|
undefined;
|
||||||
|
private runningApplyRemoteChangesLocally: Promise<void> | undefined =
|
||||||
|
undefined;
|
||||||
|
|
||||||
public constructor(
|
public constructor(
|
||||||
private readonly database: Database,
|
private readonly database: Database,
|
||||||
|
|
@ -78,7 +81,7 @@ export class Syncer {
|
||||||
);
|
);
|
||||||
}
|
}
|
||||||
|
|
||||||
public async syncRemotelyUpdatedFile(
|
private async syncRemotelyUpdatedFile(
|
||||||
remoteVersion: components["schemas"]["DocumentVersionWithoutContent"]
|
remoteVersion: components["schemas"]["DocumentVersionWithoutContent"]
|
||||||
): Promise<void> {
|
): Promise<void> {
|
||||||
await this.syncQueue.add(async () =>
|
await this.syncQueue.add(async () =>
|
||||||
|
|
@ -87,13 +90,6 @@ export class Syncer {
|
||||||
}
|
}
|
||||||
|
|
||||||
public async scheduleSyncForOfflineChanges(): Promise<void> {
|
public async scheduleSyncForOfflineChanges(): Promise<void> {
|
||||||
if (this.isRunningOfflineSync) {
|
|
||||||
Logger.getInstance().warn(
|
|
||||||
"Uploading local changes is already in progress, skipping"
|
|
||||||
);
|
|
||||||
return;
|
|
||||||
}
|
|
||||||
|
|
||||||
if (!this.settings.getSettings().isSyncEnabled) {
|
if (!this.settings.getSettings().isSyncEnabled) {
|
||||||
Logger.getInstance().debug(
|
Logger.getInstance().debug(
|
||||||
`Syncing is disabled, not uploading local changes`
|
`Syncing is disabled, not uploading local changes`
|
||||||
|
|
@ -101,9 +97,31 @@ export class Syncer {
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
|
|
||||||
this.isRunningOfflineSync = true;
|
if (this.runningScheduleSyncForOfflineChanges != null) {
|
||||||
|
Logger.getInstance().debug(
|
||||||
|
"Uploading local changes is already in progress"
|
||||||
|
);
|
||||||
|
return this.runningScheduleSyncForOfflineChanges;
|
||||||
|
}
|
||||||
|
|
||||||
try {
|
try {
|
||||||
|
this.runningScheduleSyncForOfflineChanges =
|
||||||
|
this.internalScheduleSyncForOfflineChanges();
|
||||||
|
await this.runningScheduleSyncForOfflineChanges;
|
||||||
|
Logger.getInstance().info(
|
||||||
|
`All local changes have been applied remotely`
|
||||||
|
);
|
||||||
|
} catch (e) {
|
||||||
|
Logger.getInstance().error(
|
||||||
|
`Not all local changes have been applied remotely: ${e}`
|
||||||
|
);
|
||||||
|
throw e;
|
||||||
|
} finally {
|
||||||
|
this.runningScheduleSyncForOfflineChanges = undefined;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
private async internalScheduleSyncForOfflineChanges(): Promise<void> {
|
||||||
const allLocalFiles = await this.operations.listAllFiles();
|
const allLocalFiles = await this.operations.listAllFiles();
|
||||||
let locallyDeletedFiles = [
|
let locallyDeletedFiles = [
|
||||||
...this.database.getDocuments().entries()
|
...this.database.getDocuments().entries()
|
||||||
|
|
@ -112,8 +130,7 @@ export class Syncer {
|
||||||
await Promise.all(
|
await Promise.all(
|
||||||
allLocalFiles.map(async (relativePath) =>
|
allLocalFiles.map(async (relativePath) =>
|
||||||
this.syncQueue.add(async () => {
|
this.syncQueue.add(async () => {
|
||||||
const metadata =
|
const metadata = this.database.getDocument(relativePath);
|
||||||
this.database.getDocument(relativePath);
|
|
||||||
|
|
||||||
// If there's no metadata, it must be a new file
|
// If there's no metadata, it must be a new file
|
||||||
if (!metadata) {
|
if (!metadata) {
|
||||||
|
|
@ -129,8 +146,7 @@ export class Syncer {
|
||||||
);
|
);
|
||||||
if (originalFile !== undefined) {
|
if (originalFile !== undefined) {
|
||||||
// `originalFile` hasn't been deleted but it got moved instead
|
// `originalFile` hasn't been deleted but it got moved instead
|
||||||
locallyDeletedFiles =
|
locallyDeletedFiles = locallyDeletedFiles.filter(
|
||||||
locallyDeletedFiles.filter(
|
|
||||||
(item) => item != originalFile
|
(item) => item != originalFile
|
||||||
);
|
);
|
||||||
|
|
||||||
|
|
@ -192,16 +208,64 @@ export class Syncer {
|
||||||
return this.internalSyncLocallyDeletedFile(relativePath);
|
return this.internalSyncLocallyDeletedFile(relativePath);
|
||||||
})
|
})
|
||||||
);
|
);
|
||||||
|
}
|
||||||
|
|
||||||
|
public async applyRemoteChangesLocally(): Promise<void> {
|
||||||
|
if (!this.settings.getSettings().isSyncEnabled) {
|
||||||
|
Logger.getInstance().debug(
|
||||||
|
`Syncing is disabled, not fetching remote changes`
|
||||||
|
);
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
|
||||||
|
if (this.runningApplyRemoteChangesLocally != null) {
|
||||||
|
Logger.getInstance().debug(
|
||||||
|
"Applying remote changes locally is already in progress"
|
||||||
|
);
|
||||||
|
return this.runningApplyRemoteChangesLocally;
|
||||||
|
}
|
||||||
|
|
||||||
|
try {
|
||||||
|
this.runningApplyRemoteChangesLocally =
|
||||||
|
this.internalApplyRemoteChangesLocally();
|
||||||
|
await this.runningApplyRemoteChangesLocally;
|
||||||
Logger.getInstance().info(
|
Logger.getInstance().info(
|
||||||
`All local changes have been applied remotely`
|
"All remote changes have been applied locally"
|
||||||
);
|
);
|
||||||
} catch (e) {
|
} catch (e) {
|
||||||
Logger.getInstance().error(
|
Logger.getInstance().error(
|
||||||
`Not all local changes have been applied remotely: ${e}`
|
`Failed to apply remote changes locally: ${e}`
|
||||||
);
|
);
|
||||||
|
throw e;
|
||||||
} finally {
|
} finally {
|
||||||
this.isRunningOfflineSync = false;
|
this.runningApplyRemoteChangesLocally = undefined;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
private async internalApplyRemoteChangesLocally(): Promise<void> {
|
||||||
|
const remote = await this.syncService.getAll(
|
||||||
|
this.database.getLastSeenUpdateId()
|
||||||
|
);
|
||||||
|
|
||||||
|
if (remote.latestDocuments.length === 0) {
|
||||||
|
Logger.getInstance().debug("No remote changes to apply");
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
|
||||||
|
Logger.getInstance().info("Applying remote changes locally");
|
||||||
|
|
||||||
|
await Promise.all(
|
||||||
|
remote.latestDocuments.map(async (remoteDocument) =>
|
||||||
|
this.syncRemotelyUpdatedFile(remoteDocument)
|
||||||
|
)
|
||||||
|
);
|
||||||
|
|
||||||
|
const lastSeenUpdateId = this.database.getLastSeenUpdateId();
|
||||||
|
if (
|
||||||
|
lastSeenUpdateId === undefined ||
|
||||||
|
remote.lastUpdateId > lastSeenUpdateId
|
||||||
|
) {
|
||||||
|
await this.database.setLastSeenUpdateId(remote.lastUpdateId);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue