+1
-1
@@ -1,6 +1,6 @@
|
||||
import * as plugos from "../plugos/types.ts";
|
||||
import { EndpointHookT } from "../plugos/hooks/endpoint.ts";
|
||||
import { CronHookT } from "../plugos/hooks/cron.deno.ts";
|
||||
import { CronHookT } from "../plugos/hooks/cron.ts";
|
||||
import { EventHookT } from "../plugos/hooks/event.ts";
|
||||
import { CommandHookT } from "../web/hooks/command.ts";
|
||||
import { SlashCommandHookT } from "../web/hooks/slash_command.ts";
|
||||
|
||||
@@ -130,9 +130,12 @@ Deno.test("Test store", async () => {
|
||||
ternary,
|
||||
new Map<string, SyncStatusItem>(),
|
||||
);
|
||||
console.log("N ops", await sync2.syncFiles());
|
||||
console.log(
|
||||
"N ops",
|
||||
await sync2.syncFiles(SpaceSync.primaryConflictResolver),
|
||||
);
|
||||
await sleep(2);
|
||||
assertEquals(await sync2.syncFiles(), 0);
|
||||
assertEquals(await sync2.syncFiles(SpaceSync.primaryConflictResolver), 0);
|
||||
|
||||
await Deno.remove(primaryPath, { recursive: true });
|
||||
await Deno.remove(secondaryPath, { recursive: true });
|
||||
|
||||
+172
-155
@@ -30,7 +30,7 @@ export class SpaceSync {
|
||||
) {}
|
||||
|
||||
async syncFiles(
|
||||
conflictResolver?: (
|
||||
conflictResolver: (
|
||||
name: string,
|
||||
snapshot: Map<string, SyncStatusItem>,
|
||||
primarySpace: SpacePrimitives,
|
||||
@@ -65,160 +65,12 @@ export class SpaceSync {
|
||||
|
||||
this.logger.log("info", "Iterating over all files");
|
||||
for (const name of allFilesToProcess) {
|
||||
if (
|
||||
primaryFileMap.has(name) && !secondaryFileMap.has(name) &&
|
||||
!this.snapshot.has(name)
|
||||
) {
|
||||
// New file, created on primary, copy from primary to secondary
|
||||
this.logger.log(
|
||||
"info",
|
||||
"New file created on primary, copying to secondary",
|
||||
name,
|
||||
);
|
||||
const { data } = await this.primary.readFile(name, "arraybuffer");
|
||||
const writtenMeta = await this.secondary.writeFile(
|
||||
name,
|
||||
"arraybuffer",
|
||||
data,
|
||||
);
|
||||
this.snapshot.set(name, [
|
||||
primaryFileMap.get(name)!,
|
||||
writtenMeta.lastModified,
|
||||
]);
|
||||
operations++;
|
||||
} else if (
|
||||
secondaryFileMap.has(name) && !primaryFileMap.has(name) &&
|
||||
!this.snapshot.has(name)
|
||||
) {
|
||||
// New file, created on secondary, copy from secondary to primary
|
||||
this.logger.log(
|
||||
"info",
|
||||
"New file created on secondary, copying from secondary to primary",
|
||||
name,
|
||||
);
|
||||
const { data } = await this.secondary.readFile(name, "arraybuffer");
|
||||
const writtenMeta = await this.primary.writeFile(
|
||||
name,
|
||||
"arraybuffer",
|
||||
data,
|
||||
);
|
||||
this.snapshot.set(name, [
|
||||
writtenMeta.lastModified,
|
||||
secondaryFileMap.get(name)!,
|
||||
]);
|
||||
operations++;
|
||||
} else if (
|
||||
primaryFileMap.has(name) && this.snapshot.has(name) &&
|
||||
!secondaryFileMap.has(name)
|
||||
) {
|
||||
// File deleted on B
|
||||
this.logger.log(
|
||||
"info",
|
||||
"File deleted on secondary, deleting from primary",
|
||||
name,
|
||||
);
|
||||
await this.primary.deleteFile(name);
|
||||
this.snapshot.delete(name);
|
||||
operations++;
|
||||
} else if (
|
||||
secondaryFileMap.has(name) && this.snapshot.has(name) &&
|
||||
!primaryFileMap.has(name)
|
||||
) {
|
||||
// File deleted on A
|
||||
this.logger.log(
|
||||
"info",
|
||||
"File deleted on primary, deleting from secondary",
|
||||
name,
|
||||
);
|
||||
await this.secondary.deleteFile(name);
|
||||
this.snapshot.delete(name);
|
||||
operations++;
|
||||
} else if (
|
||||
this.snapshot.has(name) && !primaryFileMap.has(name) &&
|
||||
!secondaryFileMap.has(name)
|
||||
) {
|
||||
// File deleted on both sides, :shrug:
|
||||
this.logger.log(
|
||||
"info",
|
||||
"File deleted on both ends, deleting from status",
|
||||
name,
|
||||
);
|
||||
this.snapshot.delete(name);
|
||||
operations++;
|
||||
} else if (
|
||||
primaryFileMap.has(name) && secondaryFileMap.has(name) &&
|
||||
this.snapshot.get(name) &&
|
||||
primaryFileMap.get(name) !== this.snapshot.get(name)![0] &&
|
||||
secondaryFileMap.get(name) === this.snapshot.get(name)![1]
|
||||
) {
|
||||
// File has changed on primary, but not secondary: copy from primary to secondary
|
||||
this.logger.log(
|
||||
"info",
|
||||
"File changed on primary, copying to secondary",
|
||||
name,
|
||||
);
|
||||
const { data } = await this.primary.readFile(name, "arraybuffer");
|
||||
const writtenMeta = await this.secondary.writeFile(
|
||||
name,
|
||||
"arraybuffer",
|
||||
data,
|
||||
);
|
||||
this.snapshot.set(name, [
|
||||
primaryFileMap.get(name)!,
|
||||
writtenMeta.lastModified,
|
||||
]);
|
||||
operations++;
|
||||
} else if (
|
||||
primaryFileMap.has(name) && secondaryFileMap.has(name) &&
|
||||
this.snapshot.get(name) &&
|
||||
secondaryFileMap.get(name) !== this.snapshot.get(name)![1] &&
|
||||
primaryFileMap.get(name) === this.snapshot.get(name)![0]
|
||||
) {
|
||||
// File has changed on secondary, but not primary: copy from secondary to primary
|
||||
const { data } = await this.secondary.readFile(name, "arraybuffer");
|
||||
const writtenMeta = await this.primary.writeFile(
|
||||
name,
|
||||
"arraybuffer",
|
||||
data,
|
||||
);
|
||||
this.snapshot.set(name, [
|
||||
writtenMeta.lastModified,
|
||||
secondaryFileMap.get(name)!,
|
||||
]);
|
||||
operations++;
|
||||
} else if (
|
||||
( // File changed on both ends, but we don't have any info in the snapshot (resync scenario?): have to run through conflict handling
|
||||
primaryFileMap.has(name) && secondaryFileMap.has(name) &&
|
||||
!this.snapshot.has(name)
|
||||
) ||
|
||||
( // File changed on both ends, CONFLICT!
|
||||
primaryFileMap.has(name) && secondaryFileMap.has(name) &&
|
||||
this.snapshot.get(name) &&
|
||||
secondaryFileMap.get(name) !== this.snapshot.get(name)![1] &&
|
||||
primaryFileMap.get(name) !== this.snapshot.get(name)![0]
|
||||
)
|
||||
) {
|
||||
this.logger.log(
|
||||
"info",
|
||||
"File changed on both ends, potential conflict",
|
||||
name,
|
||||
);
|
||||
if (conflictResolver) {
|
||||
operations += await conflictResolver(
|
||||
name,
|
||||
this.snapshot,
|
||||
this.primary,
|
||||
this.secondary,
|
||||
this.logger,
|
||||
);
|
||||
} else {
|
||||
throw Error(
|
||||
`Sync conflict for ${name} with no conflict resolver specified`,
|
||||
);
|
||||
}
|
||||
} else {
|
||||
// Nothing needs to happen
|
||||
}
|
||||
operations += await this.syncFile(
|
||||
name,
|
||||
primaryFileMap.get(name),
|
||||
secondaryFileMap.get(name),
|
||||
conflictResolver,
|
||||
);
|
||||
}
|
||||
} catch (e: any) {
|
||||
this.logger.log("error", "Sync error:", e.message);
|
||||
@@ -229,6 +81,171 @@ export class SpaceSync {
|
||||
return operations;
|
||||
}
|
||||
|
||||
async syncFile(
|
||||
name: string,
|
||||
primaryHash: SyncHash | undefined,
|
||||
secondaryHash: SyncHash | undefined,
|
||||
conflictResolver: (
|
||||
name: string,
|
||||
snapshot: Map<string, SyncStatusItem>,
|
||||
primarySpace: SpacePrimitives,
|
||||
secondarySpace: SpacePrimitives,
|
||||
logger: Logger,
|
||||
) => Promise<number>,
|
||||
): Promise<number> {
|
||||
let operations = 0;
|
||||
|
||||
if (
|
||||
primaryHash && !secondaryHash &&
|
||||
!this.snapshot.has(name)
|
||||
) {
|
||||
// New file, created on primary, copy from primary to secondary
|
||||
this.logger.log(
|
||||
"info",
|
||||
"New file created on primary, copying to secondary",
|
||||
name,
|
||||
);
|
||||
const { data } = await this.primary.readFile(name, "arraybuffer");
|
||||
const writtenMeta = await this.secondary.writeFile(
|
||||
name,
|
||||
"arraybuffer",
|
||||
data,
|
||||
);
|
||||
this.snapshot.set(name, [
|
||||
primaryHash,
|
||||
writtenMeta.lastModified,
|
||||
]);
|
||||
operations++;
|
||||
} else if (
|
||||
secondaryHash && !primaryHash &&
|
||||
!this.snapshot.has(name)
|
||||
) {
|
||||
// New file, created on secondary, copy from secondary to primary
|
||||
this.logger.log(
|
||||
"info",
|
||||
"New file created on secondary, copying from secondary to primary",
|
||||
name,
|
||||
);
|
||||
const { data } = await this.secondary.readFile(name, "arraybuffer");
|
||||
const writtenMeta = await this.primary.writeFile(
|
||||
name,
|
||||
"arraybuffer",
|
||||
data,
|
||||
);
|
||||
this.snapshot.set(name, [
|
||||
writtenMeta.lastModified,
|
||||
secondaryHash,
|
||||
]);
|
||||
operations++;
|
||||
} else if (
|
||||
primaryHash && this.snapshot.has(name) &&
|
||||
!secondaryHash
|
||||
) {
|
||||
// File deleted on B
|
||||
this.logger.log(
|
||||
"info",
|
||||
"File deleted on secondary, deleting from primary",
|
||||
name,
|
||||
);
|
||||
await this.primary.deleteFile(name);
|
||||
this.snapshot.delete(name);
|
||||
operations++;
|
||||
} else if (
|
||||
secondaryHash && this.snapshot.has(name) &&
|
||||
!primaryHash
|
||||
) {
|
||||
// File deleted on A
|
||||
this.logger.log(
|
||||
"info",
|
||||
"File deleted on primary, deleting from secondary",
|
||||
name,
|
||||
);
|
||||
await this.secondary.deleteFile(name);
|
||||
this.snapshot.delete(name);
|
||||
operations++;
|
||||
} else if (
|
||||
this.snapshot.has(name) && !primaryHash &&
|
||||
!secondaryHash
|
||||
) {
|
||||
// File deleted on both sides, :shrug:
|
||||
this.logger.log(
|
||||
"info",
|
||||
"File deleted on both ends, deleting from status",
|
||||
name,
|
||||
);
|
||||
this.snapshot.delete(name);
|
||||
operations++;
|
||||
} else if (
|
||||
primaryHash && secondaryHash &&
|
||||
this.snapshot.get(name) &&
|
||||
primaryHash !== this.snapshot.get(name)![0] &&
|
||||
secondaryHash === this.snapshot.get(name)![1]
|
||||
) {
|
||||
// File has changed on primary, but not secondary: copy from primary to secondary
|
||||
this.logger.log(
|
||||
"info",
|
||||
"File changed on primary, copying to secondary",
|
||||
name,
|
||||
);
|
||||
const { data } = await this.primary.readFile(name, "arraybuffer");
|
||||
const writtenMeta = await this.secondary.writeFile(
|
||||
name,
|
||||
"arraybuffer",
|
||||
data,
|
||||
);
|
||||
this.snapshot.set(name, [
|
||||
primaryHash,
|
||||
writtenMeta.lastModified,
|
||||
]);
|
||||
operations++;
|
||||
} else if (
|
||||
primaryHash && secondaryHash &&
|
||||
this.snapshot.get(name) &&
|
||||
secondaryHash !== this.snapshot.get(name)![1] &&
|
||||
primaryHash === this.snapshot.get(name)![0]
|
||||
) {
|
||||
// File has changed on secondary, but not primary: copy from secondary to primary
|
||||
const { data } = await this.secondary.readFile(name, "arraybuffer");
|
||||
const writtenMeta = await this.primary.writeFile(
|
||||
name,
|
||||
"arraybuffer",
|
||||
data,
|
||||
);
|
||||
this.snapshot.set(name, [
|
||||
writtenMeta.lastModified,
|
||||
secondaryHash,
|
||||
]);
|
||||
operations++;
|
||||
} else if (
|
||||
( // File changed on both ends, but we don't have any info in the snapshot (resync scenario?): have to run through conflict handling
|
||||
primaryHash && secondaryHash &&
|
||||
!this.snapshot.has(name)
|
||||
) ||
|
||||
( // File changed on both ends, CONFLICT!
|
||||
primaryHash && secondaryHash &&
|
||||
this.snapshot.get(name) &&
|
||||
secondaryHash !== this.snapshot.get(name)![1] &&
|
||||
primaryHash !== this.snapshot.get(name)![0]
|
||||
)
|
||||
) {
|
||||
this.logger.log(
|
||||
"info",
|
||||
"File changed on both ends, potential conflict",
|
||||
name,
|
||||
);
|
||||
operations += await conflictResolver(
|
||||
name,
|
||||
this.snapshot,
|
||||
this.primary,
|
||||
this.secondary,
|
||||
this.logger,
|
||||
);
|
||||
} else {
|
||||
// Nothing needs to happen
|
||||
}
|
||||
return operations;
|
||||
}
|
||||
|
||||
// Strategy: Primary wins
|
||||
public static async primaryConflictResolver(
|
||||
name: string,
|
||||
|
||||
+80
-21
@@ -10,7 +10,7 @@ export function syncSyscalls(
|
||||
system: System<any>,
|
||||
): SysCallMapping {
|
||||
return {
|
||||
"sync.sync": async (
|
||||
"sync.syncAll": async (
|
||||
_ctx,
|
||||
endpoint: SyncEndpoint,
|
||||
snapshot: Record<string, SyncStatusItem>,
|
||||
@@ -22,24 +22,7 @@ export function syncSyscalls(
|
||||
error?: string;
|
||||
}
|
||||
> => {
|
||||
const syncSpace = new HttpSpacePrimitives(
|
||||
endpoint.url,
|
||||
endpoint.user,
|
||||
endpoint.password,
|
||||
// Base64 PUTs to support mobile
|
||||
true,
|
||||
);
|
||||
// Convert from JSON to a Map
|
||||
const syncStatusMap = new Map<string, SyncStatusItem>(
|
||||
Object.entries(snapshot),
|
||||
);
|
||||
const spaceSync = new SpaceSync(
|
||||
localSpace,
|
||||
syncSpace,
|
||||
syncStatusMap,
|
||||
// Log to the "sync" plug sandbox
|
||||
system.loadedPlugs.get("sync")!.sandbox!,
|
||||
);
|
||||
const { spaceSync } = setupSync(endpoint, snapshot);
|
||||
|
||||
try {
|
||||
const operations = await spaceSync.syncFiles(
|
||||
@@ -58,19 +41,95 @@ export function syncSyscalls(
|
||||
};
|
||||
}
|
||||
},
|
||||
"sync.syncFile": async (
|
||||
_ctx,
|
||||
endpoint: SyncEndpoint,
|
||||
snapshot: Record<string, SyncStatusItem>,
|
||||
name: string,
|
||||
): Promise<
|
||||
{
|
||||
snapshot: Record<string, SyncStatusItem>;
|
||||
operations: number;
|
||||
// The reason to not just throw an Error is so that the partially updated snapshot can still be saved
|
||||
error?: string;
|
||||
}
|
||||
> => {
|
||||
const { spaceSync, remoteSpace } = setupSync(endpoint, snapshot);
|
||||
try {
|
||||
const localHash = (await localSpace.getFileMeta(name)).lastModified;
|
||||
let remoteHash: number | undefined = undefined;
|
||||
try {
|
||||
remoteHash =
|
||||
(await race([remoteSpace.getFileMeta(name), timeout(1000)]))
|
||||
.lastModified;
|
||||
} catch (e: any) {
|
||||
if (e.message.includes("File not found")) {
|
||||
// File doesn't exist remotely, that's ok
|
||||
} else {
|
||||
throw e;
|
||||
}
|
||||
}
|
||||
|
||||
const operations = await spaceSync.syncFile(
|
||||
name,
|
||||
localHash,
|
||||
remoteHash,
|
||||
SpaceSync.primaryConflictResolver,
|
||||
);
|
||||
return {
|
||||
// And convert back to JSON
|
||||
snapshot: Object.fromEntries(spaceSync.snapshot),
|
||||
operations,
|
||||
};
|
||||
} catch (e: any) {
|
||||
return {
|
||||
snapshot: Object.fromEntries(spaceSync.snapshot),
|
||||
operations: -1,
|
||||
error: e.message,
|
||||
};
|
||||
}
|
||||
},
|
||||
"sync.check": async (_ctx, endpoint: SyncEndpoint): Promise<void> => {
|
||||
const syncSpace = new HttpSpacePrimitives(
|
||||
endpoint.url,
|
||||
endpoint.user,
|
||||
endpoint.password,
|
||||
);
|
||||
// Let's just fetch the file list to see if it works with a timeout of 5s
|
||||
// Let's just fetch metadata for the SETTINGS.md file (which should always exist)
|
||||
try {
|
||||
await race([syncSpace.fetchFileList(), timeout(5000)]);
|
||||
await race([
|
||||
syncSpace.getFileMeta("SETTINGS.md"),
|
||||
timeout(2000),
|
||||
]);
|
||||
} catch (e: any) {
|
||||
console.error("Sync check failure", e.message);
|
||||
throw e;
|
||||
}
|
||||
},
|
||||
};
|
||||
|
||||
function setupSync(
|
||||
endpoint: SyncEndpoint,
|
||||
snapshot: Record<string, SyncStatusItem>,
|
||||
) {
|
||||
const remoteSpace = new HttpSpacePrimitives(
|
||||
endpoint.url,
|
||||
endpoint.user,
|
||||
endpoint.password,
|
||||
// Base64 PUTs to support mobile
|
||||
true,
|
||||
);
|
||||
// Convert from JSON to a Map
|
||||
const syncStatusMap = new Map<string, SyncStatusItem>(
|
||||
Object.entries(snapshot),
|
||||
);
|
||||
const spaceSync = new SpaceSync(
|
||||
localSpace,
|
||||
remoteSpace,
|
||||
syncStatusMap,
|
||||
// Log to the "sync" plug sandbox
|
||||
system.loadedPlugs.get("sync")!.sandbox!,
|
||||
);
|
||||
return { spaceSync, remoteSpace };
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user