diff --git a/src/app/storage/sync/engine.ts b/src/app/storage/sync/engine.ts new file mode 100644 index 000000000..37df9da96 --- /dev/null +++ b/src/app/storage/sync/engine.ts @@ -0,0 +1,261 @@ +import { IS_BROWSER } from '@open-pencil/core/constants' + +import { + activeStorageProviderID, + createActiveStorageAdapter, + storagePreferencesComplete +} from '@/app/integrations/storage' +import { getLocalCanvasStore } from '@/app/storage/local-store' +import { getOutbox } from '@/app/storage/sync/outbox' +import { setUploadProgress } from '@/app/storage/sync/progress' +import { setPendingSyncCount, setSyncUi } from '@/app/storage/sync/status' +import type { OutboxJob } from '@/app/storage/sync/types' + +const MAX_ATTEMPTS = 8 +const BASE_BACKOFF_MS = 1500 +const MAX_BACKOFF_MS = 60_000 + +let pumping = false +let wakeTimer: ReturnType | null = null +let onlineBound = false + +function isOnline(): boolean { + if (typeof navigator === 'undefined') return true + return navigator.onLine +} + +function backoffMs(attempts: number): number { + const exp = Math.min(MAX_BACKOFF_MS, BASE_BACKOFF_MS * 2 ** Math.max(0, attempts - 1)) + const jitter = Math.floor(exp * 0.2 * ((crypto.getRandomValues(new Uint8Array(1))[0] ?? 0) / 255)) + return exp + jitter +} + +function isPermanentError(error: unknown): boolean { + if (!(error instanceof Error)) return false + const msg = error.message.toLowerCase() + return ( + msg.includes('403') || + msg.includes('401') || + msg.includes('access denied') || + msg.includes('invalid access key') || + msg.includes('not configured') + ) +} + +async function runJob(job: OutboxJob): Promise { + if (!storagePreferencesComplete(activeStorageProviderID.value)) { + throw new Error('Storage is not configured') + } + const adapter = createActiveStorageAdapter() + const store = getLocalCanvasStore() + const meta = await store.getMeta(job.canvasId) + + if (job.type === 'deleteCanvas') { + await adapter.deleteDocument(job.canvasId) + // Keep the tombstoned row: reconcile purges it once the remote listing + // confirms the object is gone. Removing it here opened a race where a + // concurrent reconcile re-seeded the canvas from a stale remote listing. + await store.updateMeta(job.canvasId, { syncStatus: 'synced', lastSyncError: null }) + return + } + + if (!meta || meta.tombstoned) { + // Nothing to put + return + } + + if (job.type === 'putCanvas') { + // Superseded by a newer local revision already on disk + if (meta.revision > job.revision) return + if (!meta.hasFig) return + const fig = await store.readFig(job.canvasId) + if (!fig || fig.byteLength === 0) throw new Error('Local document missing for sync') + setUploadProgress(job.canvasId, 0) + try { + await adapter.putDocument( + job.canvasId, + fig, + { + name: meta.name, + updatedAt: meta.updatedAt + }, + ({ transferredBytes, totalBytes }) => { + if (totalBytes) setUploadProgress(job.canvasId, transferredBytes / totalBytes) + } + ) + } finally { + setUploadProgress(job.canvasId, null) + } + // Only mark synced if still on this revision and no other pending work for newer rev + const latest = await store.getMeta(job.canvasId) + if (latest && latest.revision === job.revision && !latest.tombstoned) { + await store.updateMeta(job.canvasId, { + syncStatus: 'synced', + lastSyncedAt: new Date().toISOString(), + lastSyncError: null + }) + } + return + } + + // Remaining job type: putThumb + if (!adapter.putThumbnail) return + const thumb = await store.readThumb(job.canvasId) + if (!thumb) return + await adapter.putThumbnail(job.canvasId, thumb) +} + +async function pumpOnce(): Promise { + const outbox = getOutbox() + const jobs = await outbox.list() + setPendingSyncCount(jobs.length) + + if (jobs.length === 0) { + if (isOnline()) setSyncUi('idle') + return + } + + if (!isOnline()) { + setSyncUi('offline') + scheduleWake(5000) + return + } + + setSyncUi('syncing') + const now = Date.now() + // Single-flight globally for simplicity (large figs) + const job = jobs.find((j) => j.nextAttemptAt <= now) + if (!job) { + const nextAt = Math.min(...jobs.map((j) => j.nextAttemptAt)) + scheduleWake(Math.max(250, nextAt - now)) + return + } + + try { + await runJob(job) + await outbox.remove(job.id) + const remaining = await outbox.list() + setPendingSyncCount(remaining.length) + if (remaining.length === 0) setSyncUi('idle') + else scheduleWake(50) + } catch (error) { + const attempts = job.attempts + 1 + const permanent = isPermanentError(error) || attempts >= MAX_ATTEMPTS + const message = error instanceof Error ? error.message : String(error) + console.warn('[Storage sync] job failed:', job.type, job.canvasId, message) + + if (permanent) { + // A failed thumbnail upload must not poison the document's sync status — + // only canvas/delete jobs reflect into the meta row. + if (job.type !== 'putThumb') { + await getLocalCanvasStore().updateMeta(job.canvasId, { + syncStatus: 'error', + lastSyncError: message + }) + setSyncUi('error', message.slice(0, 120)) + } else { + // Keep a record without touching syncStatus so the stale remote + // thumbnail is at least diagnosable. + await getLocalCanvasStore().updateMeta(job.canvasId, { lastSyncError: message }) + } + await outbox.remove(job.id) + const remaining = await outbox.list() + setPendingSyncCount(remaining.length) + if (remaining.length > 0) scheduleWake(1000) + else if (job.type === 'putThumb') setSyncUi('idle') + return + } + + const updated: OutboxJob = { + ...job, + attempts, + nextAttemptAt: Date.now() + backoffMs(attempts) + } + await outbox.update(updated) + if (job.type !== 'putThumb') { + await getLocalCanvasStore().updateMeta(job.canvasId, { + syncStatus: 'pending', + lastSyncError: message + }) + } + // Wake for the next ready job across the whole queue — not this job's + // full backoff, which starved other jobs that were ready sooner. + const all = await outbox.list() + const nextAt = Math.min(...all.map((j) => j.nextAttemptAt)) + scheduleWake(Math.max(250, nextAt - Date.now())) + } +} + +function scheduleWake(ms: number) { + if (wakeTimer != null) clearTimeout(wakeTimer) + wakeTimer = setTimeout(() => { + wakeTimer = null + void kickSyncEngine() + }, ms) +} + +function ensureOnlineListeners() { + if (onlineBound || !IS_BROWSER) return + onlineBound = true + window.addEventListener('online', () => { + setSyncUi('syncing') + void kickSyncEngine() + }) + window.addEventListener('offline', () => { + setSyncUi('offline') + }) +} + +/** Start or continue draining the outbox. Safe to call often. */ +export async function kickSyncEngine(): Promise { + ensureOnlineListeners() + if (pumping) return + pumping = true + let pumpFailed = false + try { + // Drain a few jobs per kick to avoid long tight loops blocking the tab. + for (let i = 0; i < 3; i++) { + const before = (await getOutbox().list()).length + await pumpOnce() + const after = (await getOutbox().list()).length + if (after === 0 || after >= before) break + } + } catch (error) { + // Never let an escaped rejection strand the queue — retry shortly. + pumpFailed = true + console.warn('[Storage sync] pump failed:', error) + scheduleWake(5000) + } finally { + pumping = false + } + // A job enqueued mid-pump can slip past the loop's exit check while its + // kick was swallowed by the pumping guard — re-wake if work is already due. + // (Skip offline — pumpOnce owns those wakes — and errors, which keep their + // 5s backoff; re-waking would clobber it into a tight retry loop.) + if (pumpFailed || !isOnline()) return + const jobs = await getOutbox().list() + if (jobs.some((job) => job.nextAttemptAt <= Date.now())) scheduleWake(250) +} + +export async function enqueuePutCanvas(canvasId: string, revision: number): Promise { + await getOutbox().enqueue({ canvasId, type: 'putCanvas', revision }) + void kickSyncEngine() +} + +export async function enqueuePutThumb(canvasId: string, revision: number): Promise { + await getOutbox().enqueue({ canvasId, type: 'putThumb', revision }) + void kickSyncEngine() +} + +export async function enqueueDeleteCanvas(canvasId: string): Promise { + await getOutbox().enqueue({ canvasId, type: 'deleteCanvas', revision: 0 }) + void kickSyncEngine() +} + +/** After credentials cleared — drop local mirror + outbox (optional safety). */ +export async function clearStorageLocalMirror(): Promise { + await getLocalCanvasStore().clearAll() + await getOutbox().clear() + setPendingSyncCount(0) + setSyncUi('idle') +} diff --git a/src/app/storage/sync/index.ts b/src/app/storage/sync/index.ts new file mode 100644 index 000000000..678bc70ef --- /dev/null +++ b/src/app/storage/sync/index.ts @@ -0,0 +1,24 @@ +export { + clearStorageLocalMirror, + enqueueDeleteCanvas, + enqueuePutCanvas, + enqueuePutThumb, + kickSyncEngine +} from './engine' +export { createMemoryOutbox, getOutbox, resetOutboxForTests } from './outbox' +export { setUploadProgress, uploadProgressByCanvas } from './progress' +export { + pendingSyncCount, + setPendingSyncCount, + setSyncUi, + syncStatusLabel, + syncUiDetail, + syncUiState +} from './status' +export { + makeJobId, + supersedePutCanvasJobs, + type OutboxJob, + type OutboxJobType, + type SyncUiState +} from './types' diff --git a/src/app/storage/sync/progress.ts b/src/app/storage/sync/progress.ts new file mode 100644 index 000000000..faac9999c --- /dev/null +++ b/src/app/storage/sync/progress.ts @@ -0,0 +1,14 @@ +import { ref } from 'vue' + +/** + * Per-canvas upload progress (0..1) while a putCanvas job is in flight. + * Drives the card badge fill on the Cloud Workspace grid. + */ +export const uploadProgressByCanvas = ref>(new Map()) + +export function setUploadProgress(canvasId: string, fraction: number | null) { + const next = new Map(uploadProgressByCanvas.value) + if (fraction == null) next.delete(canvasId) + else next.set(canvasId, Math.max(0, Math.min(1, fraction))) + uploadProgressByCanvas.value = next +} diff --git a/src/app/storage/sync/status.ts b/src/app/storage/sync/status.ts new file mode 100644 index 000000000..4dc0ea567 --- /dev/null +++ b/src/app/storage/sync/status.ts @@ -0,0 +1,30 @@ +import { computed, ref } from 'vue' + +import type { SyncUiState } from '@/app/storage/sync/types' + +/** Global subtle sync status for UI chips. */ +export const syncUiState = ref('idle') +export const syncUiDetail = ref(null) +export const pendingSyncCount = ref(0) + +export const syncStatusLabel = computed(() => { + switch (syncUiState.value) { + case 'syncing': + return syncUiDetail.value ?? 'Syncing…' + case 'offline': + return 'Offline · will sync' + case 'error': + return syncUiDetail.value ?? 'Sync failed' + default: + return null + } +}) + +export function setSyncUi(state: SyncUiState, detail: string | null = null) { + syncUiState.value = state + syncUiDetail.value = detail +} + +export function setPendingSyncCount(count: number) { + pendingSyncCount.value = count +}