Compare commits

...
Sign in to create a new pull request.

17 commits

168 changed files with 20441 additions and 17396 deletions

View file

@ -11,5 +11,6 @@ indent_style = space
indent_size = 4 indent_size = 4
tab_width = 4 tab_width = 4
[*.{yml,yaml}] [*.{yml,yaml,md}]
indent_size = 2 indent_size = 2
tab_width = 2

View file

@ -6,7 +6,7 @@ on:
pull_request: pull_request:
branches: ["main"] branches: ["main"]
schedule: schedule:
- cron: '*/30 * * * *' - cron: '0 * * * *'
concurrency: concurrency:
group: e2e-tests group: e2e-tests

View file

@ -1,7 +1,7 @@
{ {
"printWidth": 120, "printWidth": 120,
"tabWidth": 4, "tabWidth": 4,
"useTabs": true, "useTabs": false,
"semi": false, "semi": false,
"singleQuote": false, "singleQuote": false,
"trailingComma": "none", "trailingComma": "none",

View file

@ -125,7 +125,7 @@ sequenceDiagram
``` ```
┌─────────┐ ┌─────────┐
│ Client │ │ Client │
└───┬────┘ └───┬─-───┘
│ 1. Detect file change │ 1. Detect file change
├─► 2. Read file content ├─► 2. Read file content
@ -142,7 +142,7 @@ sequenceDiagram
┌─────────┐ ┌─────────┐
│ Server │ │ Server │
└───┬────┘ └───┬────-
│ 4. Validate message │ 4. Validate message
├─► 5. Check permissions ├─► 5. Check permissions
@ -163,7 +163,7 @@ sequenceDiagram
``` ```
┌─────────┐ ┌─────────┐
│ Server │ │ Server │
└───┬────┘ └───┬─-───┘
│ 1. File updated by another client │ 1. File updated by another client
├─► 2. Broadcast notification ├─► 2. Broadcast notification
@ -176,7 +176,7 @@ sequenceDiagram
┌─────────┐ ┌─────────┐
│ Client │ │ Client │
└───┬────┘ └───┬─-───┘
│ 3. Receive notification │ 3. Receive notification
├─► 4. Request file download ├─► 4. Request file download
@ -189,7 +189,7 @@ sequenceDiagram
┌─────────┐ ┌─────────┐
│ Server │ │ Server │
└───┬────┘ └───┬─=───┘
│ 5. Retrieve from database │ 5. Retrieve from database
└─► 6. Send file content └─► 6. Send file content
@ -203,7 +203,7 @@ sequenceDiagram
┌─────────┐ ┌─────────┐
│ Client │ │ Client │
└────┬─── └───-─┬───┘
│ 7. Write to filesystem │ 7. Write to filesystem
└─► 8. Update local metadata └─► 8. Update local metadata

View file

@ -55,6 +55,25 @@ export default [
message: "Use replaceAll instead of replace to replace all occurrences of a substring." message: "Use replaceAll instead of replace to replace all occurrences of a substring."
} }
], ],
"no-restricted-syntax": [
"error",
{
selector: "CallExpression[callee.property.name='splice'][arguments.length=2][arguments.1.type='Literal'][arguments.1.value=1]",
message: "Use `removeFromArray(array, item)` instead of manually using indexOf + splice(index, 1). Import from 'sync-client/src/utils/remove-from-array'."
},
{
selector: "CallExpression[callee.property.name='filter'] > ArrowFunctionExpression[body.type='BinaryExpression'][body.operator='!==']",
message: "Use `removeFromArray(array, item)` instead of filter(x => x !== item) for better performance. Import from 'sync-client/src/utils/remove-from-array'."
},
{
selector: "CallExpression[callee.property.name='filter'] > ArrowFunctionExpression > BlockStatement > ReturnStatement > BinaryExpression[operator='!==']",
message: "Use `removeFromArray(array, item)` instead of filter(x => { return x !== item }) for better performance. Import from 'sync-client/src/utils/remove-from-array'."
},
{
selector: "CallExpression[callee.property.name='filter'] > FunctionExpression[body.type='BlockStatement'] > BlockStatement > ReturnStatement > BinaryExpression[operator='!==']",
message: "Use `removeFromArray(array, item)` instead of filter(function(x) { return x !== item }) for better performance. Import from 'sync-client/src/utils/remove-from-array'."
}
],
"unused-imports/no-unused-vars": [ "unused-imports/no-unused-vars": [
"warn", "warn",
{ {

View file

@ -9,10 +9,11 @@
"scripts": { "scripts": {
"dev": "webpack watch --mode development", "dev": "webpack watch --mode development",
"build": "webpack --mode production", "build": "webpack --mode production",
"test": "tsx --test src/args.test.ts src/node-filesystem.test.ts" "test": "tsx --test 'src/**/*.test.ts'"
}, },
"dependencies": { "dependencies": {
"commander": "^14.0.2" "commander": "^14.0.2",
"watcher": "^2.3.1"
}, },
"devDependencies": { "devDependencies": {
"@types/node": "^24.8.1", "@types/node": "^24.8.1",

View file

@ -39,6 +39,10 @@ async function main(): Promise<void> {
const args = parseArgs(process.argv); const args = parseArgs(process.argv);
const absolutePath = path.resolve(args.localPath); const absolutePath = path.resolve(args.localPath);
if (!fsSync.existsSync(absolutePath)) {
fsSync.mkdirSync(absolutePath, { recursive: true });
}
try { try {
const stats = await fs.stat(absolutePath); const stats = await fs.stat(absolutePath);
if (!stats.isDirectory()) { if (!stats.isDirectory()) {
@ -153,7 +157,7 @@ async function main(): Promise<void> {
} }
// Add colored log formatter with level filtering // Add colored log formatter with level filtering
client.logger.addOnMessageListener((logLine) => { client.logger.onLogEmitted.add((logLine) => {
// Only show messages at or above the configured log level // Only show messages at or above the configured log level
if (LOG_LEVEL_ORDER[logLine.level] >= LOG_LEVEL_ORDER[args.logLevel]) { if (LOG_LEVEL_ORDER[logLine.level] >= LOG_LEVEL_ORDER[args.logLevel]) {
console.log(formatLogLine(logLine)); console.log(formatLogLine(logLine));
@ -164,14 +168,14 @@ async function main(): Promise<void> {
const fileWatcher = new FileWatcher(absolutePath, client); const fileWatcher = new FileWatcher(absolutePath, client);
client.addWebSocketStatusChangeListener(() => { client.onWebSocketStatusChanged.add(() => {
const isConnected = client.isWebSocketConnected; const isConnected = client.isWebSocketConnected;
client.logger.info( client.logger.info(
`WebSocket status changed: ${isConnected ? "connected" : "disconnected"}` `WebSocket status changed: ${isConnected ? "connected" : "disconnected"}`
); );
}); });
client.addRemainingSyncOperationsListener((remaining) => { client.onRemainingOperationsCountChanged.add((remaining) => {
if (remaining === 0) { if (remaining === 0) {
client.logger.info("All sync operations completed"); client.logger.info("All sync operations completed");
} else { } else {

View file

@ -1,9 +1,9 @@
import * as fs from "fs"; import Watcher from "watcher";
import * as path from "path"; import * as path from "path";
import type { SyncClient, RelativePath } from "sync-client"; import type { SyncClient, RelativePath } from "sync-client";
export class FileWatcher { export class FileWatcher {
private watcher: fs.FSWatcher | undefined; private watcher: Watcher | undefined;
private isRunning = false; private isRunning = false;
public constructor( public constructor(
@ -18,25 +18,31 @@ export class FileWatcher {
this.isRunning = true; this.isRunning = true;
this.watcher = fs.watch( this.watcher = new Watcher(this.basePath, {
this.basePath, recursive: true,
{ recursive: true }, renameDetection: true,
(eventType, filename) => { renameTimeout: 125,
if (filename === null || filename.length === 0) { ignoreInitial: true
return; });
}
// Convert to forward slashes for consistency this.watcher.on("add", (filePath: string) => {
const relativePath = this.toUnixPath(filename); this.handleCreate(this.toRelativePath(filePath));
});
if (eventType === "rename") { this.watcher.on("change", (filePath: string) => {
this.handleRenameOrDelete(relativePath); this.handleChange(this.toRelativePath(filePath));
} else { });
// Must be "change" event
this.handleChange(relativePath); this.watcher.on("unlink", (filePath: string) => {
} this.handleDelete(this.toRelativePath(filePath));
} });
this.watcher.on("rename", (oldPath: string, newPath: string) => {
this.handleRename(
this.toRelativePath(oldPath),
this.toRelativePath(newPath)
); );
});
this.client.logger.info("File watcher started"); this.client.logger.info("File watcher started");
} }
@ -50,44 +56,53 @@ export class FileWatcher {
this.client.logger.info("File watcher stopped"); this.client.logger.info("File watcher stopped");
} }
private handleCreate(relativePath: RelativePath): void {
this.client
.syncLocallyCreatedFile(relativePath)
.catch((err: unknown) => {
this.client.logger.error(
`Failed to sync created file ${relativePath}: ${this.formatError(err)}`
);
});
}
private handleChange(relativePath: RelativePath): void { private handleChange(relativePath: RelativePath): void {
this.client this.client
.syncLocallyUpdatedFile({ relativePath }) .syncLocallyUpdatedFile({ relativePath })
.catch((err: unknown) => { .catch((err: unknown) => {
this.client.logger.error( this.client.logger.error(
`Failed to sync updated file ${relativePath}: ${err instanceof Error ? err.message : String(err)}` `Failed to sync updated file ${relativePath}: ${this.formatError(err)}`
); );
}); });
} }
private handleRenameOrDelete(relativePath: RelativePath): void { private handleDelete(relativePath: RelativePath): void {
const fullPath = path.join(this.basePath, relativePath);
fs.access(fullPath, fs.constants.F_OK, (accessError) => {
if (accessError) {
this.client this.client
.syncLocallyDeletedFile(relativePath) .syncLocallyDeletedFile(relativePath)
.catch((deleteErr: unknown) => { .catch((err: unknown) => {
this.client.logger.error( this.client.logger.error(
`Failed to sync deleted file ${relativePath}: ${deleteErr instanceof Error ? deleteErr.message : String(deleteErr)}` `Failed to sync deleted file ${relativePath}: ${this.formatError(err)}`
); );
}); });
} else {
fs.stat(fullPath, (statErr, stats) => {
if (statErr !== null || !stats.isFile()) {
return;
} }
private handleRename(oldPath: RelativePath, newPath: RelativePath): void {
this.client.logger.info(`File renamed: ${oldPath} -> ${newPath}`);
this.client this.client
.syncLocallyCreatedFile(relativePath) .syncLocallyUpdatedFile({
.catch((createErr: unknown) => { oldPath,
relativePath: newPath
})
.catch((err: unknown) => {
this.client.logger.error( this.client.logger.error(
`Failed to sync created file ${relativePath}: ${createErr instanceof Error ? createErr.message : String(createErr)}` `Failed to sync renamed file ${oldPath} -> ${newPath}: ${this.formatError(err)}`
); );
}); });
});
} }
});
private toRelativePath(absolutePath: string): RelativePath {
const relative = path.relative(this.basePath, absolutePath);
return this.toUnixPath(relative);
} }
/** /**
@ -99,4 +114,8 @@ export class FileWatcher {
} }
return nativePath; return nativePath;
} }
private formatError(err: unknown): string {
return err instanceof Error ? err.message : String(err);
}
} }

View file

@ -18,5 +18,7 @@
"declarationMap": true, "declarationMap": true,
"sourceMap": true "sourceMap": true
}, },
"exclude": ["dist"] "exclude": [
"dist"
]
} }

View file

@ -0,0 +1 @@

View file

@ -85,8 +85,3 @@ If you have multiple URLs, you can also do:
## API Documentation ## API Documentation
See https://github.com/obsidianmd/obsidian-api See https://github.com/obsidianmd/obsidian-api

View file

@ -136,10 +136,7 @@ export default class VaultLinkPlugin extends Plugin {
...(IS_DEBUG_BUILD ...(IS_DEBUG_BUILD
? { ? {
fetch: debugging.slowFetchFactory(1), fetch: debugging.slowFetchFactory(1),
webSocket: debugging.slowWebSocketFactory( webSocket: debugging.slowWebSocketFactory(1, new Logger())
1,
new Logger()
)
} }
: {}) : {})
}); });
@ -174,7 +171,7 @@ export default class VaultLinkPlugin extends Plugin {
this.registerEditorExtension([remoteCursorsTheme, remoteCursorsPlugin]); this.registerEditorExtension([remoteCursorsTheme, remoteCursorsPlugin]);
client.addRemoteCursorsUpdateListener((cursors) => { client.onRemoteCursorsUpdated.add((cursors) => {
RemoteCursorsPluginValue.setCursors(cursors, this.app); RemoteCursorsPluginValue.setCursors(cursors, this.app);
renderCursorsInFileExplorer(cursors, this.app); renderCursorsInFileExplorer(cursors, this.app);
}); });

View file

@ -24,7 +24,7 @@ export class HistoryView extends ItemView {
super(leaf); super(leaf);
this.icon = HistoryView.ICON; this.icon = HistoryView.ICON;
this.client.addSyncHistoryUpdateListener(async () => this.client.onSyncHistoryUpdated.add(async () =>
this.updateView().catch((error: unknown) => { this.updateView().catch((error: unknown) => {
this.client.logger.error( this.client.logger.error(
`Failed to update history view: ${error}` `Failed to update history view: ${error}`

View file

@ -21,7 +21,7 @@ export class LogsView extends ItemView {
) { ) {
super(leaf); super(leaf);
this.icon = LogsView.ICON; this.icon = LogsView.ICON;
this.client.logger.addOnMessageListener(() => { this.client.logger.onLogEmitted.add(() => {
this.updateView(); this.updateView();
}); });
} }

View file

@ -41,8 +41,7 @@ export class SyncSettingsTab extends PluginSettingTab {
this.editedToken = this.syncClient.getSettings().token; this.editedToken = this.syncClient.getSettings().token;
this.editedVaultName = this.syncClient.getSettings().vaultName; this.editedVaultName = this.syncClient.getSettings().vaultName;
this.syncClient.addOnSettingsChangeListener( this.syncClient.onSettingsChanged.add((newSettings, oldSettings) => {
(newSettings, oldSettings) => {
let hasChanged = false; let hasChanged = false;
if (newSettings.remoteUri !== oldSettings.remoteUri) { if (newSettings.remoteUri !== oldSettings.remoteUri) {
@ -63,8 +62,7 @@ export class SyncSettingsTab extends PluginSettingTab {
if (hasChanged) { if (hasChanged) {
this.display(); this.display();
} }
} });
);
} }
private get isApplyingChanges(): boolean { private get isApplyingChanges(): boolean {

View file

@ -14,19 +14,19 @@ export class StatusBar {
private readonly syncClient: SyncClient private readonly syncClient: SyncClient
) { ) {
this.statusBarItem = plugin.addStatusBarItem(); this.statusBarItem = plugin.addStatusBarItem();
this.syncClient.addSyncHistoryUpdateListener((status) => { this.syncClient.onSyncHistoryUpdated.add((status) => {
this.lastHistoryStats = status; this.lastHistoryStats = status;
this.updateStatus(); this.updateStatus();
}); });
this.syncClient.addRemainingSyncOperationsListener( this.syncClient.onRemainingOperationsCountChanged.add(
(remainingOperations) => { (remainingOperations) => {
this.lastRemaining = remainingOperations; this.lastRemaining = remainingOperations;
this.updateStatus(); this.updateStatus();
} }
); );
this.syncClient.addOnSettingsChangeListener(() => { this.syncClient.onSettingsChanged.add(() => {
this.updateStatus(); this.updateStatus();
}); });
} }

View file

@ -5,34 +5,35 @@ import type {
NetworkConnectionStatus, NetworkConnectionStatus,
SyncClient SyncClient
} from "sync-client"; } from "sync-client";
import { utils } from "sync-client";
export class StatusDescription { export class StatusDescription {
private lastHistoryStats: HistoryStats | undefined; private lastHistoryStats: HistoryStats | undefined;
private lastRemaining: number | undefined; private lastRemaining: number | undefined;
private lastConnectionState: NetworkConnectionStatus | undefined; private lastConnectionState: NetworkConnectionStatus | undefined;
private statusChangeListeners: (() => unknown)[] = []; private readonly statusChangeListeners: (() => unknown)[] = [];
public constructor(private readonly syncClient: SyncClient) { public constructor(private readonly syncClient: SyncClient) {
void this.updateConnectionState(); void this.updateConnectionState();
syncClient.addSyncHistoryUpdateListener((status) => { syncClient.onSyncHistoryUpdated.add((status) => {
this.lastHistoryStats = status; this.lastHistoryStats = status;
this.updateDescription(); this.updateDescription();
}); });
this.syncClient.addRemainingSyncOperationsListener( this.syncClient.onRemainingOperationsCountChanged.add(
(remainingOperations) => { (remainingOperations) => {
this.lastRemaining = remainingOperations; this.lastRemaining = remainingOperations;
this.updateDescription(); this.updateDescription();
} }
); );
this.syncClient.addWebSocketStatusChangeListener(async () => this.syncClient.onWebSocketStatusChanged.add(async () =>
this.updateConnectionState() this.updateConnectionState()
); );
this.syncClient.addOnSettingsChangeListener(async () => this.syncClient.onSettingsChanged.add(async () =>
this.updateConnectionState() this.updateConnectionState()
); );
} }
@ -46,9 +47,7 @@ export class StatusDescription {
this.statusChangeListeners.push(listener); this.statusChangeListeners.push(listener);
} }
public removeStatusChangeListener(listener: () => unknown): void { public removeStatusChangeListener(listener: () => unknown): void {
this.statusChangeListeners = this.statusChangeListeners.filter( utils.removeFromArray(this.statusChangeListeners, listener);
(l) => l !== listener
);
} }
public renderStatusDescription(container: HTMLElement): void { public renderStatusDescription(container: HTMLElement): void {

File diff suppressed because it is too large Load diff

View file

@ -10,7 +10,7 @@
"prettier": { "prettier": {
"trailingComma": "none", "trailingComma": "none",
"tabWidth": 4, "tabWidth": 4,
"useTabs": true, "useTabs": false,
"endOfLine": "lf" "endOfLine": "lf"
}, },
"scripts": { "scripts": {
@ -22,6 +22,7 @@
}, },
"devDependencies": { "devDependencies": {
"concurrently": "^9.2.1", "concurrently": "^9.2.1",
"eclint": "^2.8.1",
"eslint": "9.38.0", "eslint": "9.38.0",
"eslint-plugin-unused-imports": "^4.1.4", "eslint-plugin-unused-imports": "^4.1.4",
"npm-check-updates": "^19.1.1", "npm-check-updates": "^19.1.1",

View file

@ -10,7 +10,7 @@
"scripts": { "scripts": {
"dev": "webpack watch --mode development", "dev": "webpack watch --mode development",
"build": "webpack --mode production", "build": "webpack --mode production",
"test": "tsx --test src/**/*.test.ts" "test": "tsx --test 'src/**/*.test.ts'"
}, },
"devDependencies": { "devDependencies": {
"byte-base64": "^1.1.0", "byte-base64": "^1.1.0",

View file

@ -5,6 +5,7 @@ import { slowWebSocketFactory } from "./utils/debugging/slow-web-socket-factory"
import { getRandomColor } from "./utils/get-random-color"; import { getRandomColor } from "./utils/get-random-color";
import { lineAndColumnToPosition } from "./utils/line-and-column-to-position"; import { lineAndColumnToPosition } from "./utils/line-and-column-to-position";
import { positionToLineAndColumn } from "./utils/position-to-line-and-column"; import { positionToLineAndColumn } from "./utils/position-to-line-and-column";
import { removeFromArray } from "./utils/remove-from-array";
export { export {
SyncType, SyncType,
@ -43,5 +44,6 @@ export const utils = {
getRandomColor, getRandomColor,
positionToLineAndColumn, positionToLineAndColumn,
lineAndColumnToPosition, lineAndColumnToPosition,
awaitAll awaitAll,
removeFromArray
}; };

View file

@ -2,6 +2,7 @@ import type { Logger } from "../tracing/logger";
import { EMPTY_HASH } from "../utils/hash"; import { EMPTY_HASH } from "../utils/hash";
import { CoveredValues } from "../utils/data-structures/min-covered"; import { CoveredValues } from "../utils/data-structures/min-covered";
import { awaitAll } from "../utils/await-all"; import { awaitAll } from "../utils/await-all";
import { removeFromArray } from "../utils/remove-from-array";
export type VaultUpdateId = number; export type VaultUpdateId = number;
export type DocumentId = string; export type DocumentId = string;
@ -93,6 +94,7 @@ export class Database {
public get resolvedDocuments(): DocumentRecord[] { public get resolvedDocuments(): DocumentRecord[] {
const paths = new Map<string, DocumentRecord[]>(); const paths = new Map<string, DocumentRecord[]>();
this.documents this.documents
// eslint-disable-next-line no-restricted-syntax -- Type narrowing, not removing a specific item
.filter(({ metadata }) => metadata !== undefined) .filter(({ metadata }) => metadata !== undefined)
.forEach((record) => .forEach((record) =>
paths.set(record.relativePath, [ paths.set(record.relativePath, [
@ -151,12 +153,12 @@ export class Database {
return; return;
} }
entry.updates = entry.updates.filter((update) => update !== promise); removeFromArray(entry.updates, promise);
// No need to save as Promises don't get serialized // No need to save as Promises don't get serialized
} }
public removeDocument(find: DocumentRecord): void { public removeDocument(find: DocumentRecord): void {
this.documents = this.documents.filter((document) => document !== find); removeFromArray(this.documents, find);
this.saveInTheBackground(); this.saveInTheBackground();
} }

View file

@ -1,6 +1,6 @@
import type { Logger } from "../tracing/logger"; import type { Logger } from "../tracing/logger";
import { awaitAll } from "../utils/await-all";
import { Lock } from "../utils/data-structures/locks"; import { Lock } from "../utils/data-structures/locks";
import { EventListeners } from "../utils/data-structures/event-listeners";
export interface SyncSettings { export interface SyncSettings {
remoteUri: string; remoteUri: string;
@ -33,14 +33,13 @@ export const DEFAULT_SETTINGS: SyncSettings = {
}; };
export class Settings { export class Settings {
public readonly onSettingsChanged = new EventListeners<
(newSettings: SyncSettings, oldSettings: SyncSettings) => unknown
>();
private settings: SyncSettings; private settings: SyncSettings;
private readonly lock: Lock = new Lock(); private readonly lock: Lock = new Lock();
private readonly onSettingsChangeHandlers: ((
newSettings: SyncSettings,
oldSettings: SyncSettings
) => unknown)[] = [];
public constructor( public constructor(
private readonly logger: Logger, private readonly logger: Logger,
initialState: Partial<SyncSettings> | undefined, initialState: Partial<SyncSettings> | undefined,
@ -60,21 +59,6 @@ export class Settings {
return this.settings; return this.settings;
} }
public addOnSettingsChangeListener(
listener: (settings: SyncSettings, oldSettings: SyncSettings) => unknown
): void {
this.onSettingsChangeHandlers.push(listener);
}
public removeOnSettingsChangeListener(
listener: (settings: SyncSettings, oldSettings: SyncSettings) => unknown
): void {
const index = this.onSettingsChangeHandlers.indexOf(listener);
if (index !== -1) {
this.onSettingsChangeHandlers.splice(index, 1);
}
}
public async setSetting<T extends keyof SyncSettings>( public async setSetting<T extends keyof SyncSettings>(
key: T, key: T,
value: SyncSettings[T] value: SyncSettings[T]
@ -95,14 +79,9 @@ export class Settings {
...value ...value
}; };
await awaitAll( await this.onSettingsChanged.triggerAsync(
this.onSettingsChangeHandlers this.settings,
.map((handler) => { oldSettings
return handler(this.settings, oldSettings);
})
.filter((result): result is Promise<unknown> => {
return result instanceof Promise;
})
); );
await this.save(); await this.save();

View file

@ -122,7 +122,7 @@ describe("WebSocketManager", () => {
MockWebSocket as unknown as typeof WebSocket MockWebSocket as unknown as typeof WebSocket
); );
manager.addRemoteVaultUpdateListener(async () => { manager.onRemoteVaultUpdateReceived.add(async () => {
await new Promise((resolve) => setTimeout(resolve, 10)); await new Promise((resolve) => setTimeout(resolve, 10));
}); });
manager.start(); manager.start();
@ -152,7 +152,7 @@ describe("WebSocketManager", () => {
MockWebSocket as unknown as typeof WebSocket MockWebSocket as unknown as typeof WebSocket
); );
manager.addRemoteCursorsUpdateListener(async () => { manager.onRemoteCursorsUpdateReceived.add(async () => {
await new Promise((resolve) => setTimeout(resolve, 10)); await new Promise((resolve) => setTimeout(resolve, 10));
}); });
manager.start(); manager.start();
@ -227,7 +227,7 @@ describe("WebSocketManager", () => {
); );
let statusChangeCount = 0; let statusChangeCount = 0;
manager.addWebSocketStatusChangeListener(() => { manager.onWebSocketStatusChanged.add(() => {
statusChangeCount++; statusChangeCount++;
}); });
@ -269,7 +269,7 @@ describe("WebSocketManager", () => {
resolveListener = resolve; resolveListener = resolve;
}); });
manager.addRemoteVaultUpdateListener(async () => { manager.onRemoteVaultUpdateReceived.add(async () => {
await listenerPromise; await listenerPromise;
}); });

View file

@ -6,21 +6,23 @@ import type { CursorPositionFromClient } from "./types/CursorPositionFromClient"
import type { ClientCursors } from "./types/ClientCursors"; import type { ClientCursors } from "./types/ClientCursors";
import { createPromise } from "../utils/create-promise"; import { createPromise } from "../utils/create-promise";
import type { WebSocketVaultUpdate } from "./types/WebSocketVaultUpdate"; import type { WebSocketVaultUpdate } from "./types/WebSocketVaultUpdate";
import { awaitAll } from "../utils/await-all";
import { WEBSOCKET_DISCONNECT_TIMEOUT_IN_S } from "../consts"; import { WEBSOCKET_DISCONNECT_TIMEOUT_IN_S } from "../consts";
import { removeFromArray } from "../utils/remove-from-array";
import { EventListeners } from "../utils/data-structures/event-listeners";
import { awaitAll } from "../utils/await-all";
export class WebSocketManager { export class WebSocketManager {
private readonly webSocketStatusChangeListeners: (( public readonly onWebSocketStatusChanged = new EventListeners<
isConnected: boolean (isConnected: boolean) => unknown
) => unknown)[] = []; >();
private readonly remoteVaultUpdateListeners: (( public readonly onRemoteVaultUpdateReceived = new EventListeners<
update: WebSocketVaultUpdate (update: WebSocketVaultUpdate) => Promise<void>
) => Promise<void>)[] = []; >();
private readonly remoteCursorsUpdateListeners: (( public readonly onRemoteCursorsUpdateReceived = new EventListeners<
cursors: ClientCursors[] (cursors: ClientCursors[]) => Promise<void>
) => Promise<void>)[] = []; >();
private isStopped = true; private isStopped = true;
private resolveDisconnectingPromise: null | (() => unknown) = null; private resolveDisconnectingPromise: null | (() => unknown) = null;
@ -59,24 +61,6 @@ export class WebSocketManager {
); );
} }
public addWebSocketStatusChangeListener(
listener: (isConnected: boolean) => unknown
): void {
this.webSocketStatusChangeListeners.push(listener);
}
public addRemoteCursorsUpdateListener(
listener: (cursors: ClientCursors[]) => Promise<void>
): void {
this.remoteCursorsUpdateListeners.push(listener);
}
public addRemoteVaultUpdateListener(
listener: (update: WebSocketVaultUpdate) => Promise<void>
): void {
this.remoteVaultUpdateListeners.push(listener);
}
public start(): void { public start(): void {
this.isStopped = false; this.isStopped = false;
this.initializeWebSocket(); this.initializeWebSocket();
@ -205,9 +189,7 @@ export class WebSocketManager {
this.webSocket.onopen = (): void => { this.webSocket.onopen = (): void => {
this.logger.info("WebSocket connection opened"); this.logger.info("WebSocket connection opened");
this.webSocketStatusChangeListeners.forEach((listener) => this.onWebSocketStatusChanged.trigger(true);
listener(true)
);
}; };
this.webSocket.onmessage = (event): void => { this.webSocket.onmessage = (event): void => {
@ -227,12 +209,10 @@ export class WebSocketManager {
); );
}) })
.finally(() => { .finally(() => {
const index = this.outstandingPromises.indexOf( removeFromArray(
this.outstandingPromises,
messageHandlingPromise messageHandlingPromise
); );
if (index !== -1) {
void this.outstandingPromises.splice(index, 1); // ignore the returned promise
}
}); });
void this.outstandingPromises.push(messageHandlingPromise); // ignore the returned promise void this.outstandingPromises.push(messageHandlingPromise); // ignore the returned promise
@ -247,9 +227,7 @@ export class WebSocketManager {
this.logger.warn( this.logger.warn(
`WebSocket closed with code ${event.code} (${event.reason == "" ? "unknown reason" : event.reason})` `WebSocket closed with code ${event.code} (${event.reason == "" ? "unknown reason" : event.reason})`
); );
this.webSocketStatusChangeListeners.forEach((listener) => this.onWebSocketStatusChanged.trigger(false);
listener(false)
);
if (this.isStopped) { if (this.isStopped) {
this.resolveDisconnectingPromise?.(); this.resolveDisconnectingPromise?.();
@ -267,15 +245,7 @@ export class WebSocketManager {
message: WebSocketServerMessage message: WebSocketServerMessage
): Promise<void> { ): Promise<void> {
if (message.type === "vaultUpdate") { if (message.type === "vaultUpdate") {
await awaitAll( await this.onRemoteVaultUpdateReceived.triggerAsync(message);
this.remoteVaultUpdateListeners.map(async (listener) => {
await listener(message).catch((error: unknown) => {
this.logger.error(
`Error in vault update listener: ${String(error)}`
);
});
})
);
// eslint-disable-next-line @typescript-eslint/no-unnecessary-condition // eslint-disable-next-line @typescript-eslint/no-unnecessary-condition
} else if (message.type === "cursorPositions") { } else if (message.type === "cursorPositions") {
@ -283,14 +253,8 @@ export class WebSocketManager {
`Received cursor positions for ${JSON.stringify(message.clients)}` `Received cursor positions for ${JSON.stringify(message.clients)}`
); );
await awaitAll( await this.onRemoteCursorsUpdateReceived.triggerAsync(
this.remoteCursorsUpdateListeners.map(async (listener) => { message.clients
await listener(message.clients).catch((error: unknown) => {
this.logger.error(
`Error in cursor positions listener: ${String(error)}`
);
});
})
); );
} else { } else {
this.logger.warn( this.logger.warn(

View file

@ -26,6 +26,7 @@ import { FixedSizeDocumentCache } from "./utils/data-structures/fix-sized-cache"
import { setUpTelemetry } from "./utils/set-up-telemetry"; import { setUpTelemetry } from "./utils/set-up-telemetry";
import { DIFF_CACHE_SIZE_MB } from "./consts"; import { DIFF_CACHE_SIZE_MB } from "./consts";
import { ServerConfig } from "./services/server-config"; import { ServerConfig } from "./services/server-config";
import type { EventListeners } from "./utils/data-structures/event-listeners";
export class SyncClient { export class SyncClient {
private hasStartedOfflineSync = false; private hasStartedOfflineSync = false;
@ -62,6 +63,42 @@ export class SyncClient {
public get isWebSocketConnected(): boolean { public get isWebSocketConnected(): boolean {
return this.webSocketManager.isWebSocketConnected; return this.webSocketManager.isWebSocketConnected;
} }
public get onSyncHistoryUpdated(): EventListeners<
(stats: HistoryStats) => unknown
> {
this.checkIfDestroyed("onSyncHistoryUpdated getter");
return this.history.onHistoryUpdated;
}
public get onSettingsChanged(): EventListeners<
(newSettings: SyncSettings, oldSettings: SyncSettings) => unknown
> {
this.checkIfDestroyed("onSettingsChanged getter");
return this.settings.onSettingsChanged;
}
public get onRemainingOperationsCountChanged(): EventListeners<
(remainingOperationsCount: number) => unknown
> {
this.checkIfDestroyed("onRemainingOperationsCountChanged getter");
return this.syncer.onRemainingOperationsCountChanged;
}
public get onWebSocketStatusChanged(): EventListeners<
(isConnected: boolean) => unknown
> {
this.checkIfDestroyed("onWebSocketStatusChanged getter");
return this.webSocketManager.onWebSocketStatusChanged;
}
public get onRemoteCursorsUpdated(): EventListeners<
(cursors: MaybeOutdatedClientCursors[]) => unknown
> {
this.checkIfDestroyed("onRemoteCursorsUpdated getter");
return this.cursorTracker.onRemoteCursorsUpdated;
}
public static async create({ public static async create({
fs, fs,
persistence, persistence,
@ -122,7 +159,7 @@ export class SyncClient {
settings.getSettings().isSyncEnabled, settings.getSettings().isSyncEnabled,
logger logger
); );
settings.addOnSettingsChangeListener((newSettings, oldSettings) => { settings.onSettingsChanged.add((newSettings, oldSettings) => {
if (oldSettings.isSyncEnabled != newSettings.isSyncEnabled) { if (oldSettings.isSyncEnabled != newSettings.isSyncEnabled) {
fetchController.canFetch = newSettings.isSyncEnabled; fetchController.canFetch = newSettings.isSyncEnabled;
} }
@ -221,15 +258,13 @@ export class SyncClient {
this.unloadTelemetry = setUpTelemetry(); this.unloadTelemetry = setUpTelemetry();
} }
this.logger.addOnMessageListener((log): void => { this.logger.onLogEmitted.add((log): void => {
if (log.level === LogLevel.ERROR && Sentry.isInitialized()) { if (log.level === LogLevel.ERROR && Sentry.isInitialized()) {
Sentry.captureMessage(log.message); Sentry.captureMessage(log.message);
} }
}); });
this.settings.addOnSettingsChangeListener( this.settings.onSettingsChanged.add(this.onSettingsChange.bind(this));
this.onSettingsChange.bind(this)
);
if (this.settings.getSettings().isSyncEnabled) { if (this.settings.getSettings().isSyncEnabled) {
this.logger.info("Starting SyncClient"); this.logger.info("Starting SyncClient");
@ -273,14 +308,6 @@ export class SyncClient {
return this.history.entries; return this.history.entries;
} }
public addSyncHistoryUpdateListener(
listener: (stats: HistoryStats) => unknown
): void {
this.checkIfDestroyed("addSyncHistoryUpdateListener");
this.history.addSyncHistoryUpdateListener(listener);
}
/** /**
* Wait for the in-flight operations to finish, reset all tracking, * Wait for the in-flight operations to finish, reset all tracking,
* and the local database but retain the settings. * and the local database but retain the settings.
@ -325,28 +352,6 @@ export class SyncClient {
await this.settings.setSettings(value); await this.settings.setSettings(value);
} }
public addOnSettingsChangeListener(
listener: (settings: SyncSettings, oldSettings: SyncSettings) => unknown
): void {
this.checkIfDestroyed("addOnSettingsChangeListener");
this.settings.addOnSettingsChangeListener(listener);
}
public addRemainingSyncOperationsListener(
listener: (remainingOperations: number) => unknown
): void {
this.checkIfDestroyed("addRemainingSyncOperationsListener");
this.syncer.addRemainingOperationsListener(listener);
}
public addWebSocketStatusChangeListener(listener: () => unknown): void {
this.checkIfDestroyed("addWebSocketStatusChangeListener");
this.webSocketManager.addWebSocketStatusChangeListener(listener);
}
public async syncLocallyCreatedFile( public async syncLocallyCreatedFile(
relativePath: RelativePath relativePath: RelativePath
): Promise<void> { ): Promise<void> {
@ -412,14 +417,6 @@ export class SyncClient {
await this.cursorTracker.sendLocalCursorsToServer(documentToCursors); await this.cursorTracker.sendLocalCursorsToServer(documentToCursors);
} }
public addRemoteCursorsUpdateListener(
listener: (cursors: MaybeOutdatedClientCursors[]) => unknown
): void {
this.checkIfDestroyed("addRemoteCursorsUpdateListener");
this.cursorTracker.addRemoteCursorsUpdateListener(listener);
}
public async waitUntilFinished(): Promise<void> { public async waitUntilFinished(): Promise<void> {
this.checkIfDestroyed("waitUntilIdle"); this.checkIfDestroyed("waitUntilIdle");
await this.syncer.waitUntilFinished(); await this.syncer.waitUntilFinished();

View file

@ -9,12 +9,19 @@ import { DocumentUpToDateness } from "../types/document-up-to-dateness";
import { hash } from "../utils/hash"; import { hash } from "../utils/hash";
import type { FileChangeNotifier } from "./file-change-notifier"; import type { FileChangeNotifier } from "./file-change-notifier";
import { Lock } from "../utils/data-structures/locks"; import { Lock } from "../utils/data-structures/locks";
import { EventListeners } from "../utils/data-structures/event-listeners";
// Cursor positions are updated separately from documents. However, a given cursor position is only // Cursor positions are updated separately from documents. However, a given cursor position is only
// valid within a certain version of the document it belongs to. This class tracks previous and the latest // valid within a certain version of the document it belongs to. This class tracks previous and the latest
// known remote cursor positions, and for each document, tries to return the latest cursor positions that are // known remote cursor positions, and for each document, tries to return the latest cursor positions that are
// not from the future. // not from the future.
export class CursorTracker { export class CursorTracker {
// The returned position may be accurate, if it matches the document version, or outdated, in which case
// the client has to heuristically guess it's current position based on the local edits.
public readonly onRemoteCursorsUpdated = new EventListeners<
(cursors: MaybeOutdatedClientCursors[]) => unknown
>();
private readonly updateLock = new Lock(); private readonly updateLock = new Lock();
private knownRemoteCursors: (ClientCursors & { private knownRemoteCursors: (ClientCursors & {
@ -31,7 +38,7 @@ export class CursorTracker {
private readonly fileOperations: FileOperations, private readonly fileOperations: FileOperations,
private readonly fileChangeNotifier: FileChangeNotifier private readonly fileChangeNotifier: FileChangeNotifier
) { ) {
this.webSocketManager.addRemoteCursorsUpdateListener( this.webSocketManager.onRemoteCursorsUpdateReceived.add(
async (clientCursors) => { async (clientCursors) => {
await this.updateLock.withLock(async () => { await this.updateLock.withLock(async () => {
// The latest message will contain all active clients, so we can delete the ones // The latest message will contain all active clients, so we can delete the ones
@ -58,10 +65,14 @@ export class CursorTracker {
this.knownRemoteCursors = updatedKnownRemoteCursors; this.knownRemoteCursors = updatedKnownRemoteCursors;
}); });
this.onRemoteCursorsUpdated.trigger(
this.getRelevantAndPruneKnownClientCursors()
);
} }
); );
this.fileChangeNotifier.addFileChangeListener(async (relativePath) => this.fileChangeNotifier.onFileChanged.add(async (relativePath) =>
this.updateLock.withLock(async () => { this.updateLock.withLock(async () => {
for (const clientCursor of this.knownRemoteCursors) { for (const clientCursor of this.knownRemoteCursors) {
if ( if (
@ -144,19 +155,6 @@ export class CursorTracker {
this.webSocketManager.updateLocalCursors({ documentsWithCursors }); this.webSocketManager.updateLocalCursors({ documentsWithCursors });
} }
// The returned position may be accurate, if it matches the document version, or outdated, in which case
// the client has to heuristically guess it's current position based on the local edits.
public addRemoteCursorsUpdateListener(
listener: (cursors: MaybeOutdatedClientCursors[]) => unknown
): void {
// CursorTracker registers its own event listener in the constructor so it must have been called before this
this.webSocketManager.addRemoteCursorsUpdateListener(async () => {
await this.updateLock.withLock(() =>
listener(this.getRelevantAndPruneKnownClientCursors())
);
});
}
public reset(): void { public reset(): void {
this.knownRemoteCursors = []; this.knownRemoteCursors = [];
this.lastLocalCursorState = []; this.lastLocalCursorState = [];

View file

@ -1,24 +1,12 @@
import type { RelativePath } from "../persistence/database"; import type { RelativePath } from "../persistence/database";
import { EventListeners } from "../utils/data-structures/event-listeners";
export class FileChangeNotifier { export class FileChangeNotifier {
private readonly listeners: ((filePath: RelativePath) => unknown)[] = []; public readonly onFileChanged = new EventListeners<
(filePath: RelativePath) => unknown
public addFileChangeListener( >();
listener: (filePath: RelativePath) => unknown
): void {
this.listeners.push(listener);
}
public removeFileChangeListener(
listener: (filePath: RelativePath) => unknown
): void {
const index = this.listeners.indexOf(listener);
if (index !== -1) {
this.listeners.splice(index, 1);
}
}
public notifyOfFileChange(filePath: RelativePath): void { public notifyOfFileChange(filePath: RelativePath): void {
this.listeners.forEach((listener) => listener(filePath)); this.onFileChanged.trigger(filePath);
} }
} }

View file

@ -21,18 +21,21 @@ import type { WebSocketVaultUpdate } from "../services/types/WebSocketVaultUpdat
import type { WebSocketManager } from "../services/websocket-manager"; import type { WebSocketManager } from "../services/websocket-manager";
import type { WebSocketClientMessage } from "../services/types/WebSocketClientMessage"; import type { WebSocketClientMessage } from "../services/types/WebSocketClientMessage";
import { awaitAll } from "../utils/await-all"; import { awaitAll } from "../utils/await-all";
import { EventListeners } from "../utils/data-structures/event-listeners";
export class Syncer { export class Syncer {
public readonly onRemainingOperationsCountChanged = new EventListeners<
(remainingOperations: number) => unknown
>();
private readonly remoteDocumentsLock: Locks<DocumentId>; private readonly remoteDocumentsLock: Locks<DocumentId>;
private readonly remainingOperationsListeners: ((
remainingOperations: number
) => unknown)[] = [];
// FIFO to limit the number of concurrent sync operations // FIFO to limit the number of concurrent sync operations
private readonly syncQueue: PQueue; private readonly syncQueue: PQueue;
private _isFirstSyncComplete = false; private _isFirstSyncComplete = false;
private runningScheduleSyncForOfflineChanges: Promise<void> | undefined; private runningScheduleSyncForOfflineChanges: Promise<void> | undefined;
private previousRemainingOperationsCount = 0;
public constructor( public constructor(
private readonly deviceId: string, private readonly deviceId: string,
@ -50,27 +53,28 @@ export class Syncer {
this.remoteDocumentsLock = new Locks<DocumentId>(this.logger); this.remoteDocumentsLock = new Locks<DocumentId>(this.logger);
settings.addOnSettingsChangeListener((newSettings, oldSettings) => { settings.onSettingsChanged.add((newSettings, oldSettings) => {
if (newSettings.syncConcurrency !== oldSettings.syncConcurrency) { if (newSettings.syncConcurrency !== oldSettings.syncConcurrency) {
this.syncQueue.concurrency = newSettings.syncConcurrency; this.syncQueue.concurrency = newSettings.syncConcurrency;
} }
}); });
this.syncQueue.on("active", () => { this.syncQueue.on("active", () => {
this.remainingOperationsListeners.forEach((listener) => { if (this.previousRemainingOperationsCount !== this.syncQueue.size) {
listener(this.syncQueue.size); this.previousRemainingOperationsCount = this.syncQueue.size;
}); this.onRemainingOperationsCountChanged.trigger(
this.syncQueue.size
);
}
}); });
this.webSocketManager.addWebSocketStatusChangeListener( this.webSocketManager.onWebSocketStatusChanged.add((isConnected) => {
(isConnected) => {
if (isConnected) { if (isConnected) {
// The JS WebSocket API doesn't support setting headers, so we have to send the token as a message // The JS WebSocket API doesn't support setting headers, so we have to send the token as a message
this.sendHandshakeMessage(); this.sendHandshakeMessage();
} }
} });
); this.webSocketManager.onRemoteVaultUpdateReceived.add(
this.webSocketManager.addRemoteVaultUpdateListener(
this.syncRemotelyUpdatedFile.bind(this) this.syncRemotelyUpdatedFile.bind(this)
); );
} }
@ -79,12 +83,6 @@ export class Syncer {
return this._isFirstSyncComplete; return this._isFirstSyncComplete;
} }
public addRemainingOperationsListener(
listener: (remainingOperations: number) => unknown
): void {
this.remainingOperationsListeners.push(listener);
}
public async syncLocallyCreatedFile( public async syncLocallyCreatedFile(
relativePath: RelativePath relativePath: RelativePath
): Promise<void> { ): Promise<void> {
@ -444,11 +442,13 @@ 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
/* eslint-disable no-restricted-syntax -- Comparing by property, not direct equality */
locallyPossiblyDeletedFiles = locallyPossiblyDeletedFiles =
locallyPossiblyDeletedFiles.filter( locallyPossiblyDeletedFiles.filter(
(item) => (item) =>
item.relativePath !== originalFile.relativePath item.relativePath !== originalFile.relativePath
); );
/* eslint-enable no-restricted-syntax */
this.logger.debug( this.logger.debug(
`Document '${originalFile.relativePath}' was not found under its current path in the database but was found under a different path (${relativePath}), scheduling sync to move it` `Document '${originalFile.relativePath}' was not found under its current path in the database but was found under a different path (${relativePath}), scheduling sync to move it`

View file

@ -52,7 +52,7 @@ export class UnrestrictedSyncer {
this.logger this.logger
); );
this.settings.addOnSettingsChangeListener((newSettings) => { this.settings.onSettingsChanged.add((newSettings) => {
this.ignorePatterns = globsToRegexes( this.ignorePatterns = globsToRegexes(
newSettings.ignorePatterns, newSettings.ignorePatterns,
this.logger this.logger

View file

@ -1,4 +1,5 @@
import { MAX_LOG_MESSAGE_COUNT } from "../consts"; import { MAX_LOG_MESSAGE_COUNT } from "../consts";
import { EventListeners } from "../utils/data-structures/event-listeners";
export enum LogLevel { export enum LogLevel {
DEBUG = "DEBUG", DEBUG = "DEBUG",
@ -23,14 +24,11 @@ export class LogLine {
} }
export class Logger { export class Logger {
private readonly messages: LogLine[] = []; public readonly onLogEmitted = new EventListeners<
private readonly onMessageListeners: ((message: LogLine) => unknown)[] = []; (message: LogLine) => unknown
>();
public constructor( private readonly messages: LogLine[] = [];
...onMessageListeners: ((message: LogLine) => unknown)[]
) {
this.onMessageListeners = onMessageListeners;
}
public debug(message: string): void { public debug(message: string): void {
this.pushMessage(message, LogLevel.DEBUG); this.pushMessage(message, LogLevel.DEBUG);
@ -56,19 +54,6 @@ export class Logger {
); );
} }
public addOnMessageListener(listener: (message: LogLine) => unknown): void {
this.onMessageListeners.push(listener);
}
public removeOnMessageListener(
listener: (message: LogLine) => unknown
): void {
const index = this.onMessageListeners.indexOf(listener);
if (index !== -1) {
this.onMessageListeners.splice(index, 1);
}
}
public reset(): void { public reset(): void {
this.messages.length = 0; this.messages.length = 0;
this.debug("Logger has been reset"); this.debug("Logger has been reset");
@ -82,8 +67,6 @@ export class Logger {
this.messages.shift(); this.messages.shift();
} }
this.onMessageListeners.forEach((listener) => { this.onLogEmitted.trigger(logLine);
listener(logLine);
});
} }
} }

View file

@ -4,6 +4,8 @@ import {
} from "../consts"; } from "../consts";
import type { RelativePath } from "../persistence/database"; import type { RelativePath } from "../persistence/database";
import type { Logger } from "./logger"; import type { Logger } from "./logger";
import { removeFromArray } from "../utils/remove-from-array";
import { EventListeners } from "../utils/data-structures/event-listeners";
export interface SyncCreateDetails { export interface SyncCreateDetails {
type: SyncType.CREATE; type: SyncType.CREATE;
@ -68,11 +70,11 @@ export interface HistoryStats {
} }
export class SyncHistory { export class SyncHistory {
private _entries: HistoryEntry[] = []; public readonly onHistoryUpdated = new EventListeners<
(status: HistoryStats) => unknown
>();
private readonly syncHistoryUpdateListeners: (( private readonly _entries: HistoryEntry[] = [];
status: HistoryStats
) => unknown)[] = [];
private status: HistoryStats = { private status: HistoryStats = {
success: 0, success: 0,
@ -99,7 +101,7 @@ export class SyncHistory {
const candidate = this.findSimilarRecentUpdateEntry(historyEntry); const candidate = this.findSimilarRecentUpdateEntry(historyEntry);
if (candidate !== undefined) { if (candidate !== undefined) {
this._entries = this._entries.filter((e) => e !== candidate); removeFromArray(this._entries, candidate);
} }
// Insert the entry at the beginning // Insert the entry at the beginning
@ -112,31 +114,13 @@ export class SyncHistory {
this.updateSuccessCount(historyEntry); this.updateSuccessCount(historyEntry);
} }
public addSyncHistoryUpdateListener(
listener: (stats: HistoryStats) => unknown
): void {
this.syncHistoryUpdateListeners.push(listener);
listener({ ...this.status });
}
public removeSyncHistoryUpdateListener(
listener: (stats: HistoryStats) => unknown
): void {
const index = this.syncHistoryUpdateListeners.indexOf(listener);
if (index !== -1) {
this.syncHistoryUpdateListeners.splice(index, 1);
}
}
public reset(): void { public reset(): void {
this._entries.length = 0; this._entries.length = 0;
this.status = { this.status = {
success: 0, success: 0,
error: 0 error: 0
}; };
this.syncHistoryUpdateListeners.forEach((listener) => { this.onHistoryUpdated.trigger(this.status);
listener(this.status);
});
} }
private findSimilarRecentUpdateEntry( private findSimilarRecentUpdateEntry(
@ -178,8 +162,6 @@ export class SyncHistory {
break; break;
} }
this.syncHistoryUpdateListeners.forEach((listener) => { this.onHistoryUpdated.trigger(this.status);
listener(this.status);
});
} }
} }

View file

@ -0,0 +1,150 @@
import { describe, it } from "node:test";
import assert from "node:assert";
import { EventListeners } from "./event-listeners";
describe("EventListeners", () => {
it("should add & remove listeners", () => {
const listeners = new EventListeners<() => void>();
// eslint-disable-next-line @typescript-eslint/no-empty-function
const listener = (): void => {};
listeners.add(listener);
assert.strictEqual(listeners.count, 1);
const removed = listeners.remove(listener);
assert.strictEqual(removed, true);
assert.strictEqual(listeners.count, 0);
});
it("should remove listeners using unsubscribe function", () => {
const listeners = new EventListeners<() => void>();
// eslint-disable-next-line @typescript-eslint/no-empty-function
const listener = (): void => {};
const unsubscribe = listeners.add(listener);
unsubscribe();
assert.strictEqual(listeners.count, 0);
});
it("should return false when removing non-existent listener", () => {
const listeners = new EventListeners<() => void>();
// eslint-disable-next-line @typescript-eslint/no-empty-function
const listener = (): void => {};
const removed = listeners.remove(listener);
assert.strictEqual(removed, false);
});
it("should handle multiple listeners", () => {
const listeners = new EventListeners<() => void>();
// eslint-disable-next-line @typescript-eslint/no-empty-function
const listener1 = (): void => {};
// eslint-disable-next-line @typescript-eslint/no-empty-function
const listener2 = (): void => {};
// eslint-disable-next-line @typescript-eslint/no-empty-function
const listener3 = (): void => {};
listeners.add(listener1);
listeners.add(listener2);
listeners.add(listener3);
assert.strictEqual(listeners.count, 3);
listeners.remove(listener2);
assert.strictEqual(listeners.count, 2);
});
it("should trigger all listeners synchronously", () => {
const listeners = new EventListeners<(value: string) => void>();
const calls: string[] = [];
listeners.add((value) => calls.push(`listener1-${value}`));
listeners.add((value) => calls.push(`listener2-${value}`));
listeners.trigger("test");
assert.deepStrictEqual(calls, ["listener1-test", "listener2-test"]);
});
it("should trigger listeners with multiple arguments", () => {
const listeners = new EventListeners<
(a: number, b: string, c: boolean) => void
>();
const calls: [number, string, boolean][] = [];
listeners.add((a, b, c) => calls.push([a, b, c]));
listeners.trigger(42, "hello", true);
assert.deepStrictEqual(calls, [[42, "hello", true]]);
});
it("should not trigger removed listeners", () => {
const listeners = new EventListeners<() => void>();
let count1 = 0;
let count2 = 0;
const listener1 = (): void => {
count1++;
};
const listener2 = (): void => {
count2++;
};
listeners.add(listener1);
const unsubscribe = listeners.add(listener2);
unsubscribe();
listeners.trigger();
assert.strictEqual(count1, 1);
assert.strictEqual(count2, 0);
});
it("should trigger all listeners and await promises", async () => {
const listeners = new EventListeners<
(value: string) => Promise<void> | void
>();
const results: string[] = [];
listeners.add(async (value) => {
await new Promise((resolve) => setTimeout(resolve, 10));
results.push(`async1-${value}`);
});
listeners.add((value) => {
results.push(`sync-${value}`);
});
listeners.add(async (value) => {
await new Promise((resolve) => setTimeout(resolve, 5));
results.push(`async2-${value}`);
});
await listeners.triggerAsync("test");
assert.ok(results.includes("async1-test"));
assert.ok(results.includes("sync-test"));
assert.ok(results.includes("async2-test"));
assert.strictEqual(results.length, 3);
});
it("should not trigger cleared listeners", () => {
const listeners = new EventListeners<() => void>();
let called = false;
const listener = (): void => {
called = true;
};
listeners.add(listener);
listeners.clear();
assert.strictEqual(listeners.count, 0);
listeners.trigger();
assert.strictEqual(called, false);
});
});

View file

@ -0,0 +1,71 @@
import { removeFromArray } from "../remove-from-array";
import { awaitAll } from "../await-all";
/**
* A utility class for managing event listeners with type-safe add/remove operations.
*/
// eslint-disable-next-line @typescript-eslint/no-explicit-any
export class EventListeners<TListener extends (...args: any[]) => any> {
private readonly listeners: TListener[] = [];
public get count(): number {
return this.listeners.length;
}
/**
* Adds a new listener to the collection.
*
* @param listener The listener callback to add
* @returns An unsubscribe function that removes this listener when called
*/
public add(listener: TListener): () => void {
this.listeners.push(listener);
return () => this.remove(listener);
}
/**
* Removes a listener from the collection.
*
* @param listener The listener callback to remove
* @returns true if the listener was found and removed, false otherwise
*/
public remove(listener: TListener): boolean {
return removeFromArray(this.listeners, listener);
}
/**
* Triggers all listeners synchronously with the provided arguments.
* Any returned promises are ignored. Use triggerAsync() to await them.
*
* @param args The arguments to pass to each listener
*/
public trigger(...args: Parameters<TListener>): void {
this.listeners.forEach((listener) => {
listener(...args);
});
}
/**
* Triggers all listeners and awaits any promises they return.
* Synchronous listeners are called immediately, and any async listeners
* are awaited in parallel.
*
* @param args The arguments to pass to each listener
*/
public async triggerAsync(...args: Parameters<TListener>): Promise<void> {
await awaitAll(
this.listeners
.map((listener) => {
// eslint-disable-next-line @typescript-eslint/no-unsafe-return
return listener(...args);
})
.filter((result): result is Promise<unknown> => {
return result instanceof Promise;
})
);
}
public clear(): void {
this.listeners.length = 0;
}
}

View file

@ -3,7 +3,7 @@ import type { LogLine } from "../../tracing/logger";
import { LogLevel } from "../../tracing/logger"; import { LogLevel } from "../../tracing/logger";
export function logToConsole(client: SyncClient): void { export function logToConsole(client: SyncClient): void {
client.logger.addOnMessageListener((logLine: LogLine) => { client.logger.onLogEmitted.add((logLine: LogLine) => {
const formatted = `${logLine.timestamp.toISOString()} ${logLine.level} ${logLine.message}`; const formatted = `${logLine.timestamp.toISOString()} ${logLine.level} ${logLine.message}`;
switch (logLine.level) { switch (logLine.level) {

View file

@ -2,7 +2,8 @@ import { makeRe } from "minimatch";
import type { Logger } from "../tracing/logger"; import type { Logger } from "../tracing/logger";
export function globsToRegexes(globs: string[], logger: Logger): RegExp[] { export function globsToRegexes(globs: string[], logger: Logger): RegExp[] {
return globs return (
globs
.map((pattern) => { .map((pattern) => {
const result = makeRe(pattern, { const result = makeRe(pattern, {
dot: true dot: true
@ -14,5 +15,7 @@ export function globsToRegexes(globs: string[], logger: Logger): RegExp[] {
} }
return result; return result;
}) })
.filter((pattern) => pattern !== false); // eslint-disable-next-line no-restricted-syntax -- Filtering out false values, not removing a specific item
.filter((pattern) => pattern !== false)
);
} }

View file

@ -0,0 +1,17 @@
/**
* Efficiently removes a specific item from an array by modifying it in place.
* This is more efficient than using `.filter(item => item !== toRemove)` as it avoids creating a new array
*
* @param array The array to modify
* @param item The item to remove
* @returns true if the item was found and removed, false otherwise
*/
export function removeFromArray<T>(array: T[], item: T): boolean {
const index = array.indexOf(item);
if (index !== -1) {
// eslint-disable-next-line no-restricted-syntax -- This is the implementation of the helper itself
array.splice(index, 1);
return true;
}
return false;
}

View file

@ -8,7 +8,7 @@
"scripts": { "scripts": {
"dev": "webpack watch --mode development", "dev": "webpack watch --mode development",
"build": "webpack --mode production", "build": "webpack --mode production",
"test": "tsx --test src/**/*.test.ts" "test": "tsx --test 'src/**/*.test.ts'"
}, },
"devDependencies": { "devDependencies": {
"@types/node": "^24.8.1", "@types/node": "^24.8.1",

View file

@ -15,7 +15,7 @@ export class MockAgent extends MockClient {
private readonly pendingActions: Promise<unknown>[] = []; private readonly pendingActions: Promise<unknown>[] = [];
// The renamed file finding algorithm isn't too smart so we can't both update and rename the same file // The renamed file finding algorithm isn't too smart so we can't both update and rename the same file
private doNotTouchWhileOffline: string[] = []; private readonly doNotTouchWhileOffline: string[] = [];
public constructor( public constructor(
initialSettings: Partial<SyncSettings>, initialSettings: Partial<SyncSettings>,
@ -42,7 +42,7 @@ export class MockAgent extends MockClient {
"Connection check failed" "Connection check failed"
); );
this.client.logger.addOnMessageListener((logLine: LogLine) => { this.client.logger.onLogEmitted.add((logLine: LogLine) => {
const state = this.client.getSettings().isSyncEnabled const state = this.client.getSettings().isSyncEnabled
? "(online) " ? "(online) "
: "(offline)"; : "(offline)";
@ -54,9 +54,9 @@ export class MockAgent extends MockClient {
); );
if (historyEntry) { if (historyEntry) {
this.doNotTouchWhileOffline = utils.removeFromArray(
this.doNotTouchWhileOffline.filter( this.doNotTouchWhileOffline,
(file) => file !== historyEntry[1] historyEntry[1]
); );
} }
switch (logLine.level) { switch (logLine.level) {

View file

@ -2,7 +2,6 @@
set -e set -e
# Parse arguments
FIX_MODE=false FIX_MODE=false
if [[ "$1" == "--fix" ]]; then if [[ "$1" == "--fix" ]]; then
FIX_MODE=true FIX_MODE=true
@ -32,10 +31,18 @@ if [[ "$FIX_MODE" == true ]]; then
else else
npm ci npm ci
fi fi
cd ..
cd frontend
npm run build npm run build
npm run test npm run test
npm run lint npm run lint
# Use git ls-files to only check tracked files, respecting .gitignore
# We always run in fix mode and then check with git status
git ls-files | xargs npx eclint fix
if [[ "$FIX_MODE" == false ]] && [[ $(git status --porcelain) ]]; then if [[ "$FIX_MODE" == false ]] && [[ $(git status --porcelain) ]]; then
git status --porcelain git status --porcelain
echo "Failing CI because the working directory is not clean after linting" echo "Failing CI because the working directory is not clean after linting"

View file

@ -109,4 +109,3 @@ while true; do
sleep 0.2 sleep 0.2
done done

View file

@ -11,5 +11,6 @@ cd -
cp -r sync-server/bindings/* frontend/sync-client/src/services/types/ cp -r sync-server/bindings/* frontend/sync-client/src/services/types/
cd frontend cd frontend
npm run lint || npx prettier --write sync-client/src/services/types/*.ts npm run lint
git ls-files | xargs npx eclint fix
cd - cd -

11
sync-server/Cargo.lock generated
View file

@ -375,16 +375,6 @@ dependencies = [
"clap_derive", "clap_derive",
] ]
[[package]]
name = "clap-verbosity-flag"
version = "3.0.3"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "eeab6a5cdfc795a05538422012f20a5496f050223c91be4e5420bfd13c641fb1"
dependencies = [
"clap",
"log",
]
[[package]] [[package]]
name = "clap_builder" name = "clap_builder"
version = "4.5.38" version = "4.5.38"
@ -2143,7 +2133,6 @@ dependencies = [
"bimap", "bimap",
"chrono", "chrono",
"clap", "clap",
"clap-verbosity-flag",
"futures", "futures",
"humantime-serde", "humantime-serde",
"log", "log",

View file

@ -30,7 +30,6 @@ clap = { version = "4.5.38", features = ["derive"] }
futures = "0.3.31" futures = "0.3.31"
serde_yaml = "0.9.34" serde_yaml = "0.9.34"
serde_json = "1.0.140" serde_json = "1.0.140"
clap-verbosity-flag = "3.0.3"
bimap = "0.6.3" bimap = "0.6.3"
ts-rs = { version = "10.1", features = ["uuid-impl", "chrono-impl"] } ts-rs = { version = "10.1", features = ["uuid-impl", "chrono-impl"] }
base64 = "0.22.1" base64 = "0.22.1"

View file

@ -30,3 +30,4 @@ users:
logging: logging:
log_directory: logs log_directory: logs
log_rotation: 7days log_rotation: 7days
log_level: info

View file

@ -57,6 +57,7 @@ impl Database {
let mut connection_pools = std::collections::HashMap::new(); let mut connection_pools = std::collections::HashMap::new();
info!("Applying pending database migrations");
let mut entries = tokio::fs::read_dir(&config.databases_directory_path).await?; let mut entries = tokio::fs::read_dir(&config.databases_directory_path).await?;
while let Some(entry) = entries.next_entry().await? { while let Some(entry) = entries.next_entry().await? {
if !entry.file_name().to_string_lossy().ends_with(".sqlite") { if !entry.file_name().to_string_lossy().ends_with(".sqlite") {
@ -319,7 +320,7 @@ impl Database {
device_id, device_id,
has_been_merged has_been_merged
from latest_document_versions from latest_document_versions
where relative_path = ? where relative_path = ? and is_deleted = false
order by vault_update_id desc -- `latest_document_versions` only contains a single latest version of each document, however, order by vault_update_id desc -- `latest_document_versions` only contains a single latest version of each document, however,
-- multiple documents can have the same `relative_path`, if they have been deleted. That's -- multiple documents can have the same `relative_path`, if they have been deleted. That's
-- why we only care about the latest version of the document with the given relative path. -- why we only care about the latest version of the document with the given relative path.

View file

@ -1,7 +1,6 @@
use std::ffi::OsString; use std::ffi::OsString;
use clap::Parser; use clap::Parser;
use clap_verbosity_flag::{InfoLevel, Verbosity};
use crate::cli::color_when::ColorWhen; use crate::cli::color_when::ColorWhen;
@ -12,9 +11,6 @@ pub struct Args {
#[arg(index = 1)] #[arg(index = 1)]
pub config_path: Option<OsString>, pub config_path: Option<OsString>,
#[command(flatten)]
pub verbose: Verbosity<InfoLevel>,
#[arg( #[arg(
long, long,
value_name = "WHEN", value_name = "WHEN",

View file

@ -3,7 +3,10 @@ use std::time::Duration;
use log::debug; use log::debug;
use serde::{Deserialize, Serialize}; use serde::{Deserialize, Serialize};
use crate::consts::{DEFAULT_LOG_DIRECTORY, DEFAULT_LOG_ROTATION_INTERVAL}; use crate::{
consts::{DEFAULT_LOG_DIRECTORY, DEFAULT_LOG_LEVEL, DEFAULT_LOG_ROTATION_INTERVAL},
utils::log_level::LogLevel,
};
#[derive(Debug, Deserialize, Serialize, Clone)] #[derive(Debug, Deserialize, Serialize, Clone)]
pub struct LoggingConfig { pub struct LoggingConfig {
@ -12,6 +15,9 @@ pub struct LoggingConfig {
#[serde(default = "default_log_rotation", with = "humantime_serde")] #[serde(default = "default_log_rotation", with = "humantime_serde")]
pub log_rotation: Duration, pub log_rotation: Duration,
#[serde(default = "default_log_level")]
pub log_level: LogLevel,
} }
impl Default for LoggingConfig { impl Default for LoggingConfig {
@ -19,6 +25,7 @@ impl Default for LoggingConfig {
Self { Self {
log_directory: default_log_directory(), log_directory: default_log_directory(),
log_rotation: default_log_rotation(), log_rotation: default_log_rotation(),
log_level: default_log_level(),
} }
} }
} }
@ -32,3 +39,8 @@ fn default_log_rotation() -> Duration {
debug!("Using default log rotation: {DEFAULT_LOG_ROTATION_INTERVAL:?}"); debug!("Using default log rotation: {DEFAULT_LOG_ROTATION_INTERVAL:?}");
DEFAULT_LOG_ROTATION_INTERVAL DEFAULT_LOG_ROTATION_INTERVAL
} }
fn default_log_level() -> LogLevel {
debug!("Using default log level: Info");
DEFAULT_LOG_LEVEL
}

View file

@ -1,5 +1,7 @@
use std::time::Duration; use std::time::Duration;
use crate::utils::log_level::LogLevel;
pub const DEFAULT_CONFIG_PATH: &str = "config.yml"; pub const DEFAULT_CONFIG_PATH: &str = "config.yml";
pub const DEFAULT_DATABASES_DIRECTORY_PATH: &str = "databases"; pub const DEFAULT_DATABASES_DIRECTORY_PATH: &str = "databases";
@ -14,6 +16,7 @@ pub const DEFAULT_MAX_CLIENTS_PER_VAULT: usize = 256;
pub const DEFAULT_LOG_DIRECTORY: &str = "logs"; pub const DEFAULT_LOG_DIRECTORY: &str = "logs";
pub const DEFAULT_LOG_ROTATION_INTERVAL: Duration = Duration::from_secs(60 * 60 * 24); // 1 day pub const DEFAULT_LOG_ROTATION_INTERVAL: Duration = Duration::from_secs(60 * 60 * 24); // 1 day
pub const DEFAULT_LOG_LEVEL: LogLevel = LogLevel::Info;
pub const DEFAULT_MERGEABLE_FILE_EXTENSIONS: &[&str] = &["md", "txt"]; pub const DEFAULT_MERGEABLE_FILE_EXTENSIONS: &[&str] = &["md", "txt"];

View file

@ -60,14 +60,7 @@ fn set_up_logging(
args: &Args, args: &Args,
logging_config: &config::logging_config::LoggingConfig, logging_config: &config::logging_config::LoggingConfig,
) -> Result<(), SyncServerError> { ) -> Result<(), SyncServerError> {
let level_filter = match args.verbose.log_level_filter() { let level_filter = logging_config.log_level.as_tracing_level();
// We don't want to allow disabling all logging
log::LevelFilter::Off | log::LevelFilter::Error => tracing::Level::ERROR,
log::LevelFilter::Warn => tracing::Level::WARN,
log::LevelFilter::Info => tracing::Level::INFO,
log::LevelFilter::Debug => tracing::Level::DEBUG,
log::LevelFilter::Trace => tracing::Level::TRACE,
};
let env_filter = EnvFilter::builder() let env_filter = EnvFilter::builder()
.with_default_directive(level_filter.into()) .with_default_directive(level_filter.into())
@ -77,7 +70,7 @@ fn set_up_logging(
let use_colors = args.color.use_colors(); let use_colors = args.color.use_colors();
let is_debug_mode = args.verbose.log_level_filter() >= log::LevelFilter::Debug; let is_debug_mode = logging_config.log_level.is_debug_or_trace();
let file_appender = RotatingFileWriter::new( let file_appender = RotatingFileWriter::new(
&logging_config.log_directory, &logging_config.log_directory,

View file

@ -2,6 +2,7 @@ pub mod dedup_paths;
pub mod find_first_available_path; pub mod find_first_available_path;
pub mod is_binary; pub mod is_binary;
pub mod is_file_type_mergable; pub mod is_file_type_mergable;
pub mod log_level;
pub mod normalize; pub mod normalize;
pub mod rotating_file_writer; pub mod rotating_file_writer;
pub mod sanitize_path; pub mod sanitize_path;

View file

@ -0,0 +1,27 @@
use serde::{Deserialize, Serialize};
#[derive(Debug, Deserialize, Serialize, Clone, Copy, PartialEq, Eq)]
#[serde(rename_all = "lowercase")]
pub enum LogLevel {
Error,
Warn,
Info,
Debug,
Trace,
}
impl LogLevel {
pub fn as_tracing_level(self) -> tracing::Level {
match self {
Self::Error => tracing::Level::ERROR,
Self::Warn => tracing::Level::WARN,
Self::Info => tracing::Level::INFO,
Self::Debug => tracing::Level::DEBUG,
Self::Trace => tracing::Level::TRACE,
}
}
pub fn is_debug_or_trace(self) -> bool {
matches!(self, Self::Debug | Self::Trace)
}
}

View file

@ -173,6 +173,7 @@ mod tests {
#[test] #[test]
fn test_write_creates_log_file_and_directory() { fn test_write_creates_log_file_and_directory() {
let temp_dir = std::env::temp_dir().join("test_write_creates_log_file_and_directory"); let temp_dir = std::env::temp_dir().join("test_write_creates_log_file_and_directory");
let _ = fs::remove_dir_all(&temp_dir);
let mut writer = let mut writer =
RotatingFileWriter::new(&temp_dir, "test", Duration::from_secs(3600)).unwrap(); RotatingFileWriter::new(&temp_dir, "test", Duration::from_secs(3600)).unwrap();
@ -195,6 +196,7 @@ mod tests {
#[test] #[test]
fn test_rotation_after_duration() { fn test_rotation_after_duration() {
let temp_dir = std::env::temp_dir().join("test_rotation_after_duration"); let temp_dir = std::env::temp_dir().join("test_rotation_after_duration");
let _ = fs::remove_dir_all(&temp_dir);
// Use a very short rotation duration // Use a very short rotation duration
// Note: We need to wait at least 1 second between rotations since // Note: We need to wait at least 1 second between rotations since
@ -227,6 +229,7 @@ mod tests {
fn test_calculate_next_rotation_time_no_existing_logs() { fn test_calculate_next_rotation_time_no_existing_logs() {
let temp_dir = let temp_dir =
std::env::temp_dir().join("test_calculate_next_rotation_time_no_existing_logs"); std::env::temp_dir().join("test_calculate_next_rotation_time_no_existing_logs");
let _ = fs::remove_dir_all(&temp_dir);
fs::create_dir_all(&temp_dir).unwrap(); fs::create_dir_all(&temp_dir).unwrap();
@ -248,6 +251,7 @@ mod tests {
fn test_calculate_next_rotation_time_with_existing_log() { fn test_calculate_next_rotation_time_with_existing_log() {
let temp_dir = let temp_dir =
std::env::temp_dir().join("test_calculate_next_rotation_time_with_existing_log"); std::env::temp_dir().join("test_calculate_next_rotation_time_with_existing_log");
let _ = fs::remove_dir_all(&temp_dir);
fs::create_dir_all(&temp_dir).unwrap(); fs::create_dir_all(&temp_dir).unwrap();
@ -286,6 +290,7 @@ mod tests {
#[test] #[test]
fn test_picks_latest_log_file() { fn test_picks_latest_log_file() {
let temp_dir = std::env::temp_dir().join("test_picks_latest_log_file"); let temp_dir = std::env::temp_dir().join("test_picks_latest_log_file");
let _ = fs::remove_dir_all(&temp_dir);
fs::create_dir_all(&temp_dir).unwrap(); fs::create_dir_all(&temp_dir).unwrap();
@ -320,6 +325,7 @@ mod tests {
#[test] #[test]
fn test_ignores_malformed_filenames() { fn test_ignores_malformed_filenames() {
let temp_dir = std::env::temp_dir().join("test_ignores_malformed_filenames"); let temp_dir = std::env::temp_dir().join("test_ignores_malformed_filenames");
let _ = fs::remove_dir_all(&temp_dir);
fs::create_dir_all(&temp_dir).unwrap(); fs::create_dir_all(&temp_dir).unwrap();