diff --git a/apps/obsidian-connector/src/main.ts b/apps/obsidian-connector/src/main.ts index 6b919981..1fd8eb06 100644 --- a/apps/obsidian-connector/src/main.ts +++ b/apps/obsidian-connector/src/main.ts @@ -1,6 +1,7 @@ -import { Notice, Plugin } from "obsidian" +import { Notice, Plugin, type TAbstractFile, TFile } from "obsidian" import { configure, createConnection } from "./api" import { SupermemorySettingTab } from "./settings" +import { SyncEngine } from "./sync" import type { SupermemorySettings } from "./types" const DEFAULT_SETTINGS: SupermemorySettings = { @@ -15,6 +16,7 @@ const DEFAULT_SETTINGS: SupermemorySettings = { export default class SupermemoryPlugin extends Plugin { settings: SupermemorySettings = DEFAULT_SETTINGS + private syncEngine: SyncEngine | null = null async onload() { await this.loadSettings() @@ -31,6 +33,37 @@ export default class SupermemoryPlugin extends Plugin { name: "Sync vault now", callback: () => this.syncVault(), }) + + if (this.settings.syncOnSave) { + this.registerVaultEvents() + } + + if (this.settings.syncOnStartup && this.settings.apiKey) { + this.app.workspace.onLayoutReady(() => this.syncVault()) + } + } + + private registerVaultEvents() { + this.registerEvent( + this.app.vault.on("modify", (file: TAbstractFile) => { + if (file instanceof TFile) this.syncEngine?.onFileChange(file) + }), + ) + this.registerEvent( + this.app.vault.on("create", (file: TAbstractFile) => { + if (file instanceof TFile) this.syncEngine?.onFileChange(file) + }), + ) + this.registerEvent( + this.app.vault.on("delete", (file: TAbstractFile) => { + if (file instanceof TFile) this.syncEngine?.onFileDelete(file) + }), + ) + this.registerEvent( + this.app.vault.on("rename", (file: TAbstractFile, oldPath: string) => { + if (file instanceof TFile) this.syncEngine?.onFileRename(file, oldPath) + }), + ) } async ensureConnection(): Promise { @@ -66,7 +99,12 @@ export default class SupermemoryPlugin extends Plugin { private async syncVault() { const connectionId = await this.ensureConnection() if (!connectionId) return - new Notice("Supermemory: sync starting...") + + if (!this.syncEngine) { + this.syncEngine = new SyncEngine(this.app, connectionId) + } + + await this.syncEngine.fullSync() } async loadSettings() { diff --git a/apps/obsidian-connector/src/sync.ts b/apps/obsidian-connector/src/sync.ts new file mode 100644 index 00000000..aca97971 --- /dev/null +++ b/apps/obsidian-connector/src/sync.ts @@ -0,0 +1,174 @@ +import { type App, Notice, TFile, parseFrontMatterEntry, parseFrontMatterTags } from "obsidian" +import { pushDeletions, pushNotes, type NotePayload } from "./api" + +const BATCH_SIZE = 50 +const DEBOUNCE_MS = 3000 + +interface SyncState { + syncing: boolean + lastFullSync: number +} + +export class SyncEngine { + private app: App + private connectionId: string + private state: SyncState = { syncing: false, lastFullSync: 0 } + private pendingChanges: Map = new Map() + private debounceTimer: ReturnType | null = null + + constructor(app: App, connectionId: string) { + this.app = app + this.connectionId = connectionId + } + + async fullSync(): Promise<{ queued: number; failed: number }> { + if (this.state.syncing) { + new Notice("Supermemory: sync already in progress.") + return { queued: 0, failed: 0 } + } + + this.state.syncing = true + let totalQueued = 0 + let totalFailed = 0 + + try { + const files = this.app.vault.getMarkdownFiles() + const batches = this.chunk(files, BATCH_SIZE) + + for (const batch of batches) { + const notes = await Promise.all(batch.map((f) => this.fileToPayload(f))) + const valid = notes.filter((n): n is NotePayload => n !== null) + if (valid.length === 0) continue + + try { + const result = await pushNotes(this.connectionId, valid) + totalQueued += result.queuedCount + } catch { + totalFailed += valid.length + } + } + + this.state.lastFullSync = Date.now() + new Notice( + `Supermemory: synced ${totalQueued} note${totalQueued !== 1 ? "s" : ""}${totalFailed > 0 ? `, ${totalFailed} failed` : ""}.`, + ) + } finally { + this.state.syncing = false + } + + return { queued: totalQueued, failed: totalFailed } + } + + onFileChange(file: TFile) { + if (!(file instanceof TFile) || file.extension !== "md") return + this.pendingChanges.set(file.path, "upsert") + this.schedulePush() + } + + onFileDelete(file: TFile) { + if (!(file instanceof TFile) || file.extension !== "md") return + this.pendingChanges.set(file.path, "delete") + this.schedulePush() + } + + onFileRename(file: TFile, oldPath: string) { + if (!(file instanceof TFile) || file.extension !== "md") return + this.pendingChanges.set(oldPath, "delete") + this.pendingChanges.set(file.path, "upsert") + this.schedulePush() + } + + private schedulePush() { + if (this.debounceTimer) clearTimeout(this.debounceTimer) + this.debounceTimer = setTimeout(() => this.flushPending(), DEBOUNCE_MS) + } + + private async flushPending() { + if (this.pendingChanges.size === 0) return + if (this.state.syncing) { + this.schedulePush() + return + } + + this.state.syncing = true + const changes = new Map(this.pendingChanges) + this.pendingChanges.clear() + + try { + const upserts: string[] = [] + const deletes: string[] = [] + + for (const [path, action] of changes) { + if (action === "upsert") upserts.push(path) + else deletes.push(path) + } + + if (deletes.length > 0) { + await pushDeletions( + this.connectionId, + deletes.map((path) => ({ path })), + ) + } + + if (upserts.length > 0) { + const files = upserts + .map((path) => this.app.vault.getAbstractFileByPath(path)) + .filter((f): f is TFile => f instanceof TFile) + + const batches = this.chunk(files, BATCH_SIZE) + for (const batch of batches) { + const notes = await Promise.all( + batch.map((f) => this.fileToPayload(f)), + ) + const valid = notes.filter((n): n is NotePayload => n !== null) + if (valid.length > 0) { + await pushNotes(this.connectionId, valid) + } + } + } + } catch { + // Re-queue on failure so the next debounce retries + for (const [path, action] of changes) { + if (!this.pendingChanges.has(path)) { + this.pendingChanges.set(path, action) + } + } + this.schedulePush() + } finally { + this.state.syncing = false + } + } + + private async fileToPayload(file: TFile): Promise { + try { + const content = await this.app.vault.cachedRead(file) + const cache = this.app.metadataCache.getFileCache(file) + let frontmatter: Record | undefined + + if (cache?.frontmatter) { + frontmatter = { ...cache.frontmatter } + delete frontmatter.position + const tags = parseFrontMatterTags(cache.frontmatter) + if (tags) frontmatter.tags = tags + } + + return { + path: file.path, + content, + title: file.basename, + mtime: file.stat.mtime, + frontmatter, + } + } catch { + return null + } + } + + private chunk(arr: T[], size: number): T[][] { + const result: T[][] = [] + for (let i = 0; i < arr.length; i += size) { + result.push(arr.slice(i, i + size)) + } + return result + } +}