How the applyChanges() Function Works on the Durable Object Side in Cloudflare Computer
The applyChangesSync() function in Cloudflare Computer's Durable Object acts as the authoritative apply engine, processing batched file system changes inside a SQLite transaction to ensure atomic synchronization between container-side clients and the DO's virtual filesystem.
In the cloudflare/computer repository, the Durable Object (DO) serves as the authoritative sync server for container-side SQLite databases. When the in-container computerd process pushes file system modifications, the DO receives these changes through the SyncRPC interface and applies them via the applyChangesSync() function located in the dofs package.
RPC Push Handling and Transactional Entry Point
The entry point for remote changes resides in packages/rpc/src/server.ts, where the push() method handles incoming synchronization requests from peers. This method receives a ReadableStream<ChangeEntry> plus metadata about the sender's revision, then orchestrates the atomic apply process.
// packages/rpc/src/server.ts
async push(input: {
senderRev: number;
changes: ReadableStream<ChangeEntry>;
}): Promise<{ rev: number; appliedPushCursor: ChangeCursor }> {
// 1️⃣ Collect the streamed ChangeEntry objects into an array
const entries: ChangeEntry[] = [];
const reader = input.changes.getReader();
…
// 2️⃣ Determine whether the sender is a peer (senderRev > 0) or a local writer
const isPeer = input.senderRev > 0;
// 3️⃣ Run the whole batch inside a single SQLite transaction.
// applyChangesSync() performs the per‑entry work.
this.db.transactionSync(() => {
applyChangesSync(this.db, entries, new Map(), {
source: isPeer ? "upstream" : "local",
});
// 4️⃣ If the sender was a peer, update the fetch cursor so the next pull
// sees the applied push rev.
if (isPeer) {
const nextCursor = { rev: input.senderRev, path: null };
if (compareChangeCursors(nextCursor, readFetchCursor(this.db)) > 0) {
writeFetchCursor(this.db, nextCursor);
}
}
});
…
return {
rev: currentRev(this.db),
appliedPushCursor: { rev: input.senderRev, path: null },
};
}
Detecting Source and Managing Cursors
The implementation distinguishes between peer and local writers using the senderRev parameter. When senderRev > 0, the source is marked as "upstream", which signals applyChangesSync() not to advance the local pushRev watermark. After successful application, the DO advances its fetch cursor to the sender's revision, ensuring the remote peer will not re-push identical entries in subsequent synchronization cycles.
Inside applyChangesSync: The Core Apply Engine
The synchronous apply logic lives in packages/dofs/src/sync/apply.ts. This function iterates over each ChangeEntry, performs deduplication checks, enforces mount constraints, and dispatches to specialized handlers based on entry kind (file, directory, symlink, or delete).
// packages/dofs/src/sync/apply.ts
export function applyChangesSync(
db: Database,
entries: readonly ChangeEntry[],
objects: Map<string, Uint8Array>,
options: ApplyOptions = {},
): ApplyResult {
// (batch size limits omitted for brevity)
for (const entry of entries) {
// ① Upstream deduplication – skip entries already present
if (options.source === "upstream" && entry.kind !== "delete") {
if (alreadyApplied(db, entry)) continue;
}
// ② Read‑only mount guard – drop entries that target a read‑only mount
const blockingRoot = readOnlyRootFor(db, entry.path);
if (blockingRoot !== undefined) {
skipped.push({ … });
continue;
}
// ③ Dispatch based on entry kind
if (entry.kind === "delete") {
rm(db, entry.path, { recursive: true, force: true });
} else if (entry.kind === "dir") {
applyDirectoryEntry(db, entry);
} else if (entry.kind === "symlink") {
removeReplaceableFinalEntry(db, entry.path, "symlink");
symlink(db, entry.target, entry.path, () => entry.mtime);
} else { // file
const total = applyFileEntry(db, entry, objects);
// `applyFileEntry` stages missing blobs, validates chunk sizes,
// removes any conflicting inode, then links the staged chunks.
}
applied++;
// (batch‑size accounting omitted)
}
return { applied, skipped };
}
Deduplication and Read-Only Mount Protection
The function employs upstream deduplication to maintain efficiency. When options.source equals "upstream", the alreadyApplied() helper (lines 59-82 in apply.ts) compares the live inode's manifest hash against the incoming entry, skipping unnecessary writes for identical data.
For read-only mounts, readOnlyRootFor(db, entry.path) checks if the target path falls under a protected mount root. Entries targeting these paths are added to the skipped array and never written to the database, preserving the authoritativeness of the mounted content.
Directory, Symlink, and Delete Operations
The dispatch logic handles three distinct entry types:
- Directory entries trigger
applyDirectoryEntry(), which creates new directories viamkdir()or replaces existing non-directory inodes by removing the subtree first. - Symlink entries invoke
removeReplaceableFinalEntry()to clear conflicting files, followed bysymlink()to create the new symbolic link with the specified modification time. - Delete entries execute
rm()with recursive and force flags, tolerating "already gone" conditions without error.
File Entry Processing and Blob Staging
File entries receive the most complex handling through applyFileEntry() (around line 408 in apply.ts). This function:
- Validates chunk windows using
assertChunkWindows()to ensure data integrity. - Stages missing blobs by checking
stagedBlobSize()for each chunk. Missing payloads are retrieved from theobjectsMap (populated from the RPC request) and written viastageBlob(). - Removes conflicting entries using
removeReplaceableFinalEntry()to clear the path. - Links staged chunks via
linkStagedChunksSync()to construct the final file vnode with proper mode and modification time.
The function returns the total bytes written, which the caller uses for batch-size accounting and progress tracking.
Transactional Guarantees and Synchronous Execution
The DO uses applyChangesSync() rather than its asynchronous counterpart to guarantee atomic batch processing. Because the push RPC receives a complete batch from the container, wrapping the call in this.db.transactionSync() ensures that any failure during entry processing rolls back all prior changes in the batch. This prevents partial state corruption where some files update while others fail.
The async variant (applyChanges()) exists for the pull side, where streaming large datasets requires entry-by-entry processing without buffering the entire set into memory.
Summary
- The
push()RPC handler inpackages/rpc/src/server.tsbuffers incomingChangeEntrystreams and invokesapplyChangesSync()inside a SQLite transaction. applyChangesSync()inpackages/dofs/src/sync/apply.tsiterates over entries, skipping duplicates viaalreadyApplied()and filtering read-only mounts viareadOnlyRootFor().- Entry kinds dispatch to specialized handlers:
applyDirectoryEntry()for directories,symlink()for links,rm()for deletions, andapplyFileEntry()for files. - File processing stages missing blob chunks from the provided
objectsMap before linking them into the virtual filesystem. - Synchronous execution ensures atomicity; the entire batch commits or rolls back as a single unit, maintaining consistency between the container and Durable Object.
Frequently Asked Questions
What is the difference between applyChanges() and applyChangesSync() in Cloudflare Computer?
applyChangesSync() is the synchronous variant used on the Durable Object side when receiving complete batches via RPC, allowing the entire operation to run inside a transactionSync() block for atomicity. The asynchronous applyChanges() variant is used during pull operations where entries stream in continuously and memory efficiency requires processing one entry at a time without buffering the entire batch.
How does the Durable Object handle duplicate changes from upstream peers?
The DO calls alreadyApplied() (defined around lines 59-82 in packages/dofs/src/sync/apply.ts) for each incoming entry when the source is "upstream". This function compares the entry's manifest hash against the current inode state; if they match, the entry is skipped, preventing redundant writes and breaking potential infinite sync loops between peers.
What happens when a change targets a read-only mount?
When readOnlyRootFor(db, entry.path) returns a matching mount root, the entry is added to the skipped array and omitted from the database write. This preserves the immutability of read-only mounts, ensuring that remote changes cannot override authoritative content mounted from external sources.
How are large file chunks handled during the apply process?
The applyFileEntry() function processes files chunk-by-chunk, validating window boundaries with assertChunkWindows() and verifying sizes with assertChunkSize(). Missing chunks are staged to the database via stageBlob() using byte payloads from the objects Map provided in the RPC call, then linked into the filesystem atomically using linkStagedChunksSync().
Have a question about this repo?
These articles cover the highlights, but your codebase questions are specific. Give your agent direct access to the source. Share this with your agent to get started:
curl -s "https://instagit.com/install.md" Maintain an open-source project? Get it listed too →