Compare commits

...
Author SHA1 Message Date
Shoubhit Dash cb039c3895 perf(core): keep the clean-capture shortcut when the source ignores untracked paths
- Without exclude mirroring, a path ignored only by the source's info/exclude is listed as untracked on every scan, which defeated the clean-capture memo and cost check-ignore, update-index, and write-tree on every step. When the index is unchanged and nothing tracked changed, one check-ignore confirms every untracked path is ignored and returns the remembered tree; the result is reused on a miss.
- Only tracked paths that became ignored need --force-remove; untracked paths are never in the index.
2026-10-06 03:59:57 +05:30
Shoubhit Dash e293f57547 fix(core): close snapshot races found in review
- Scan the temporary index that the capture writes, never the shared one, so a concurrent process cannot roll entries back into a snapshot.
- Sweep abandoned temporary indexes by the creation time in their name; a hard link keeps the store index's old mtime while in use.
- Rebuild the index only when Git itself reports it unreadable, and install the source index atomically instead of deleting first.
- Record a tracked file replaced by an embedded repository as a gitlink.
- Drop exclude mirroring, which changed how source negations apply.
- Keep step cancellation responsive while the start snapshot is captured.
- Restore checkouts before removals so one failed removal cannot block the rest.
- Rename internals: objects.pack (not compact), withTemporaryIndex, lastCaptures, attemptCapture, canBatchRestore, pendingSnapshot; fold the capture's seed into ignores.
2026-10-05 14:26:49 +05:30
Shoubhit Dash f5c246c11d test(core): add snapshot parity and benchmark scripts 2026-10-05 11:46:42 +05:30
Shoubhit Dash 673360d6b9 perf(core): speed up and harden snapshot capture 2026-10-05 11:46:42 +05:30
14 changed files with 2553 additions and 124 deletions

No files matched your search

@@ -0,0 +1,66 @@
/**
* Copy existing snapshot stores and run the real object packing on the copies. Verifies
* that every locally stored object survives and reports size, time, and peak RSS
* of the Git child. The originals are only read.
*
* bun run script/benchmark-snapshot-pack.ts <store>...
*/
import { $ } from "bun"
import fs from "fs/promises"
import os from "os"
import path from "path"
import { Effect, Logger } from "effect"
import { AppNodeBuilder } from "../src/effect/app-node-builder"
import { Git } from "../src/git"
import { AbsolutePath } from "../src/schema"
const root = path.join(process.env.SNAPSHOT_BENCH_ROOT ?? os.tmpdir(), "opencode-snapshot-pack")
await fs.mkdir(root, { recursive: true })
// Local objects only: alternates are hidden so borrowed source objects do not count.
async function objects(store: string) {
const alternates = path.join(store, "objects", "info", "alternates")
const borrowed = await fs.readFile(alternates, "utf8").catch(() => undefined)
if (borrowed !== undefined) await fs.rm(alternates)
const listed =
await $`git --git-dir ${store} --work-tree ${store} cat-file --batch-all-objects ${"--batch-check=%(objectname)"}`
.quiet()
.text()
if (borrowed !== undefined) await fs.writeFile(alternates, borrowed)
return listed.split("\n").filter(Boolean).toSorted()
}
const size = async (directory: string) => Number((await $`du -sk ${directory}`.quiet().text()).split("\t")[0])
const loose = async (store: string) =>
(await $`find ${path.join(store, "objects")} -path '*/objects/??/*' -type f`.quiet().text())
.split("\n")
.filter(Boolean).length
let failed = false
for (const source of process.argv.slice(2)) {
const copy = await fs.mkdtemp(path.join(root, "store-"))
await $`cp -R ${source}/. ${copy}`.quiet()
const before = { objects: await objects(copy), size: await size(copy), loose: await loose(copy) }
const started = performance.now()
await Effect.runPromise(
Effect.gen(function* () {
const git = yield* Git.Service
yield* git.objects.pack(
new Git.Repository({
worktree: AbsolutePath.make(copy),
gitDirectory: AbsolutePath.make(copy),
commonDirectory: AbsolutePath.make(copy),
}),
)
}).pipe(Effect.provide(AppNodeBuilder.build(Git.node)), Effect.provide(Logger.layer([]))),
)
const elapsed = performance.now() - started
const after = { objects: await objects(copy), size: await size(copy), loose: await loose(copy) }
const same = before.objects.join("\n") === after.objects.join("\n")
if (!same) failed = true
console.log(
`${path.basename(source)} ${before.size} KiB -> ${after.size} KiB loose ${before.loose} -> ${after.loose} objects ${before.objects.length} -> ${after.objects.length} ${same ? "identical" : "MISMATCH"} ${elapsed.toFixed(0)} ms`,
)
await fs.rm(copy, { recursive: true, force: true })
}
if (failed) process.exit(1)
@@ -0,0 +1,194 @@
/**
* Drives the real SessionStep.attempt through a scripted agent trace against a real
* snapshot store, then repeats it with snapshots disabled. The difference is the
* wall-clock cost snapshots add to the session.
*
* bun run script/benchmark-snapshot-session.ts [--fixture opencode] [--ttft 400] [--rounds 2]
*
* Uses the fixtures generated by benchmark-snapshot.ts.
*/
import { $ } from "bun"
import fs from "fs/promises"
import os from "os"
import path from "path"
import { LanguageModel, LLM, LLMEvent } from "@opencode/ai"
import { OpenAIChat } from "@opencode/ai/protocols/openai-chat"
import { TestLLM } from "@opencode/ai/testing"
import { Global } from "@opencode/util/global"
import { AppProcess } from "@opencode/util/process"
import { LayerNode } from "@opencode/util/effect/layer-node"
import { Money } from "@opencode/schema/money"
import { Effect, Layer, Logger, Stream } from "effect"
import { Agent } from "../src/agent"
import { Bus } from "../src/bus"
import { Database } from "../src/database/database"
import { AppNodeBuilder } from "../src/effect/app-node-builder"
import { Location } from "../src/location"
import { Project } from "../src/project"
import { ProjectTable } from "../src/project/sql"
import { AbsolutePath } from "../src/schema"
import { Session } from "../src/session"
import { SessionMessage } from "../src/session/message"
import { SessionProjector } from "../src/session/projector"
import { SessionRunnerModel } from "../src/session/runner/model"
import { SessionStep } from "../src/session/runner/step"
import { SessionTable } from "../src/session/sql"
import { Snapshot } from "../src/snapshot"
import { ToolOutput } from "../src/tool-output"
const args = process.argv.slice(2)
const flag = (name: string) => {
const index = args.indexOf(`--${name}`)
return index === -1 ? undefined : args[index + 1]
}
const fixture = flag("fixture") ?? "opencode"
const ttft = Number(flag("ttft") ?? 400)
const rounds = Number(flag("rounds") ?? 2)
const root = path.join(process.env.SNAPSHOT_BENCH_ROOT ?? os.tmpdir(), "opencode-snapshot-bench")
const directory = path.join(root, "fixtures", fixture)
// A read-heavy coding loop: explore, edit, run checks, answer.
const trace = [
"read",
"grep",
"read",
"glob",
"read",
"edit",
"shell",
"read",
"edit",
"shell",
"grep",
"read",
"edit",
"shell",
"text",
] as const
let spawns = 0
const countingProcess = Layer.effect(
AppProcess.Service,
Effect.gen(function* () {
const real = yield* AppProcess.Service
return AppProcess.Service.of({
...real,
run: (command, options) => {
spawns++
return real.run(command, options)
},
})
}),
).pipe(Layer.provide(AppNodeBuilder.build(AppProcess.node)))
const snapshotLayer = (data: string) =>
AppNodeBuilder.build(Snapshot.node, [
Location.node.replace(Location.boundNode(Location.Ref.make({ directory: AbsolutePath.make(directory) }))),
Global.node.replace(Global.layerWith({ data, config: path.join(data, "config") })),
AppProcess.node.replace(countingProcess),
])
const files = (await $`git ls-files -z -- '*.ts'`.cwd(directory).text()).split("\0").filter(Boolean).slice(0, 50)
const model = SessionRunnerModel.resolved(
LanguageModel.make({ id: "bench-model", provider: "test", route: OpenAIChat.route }),
{
capabilities: { tools: true, input: ["text"], output: ["text"] },
limit: { context: 100_000, output: 1_000 },
cost: [
{
input: Money.USDPerMillionTokens.make(1),
output: Money.USDPerMillionTokens.make(2),
cache: { read: Money.USDPerMillionTokens.make(0.1), write: Money.USDPerMillionTokens.make(0.5) },
},
],
},
)
const delayed = (response: ReturnType<typeof TestLLM.stop>) =>
Stream.unwrap(Effect.sleep(ttft).pipe(Effect.as(Stream.fromIterable(response as Iterable<LLMEvent>))))
const run = (snapshots: Layer.Layer<Snapshot.Service>) =>
Effect.gen(function* () {
const db = (yield* Database.Service).db
const llm = yield* TestLLM.Test
const sessionID = Session.ID.create()
yield* db
.insert(ProjectTable)
.values({ id: Project.ID.global, worktree: AbsolutePath.make(directory), sandboxes: [] })
.onConflictDoNothing()
.run()
yield* db
.insert(SessionTable)
.values({ id: sessionID, project_id: Project.ID.global, slug: "bench", directory, version: "bench" })
.run()
const steps = yield* SessionStep.make
const snapshot = yield* Snapshot.Service
// Warm the snapshot store so the first step does not pay repository creation.
yield* snapshot.capture()
const before = spawns
const start = performance.now()
let edits = 0
for (let round = 0; round < rounds; round++)
for (const kind of trace) {
const call = `call-${round}-${edits}-${kind}-${Math.random().toString(36).slice(2)}`
yield* llm.push(
delayed(
kind === "text"
? TestLLM.text("Done.", `text-${call}`)
: TestLLM.tool(call, kind, { path: files[edits % files.length] }),
),
)
yield* steps.attempt({
isLocationClosed: () => false,
sessionID,
assistantMessageID: SessionMessage.ID.create(),
agent: Agent.defaultID,
model,
prepared: {
retry: () => Effect.void,
request: LLM.request({ model: model.model, prompt: "bench" }),
options: {},
executeTool: (input) =>
Effect.gen(function* () {
if (input.call.name === "edit")
yield* Effect.promise(() =>
Bun.write(path.join(directory, files[edits++ % files.length]!), `// bench ${edits}\n`),
)
return { content: [{ type: "text" as const, text: "ok" }] }
}),
},
retry: () => Effect.succeed({ retry: false as const }),
recoverContinuation: false,
recoverOverflow: Effect.succeed(false),
})
}
return { ms: performance.now() - start, spawns: spawns - before }
}).pipe(
Effect.provide(snapshots),
Effect.provide(
Layer.merge(
AppNodeBuilder.build(LayerNode.group([Database.node, Bus.node, SessionProjector.node, ToolOutput.node]), [
Bus.node.replace(Bus.configured({ persist: true })),
]),
TestLLM.testLayer(),
),
),
Effect.provide(Logger.layer([])),
)
await $`git -c core.fsmonitor=false reset -q --hard`.cwd(directory).quiet()
const data = await fs.mkdtemp(path.join(root, "session-data-"))
const noop = await Effect.runPromise(run(Snapshot.noopLayer))
await $`git -c core.fsmonitor=false reset -q --hard`.cwd(directory).quiet()
const real = await Effect.runPromise(run(snapshotLayer(data)))
await $`git -c core.fsmonitor=false reset -q --hard`.cwd(directory).quiet()
await fs.rm(data, { recursive: true, force: true })
const count = trace.length * rounds
console.log(`fixture ${fixture}, ${count} steps, simulated time to first event ${ttft} ms`)
console.log(`no snapshots ${noop.ms.toFixed(0).padStart(7)} ms`)
console.log(`snapshots ${real.ms.toFixed(0).padStart(7)} ms git ${real.spawns}`)
console.log(
`overhead ${(real.ms - noop.ms).toFixed(0).padStart(7)} ms ${((real.ms - noop.ms) / count).toFixed(1)} ms/step ${(real.spawns / count).toFixed(1)} git/step`,
)
@@ -0,0 +1,28 @@
// Child process for the cross-process snapshot probe: edits its own file and captures repeatedly,
// printing the number of failed captures.
import path from "path"
import { Effect, Logger } from "effect"
import { AppNodeBuilder } from "../src/effect/app-node-builder"
import { Location } from "../src/location"
import { AbsolutePath } from "../src/schema"
import { Snapshot } from "../src/snapshot"
import { Global } from "@opencode/util/global"
const [data, directory, id, count] = process.argv.slice(2)
const layer = AppNodeBuilder.build(Snapshot.node, [
Location.node.replace(Location.boundNode(Location.Ref.make({ directory: AbsolutePath.make(directory!) }))),
Global.node.replace(Global.layerWith({ data: data!, config: path.join(data!, "config") })),
])
const failures = await Effect.runPromise(
Effect.gen(function* () {
const snapshot = yield* Snapshot.Service
let failed = 0
for (let index = 0; index < Number(count); index++) {
yield* Effect.promise(() => Bun.write(path.join(directory!, `worker-${id}.txt`), `${index}\n`))
if ((yield* snapshot.capture()) === undefined) failed++
}
return failed
}).pipe(Effect.provide(layer), Effect.provide(Logger.layer([]))),
)
console.log(failures)
+385
View File
@@ -0,0 +1,385 @@
/**
* Snapshot benchmark + robustness probe.
*
* bun run script/benchmark-snapshot.ts [--fixtures small,medium,large,opencode] [--iterations 10] [--json out.json]
*
* Fixtures are generated once under $TMPDIR/opencode-snapshot-bench and reset before each run.
* Every timed scenario also reports how many git processes it spawned.
*/
import { $ } from "bun"
import fs from "fs/promises"
import os from "os"
import path from "path"
import { Effect, Layer, Logger } from "effect"
import { AppNodeBuilder } from "../src/effect/app-node-builder"
import { Git } from "../src/git"
import { Location } from "../src/location"
import { AbsolutePath, RelativePath } from "../src/schema"
import { Snapshot } from "../src/snapshot"
import { Global } from "@opencode/util/global"
import { AppProcess } from "@opencode/util/process"
const args = process.argv.slice(2)
const flag = (name: string) => {
const index = args.indexOf(`--${name}`)
return index === -1 ? undefined : args[index + 1]
}
const iterations = Number(flag("iterations") ?? 10)
const selected = (flag("fixtures") ?? "small,medium,large,opencode").split(",")
const jsonOut = flag("json")
const root = path.join(process.env.SNAPSHOT_BENCH_ROOT ?? os.tmpdir(), "opencode-snapshot-bench")
type Fixture = { name: string; tracked: number; ignored: number; dirs: number }
const fixtures: Record<string, Fixture> = {
small: { name: "small", tracked: 1_000, ignored: 0, dirs: 20 },
medium: { name: "medium", tracked: 20_000, ignored: 20_000, dirs: 400 },
large: { name: "large", tracked: 100_000, ignored: 50_000, dirs: 2_000 },
}
// ---------- fixture generation ----------
async function writeMany(files: Array<[string, string]>) {
const dirs = new Set(files.map(([file]) => path.dirname(file)))
await Promise.all([...dirs].map((dir) => fs.mkdir(dir, { recursive: true })))
for (let index = 0; index < files.length; index += 512)
await Promise.all(files.slice(index, index + 512).map(([file, content]) => Bun.write(file, content)))
}
async function generate(fixture: Fixture) {
const dir = path.join(root, "fixtures", fixture.name)
if (await Bun.file(path.join(dir, ".bench-ready")).exists()) return dir
await fs.rm(dir, { recursive: true, force: true })
console.error(`generating fixture ${fixture.name} (${fixture.tracked} tracked, ${fixture.ignored} ignored)`)
const tracked = Array.from({ length: fixture.tracked }, (_, index) => {
const file = path.join(dir, "src", `d${index % fixture.dirs}`, `f${index}.ts`)
return [file, `export const value${index} = ${index}\n`.repeat(8)] as [string, string]
})
const ignored = Array.from({ length: fixture.ignored }, (_, index) => {
const file = path.join(dir, "node_modules", `pkg${index % 500}`, `f${index}.js`)
return [file, `module.exports = ${index}\n`] as [string, string]
})
await writeMany([...tracked, ...ignored, [path.join(dir, ".gitignore"), "node_modules\ndist\n"]])
await gitInit(dir)
await Bun.write(path.join(dir, ".bench-ready"), "")
return dir
}
async function cloneOpencode() {
const dir = path.join(root, "fixtures", "opencode")
if (await Bun.file(path.join(dir, ".bench-ready")).exists()) return dir
await fs.rm(dir, { recursive: true, force: true })
const source = (await $`git rev-parse --show-toplevel`.cwd(import.meta.dir).text()).trim()
console.error(`cloning ${source} into fixture opencode`)
await $`git clone --quiet --local --no-hardlinks ${source} ${dir}`.quiet()
// A realistic ignored dependency tree without paying for a full install.
await writeMany(
Array.from({ length: 30_000 }, (_, index) => [
path.join(dir, "node_modules", `pkg${index % 700}`, `f${index}.js`),
`module.exports = ${index}\n`,
]),
)
await Bun.write(path.join(dir, ".git", "info", "exclude"), ".bench-ready\n")
await Bun.write(path.join(dir, ".bench-ready"), "")
return dir
}
async function gitInit(dir: string) {
await $`git init -q`.cwd(dir).quiet()
await Bun.write(path.join(dir, ".git", "info", "exclude"), ".bench-ready\n")
await $`git -c core.fsmonitor=false add -A`.cwd(dir).quiet()
await $`git -c user.email=bench@opencode.test -c user.name=Bench commit -q --no-gpg-sign -m initial`.cwd(dir).quiet()
}
async function reset(dir: string) {
await $`git -c core.fsmonitor=false reset -q --hard`.cwd(dir).quiet()
await $`git -c core.fsmonitor=false clean -qfd`.cwd(dir).quiet()
}
// ---------- harness ----------
let spawns = 0
const countingProcess = Layer.effect(
AppProcess.Service,
Effect.gen(function* () {
const real = yield* AppProcess.Service
return AppProcess.Service.of({
...real,
run: (command, options) => {
spawns++
return real.run(command, options)
},
})
}),
).pipe(Layer.provide(AppNodeBuilder.build(AppProcess.node)))
function snapshotLayer(data: string, directory: string) {
return AppNodeBuilder.build(Snapshot.node, [
Location.node.replace(Location.boundNode(Location.Ref.make({ directory: AbsolutePath.make(directory) }))),
Global.node.replace(Global.layerWith({ data, config: path.join(data, "config") })),
AppProcess.node.replace(countingProcess),
])
}
type Sample = { ms: number; spawns: number }
const results: Array<{ fixture: string; scenario: string; samples: Sample[] }> = []
const time = <A, E, R>(effect: Effect.Effect<A, E, R>) =>
Effect.gen(function* () {
const before = spawns
const start = performance.now()
const value = yield* effect
return { value, sample: { ms: performance.now() - start, spawns: spawns - before } }
})
function record(fixture: string, scenario: string, samples: Sample[]) {
results.push({ fixture, scenario, samples })
const sorted = samples.map((sample) => sample.ms).toSorted((a, b) => a - b)
const p = (value: number) => sorted[Math.min(Math.ceil(sorted.length * value) - 1, sorted.length - 1)] ?? 0
const mean = sorted.reduce((total, value) => total + value, 0) / sorted.length
const spawned = samples.reduce((total, sample) => total + sample.spawns, 0) / samples.length
console.log(
`${fixture.padEnd(9)} ${scenario.padEnd(22)} p50 ${p(0.5).toFixed(1).padStart(8)} ms mean ${mean
.toFixed(1)
.padStart(8)} ms max ${(sorted.at(-1) ?? 0).toFixed(1).padStart(8)} ms git ${spawned.toFixed(1).padStart(5)}`,
)
}
const edit = (dir: string, files: string[], tag: string) =>
Effect.promise(() => Promise.all(files.map((file) => Bun.write(path.join(dir, file), `// ${tag}\n${file}\n`))))
function trackedSample(dir: string, count: number) {
return $`git ls-files -z`
.cwd(dir)
.text()
.then((text) => {
const files = text.split("\0").filter((file) => file.endsWith(".ts") || file.endsWith(".md"))
const stride = Math.max(1, Math.floor(files.length / count))
return Array.from({ length: count }, (_, index) => files[(index * stride) % files.length]!)
})
}
const repeat = <A, E, R>(count: number, effect: (index: number) => Effect.Effect<Sample, E, R>) =>
Effect.forEach(
Array.from({ length: count }, (_, index) => index),
effect,
)
async function bench(name: string, dir: string) {
await reset(dir)
const data = await fs.mkdtemp(path.join(root, "data-"))
const sample = await trackedSample(dir, 100)
const program = Effect.gen(function* () {
const snapshot = yield* Snapshot.Service
const capture = snapshot
.capture()
.pipe(Effect.flatMap((id) => (id ? Effect.succeed(id) : Effect.die(new Error("capture returned undefined")))))
const cold = yield* time(capture)
record(name, "capture cold", [cold.sample])
record(
name,
"capture clean",
yield* repeat(iterations, () => time(capture).pipe(Effect.map((result) => result.sample))),
)
record(
name,
"capture edit 1",
yield* repeat(iterations, (index) =>
edit(dir, [sample[0]!], `edit-${index}`).pipe(
Effect.andThen(time(capture)),
Effect.map((result) => result.sample),
),
),
)
record(
name,
"capture add 1",
yield* repeat(iterations, (index) =>
edit(dir, [`src/new-${index}.ts`], "new").pipe(
Effect.andThen(time(capture)),
Effect.map((result) => result.sample),
),
),
)
const before = yield* capture
const edited = yield* repeat(Math.max(3, Math.ceil(iterations / 2)), (index) =>
edit(dir, sample, `bulk-${index}`).pipe(
Effect.andThen(time(capture)),
Effect.map((result) => result.sample),
),
)
record(name, "capture edit 100", edited)
const after = yield* capture
record(
name,
"files (100 changed)",
yield* repeat(iterations, () =>
time(snapshot.files({ from: before, to: after })).pipe(Effect.map((result) => result.sample)),
),
)
record(
name,
"diff (100 changed)",
yield* repeat(iterations, () =>
time(snapshot.diff({ from: before, to: after })).pipe(Effect.map((result) => result.sample)),
),
)
const plan = new Map(sample.map((file) => [RelativePath.make(file), before] as const))
record(
name,
"restore 100",
yield* repeat(Math.max(3, Math.ceil(iterations / 2)), (index) =>
edit(dir, sample, `restore-${index}`).pipe(
Effect.andThen(time(snapshot.restore({ files: plan }))),
Effect.map((result) => result.sample),
),
),
)
// A read-mostly agent loop under the step policy: capture at attempt start, capture at settlement,
// then list changed files. One step in five edits two files.
const steps = yield* repeat(20, (index) =>
Effect.gen(function* () {
const before = spawns
const start = performance.now()
const from = yield* capture
if (index % 5 === 0) yield* edit(dir, [sample[index]!, sample[index + 1]!], `step-${index}`)
const to = yield* capture
if (from !== to) yield* snapshot.files({ from, to })
return { ms: performance.now() - start, spawns: spawns - before }
}),
)
record(name, "step loop (per step)", steps)
}).pipe(Effect.provide(snapshotLayer(data, dir)))
await Effect.runPromise(program.pipe(Effect.provide(Logger.layer([]))))
await reset(dir)
const store = (await $`find ${path.join(data, "snapshot")} -mindepth 2 -maxdepth 2 -type d`.text()).trim()
const before = Number((await $`du -sk ${store}`.text()).split("\t")[0])
const started = performance.now()
await Effect.runPromise(
Effect.gen(function* () {
const git = yield* Git.Service
// The baseline worktree predates object packing.
if (!("objects" in git)) return
yield* git.objects.pack(
new Git.Repository({
worktree: AbsolutePath.make(dir),
gitDirectory: AbsolutePath.make(store),
commonDirectory: AbsolutePath.make(store),
}),
)
}).pipe(Effect.provide(AppNodeBuilder.build(Git.node)), Effect.provide(Logger.layer([]))),
)
const after = Number((await $`du -sk ${store}`.text()).split("\t")[0])
console.log(
`${name.padEnd(9)} packing ${before} KiB -> ${after} KiB in ${(performance.now() - started).toFixed(0)} ms`,
)
const size = await $`du -sk ${data}`.text()
const loose = (await $`find ${data}/snapshot -path '*/objects/??/*' -type f`.text()).split("\n").filter(Boolean)
const files = (await $`find ${data}/snapshot -maxdepth 3 -type f`.text()).split("\n").filter(Boolean)
console.log(
`${name.padEnd(9)} snapshot store ${size.split("\t")[0]} KiB, ${loose.length} loose objects, top-level files: ${files.map((file) => path.basename(file)).join(" ")}`,
)
await fs.rm(data, { recursive: true, force: true })
}
// ---------- robustness probes ----------
async function probe(label: string, run: () => Promise<boolean>) {
const ok = await run().catch((error) => {
console.error(error)
return false
})
console.log(`probe ${label.padEnd(44)} ${ok ? "PASS" : "FAIL"}`)
results.push({ fixture: "probe", scenario: label, samples: [{ ms: ok ? 1 : 0, spawns: 0 }] })
}
async function withRepo<A>(setup: (dir: string) => Promise<void>, body: (dir: string, data: string) => Promise<A>) {
const dir = await fs.realpath(await fs.mkdtemp(path.join(root, "probe-")))
const project = path.join(dir, "project")
await fs.mkdir(project)
await setup(project)
await gitInit(project)
const result = await body(project, path.join(dir, "data"))
await fs.rm(dir, { recursive: true, force: true })
return result
}
const captureOnce = (data: string, directory: string) =>
Effect.runPromise(
Effect.gen(function* () {
const snapshot = yield* Snapshot.Service
return yield* snapshot.capture()
}).pipe(Effect.provide(snapshotLayer(data, directory)), Effect.provide(Logger.layer([]))),
)
async function snapshotGitDir(data: string) {
const projects = await fs.readdir(path.join(data, "snapshot"))
const project = path.join(data, "snapshot", projects[0]!)
return path.join(project, (await fs.readdir(project))[0]!)
}
async function probes() {
await probe("capture self-heals a zeroed index", () =>
withRepo(
(dir) => Bun.write(path.join(dir, "a.txt"), "a\n").then(() => {}),
async (dir, data) => {
await captureOnce(data, dir)
const gitDir = await snapshotGitDir(data)
await Bun.write(path.join(gitDir, "index"), new Uint8Array(1024))
await Bun.write(path.join(dir, "a.txt"), "b\n")
return (await captureOnce(data, dir)) !== undefined
},
),
)
await probe("capture survives a stale index.lock", () =>
withRepo(
(dir) => Bun.write(path.join(dir, "a.txt"), "a\n").then(() => {}),
async (dir, data) => {
await captureOnce(data, dir)
const gitDir = await snapshotGitDir(data)
await Bun.write(path.join(gitDir, "index.lock"), "")
await Bun.write(path.join(dir, "a.txt"), "b\n")
return (await captureOnce(data, dir)) !== undefined
},
),
)
await probe("capture works in a `..scope` directory", () =>
withRepo(
(dir) => Bun.write(path.join(dir, "..scope", "a.txt"), "a\n").then(() => {}),
async (dir, data) => (await captureOnce(data, path.join(dir, "..scope"))) !== undefined,
),
)
await probe("two processes capture concurrently (4x25)", () =>
withRepo(
(dir) => writeMany(Array.from({ length: 2000 }, (_, i) => [path.join(dir, `d${i % 20}`, `f${i}.txt`), `${i}\n`])),
async (dir, data) => {
await captureOnce(data, dir)
const worker = path.join(import.meta.dir, "benchmark-snapshot-worker.ts")
const outputs = await Promise.all(
Array.from({ length: 4 }, (_, id) => $`bun run ${worker} ${data} ${dir} ${id} 25`.nothrow().quiet()),
)
const failures = outputs.map((output) => Number(output.stdout.toString().trim() || "25"))
console.log(` concurrent capture failures per process: ${failures.join(", ")}`)
return failures.every((count) => count === 0)
},
),
)
}
await fs.mkdir(root, { recursive: true })
for (const name of selected) {
if (name === "probes") continue
const dir = name === "opencode" ? await cloneOpencode() : await generate(fixtures[name]!)
await bench(name, dir)
}
if (selected.includes("probes") || !flag("fixtures")) await probes()
if (jsonOut) await Bun.write(jsonOut, JSON.stringify(results, null, 2))
@@ -0,0 +1,22 @@
// Compare two snapshot-parity.ts outputs and print every diverging step per scenario.
const [left, right] = await Promise.all(process.argv.slice(2, 4).map((file) => Bun.file(file!).json()))
let divergent = 0
for (const name of new Set([...Object.keys(left), ...Object.keys(right)])) {
const a: unknown[] = left[name] ?? []
const b: unknown[] = right[name] ?? []
const differences = Array.from({ length: Math.max(a.length, b.length) }, (_, index) => index).filter(
(index) => JSON.stringify(a[index]) !== JSON.stringify(b[index]),
)
if (!differences.length) {
console.log(`same ${name}`)
continue
}
divergent++
console.log(`DIFFERENT ${name}`)
for (const index of differences) {
console.log(` step ${index}`)
console.log(` base: ${JSON.stringify(a[index])}`)
console.log(` new: ${JSON.stringify(b[index])}`)
}
}
console.log(divergent ? `${divergent} divergent scenarios` : "all scenarios identical")
+20
View File
@@ -0,0 +1,20 @@
#!/bin/sh
# Run the snapshot parity corpus from two worktrees under every Git config variant and compare.
# script/snapshot-parity-run.sh <baseline-worktree> [variant...]
set -e
base="$1"
shift
here="$(cd "$(dirname "$0")/.." && pwd)"
out="${SNAPSHOT_BENCH_ROOT:-${TMPDIR:-/tmp}}/opencode-snapshot-parity-results"
mkdir -p "$out"
cp "$here/script/snapshot-parity.ts" "$base/packages/core/script/snapshot-parity.ts"
variants="${*:-user empty split-index skip-hash autocrlf-and-no-untracked-cache fsmonitor template-hook}"
status=0
for name in $variants; do
variant="$(echo "$name" | tr '-' ' ')"
(cd "$base/packages/core" && PARITY_CONFIG="$variant" bun run script/snapshot-parity.ts "$out/$name-base.json" >/dev/null 2>&1)
(cd "$here" && PARITY_CONFIG="$variant" bun run script/snapshot-parity.ts "$out/$name-new.json" >/dev/null 2>&1)
echo "=== $variant"
bun run "$here/script/snapshot-parity-compare.ts" "$out/$name-base.json" "$out/$name-new.json" | grep -v "^same" || status=1
done
exit $status
+710
View File
@@ -0,0 +1,710 @@
/**
* Differential corpus for the snapshot engine. Run the same file from two worktrees
* and compare the JSON: tree IDs, changed-file lists, diffs, and post-restore
* worktree digests must be byte-identical.
*
* bun run script/snapshot-parity.ts <result.json> [scenario...]
*/
import { $ } from "bun"
import crypto from "crypto"
import fs from "fs/promises"
import os from "os"
import path from "path"
import { Effect, Layer, Logger, ManagedRuntime } from "effect"
import { AppNodeBuilder } from "../src/effect/app-node-builder"
import { Location } from "../src/location"
import { AbsolutePath, RelativePath } from "../src/schema"
import { Snapshot } from "../src/snapshot"
import { Global } from "@opencode/util/global"
const root = path.join(process.env.SNAPSHOT_BENCH_ROOT ?? os.tmpdir(), "opencode-snapshot-parity")
// User-level Git configuration reaches every snapshot command, so the corpus runs under several.
const variants: Record<string, string | undefined> = {
user: undefined,
empty: "",
"split index": "[core]\n\tsplitIndex = true\n",
"skip hash": "[index]\n\tskipHash = true\n[feature]\n\tmanyFiles = true\n",
"autocrlf and no untracked cache": "[core]\n\tautocrlf = true\n\tuntrackedCache = false\n",
fsmonitor: "[core]\n\tfsmonitor = true\n",
"template hook": "[init]\n\ttemplateDir = TEMPLATE\n",
}
const variant = process.env.PARITY_CONFIG ?? "user"
if (variants[variant] !== undefined) {
await fs.mkdir(root, { recursive: true })
const template = path.join(root, `template-${process.pid}`)
await fs.mkdir(path.join(template, "hooks"), { recursive: true })
await fs.writeFile(path.join(template, "hooks", "post-checkout"), '#!/bin/sh\necho "checkout $3" > .hooklog\n', {
mode: 0o755,
})
// Stores are created with `git init`, so template ignore rules must keep applying to them.
await fs.mkdir(path.join(template, "info"), { recursive: true })
await fs.writeFile(path.join(template, "info", "exclude"), "template-ignored/\n")
const config = path.join(root, `gitconfig-${process.pid}`)
await fs.writeFile(config, variants[variant]!.replace("TEMPLATE", template))
process.env.GIT_CONFIG_GLOBAL = config
}
const env = {
...process.env,
GIT_ALLOW_PROTOCOL: "file",
GIT_AUTHOR_NAME: "Parity",
GIT_AUTHOR_EMAIL: "parity@opencode.test",
GIT_COMMITTER_NAME: "Parity",
GIT_COMMITTER_EMAIL: "parity@opencode.test",
GIT_AUTHOR_DATE: "2026-01-01T00:00:00Z",
GIT_COMMITTER_DATE: "2026-01-01T00:00:00Z",
}
type Context = {
readonly dir: string
readonly write: (file: string, content: string | Uint8Array) => Promise<void>
readonly remove: (file: string) => Promise<void>
readonly git: (args: string[], cwd?: string) => Promise<string>
readonly capture: (label: string, location?: string) => Promise<void>
readonly files: (from: string, to: string) => Promise<void>
readonly diff: (from: string, to: string, paths?: string[]) => Promise<void>
readonly restore: (label: string, files: Record<string, string>) => Promise<void>
readonly digest: (label: string) => Promise<void>
}
type Scenario = (ctx: Context) => Promise<void>
const commit = async (ctx: Context, message = "initial") => {
await ctx.git(["add", "-A"])
await ctx.git(["commit", "-q", "--no-gpg-sign", "--allow-empty", "-m", message])
}
const scenarios: Record<string, Scenario> = {
async basic(ctx) {
await ctx.write("a.txt", "a\n")
await ctx.write("b.txt", "b\n")
await ctx.write("dir/c.txt", "c\n")
await commit(ctx)
await ctx.capture("t0")
await ctx.write("a.txt", "a2\n")
await ctx.remove("b.txt")
await ctx.write("dir/new.txt", "new\n")
await ctx.capture("t1")
await ctx.capture("t1-again")
await ctx.files("t0", "t1")
await ctx.diff("t0", "t1")
await ctx.diff("t0", "t1", ["a.txt"])
await ctx.restore("undo", { "a.txt": "t0", "b.txt": "t0", "dir/new.txt": "t0" })
await ctx.capture("t2")
await ctx.files("t0", "t2")
},
async "rename directory"(ctx) {
for (let i = 0; i < 20; i++) await ctx.write(`src/f${i}.ts`, `${i}\n`)
await commit(ctx)
await ctx.capture("t0")
await fs.rename(path.join(ctx.dir, "src"), path.join(ctx.dir, "lib"))
await ctx.capture("t1")
await ctx.files("t0", "t1")
await ctx.diff("t0", "t1")
},
async "file becomes directory"(ctx) {
await ctx.write("a", "file\n")
await ctx.write("keep.txt", "keep\n")
await commit(ctx)
await ctx.capture("t0")
await ctx.remove("a")
await ctx.write("a/b", "nested\n")
await ctx.capture("t1")
await ctx.files("t0", "t1")
await ctx.restore("undo", { a: "t0", "a/b": "t0" })
await ctx.capture("t2")
},
async "directory becomes file"(ctx) {
await ctx.write("d/x", "x\n")
await ctx.write("d/y", "y\n")
await commit(ctx)
await ctx.capture("t0")
await ctx.remove("d")
await ctx.write("d", "file now\n")
await ctx.capture("t1")
await ctx.files("t0", "t1")
await ctx.restore("undo", { "d/x": "t0", "d/y": "t0", d: "t0" })
await ctx.capture("t2")
},
async "embedded repository"(ctx) {
await ctx.write("top.txt", "top\n")
await commit(ctx)
await ctx.capture("t0")
await ctx.write("sub/inner.txt", "inner\n")
await ctx.git(["init", "-q"], path.join(ctx.dir, "sub"))
await ctx.git(["add", "-A"], path.join(ctx.dir, "sub"))
await ctx.git(["commit", "-q", "--no-gpg-sign", "-m", "inner"], path.join(ctx.dir, "sub"))
await ctx.capture("t1")
await ctx.write("sub/inner.txt", "changed\n")
await ctx.capture("t2")
await ctx.git(["commit", "-q", "--no-gpg-sign", "-am", "inner 2"], path.join(ctx.dir, "sub"))
await ctx.capture("t3")
await ctx.files("t0", "t3")
await ctx.diff("t0", "t3")
},
async symlinks(ctx) {
await ctx.write("target.txt", "target\n")
await ctx.write("dir/inside.txt", "inside\n")
await ctx.write("plain.txt", "plain\n")
await commit(ctx)
await ctx.capture("t0")
await fs.symlink("target.txt", path.join(ctx.dir, "link"))
await fs.symlink("dir", path.join(ctx.dir, "dirlink"))
await fs.symlink("missing", path.join(ctx.dir, "dangling"))
await ctx.remove("plain.txt")
await fs.symlink("target.txt", path.join(ctx.dir, "plain.txt"))
await ctx.capture("t1")
await ctx.write("target.txt", "target 2\n")
await ctx.capture("t2")
await ctx.files("t0", "t2")
await ctx.diff("t0", "t2")
await ctx.restore("undo", { link: "t0", dirlink: "t0", dangling: "t0", "plain.txt": "t0", "target.txt": "t0" })
await ctx.capture("t3")
},
async "executable bit"(ctx) {
await ctx.write("run.sh", "echo hi\n")
await commit(ctx)
await ctx.capture("t0")
await fs.chmod(path.join(ctx.dir, "run.sh"), 0o755)
await ctx.capture("t1")
await ctx.files("t0", "t1")
await ctx.diff("t0", "t1")
await ctx.restore("undo", { "run.sh": "t0" })
await ctx.capture("t2")
},
async "special names"(ctx) {
const names = [
"with space.txt",
"caf\u00e9.txt",
"cafe\u0301-nfd.txt",
'quote"s.txt',
"-rf",
":colon.txt",
"!bang.txt",
"#hash.txt",
"app/[slug]/page.tsx",
"star*.txt",
"q?.txt",
"back\\slash.txt",
"tab\tname.txt",
"new\nline.txt",
"\u65e5\u672c\u8a9e/\u30d5\u30a1\u30a4\u30eb.txt",
]
for (const name of names) await ctx.write(name, `${name}\n`)
await commit(ctx)
await ctx.capture("t0")
for (const name of names) await ctx.write(name, `${name} changed\n`)
await ctx.write("added [x].txt", "added\n")
await ctx.capture("t1")
await ctx.files("t0", "t1")
await ctx.diff("t0", "t1")
await ctx.restore("undo", Object.fromEntries([...names, "added [x].txt"].map((name) => [name, "t0"])))
await ctx.capture("t2")
},
async "ignore rules"(ctx) {
await ctx.write(".gitignore", "*.log\nbuild/\n!keep.log\n")
await ctx.write("keep.log", "keep\n")
await ctx.write("src/a.ts", "a\n")
await ctx.write("build/forced.txt", "forced\n")
await ctx.git(["add", "-A"])
await ctx.git(["add", "-f", "build/forced.txt"])
await ctx.git(["commit", "-q", "--no-gpg-sign", "-m", "initial"])
await ctx.write(".git/info/exclude", "secret/\n*.local\n")
await ctx.capture("t0")
await ctx.write("debug.log", "ignored\n")
await ctx.write("keep.log", "keep 2\n")
await ctx.write("build/forced.txt", "forced 2\n")
await ctx.write("build/out.js", "out\n")
await ctx.write("secret/key.txt", "key\n")
await ctx.write("config.local", "local\n")
await ctx.write("template-ignored/x.txt", "template\n")
await ctx.write("src/a.ts", "a2\n")
await ctx.capture("t1")
await ctx.files("t0", "t1")
await ctx.diff("t0", "t1")
await ctx.write(".git/info/exclude", "")
await ctx.capture("t2")
await ctx.files("t1", "t2")
},
async "gitignore changes"(ctx) {
await ctx.write("gen/out.txt", "1\n")
await ctx.write("src.txt", "src\n")
await commit(ctx)
await ctx.capture("t0")
await ctx.write(".gitignore", "gen/\n")
await ctx.write("gen/out.txt", "2\n")
await ctx.capture("t1")
await ctx.files("t0", "t1")
await ctx.write(".gitignore", "")
await ctx.capture("t2")
await ctx.files("t1", "t2")
},
async "large files"(ctx) {
await ctx.write("tracked.bin", "small\n")
await commit(ctx)
await ctx.capture("t0")
await ctx.write("big.bin", new Uint8Array(3 * 1024 * 1024).fill(7))
await ctx.write("tracked.bin", new Uint8Array(3 * 1024 * 1024).fill(9))
await ctx.capture("t1")
await ctx.files("t0", "t1")
await ctx.write("big.bin", "now small\n")
await ctx.capture("t2")
await ctx.files("t1", "t2")
},
async "line endings"(ctx) {
await ctx.write(".gitattributes", "* text=auto\n*.bat eol=crlf\n")
await ctx.write("unix.txt", "one\ntwo\n")
await commit(ctx)
await ctx.capture("t0")
await ctx.write("unix.txt", "one\r\ntwo\r\nthree\r\n")
await ctx.write("run.bat", "echo\r\n")
await ctx.capture("t1")
await ctx.files("t0", "t1")
await ctx.diff("t0", "t1")
await ctx.restore("undo", { "unix.txt": "t0", "run.bat": "t0" })
await ctx.digest("after-undo")
},
async "merge conflict in source"(ctx) {
await ctx.write("c.txt", "base\n")
await commit(ctx)
await ctx.git(["checkout", "-q", "-b", "other"])
await ctx.write("c.txt", "other\n")
await commit(ctx, "other")
await ctx.git(["checkout", "-q", "-"])
await ctx.write("c.txt", "main\n")
await commit(ctx, "main")
await ctx.git(["merge", "-q", "other"]).catch(() => "")
await ctx.capture("t0")
await ctx.write("c.txt", "resolved\n")
await ctx.capture("t1")
await ctx.files("t0", "t1")
},
async "case-only rename"(ctx) {
await ctx.write("Readme.md", "readme\n")
await commit(ctx)
await ctx.capture("t0")
await fs.rename(path.join(ctx.dir, "Readme.md"), path.join(ctx.dir, "README.md"))
await ctx.capture("t1")
await ctx.files("t0", "t1")
await ctx.write("README.md", "changed\n")
await ctx.capture("t2")
await ctx.files("t0", "t2")
},
async "location subdirectory"(ctx) {
await ctx.write("pkg/a.txt", "a\n")
await ctx.write("other/b.txt", "b\n")
await commit(ctx)
await ctx.capture("t0", "pkg")
await ctx.write("pkg/a.txt", "a2\n")
await ctx.write("pkg/new.txt", "new\n")
await ctx.write("other/b.txt", "b2\n")
await ctx.capture("t1", "pkg")
await ctx.files("t0", "t1")
await ctx.capture("t2", ".")
await ctx.files("t1", "t2")
await ctx.capture("t3", "pkg")
await ctx.files("t2", "t3")
},
async "many untracked"(ctx) {
await ctx.write("seed.txt", "seed\n")
await commit(ctx)
for (let i = 0; i < 3000; i++) await ctx.write(`gen/d${i % 30}/f${i}.txt`, `${i}\n`)
await ctx.capture("t0")
await ctx.capture("t0-again")
},
async "empty directories and delete all"(ctx) {
await ctx.write("a/b/c.txt", "c\n")
await ctx.write("d.txt", "d\n")
await commit(ctx)
await ctx.capture("t0")
await fs.mkdir(path.join(ctx.dir, "empty/nested"), { recursive: true })
await ctx.capture("t1")
await ctx.remove("a")
await ctx.remove("d.txt")
await ctx.capture("t2")
await ctx.files("t0", "t2")
await ctx.restore("undo", { "a/b/c.txt": "t0", "d.txt": "t0" })
await ctx.capture("t3")
},
async "restore with overlapping paths"(ctx) {
await ctx.write("a", "file a\n")
await commit(ctx)
await ctx.capture("t0")
await ctx.remove("a")
await ctx.write("a/b", "nested b\n")
await ctx.capture("t1")
await ctx.remove("a")
await ctx.write("a", "file again\n")
await ctx.capture("t2")
await ctx.restore("order-1", { a: "t0", "a/b": "t0" })
await ctx.restore("order-2", { "a/b": "t1", a: "t1" })
await ctx.restore("order-3", { a: "t1", "a/b": "t1" })
await ctx.restore("order-4", { "a/b": "t0", a: "t2" })
},
async "restore from several trees"(ctx) {
await ctx.write("x.txt", "x0\n")
await ctx.write("y.txt", "y0\n")
await commit(ctx)
await ctx.capture("t0")
await ctx.write("x.txt", "x1\n")
await ctx.write("z.txt", "z1\n")
await ctx.capture("t1")
await ctx.write("x.txt", "x2\n")
await ctx.write("y.txt", "y2\n")
await ctx.write("z.txt", "z2\n")
await ctx.write("w.txt", "w2\n")
await ctx.restore("mixed", { "x.txt": "t1", "y.txt": "t0", "z.txt": "t0", "w.txt": "t1" })
},
async "staged changes in source"(ctx) {
await ctx.write("a.txt", "a\n")
await commit(ctx)
await ctx.write("a.txt", "staged\n")
await ctx.write("new.txt", "staged new\n")
await ctx.git(["add", "-A"])
await ctx.write("a.txt", "worktree\n")
await ctx.capture("t0")
await ctx.write("a.txt", "worktree 2\n")
await ctx.capture("t1")
await ctx.files("t0", "t1")
},
async "binary and replaced types"(ctx) {
await ctx.write("img.bin", new Uint8Array([0, 1, 2, 3, 0, 255]))
await ctx.write("f.txt", "file\n")
await commit(ctx)
await ctx.capture("t0")
await ctx.write("img.bin", new Uint8Array([0, 9, 9, 3, 0, 255, 1]))
await ctx.remove("f.txt")
await fs.symlink("img.bin", path.join(ctx.dir, "f.txt"))
await ctx.capture("t1")
await ctx.files("t0", "t1")
await ctx.diff("t0", "t1")
await ctx.remove("f.txt")
await ctx.write("f.txt", "file\n")
await ctx.capture("t2")
await ctx.files("t0", "t2")
},
async "sparse checkout source"(ctx) {
await ctx.write("in/a.txt", "a\n")
await ctx.write("out/b.txt", "b\n")
await commit(ctx)
await ctx.git(["sparse-checkout", "set", "in"])
await ctx.capture("t0")
await ctx.write("in/a.txt", "a2\n")
await ctx.write("in/new.txt", "new\n")
await ctx.capture("t1")
await ctx.files("t0", "t1")
},
async "unreadable file"(ctx) {
await ctx.write("ok.txt", "ok\n")
await commit(ctx)
await ctx.capture("t0")
await ctx.write("locked.txt", "locked\n")
await fs.chmod(path.join(ctx.dir, "locked.txt"), 0o000)
await ctx.capture("t1")
await fs.chmod(path.join(ctx.dir, "locked.txt"), 0o644)
await ctx.capture("t2")
await ctx.files("t0", "t2")
},
async "no commits yet"(ctx) {
await ctx.write("a.txt", "a\n")
await ctx.capture("t0")
await ctx.git(["add", "a.txt"])
await ctx.write("b.txt", "b\n")
await ctx.capture("t1")
await ctx.files("t0", "t1")
},
async "special files"(ctx) {
await ctx.write("a.txt", "a\n")
await commit(ctx)
await ctx.capture("t0")
await $`mkfifo ${path.join(ctx.dir, "pipe")}`.quiet()
await ctx.capture("t1")
await ctx.remove("pipe")
await ctx.capture("t2")
await ctx.files("t0", "t2")
},
async "tracked submodule"(ctx) {
const other = path.join(path.dirname(ctx.dir), "library")
await fs.mkdir(other)
await ctx.git(["init", "-q"], other)
await fs.writeFile(path.join(other, "lib.txt"), "lib\n")
await ctx.git(["add", "-A"], other)
await ctx.git(["commit", "-q", "--no-gpg-sign", "-m", "lib"], other)
await ctx.git(["-c", "protocol.file.allow=always", "submodule", "add", "-q", other, "vendor/library"])
// The absolute temporary URL differs between runs; pin it so trees are comparable.
await ctx.git(["config", "-f", ".gitmodules", "submodule.vendor/library.url", "../library"])
await commit(ctx)
await ctx.capture("t0")
await ctx.write("vendor/library/lib.txt", "changed\n")
await ctx.capture("t1")
await ctx.git(["commit", "-q", "--no-gpg-sign", "-am", "bump"], path.join(ctx.dir, "vendor/library"))
await ctx.capture("t2")
await ctx.files("t0", "t2")
await ctx.diff("t0", "t2")
},
async "index flags in source"(ctx) {
await ctx.write("assumed.txt", "a\n")
await ctx.write("skipped.txt", "s\n")
await ctx.write("plain.txt", "p\n")
await commit(ctx)
await ctx.git(["update-index", "--assume-unchanged", "assumed.txt"])
await ctx.git(["update-index", "--skip-worktree", "skipped.txt"])
await ctx.write("intent.txt", "intent\n")
await ctx.git(["add", "-N", "intent.txt"])
await ctx.capture("t0")
await ctx.write("assumed.txt", "a2\n")
await ctx.write("skipped.txt", "s2\n")
await ctx.write("intent.txt", "intent 2\n")
await ctx.write("plain.txt", "p2\n")
await ctx.capture("t1")
await ctx.files("t0", "t1")
},
async "restore key spellings"(ctx) {
await ctx.write("Readme.md", "readme\n")
await ctx.write("dir/a.txt", "a\n")
await ctx.write("dir/b.txt", "b\n")
await ctx.write("caf\u00e9.txt", "nfc\n")
await commit(ctx)
await ctx.capture("t0")
await ctx.write("Readme.md", "changed\n")
await ctx.write("dir/a.txt", "a2\n")
await ctx.write("dir/c.txt", "c2\n")
await ctx.write("caf\u00e9.txt", "nfc 2\n")
await ctx.restore("case", { "readme.md": "t0" })
await ctx.restore("directory", { dir: "t0" })
await ctx.restore("nfd", { "cafe\u0301.txt": "t0" })
await ctx.restore("dot-slash", { "./Readme.md": "t0" })
},
async "large batched restore"(ctx) {
for (let i = 0; i < 300; i++) await ctx.write(`src/m${i % 7}/f${i}.ts`, `export const v = ${i}\n`)
await commit(ctx)
await ctx.capture("t0")
for (let i = 0; i < 300; i += 2) await ctx.write(`src/m${i % 7}/f${i}.ts`, `changed ${i}\n`)
for (let i = 0; i < 50; i++) await ctx.write(`src/new/n${i}.ts`, `new ${i}\n`)
for (let i = 1; i < 300; i += 3) await ctx.remove(`src/m${i % 7}/f${i}.ts`)
await ctx.capture("t1")
await ctx.files("t0", "t1")
const changed = await Promise.resolve().then(async () => {
const files: Record<string, string> = {}
for (let i = 0; i < 300; i++) files[`src/m${i % 7}/f${i}.ts`] = "t0"
for (let i = 0; i < 50; i++) files[`src/new/n${i}.ts`] = "t0"
return files
})
await ctx.restore("all", changed)
await ctx.capture("t2")
await ctx.files("t0", "t2")
},
async "restore blocked by a file"(ctx) {
await ctx.write("a/b.txt", "b\n")
await ctx.write("z.txt", "z\n")
await commit(ctx)
await ctx.capture("t0")
await ctx.remove("a")
await ctx.write("a", "now a file\n")
await ctx.write("z.txt", "z2\n")
await ctx.restore("blocked", { "a/b.txt": "t0", "z.txt": "t0" })
},
async "location with pattern characters"(ctx) {
await ctx.write("app/[id]/page.tsx", "page\n")
await ctx.write("app/i/other.tsx", "other\n")
await commit(ctx)
await ctx.capture("t0", "app/[id]")
await ctx.write("app/[id]/page.tsx", "page 2\n")
await ctx.write("app/i/other.tsx", "other 2\n")
await ctx.capture("t1", "app/[id]")
await ctx.files("t0", "t1")
},
async "restore paths absent from both sides"(ctx) {
await ctx.write("a.txt", "a\n")
await commit(ctx)
await ctx.capture("t0")
await ctx.restore("missing", { "never.txt": "t0", "gone/deep/x.txt": "t0" })
await ctx.write("later.txt", "later\n")
await ctx.restore("remove-later", { "later.txt": "t0" })
},
}
async function run(name: string, scenario: Scenario) {
const base = await fs.realpath(await fs.mkdtemp(path.join(root, "case-")))
const dir = path.join(base, "project")
const data = path.join(base, "data")
await fs.mkdir(dir, { recursive: true })
const git = (args: string[], cwd = dir) =>
$`git -c core.fsmonitor=false -c core.splitIndex=false -c init.defaultBranch=main ${args}`
.cwd(cwd)
.env(env)
.quiet()
.text()
await git(["init", "-q"])
const log: unknown[] = []
const trees = new Map<string, Snapshot.ID | undefined>()
// One long-lived runtime per Location, as in production, so in-process caches are exercised.
const runtimes = new Map<string, ManagedRuntime.ManagedRuntime<Snapshot.Service, never>>()
const runtime = (location = ".") => {
const existing = runtimes.get(location)
if (existing) return existing
const created = ManagedRuntime.make(
Layer.provide(snapshotLayer(data, path.join(dir, location)), Logger.layer([])) as Layer.Layer<Snapshot.Service>,
)
runtimes.set(location, created)
return created
}
const use = <A>(location: string | undefined, body: (snapshot: Snapshot.Interface) => Effect.Effect<A, unknown>) =>
runtime(location)
.runPromise(
Effect.gen(function* () {
return yield* body(yield* Snapshot.Service)
}).pipe(Effect.exit),
)
.then((exit) => (exit._tag === "Success" ? { ok: exit.value } : { error: true }))
const tree = (label: string) => {
const id = trees.get(label)
if (!id) throw new Error(`no tree ${label}`)
return id
}
const ctx: Context = {
dir,
write: async (file, content) => {
await fs.mkdir(path.dirname(path.join(dir, file)), { recursive: true })
await fs.writeFile(path.join(dir, file), content)
},
remove: (file) => fs.rm(path.join(dir, file), { recursive: true, force: true }),
git,
capture: async (label, location) => {
const result = await use(location, (snapshot) => snapshot.capture())
const id = "ok" in result ? result.ok : undefined
trees.set(label, id)
log.push({ capture: label, tree: id ?? null, entries: id ? await listing(data, id) : null })
},
files: async (from, to) => {
if (!trees.get(from) || !trees.get(to)) return void log.push({ files: [from, to], skipped: true })
log.push({ files: [from, to], result: await use(undefined, (s) => s.files({ from: tree(from), to: tree(to) })) })
},
diff: async (from, to, paths) => {
if (!trees.get(from) || !trees.get(to)) return void log.push({ diff: [from, to], skipped: true })
const result = await use(undefined, (s) =>
s.diff({ from: tree(from), to: tree(to), paths: paths?.map((file) => RelativePath.make(file)) }),
)
log.push({
diff: [from, to, paths ?? null],
result:
"ok" in result && result.ok
? result.ok.map((file) => ({
...file,
patch: crypto.createHash("sha1").update(file.patch).digest("hex"),
}))
: result,
})
},
restore: async (label, files) => {
const plan = new Map(
Object.entries(files).flatMap(([file, from]) => {
const id = trees.get(from)
return id ? [[RelativePath.make(file), id] as const] : []
}),
)
const result = await use(undefined, (s) => s.restore({ files: plan }))
log.push({ restore: label, result: "ok" in result ? "ok" : result, worktree: await digest(dir) })
},
digest: async (label) => {
log.push({ digest: label, worktree: await digest(dir) })
},
}
const failure = await scenario(ctx).then(
() => undefined,
(error: unknown) => String(error),
)
if (failure) log.push({ scenarioError: failure })
await Promise.all([...runtimes.values()].map((item) => item.dispose()))
await $`chmod -R u+rwX ${base}`.quiet().nothrow()
if (process.env.PARITY_KEEP) console.error(`kept ${base}`)
if (!process.env.PARITY_KEEP) await fs.rm(base, { recursive: true, force: true })
return log
}
function snapshotLayer(data: string, directory: string) {
return AppNodeBuilder.build(Snapshot.node, [
Location.node.replace(Location.boundNode(Location.Ref.make({ directory: AbsolutePath.make(directory) }))),
Global.node.replace(Global.layerWith({ data, config: path.join(data, "config") })),
])
}
/** The recursive listing of a captured tree, read from whichever snapshot store holds it. */
async function listing(data: string, tree: string) {
const stores = (await $`find ${path.join(data, "snapshot")} -mindepth 2 -maxdepth 2 -type d`.quiet().text())
.split("\n")
.filter(Boolean)
for (const store of stores) {
const result = await $`git --git-dir ${store} ls-tree -r -z ${tree}`.quiet().nothrow()
if (result.exitCode === 0) return result.stdout.toString().split("\0").filter(Boolean)
}
return null
}
/** Paths, types, modes, and content hashes of the worktree, excluding `.git`. */
async function digest(dir: string) {
const entries: string[] = []
const walk = async (current: string) => {
for (const entry of (await fs.readdir(current, { withFileTypes: true })).toSorted((a, b) =>
a.name < b.name ? -1 : 1,
)) {
const full = path.join(current, entry.name)
const relative = path.relative(dir, full)
if (relative === ".git" || relative.startsWith(".git/")) continue
const stat = await fs.lstat(full)
if (stat.isSymbolicLink()) entries.push(`L ${relative} -> ${await fs.readlink(full)}`)
else if (stat.isDirectory()) {
entries.push(`D ${relative}`)
await walk(full)
} else {
const content = await fs.readFile(full).catch(() => Buffer.from("<unreadable>"))
entries.push(
`F ${relative} ${stat.mode & 0o111 ? "x" : "-"} ${crypto.createHash("sha1").update(content).digest("hex")}`,
)
}
}
}
await walk(dir)
return entries
}
await fs.mkdir(root, { recursive: true })
const [outputFile, ...selected] = process.argv.slice(2)
const output: Record<string, unknown> = {}
for (const [name, scenario] of Object.entries(scenarios)) {
if (selected.length && !selected.includes(name)) continue
output[name] = await run(name, scenario)
}
await Bun.write(outputFile!, JSON.stringify(output, null, 2))
+684 -106
View File
@@ -1,7 +1,7 @@
export * as Git from "./git.js"
import path from "path"
import { Context, Effect, Layer, Schema } from "effect"
import { Cause, Context, Effect, Exit, Layer, Option, Schedule, Schema } from "effect"
import { ChildProcess } from "effect/unstable/process"
import { AbsolutePath, RelativePath } from "./schema.js"
import { FSUtil } from "@opencode/util/fs-util"
@@ -30,6 +30,8 @@ const snapshotConfig = `[core]
symlinks = true
fsmonitor = false
untrackedCache = true
# A split index cannot name its shared file once manyFiles skips index checksums.
splitIndex = false
[feature]
manyFiles = true
[index]
@@ -40,6 +42,22 @@ const snapshotConfig = `[core]
export const TreeID = Schema.String.pipe(Schema.brand("Git.TreeID"))
export type TreeID = typeof TreeID.Type
const temporaryIndexPrefix = "index.opencode-"
// Like `git gc --auto`, pack once loose objects accumulate, and combine packs before lookups slow down.
const looseObjectLimit = 2048
const packCountLimit = 16
export interface CaptureInput {
readonly repository: Repository
readonly scopes: readonly RelativePath[]
/**
* Source repository whose ignore rules decide which paths are recorded. Its index also
* rebuilds an unreadable snapshot index without rehashing every tracked file.
*/
readonly ignores?: Repository
readonly maximumUntrackedFileBytes?: number
}
export class OperationError extends Schema.TaggedError<OperationError>()("Git.OperationError", {
operation: Schema.Literals([
"clone",
@@ -52,6 +70,7 @@ export class OperationError extends Schema.TaggedError<OperationError>()("Git.Op
"list_files",
"diff",
"restore",
"pack",
]),
message: Schema.String,
directory: Schema.optional(AbsolutePath),
@@ -134,12 +153,7 @@ export interface Interface {
}) => Effect.Effect<ReadonlySet<RelativePath>, OperationError>
}
readonly tree: {
readonly capture: (input: {
repository: Repository
scopes: readonly RelativePath[]
ignores?: Repository
maximumUntrackedFileBytes?: number
}) => Effect.Effect<TreeID, OperationError>
readonly capture: (input: CaptureInput) => Effect.Effect<TreeID, OperationError>
readonly write: (repository: Repository) => Effect.Effect<TreeID, OperationError>
readonly files: (input: {
repository: Repository
@@ -158,6 +172,14 @@ export interface Interface {
files: ReadonlyMap<RelativePath, TreeID>
}) => Effect.Effect<void, OperationError>
}
readonly objects: {
/**
* Pack loose objects and combine small packs without dropping any object,
* reachable or not. Returns without work below the thresholds; safe to run
* concurrently with captures and with other processes.
*/
readonly pack: (repository: Repository) => Effect.Effect<void, OperationError>
}
}
export class Service extends Context.Service<Service, Interface>()("@opencode/Git") {}
@@ -310,13 +332,13 @@ const layer = Layer.effect(
operationName: OperationError["operation"],
repository: Repository,
args: string[],
options?: { stdin?: string; env?: Record<string, string>; maxOutputBytes?: number },
options?: { stdin?: string | Uint8Array; env?: Record<string, string>; maxOutputBytes?: number; index?: string },
) {
const result = yield* proc
.run(
ChildProcess.make(gitExecutable, repositoryArgs(repository, args), {
cwd: repository.worktree,
env: options?.env,
env: options?.index ? { ...options.env, GIT_INDEX_FILE: options.index } : options?.env,
extendEnv: true,
}),
{ stdin: options?.stdin, maxOutputBytes: options?.maxOutputBytes },
@@ -342,43 +364,155 @@ const layer = Layer.effect(
})
})
/**
* A new store is initialized in a temporary sibling directory and renamed into
* place, so processes racing to create the same store never observe or write a
* half-initialized one; the loser discards its copy and adopts the winner's.
*/
const create = Effect.fn("Git.repo.create")(function* (input: {
worktree: AbsolutePath
gitDirectory: AbsolutePath
seed?: Repository
}) {
const operationError = (message: string) => (cause: unknown) =>
new OperationError({ operation: "create", directory: input.gitDirectory, message, cause })
yield* fs.ensureDir(input.gitDirectory).pipe(Effect.mapError(operationError("Failed to create Git storage")))
const repository = new Repository({
worktree: input.worktree,
gitDirectory: input.gitDirectory,
commonDirectory: input.gitDirectory,
})
yield* repositoryOperation("create", repository, ["init"])
if (yield* fs.existsSafe(path.join(input.gitDirectory, "HEAD"))) {
yield* initRepository({ ...input, directory: input.gitDirectory })
return repository
}
const temporaryDirectory = AbsolutePath.make(`${input.gitDirectory}.init-${uniqueSuffix()}`)
yield* Effect.gen(function* () {
yield* fs.writeFileString(path.join(input.gitDirectory, snapshotConfigFile), snapshotConfig)
const config = path.join(input.gitDirectory, "config")
yield* initRepository({ ...input, directory: temporaryDirectory })
const renamed = yield* fs.rename(temporaryDirectory, input.gitDirectory).pipe(Effect.exit)
if (Exit.isSuccess(renamed) || (yield* fs.existsSafe(path.join(input.gitDirectory, "HEAD")))) return
return yield* new OperationError({
operation: "create",
directory: input.gitDirectory,
message: "Failed to move Git storage into place",
cause: Cause.squash(renamed.cause),
})
}).pipe(Effect.ensuring(fs.remove(temporaryDirectory, { recursive: true, force: true }).pipe(Effect.ignore)))
return repository
})
const initRepository = Effect.fnUntraced(function* (input: {
worktree: AbsolutePath
gitDirectory: AbsolutePath
directory: AbsolutePath
seed?: Repository
}) {
const operationError = (message: string) => (cause: unknown) =>
new OperationError({ operation: "create", directory: input.gitDirectory, message, cause })
yield* fs.ensureDir(input.directory).pipe(Effect.mapError(operationError("Failed to create Git storage")))
yield* repositoryOperation(
"create",
new Repository({ worktree: input.worktree, gitDirectory: input.directory, commonDirectory: input.directory }),
["init"],
)
yield* Effect.gen(function* () {
yield* fs.writeFileString(path.join(input.directory, snapshotConfigFile), snapshotConfig)
const config = path.join(input.directory, "config")
const current = yield* fs.readFileString(config)
if (current.includes(snapshotConfigInclude)) return
yield* fs.writeFileString(config, `${current.endsWith("\n") ? "\n" : "\n\n"}${snapshotConfigInclude}`, {
flag: "a",
})
}).pipe(Effect.mapError(operationError("Failed to configure Git storage")))
if (!input.seed) return repository
if (!input.seed) return
yield* fs
.ensureDir(path.join(input.gitDirectory, "objects", "info"))
.ensureDir(path.join(input.directory, "objects", "info"))
.pipe(Effect.mapError(operationError("Failed to configure shared Git objects")))
yield* fs
.writeFileString(
path.join(input.gitDirectory, "objects", "info", "alternates"),
path.join(input.directory, "objects", "info", "alternates"),
path.join(input.seed.commonDirectory, "objects") + "\n",
)
.pipe(Effect.mapError(operationError("Failed to configure shared Git objects")))
yield* fs
.copyFile(path.join(input.seed.gitDirectory, "index"), path.join(input.gitDirectory, "index"))
.copyFile(path.join(input.seed.gitDirectory, "index"), path.join(input.directory, "index"))
.pipe(Effect.ignore)
return repository
})
/**
* Both commands only read the index, so they never take `index.lock`. Two parallel processes beat one combined
* `ls-files -m -o`, whose lstat pass is not threaded like diff-files'.
*/
const listChanges = Effect.fnUntraced(function* (repository: Repository, scope: RelativePath, index?: string) {
const list = (args: string[]) =>
repositoryOperation("refresh", repository, args, { index }).pipe(
// Embedded repositories are listed as `dir/`; update-index records them as gitlinks.
Effect.map((result) => nuls(result.text).map((file) => RelativePath.make(file.replace(/\/$/, "")))),
)
const [tracked, untracked] = yield* Effect.all(
[
list(["diff-files", "--name-only", "-z", "--", literalPathspec(scope)]),
list(["ls-files", "--others", "--exclude-standard", "-z", "--", literalPathspec(scope)]),
],
{ concurrency: 2 },
)
return { tracked, untracked }
})
const updateIndex = Effect.fnUntraced(function* (input: {
repository: Repository
changes: { tracked: readonly RelativePath[]; untracked: readonly RelativePath[] }
ignores?: Repository
maximumUntrackedFileBytes?: number
index?: string
excluded?: ReadonlySet<RelativePath>
}) {
const candidates = [...input.changes.tracked, ...input.changes.untracked]
if (!candidates.length) return { skipped: [] }
const excluded =
input.excluded ??
(input.ignores ? yield* ignored({ repository: input.ignores, paths: candidates }) : new Set<RelativePath>())
const maximum = input.maximumUntrackedFileBytes
const skipped = maximum
? (yield* Effect.forEach(
input.changes.untracked.filter((item) => !excluded.has(item)),
(item) =>
fs.stat(path.join(input.repository.worktree, item)).pipe(
Effect.map((info) => (info.type === "File" && Number(info.size) > maximum ? item : undefined)),
Effect.orElseSucceed(() => undefined),
),
{ concurrency: 8 },
)).filter((item): item is RelativePath => item !== undefined)
: []
const skip = new Set(skipped)
const added = candidates.filter((item) => !excluded.has(item) && !skip.has(item))
// Untracked paths are not in the index, so only tracked paths that are now ignored need removing.
const removed = input.changes.tracked.filter((item) => excluded.has(item))
// update-index takes literal paths, so large lists avoid add's quadratic pathspec matching.
if (removed.length)
yield* repositoryOperation("refresh", input.repository, ["update-index", "--force-remove", "-z", "--stdin"], {
stdin: removed.join("\0") + "\0",
index: input.index,
})
const add = (paths: readonly RelativePath[]) =>
repositoryOperation(
"refresh",
input.repository,
["update-index", "--add", "--remove", "--replace", "-z", "--stdin"],
{
stdin: paths.join("\0") + "\0",
index: input.index,
},
)
if (added.length) yield* add(added)
// update-index drops a tracked file replaced by an embedded repository; only a second pass records its gitlink.
const embedded = (yield* Effect.forEach(
input.changes.tracked.filter((item) => !excluded.has(item)),
(item) =>
fs
.existsSafe(path.join(input.repository.worktree, item, ".git"))
.pipe(Effect.map((found) => (found ? item : undefined))),
{ concurrency: 8 },
)).filter((item): item is RelativePath => item !== undefined)
if (embedded.length) yield* add(embedded)
return { skipped }
})
const refresh = Effect.fn("Git.index.refresh")(function* (input: {
@@ -387,52 +521,7 @@ const layer = Layer.effect(
ignores?: Repository
maximumUntrackedFileBytes?: number
}) {
const list = (args: string[]) =>
repositoryOperation("refresh", input.repository, args).pipe(Effect.map((result) => nuls(result.text)))
const [tracked, untracked] = yield* Effect.all(
[
list(["diff-files", "--name-only", "-z", "--", input.scope]),
list(["ls-files", "--others", "--exclude-standard", "-z", "--", input.scope]),
],
{ concurrency: 2 },
)
const candidates = Array.from(new Set([...tracked, ...untracked])).map((file) => RelativePath.make(file))
if (!candidates.length) return { skipped: [] }
const excluded = input.ignores
? yield* ignored({ repository: input.ignores, paths: candidates })
: new Set<RelativePath>()
const allowed = candidates.filter((item) => !excluded.has(item))
const maximum = input.maximumUntrackedFileBytes
const skipped = maximum
? (yield* Effect.forEach(
untracked.filter((item) => allowed.includes(RelativePath.make(item))),
(item) =>
fs.stat(path.join(input.repository.worktree, item)).pipe(
Effect.map((info) =>
info.type === "File" && Number(info.size) > maximum ? RelativePath.make(item) : undefined,
),
Effect.orElseSucceed(() => undefined),
),
{ concurrency: 8 },
)).filter((item): item is RelativePath => item !== undefined)
: []
const stage = allowed.filter((item) => !skipped.includes(item))
const remove = [...excluded, ...skipped]
if (remove.length)
yield* repositoryOperation(
"refresh",
input.repository,
["rm", "--cached", "-f", "--ignore-unmatch", "--pathspec-from-file=-", "--pathspec-file-nul"],
{ stdin: remove.join("\0") + "\0" },
)
if (stage.length)
yield* repositoryOperation(
"refresh",
input.repository,
["add", "--all", "--sparse", "--pathspec-from-file=-", "--pathspec-file-nul"],
{ stdin: stage.join("\0") + "\0" },
)
return { skipped }
return yield* updateIndex({ ...input, changes: yield* listChanges(input.repository, input.scope) })
})
const ignored = Effect.fn("Git.index.ignored")(function* (input: {
@@ -472,26 +561,228 @@ const layer = Layer.effect(
return new Set(nuls(result.stdout.toString("utf8")).map((file) => RelativePath.make(file)))
})
const writeTree = Effect.fn("Git.tree.write")(function* (repository: Repository) {
return TreeID.make((yield* repositoryOperation("write_tree", repository, ["write-tree"])).text.trim())
const writeTree = Effect.fn("Git.tree.write")(function* (repository: Repository, index?: string) {
const tree = (yield* repositoryOperation("write_tree", repository, ["write-tree"], { index })).text.trim()
if (/^[0-9a-f]{40,64}$/.test(tree)) return TreeID.make(tree)
return yield* new OperationError({
operation: "write_tree",
directory: repository.worktree,
message: `Invalid tree ID: ${tree}`,
})
})
const captureTree = Effect.fn("Git.tree.capture")(
(input: {
repository: Repository
scopes: readonly RelativePath[]
ignores?: Repository
maximumUntrackedFileBytes?: number
}) =>
locked(
input.repository,
Effect.gen(function* () {
yield* Effect.forEach(input.scopes, (scope) => refresh({ ...input, scope }), { discard: true })
return yield* writeTree(input.repository)
}),
const indexFile = (repository: Repository) => path.join(repository.gitDirectory, "index")
/** Git only replaces an index by renaming a new file into place, so every write changes its inode. */
const indexFingerprint = (file: string) =>
fs.stat(file).pipe(
Effect.map((info) =>
[Option.getOrUndefined(info.ino), info.size, Option.getOrUndefined(info.mtime)?.getTime()].join(":"),
),
Effect.orElseSucceed(() => undefined),
)
/** A copied index must keep the timestamp Git uses to detect racily clean entries. */
const copyIndex = Effect.fnUntraced(function* (from: string, to: string) {
const info = yield* fs.stat(from)
yield* fs.copyFile(from, to)
const mtime = Option.getOrUndefined(info.mtime)
if (mtime)
yield* fs
.utimes(
to,
Option.getOrElse(info.atime, () => mtime),
mtime,
)
.pipe(Effect.ignore)
})
/**
* Run index writes against a temporary index, then rename it over the store's
* index. Processes never contend on `index.lock`, an interrupted or killed
* writer cannot leave a partial index behind, and a stale lock is irrelevant.
* Git never writes an index in place, it writes a lock file and renames it over
* the temporary name, so a hard-linked temporary index leaves the store's index
* untouched; a copy is the fallback when linking fails. The index is only a stat
* cache over the object store; when writers race, the last rename wins and the
* next capture reconciles against the worktree.
*/
const withTemporaryIndex = <A, E, R>(repository: Repository, use: (index: string) => Effect.Effect<A, E, R>) =>
Effect.acquireUseRelease(
Effect.gen(function* () {
const current = indexFile(repository)
const index = temporaryIndex(repository)
if (!(yield* fs.existsSafe(current))) return index
const linked = yield* fs.link(current, index).pipe(
Effect.as(true),
Effect.orElseSucceed(() => false),
)
if (linked) return index
yield* copyIndex(current, index).pipe(
Effect.mapError(
(cause) =>
new OperationError({
operation: "refresh",
directory: repository.gitDirectory,
message: "Failed to prepare a temporary index",
cause,
}),
),
)
return index
}),
(index) =>
Effect.gen(function* () {
const value = yield* use(index)
const installed = yield* Effect.uninterruptible(
Effect.gen(function* () {
// Fingerprint before the rename, so a concurrent writer's index can never be paired with this tree.
const fingerprint = yield* indexFingerprint(index)
const renamed = yield* fs.rename(index, indexFile(repository)).pipe(
// Windows reports a transient sharing violation while another process reads the index.
Effect.retry({ times: 3, schedule: Schedule.spaced("20 millis") }),
Effect.as(true),
Effect.orElseSucceed(() => false),
)
return renamed ? fingerprint : undefined
}),
)
return { value, installed }
}),
(index) => fs.remove(index, { force: true }).pipe(Effect.ignore),
)
// A clean capture returns the tree last written from the exact index it scanned.
const lastCaptures = new Map<string, { readonly fingerprint: string; readonly tree: TreeID }>()
const preparedStores = new Set<string>()
/**
* Changes are listed against the temporary index they are written to, never the
* store's index, which another process may replace at any moment. Each tree is
* therefore exactly its own index plus the worktree changes against it.
*/
const attemptCapture = Effect.fnUntraced(function* (input: CaptureInput) {
const result = yield* withTemporaryIndex(input.repository, (index) =>
Effect.gen(function* () {
const scanned = yield* indexFingerprint(index)
const changes = yield* Effect.forEach(input.scopes, (scope) => listChanges(input.repository, scope, index), {
concurrency: "unbounded",
})
const last = lastCaptures.get(input.repository.gitDirectory)
const untracked = changes.flatMap((change) => change.untracked)
// Paths the source ignores are listed as untracked on every scan; they leave the index unchanged.
const excluded =
scanned && last?.fingerprint === scanned && changes.every((change) => !change.tracked.length)
? input.ignores
? yield* ignored({ repository: input.ignores, paths: untracked })
: new Set<RelativePath>()
: undefined
if (last && excluded && untracked.every((file) => excluded.has(file))) return last.tree
yield* Effect.forEach(changes, (change) => updateIndex({ ...input, changes: change, index, excluded }), {
discard: true,
})
return yield* writeTree(input.repository, index)
}),
)
if (result.installed)
lastCaptures.set(input.repository.gitDirectory, { fingerprint: result.installed, tree: result.value })
if (!result.installed) lastCaptures.delete(input.repository.gitDirectory)
return result.value
})
const captureTree = Effect.fn("Git.tree.capture")((input: CaptureInput) =>
locked(
input.repository,
Effect.gen(function* () {
if (!preparedStores.has(input.repository.gitDirectory)) {
preparedStores.add(input.repository.gitDirectory)
yield* sweepTemporaryIndexes(input.repository)
yield* ensureSnapshotConfig(input.repository)
}
return yield* attemptCapture(input).pipe(
Effect.catch((error) =>
Effect.gen(function* () {
if (!(yield* isIndexUnreadable(input.repository))) return yield* error
yield* Effect.logWarning("rebuilding unreadable snapshot index", {
directory: input.repository.gitDirectory,
message: error.message,
})
yield* rebuildIndex(input.repository, input.ignores)
return yield* attemptCapture(input)
}),
),
)
}),
),
)
/**
* Only an index Git itself cannot read is rebuilt, decided by exit status rather
* than by localized messages. Other failures, such as an unreadable worktree
* file or a process that could not start, surface unchanged.
*/
const isIndexUnreadable = Effect.fnUntraced(function* (repository: Repository) {
if (!(yield* fs.existsSafe(indexFile(repository)))) return false
return yield* repositoryOperation("refresh", repository, [
"ls-files",
"-z",
"--",
literalPathspec(".opencode-probe"),
]).pipe(Effect.match({ onSuccess: () => false, onFailure: (error) => error.cause === undefined }))
})
const rebuildIndex = Effect.fnUntraced(function* (repository: Repository, source?: Repository) {
lastCaptures.delete(repository.gitDirectory)
const temporary = temporaryIndex(repository)
// The source's stat cache and cache-tree spare recovery from rehashing every tracked file.
const copied = source
? yield* copyIndex(path.join(source.gitDirectory, "index"), temporary).pipe(
Effect.as(true),
Effect.orElseSucceed(() => false),
)
: false
yield* (
copied ? fs.rename(temporary, indexFile(repository)) : fs.remove(indexFile(repository), { force: true })
).pipe(Effect.ignore, Effect.ensuring(fs.remove(temporary, { force: true }).pipe(Effect.ignore)))
})
// Stores created by earlier releases would otherwise keep their original settings; only OpenCode's include file is rewritten, never `config`.
const ensureSnapshotConfig = Effect.fnUntraced(function* (repository: Repository) {
const file = path.join(repository.gitDirectory, snapshotConfigFile)
const current = yield* fs.readFileStringSafe(file).pipe(Effect.orElseSucceed(() => undefined))
if (current === undefined || current === snapshotConfig) return
// Concurrent Git processes must never read a partially written file.
const temporary = `${file}.${uniqueSuffix()}`
yield* fs
.writeFileString(temporary, snapshotConfig)
.pipe(
Effect.andThen(fs.rename(temporary, file)),
Effect.ignore,
Effect.ensuring(fs.remove(temporary, { force: true }).pipe(Effect.ignore)),
)
})
const temporaryIndex = (repository: Repository) =>
path.join(repository.gitDirectory, `${temporaryIndexPrefix}${Date.now()}-${uniqueSuffix()}`)
/**
* Temporary indexes left by killed processes. Age comes from the creation time in
* the name, because a hard-linked index keeps the store index's old mtime; live
* captures finish long before the cutoff.
*/
const sweepTemporaryIndexes = Effect.fnUntraced(function* (repository: Repository) {
const cutoff = Date.now() - 60 * 60 * 1000
const entries = yield* fs.readDirectory(repository.gitDirectory).pipe(Effect.orElseSucceed(() => []))
yield* Effect.forEach(
entries.filter(
(entry) =>
entry.startsWith(temporaryIndexPrefix) &&
Number(entry.slice(temporaryIndexPrefix.length).split("-")[0]) < cutoff,
),
(entry) => fs.remove(path.join(repository.gitDirectory, entry), { force: true }).pipe(Effect.ignore),
{ discard: true },
)
})
const treeFiles = Effect.fn("Git.tree.files")(function* (input: {
repository: Repository
from: TreeID
@@ -523,7 +814,7 @@ const layer = Layer.effect(
paths?: readonly RelativePath[]
}) {
if (input.paths?.length === 0) return []
const args = ["--no-renames", input.from, input.to, "--", ...(input.paths ?? [])]
const args = ["--no-renames", input.from, input.to, "--", ...(input.paths ?? []).map(literalPathspec)]
// Patch headers have no -z form: unquoted paths keep chunksByFile matching non-ASCII names.
const [names, numbers, patch] = yield* Effect.all(
[
@@ -581,7 +872,7 @@ const layer = Layer.effect(
"-z",
tree,
"--",
file,
literalPathspec(file),
])).text.replace(/\0$/, "")
if (!text) return false
if (!/^\d+\s+\w+\s+[0-9a-f]+\t/.test(text))
@@ -593,35 +884,293 @@ const layer = Layer.effect(
return true
})
const removePath = (repository: Repository, file: RelativePath) =>
fs.remove(path.join(repository.worktree, file), { recursive: true, force: true }).pipe(
Effect.mapError(
(cause) =>
new OperationError({
operation: "restore",
directory: repository.worktree,
message: `Failed to remove ${file}`,
cause,
}),
),
)
/** Chunked to stay below argument limits. */
const pathsInTree = Effect.fnUntraced(function* (
repository: Repository,
tree: TreeID,
files: readonly RelativePath[],
) {
const chunks = Array.from({ length: Math.ceil(files.length / 512) }, (_, index) =>
files.slice(index * 512, (index + 1) * 512),
)
const listed = yield* Effect.forEach(chunks, (chunk) =>
repositoryOperation("restore", repository, ["ls-tree", "-z", tree, "--", ...chunk.map(literalPathspec)]).pipe(
Effect.flatMap((result) =>
Effect.forEach(nuls(result.text), (record) => {
const match = /^\d+ \w+ [0-9a-f]+\t(.*)$/s.exec(record)
if (match) return Effect.succeed(match[1])
return Effect.fail(
new OperationError({
operation: "restore",
directory: repository.worktree,
message: `Invalid tree entry: ${record}`,
}),
)
}),
),
),
)
return new Set(listed.flat())
})
/** Batched paths cost one ls-tree and one checkout per source tree instead of two processes per file. */
const restore = Effect.fn("Git.tree.restore")(
(input: { repository: Repository; files: ReadonlyMap<RelativePath, TreeID> }) =>
locked(
input.repository,
Effect.forEach(
input.files,
([file, tree]) =>
Effect.gen(function* () {
if (yield* hasEntry(input.repository, tree, file)) {
yield* repositoryOperation("restore", input.repository, ["checkout", tree, "--", file])
return
}
yield* fs.remove(path.join(input.repository.worktree, file), { recursive: true, force: true }).pipe(
Effect.mapError(
(cause) =>
new OperationError({
operation: "restore",
directory: input.repository.worktree,
message: `Failed to remove ${file}`,
cause,
}),
),
)
}),
{ discard: true },
),
Effect.gen(function* () {
if (!input.files.size) return
if (!canBatchRestore([...input.files.keys()]))
return yield* Effect.forEach(
input.files,
([file, tree]) =>
Effect.gen(function* () {
if (!(yield* hasEntry(input.repository, tree, file)))
return yield* removePath(input.repository, file)
yield* withTemporaryIndex(input.repository, (index) =>
repositoryOperation(
"restore",
input.repository,
["checkout", tree, "--", literalPathspec(file)],
{
index,
},
),
)
}),
{ discard: true },
)
const groups = new Map<TreeID, RelativePath[]>()
input.files.forEach((tree, file) => groups.set(tree, [...(groups.get(tree) ?? []), file]))
const plan = yield* Effect.forEach(groups, ([tree, files]) =>
pathsInTree(input.repository, tree, files).pipe(
Effect.map((present) => ({
tree,
present: files.filter((file) => present.has(file)),
absent: files.filter((file) => !present.has(file)),
})),
),
)
// Checkouts go first, so one failing removal cannot stop the other files from being restored.
const checkouts = plan.filter((item) => item.present.length)
if (checkouts.length)
yield* withTemporaryIndex(input.repository, (index) =>
Effect.forEach(
checkouts,
(item) =>
repositoryOperation(
"restore",
input.repository,
["checkout", item.tree, "--pathspec-from-file=-", "--pathspec-file-nul"],
{ stdin: item.present.map(literalPathspec).join("\0") + "\0", index },
),
{ discard: true },
),
)
yield* Effect.forEach(
plan.flatMap((item) => item.absent),
(file) => removePath(input.repository, file),
{ concurrency: 16, discard: true },
)
}),
),
)
/**
* Snapshot trees have no refs, so Git's own repack and prune would treat every
* snapshot as garbage. Packing instead hands pack-objects an explicit list of
* every loose object (and, when combining, every object in the old packs), checks
* the new index lists all of them, and only then deletes the loose copies and
* the combined packs. Nothing is pruned, and objects are never delegated to the
* source repository through alternates.
*/
const packObjects = Effect.fn("Git.objects.pack")(function* (repository: Repository) {
const objects = path.join(repository.gitDirectory, "objects")
const packDirectory = path.join(objects, "pack")
const loose = yield* looseObjects(objects)
const packs = yield* localPacks(packDirectory)
const combinePacks = packs.length >= packCountLimit
if (loose.length < looseObjectLimit && !combinePacks) return
yield* withPackLock(
repository,
Effect.gen(function* () {
yield* sweepTemporaryPacks(packDirectory)
const combined = combinePacks
? yield* Effect.forEach(packs, (pack) =>
packIndex(repository, path.join(packDirectory, `${pack}.idx`)).pipe(
Effect.map((oids) => ({ pack, oids })),
),
)
: []
const objectIds = [...new Set([...loose.map((item) => item.oid), ...combined.flatMap((item) => item.oids)])]
if (!objectIds.length) return
const newPacks = (yield* repositoryOperation(
"pack",
repository,
[
"-c",
"pack.threads=2",
"-c",
"pack.windowMemory=64m",
"pack-objects",
"-q",
"--non-empty",
path.join(packDirectory, "pack"),
],
{ stdin: objectIds.join("\n") + "\n" },
)).text
.split("\n")
.map((line) => line.trim())
.filter(Boolean)
const packedIds = new Set(
(yield* Effect.forEach(newPacks, (name) =>
packIndex(repository, path.join(packDirectory, `pack-${name}.idx`)),
)).flat(),
)
const missing = objectIds.filter((oid) => !packedIds.has(oid))
if (missing.length)
return yield* new OperationError({
operation: "pack",
directory: repository.gitDirectory,
message: `Packed ${packedIds.size} objects but ${missing.length} are missing; nothing was removed`,
})
// Deletion is quick and must not stop halfway, which could strand a pack without its index.
yield* Effect.uninterruptible(
Effect.gen(function* () {
yield* Effect.forEach(loose, (item) => fs.remove(item.file, { force: true }).pipe(Effect.ignore), {
concurrency: 16,
discard: true,
})
const newPackNames = new Set(newPacks.map((name) => `pack-${name}`))
yield* Effect.forEach(
combined.filter((item) => !newPackNames.has(item.pack)),
(item) =>
// The index goes first so readers never see an index without its pack.
Effect.forEach(
[".idx", ".pack", ".rev", ".bitmap", ".mtimes"],
(extension) =>
fs
.remove(path.join(packDirectory, `${item.pack}${extension}`), { force: true })
.pipe(Effect.ignore),
{ discard: true },
),
{ discard: true },
)
}),
)
}),
)
})
const looseObjects = Effect.fnUntraced(function* (objects: string) {
const fanout = (yield* fs.readDirectory(objects).pipe(Effect.orElseSucceed(() => []))).filter((entry) =>
/^[0-9a-f]{2}$/.test(entry),
)
const listed = yield* Effect.forEach(
fanout,
(prefix) =>
fs.readDirectory(path.join(objects, prefix)).pipe(
Effect.orElseSucceed(() => []),
Effect.map((entries) =>
entries
.filter((entry) => /^[0-9a-f]{38}$|^[0-9a-f]{62}$/.test(entry))
.map((entry) => ({ oid: prefix + entry, file: path.join(objects, prefix, entry) })),
),
),
{ concurrency: 16 },
)
return listed.flat()
})
// Packs with a .keep marker or a promisor file belong to someone else's policy and are left alone.
const localPacks = Effect.fnUntraced(function* (directory: string) {
const entries = new Set(yield* fs.readDirectory(directory).pipe(Effect.orElseSucceed(() => [])))
return [...entries]
.filter((entry) => entry.startsWith("pack-") && entry.endsWith(".pack"))
.map((entry) => entry.slice(0, -".pack".length))
.filter(
(pack) => entries.has(`${pack}.idx`) && !entries.has(`${pack}.keep`) && !entries.has(`${pack}.promisor`),
)
})
const packIndex = Effect.fnUntraced(function* (repository: Repository, file: string) {
const bytes = yield* fs.readFile(file).pipe(
Effect.mapError(
(cause) =>
new OperationError({
operation: "pack",
directory: repository.gitDirectory,
message: `Failed to read ${file}`,
cause,
}),
),
)
const result = yield* repositoryOperation("pack", repository, ["show-index"], { stdin: bytes })
return result.text.split("\n").flatMap((line) => {
const oid = line.split(" ")[1]
return oid ? [oid] : []
})
})
/**
* pack-objects writes tmp_* files and renames them; leftovers mean a killed
* process. A pack without its index is unreadable to Git, so removing one is
* lossless; the age cutoff skips a pack whose index is still being renamed.
*/
const sweepTemporaryPacks = Effect.fnUntraced(function* (directory: string) {
const cutoff = Date.now() - 60 * 60 * 1000
const entries = yield* fs.readDirectory(directory).pipe(Effect.orElseSucceed(() => []))
const indexed = new Set(entries.filter((entry) => entry.endsWith(".idx")).map((entry) => entry.slice(0, -4)))
yield* Effect.forEach(
entries.filter(
(entry) =>
entry.startsWith("tmp_") ||
(entry.startsWith("pack-") && entry.endsWith(".pack") && !indexed.has(entry.slice(0, -5))),
),
(entry) =>
fs.stat(path.join(directory, entry)).pipe(
Effect.flatMap((info) =>
(Option.getOrUndefined(info.mtime)?.getTime() ?? 0) < cutoff
? fs.remove(path.join(directory, entry), { force: true })
: Effect.void,
),
Effect.ignore,
),
{ discard: true },
)
})
// Cross-process exclusion: concurrent packing would stay lossless but duplicate work and objects.
const withPackLock = <A, E, R>(repository: Repository, effect: Effect.Effect<A, E, R>) => {
const file = path.join(repository.gitDirectory, "opencode-pack.lock")
return Effect.gen(function* () {
const stale = yield* fs.stat(file).pipe(
Effect.map((info) => (Option.getOrUndefined(info.mtime)?.getTime() ?? 0) < Date.now() - 60 * 60 * 1000),
Effect.orElseSucceed(() => false),
)
if (stale) yield* fs.remove(file, { force: true }).pipe(Effect.ignore)
const acquired = yield* fs.writeFileString(file, String(process.pid), { flag: "wx" }).pipe(
Effect.as(true),
Effect.orElseSucceed(() => false),
)
if (!acquired) return
yield* effect.pipe(Effect.ensuring(fs.remove(file, { force: true }).pipe(Effect.ignore)))
})
}
const worktreeRun = Effect.fnUntraced(function* (
operation: "create" | "remove" | "list",
repository: Repository,
@@ -702,11 +1251,12 @@ const layer = Layer.effect(
index: { refresh, ignored },
tree: {
capture: captureTree,
write: writeTree,
write: (repository) => writeTree(repository),
files: treeFiles,
diff: treeDiff,
restore,
},
objects: { pack: packObjects },
})
}),
)
@@ -749,6 +1299,34 @@ function nuls(text: string) {
return text.split("\0").filter(Boolean)
}
/** Pathspec magic that matches exactly one path, so names with `*`, `?`, `[`, or a leading `:` are not patterns. */
function literalPathspec(file: string) {
return `:(literal)${file}`
}
/**
* Paths whose restore operations commute: none inside another, and every segment
* printable ASCII that no platform rewrites. Git precomposes Unicode, filesystems
* may fold case, and Windows drops trailing dots and spaces and treats `\` and
* `:` specially, so anything else takes the ordered path.
*/
function canBatchRestore(files: readonly string[]) {
const canonical = files.every((file) =>
file.split("/").every((part) => /^[\x20-\x7e]+$/.test(part) && !/[\\:]/.test(part) && !/[. ]$/.test(part)),
)
if (!canonical) return false
const folded = new Set(files.map((file) => file.toLowerCase()))
if (folded.size !== files.length) return false
return files.every((file) => {
const parts = file.toLowerCase().split("/")
return parts.slice(1).every((_, index) => !folded.has(parts.slice(0, index + 1).join("/")))
})
}
function uniqueSuffix() {
return `${process.pid}-${Math.random().toString(36).slice(2)}`
}
function resolvePath(cwd: string, value: string) {
const trimmed = value.replace(/[\r\n]+$/, "")
if (!trimmed) return cwd
@@ -20,7 +20,8 @@ type Input = {
readonly agent: Agent.ID
readonly model: Model.Ref
readonly providerMetadataKey: string
readonly snapshot?: Snapshot.ID
/** The start snapshot, awaited before `Step.Started` so its capture can overlap the provider request. */
readonly pendingSnapshot?: Effect.Effect<Snapshot.ID | undefined>
readonly started: number
readonly assistantMessageID: SessionMessage.ID
}
@@ -98,6 +99,9 @@ export const createLLMEventPublisher = (bus: Pick<Bus.Interface, "publish">, inp
let stepSettlement: StepRecord["finish"]
const startAssistant = Effect.fnUntraced(function* () {
if (stepStarted) return assistantMessageID
const snapshot = input.pendingSnapshot ? yield* input.pendingSnapshot : undefined
// Check again after the await, so the check and the mark below never straddle a yield.
if (stepStarted) return assistantMessageID
stepStarted = true
yield* bus.publish(SessionEvent.Step.Started, {
@@ -105,7 +109,7 @@ export const createLLMEventPublisher = (bus: Pick<Bus.Interface, "publish">, inp
agent: input.agent,
model: input.model,
assistantMessageID,
snapshot: input.snapshot,
snapshot,
started: input.started,
})
return assistantMessageID
+10 -2
View File
@@ -70,14 +70,16 @@ export const make = Effect.gen(function* () {
const toolOutput = yield* ToolOutput.Service
const attempt = Effect.fn("SessionStep.attempt")(function* (input: Input) {
const startSnapshot = yield* snapshots.capture()
// The start snapshot only has to exist before local tools run, which cannot happen before Step.Started,
// so it is captured while the provider request is in flight instead of delaying it.
const pendingStartSnapshot = yield* snapshots.capture().pipe(Effect.forkScoped)
const publisher = createLLMEventPublisher(bus, {
sessionID: input.sessionID,
assistantMessageID: input.assistantMessageID,
agent: input.agent,
model: input.model.ref,
providerMetadataKey: input.model.model.route.providerMetadataKey ?? input.model.model.provider,
snapshot: startSnapshot,
pendingSnapshot: Fiber.join(pendingStartSnapshot),
started: yield* Clock.currentTimeMillis,
})
const toolRuns: Array<{
@@ -105,6 +107,8 @@ export const make = Effect.gen(function* () {
Stream.runForEach((event) =>
Effect.gen(function* () {
if (overflowFailure || publisher.hasProviderError()) return
// Wait here, where cancellation still works, rather than inside the uninterruptible publish.
if (!publisher.hasStarted()) yield* Fiber.join(pendingStartSnapshot)
if (
LLMEvent.is.providerError(event) &&
isContextOverflowFailure(event) &&
@@ -138,6 +142,9 @@ export const make = Effect.gen(function* () {
const stream = yield* restore(providerStream).pipe(Effect.exit)
const streamFailure = Option.getOrUndefined(Exit.findErrorOption(stream))
const streamInterrupted = Exit.hasInterrupts(stream)
// Cancelled before the start snapshot existed: record nothing, as when the capture preceded the request.
if (streamInterrupted && !publisher.hasStarted() && !pendingStartSnapshot.pollUnsafe())
return yield* Effect.failCause(stream.cause)
if (!overflowFailure && publisher.hasStarted()) yield* publisher.streamed()
if (streamInterrupted) yield* interruptTools
const joined = yield* restore(Fiber.awaitAll(toolRuns.map((run) => run.fiber))).pipe(Effect.exit)
@@ -224,6 +231,7 @@ export const make = Effect.gen(function* () {
const record = publisher.record()
if (record.finish || record.failure) {
const startSnapshot = yield* Fiber.join(pendingStartSnapshot)
const snapshot = yield* snapshots.capture()
const files =
startSnapshot && snapshot
+23 -12
View File
@@ -2,7 +2,7 @@ export * as Snapshot from "./snapshot.js"
import { makeLocationNode } from "@opencode/util/effect/app-node"
import path from "path"
import { Context, Effect, Fiber, Layer, Schema, Scope } from "effect"
import { Clock, Context, Effect, Fiber, Layer, Schema, Scope } from "effect"
import { FileDiff } from "@opencode/schema/file-diff"
import { FSUtil } from "@opencode/util/fs-util"
import { Git } from "./git.js"
@@ -106,26 +106,37 @@ const layer = Layer.effect(
const repository = repositoryFiber.pipe(Effect.uninterruptible, Effect.flatMap(Fiber.join))
const scope = Effect.fnUntraced(function* (worktree: AbsolutePath) {
const relative = path.relative(worktree, location.directory)
if (relative.startsWith("..") || path.isAbsolute(relative))
if (!FSUtil.contains(worktree, location.directory))
return yield* new Error({ operation: "capture", message: "Location is outside the project" })
return RelativePath.make(relative.replaceAll("\\", "/") || ".")
return RelativePath.make(path.relative(worktree, location.directory).replaceAll("\\", "/") || ".")
})
const enabled = () => location.vcs?.type === "git" && state.get().enabled
// `objects.pack` takes a cross-process lock, so a run that outlasts the interval never overlaps the next.
let lastPackCheck = Number.NEGATIVE_INFINITY
const packWhenDue = Effect.fnUntraced(function* (repository: Git.Repository) {
const now = yield* Clock.currentTimeMillis
if (now - lastPackCheck < 10 * 60 * 1000) return
lastPackCheck = now
yield* git.objects.pack(repository).pipe(
Effect.catch((cause) => Effect.logWarning("failed to pack snapshot objects", { cause })),
Effect.forkIn(lifetime),
)
})
const capture = Effect.fn("Snapshot.capture")(function* () {
if (!enabled()) return undefined
return yield* Effect.gen(function* () {
const repo = yield* repository
return ID.make(
yield* git.tree.capture({
repository: repo.snapshotRepository,
scopes: [yield* scope(repo.worktree)],
ignores: repo.source,
maximumUntrackedFileBytes: 2 * 1024 * 1024,
}),
)
const tree = yield* git.tree.capture({
repository: repo.snapshotRepository,
scopes: [yield* scope(repo.worktree)],
ignores: repo.source,
maximumUntrackedFileBytes: 2 * 1024 * 1024,
})
yield* packWhenDue(repo.snapshotRepository)
return ID.make(tree)
}).pipe(
Effect.catch((cause) => Effect.logWarning("failed to capture snapshot", { cause }).pipe(Effect.as(undefined))),
)
+194 -1
View File
@@ -2,7 +2,7 @@ import { describe, expect } from "bun:test"
import { $ } from "bun"
import fs from "fs/promises"
import path from "path"
import { Effect } from "effect"
import { Effect, Exit } from "effect"
import { LayerNode } from "@opencode/util/effect/layer-node"
import { Git } from "@opencode/core/git"
import { AbsolutePath, RelativePath } from "@opencode/core/schema"
@@ -296,3 +296,196 @@ describe("Git trees", () => {
}),
)
})
describe("Git objects", () => {
it.live(
"packs loose objects and combines packs without losing or delegating any object",
() =>
Effect.gen(function* () {
const root = yield* Effect.acquireRelease(
Effect.promise(() => tmpdir()),
(dir) => Effect.promise(() => dir[Symbol.asyncDispose]()),
)
const project = path.join(root.path, "project")
yield* Effect.promise(async () => {
await fs.mkdir(project)
await initRepo(project)
await Bun.write(path.join(project, "seed.txt"), "seed\n")
await $`git add . && git commit -q -m initial`.cwd(project).quiet()
})
const git = yield* Git.Service
const source = yield* git.repo.discover(AbsolutePath.make(project))
if (!source) throw new Error("Repository not found")
const repository = yield* git.repo.create({
worktree: source.worktree,
gitDirectory: AbsolutePath.make(path.join(root.path, "storage")),
seed: source,
})
const capture = () => git.tree.capture({ repository, scopes: [RelativePath.make(".")], ignores: source })
const trees = [yield* capture()]
// Enough distinct content to cross the loose-object threshold.
yield* Effect.promise(() =>
Promise.all(
Array.from({ length: 2100 }, (_, index) =>
Bun.write(path.join(project, `gen/d${index % 20}/f${index}.txt`), `${index}\n`),
),
),
)
trees.push(yield* capture())
const objects = path.join(repository.gitDirectory, "objects")
const alternates = path.join(objects, "info", "alternates")
// Objects stored in the snapshot repository itself, excluding anything borrowed through alternates.
const everything = () =>
Effect.promise(async () => {
const borrowed = await fs.readFile(alternates, "utf8")
await fs.rm(alternates)
const listed =
await $`git --git-dir ${repository.gitDirectory} cat-file --batch-all-objects ${"--batch-check=%(objectname)"}`.text()
await fs.writeFile(alternates, borrowed)
return listed.split("\n").filter(Boolean).toSorted()
})
const loose = () => Effect.promise(() => looseObjectIDs(objects))
const before = yield* everything()
expect((yield* loose()).length).toBeGreaterThan(2048)
yield* git.objects.pack(repository)
expect(yield* loose()).toEqual([])
expect(yield* everything()).toEqual(before)
for (let batch = 0; batch < 16; batch++) {
yield* Effect.promise(() => Bun.write(path.join(project, `batch-${batch}.txt`), `${batch}\n`))
trees.push(yield* capture())
yield* Effect.promise(async () => {
const list = await looseObjectIDs(objects)
await $`git --git-dir ${repository.gitDirectory} pack-objects -q ${path.join(objects, "pack", "pack")} < ${Buffer.from(list.join("\n") + "\n")}`.quiet()
await $`git --git-dir ${repository.gitDirectory} prune-packed`.quiet()
})
}
const packs = () =>
Effect.promise(async () =>
(await fs.readdir(path.join(objects, "pack"))).filter((file) => file.endsWith(".pack")),
)
expect((yield* packs()).length).toBeGreaterThanOrEqual(16)
const merged = yield* everything()
yield* git.objects.pack(repository)
expect((yield* packs()).length).toBe(1)
expect(yield* everything()).toEqual(merged)
expect((yield* git.tree.files({ repository, from: trees[0]!, to: trees.at(-1)! })).length).toBe(2100 + 16)
}),
{ timeout: 60_000 },
)
})
describe("Git capture", () => {
it.live("applies both the store's own and the source's ignore rules and never splits the index", () =>
Effect.gen(function* () {
const root = yield* Effect.acquireRelease(
Effect.promise(() => tmpdir()),
(dir) => Effect.promise(() => dir[Symbol.asyncDispose]()),
)
const project = path.join(root.path, "project")
yield* Effect.promise(async () => {
await fs.mkdir(project)
await initRepo(project)
await Bun.write(path.join(project, "tracked.bin"), "small\n")
await $`git add . && git commit -q -m initial`.cwd(project).quiet()
await Bun.write(path.join(project, ".git", "info", "exclude"), "source-ignored/\n")
})
const git = yield* Git.Service
const source = yield* git.repo.discover(AbsolutePath.make(project))
if (!source) throw new Error("Repository not found")
const storage = AbsolutePath.make(path.join(root.path, "storage"))
const repository = yield* git.repo.create({ worktree: source.worktree, gitDirectory: storage, seed: source })
yield* Effect.promise(async () => {
// Rules an init template left in the store keep applying, and a split index would lose entries.
await Bun.write(path.join(storage, "info", "exclude"), "store-ignored/\n")
await $`git --git-dir ${storage} config core.splitIndex true`.quiet()
})
const capture = () => git.tree.capture({ repository, scopes: [RelativePath.make(".")], ignores: source })
const before = yield* capture()
yield* Effect.promise(async () => {
await Bun.write(path.join(project, "tracked.bin"), "changed\n")
await Bun.write(path.join(project, "store-ignored", "a.txt"), "a\n")
await Bun.write(path.join(project, "source-ignored", "b.txt"), "b\n")
await Bun.write(path.join(project, "kept.txt"), "kept\n")
})
const after = yield* capture()
expect(yield* git.tree.files({ repository, from: before, to: after })).toEqual([
RelativePath.make("kept.txt"),
RelativePath.make("tracked.bin"),
])
expect(yield* capture()).toBe(after)
expect(yield* Effect.promise(() => fs.readFile(path.join(storage, "info", "exclude"), "utf8"))).toStartWith(
"store-ignored/\n",
)
}),
)
})
describe("Git capture recovery", () => {
const setup = Effect.gen(function* () {
const root = yield* Effect.acquireRelease(
Effect.promise(() => tmpdir()),
(dir) => Effect.promise(() => dir[Symbol.asyncDispose]()),
)
const project = path.join(root.path, "project")
yield* Effect.promise(async () => {
await fs.mkdir(project)
await initRepo(project)
await Bun.write(path.join(project, "a.txt"), "a\n")
await Bun.write(path.join(project, "b.txt"), "b\n")
await $`git add . && git commit -q -m initial`.cwd(project).quiet()
})
const git = yield* Git.Service
const source = yield* git.repo.discover(AbsolutePath.make(project))
if (!source) throw new Error("Repository not found")
const repository = yield* git.repo.create({
worktree: source.worktree,
gitDirectory: AbsolutePath.make(path.join(root.path, "storage")),
seed: source,
})
const capture = () => git.tree.capture({ repository, scopes: [RelativePath.make(".")], ignores: source })
return { root, project, repository, capture }
})
it.live("sweeps abandoned temporary indexes by their creation time, not their mtime", () =>
Effect.gen(function* () {
const { repository, capture } = yield* setup
const twoHoursAgo = Date.now() - 2 * 60 * 60 * 1000
const live = path.join(repository.gitDirectory, `index.opencode-${Date.now()}-1-live`)
const abandoned = path.join(repository.gitDirectory, `index.opencode-${twoHoursAgo}-1-abandoned`)
yield* Effect.promise(async () => {
await Bun.write(live, "")
await Bun.write(abandoned, "")
// A hard link to an index nobody wrote for hours carries that old mtime while in use.
await fs.utimes(live, new Date(twoHoursAgo), new Date(twoHoursAgo))
})
yield* capture()
expect(yield* Effect.promise(() => Bun.file(live).exists())).toBe(true)
expect(yield* Effect.promise(() => Bun.file(abandoned).exists())).toBe(false)
}),
)
it.live("keeps the index when Git cannot even start", () =>
Effect.gen(function* () {
const { root, project, repository, capture } = yield* setup
const before = yield* capture()
const moved = path.join(root.path, "moved")
yield* Effect.promise(() => fs.rename(project, moved))
expect(Exit.isFailure(yield* capture().pipe(Effect.exit))).toBe(true)
yield* Effect.promise(() => fs.rename(moved, project))
expect(yield* Effect.promise(() => Bun.file(path.join(repository.gitDirectory, "index")).exists())).toBe(true)
expect(yield* capture()).toBe(before)
}),
)
})
async function looseObjectIDs(objects: string) {
const prefixes = (await fs.readdir(objects)).filter((entry) => /^[0-9a-f]{2}$/.test(entry))
const listed = await Promise.all(
prefixes.map(async (prefix) => (await fs.readdir(path.join(objects, prefix))).map((entry) => prefix + entry)),
)
return listed.flat()
}
@@ -30,7 +30,11 @@ const base64 = "iVBORw0KGgoAAAANSUhEUgAAAAEAAAAB"
const capture = (
providerMetadataKey = "anthropic",
options?: { readonly interruptProgress?: boolean; readonly beforeTextDelta?: Effect.Effect<void> },
options?: {
readonly interruptProgress?: boolean
readonly beforeTextDelta?: Effect.Effect<void>
readonly pendingSnapshot?: Effect.Effect<Snapshot.ID | undefined>
},
) => {
const published: Array<{ readonly type: string; readonly data: unknown }> = []
const bus: Pick<Bus.Interface, "publish"> = {
@@ -60,6 +64,7 @@ const capture = (
providerID: Provider.ID.opencode,
},
providerMetadataKey,
pendingSnapshot: options?.pendingSnapshot,
started: 0,
assistantMessageID: SessionMessage.ID.create(),
}),
@@ -580,6 +585,20 @@ test("success event data can carry provider-executed result state", () => {
expect(decoded.resultState).toMatchObject({ result: { type: "content" } })
})
test("step start waits for the pending start snapshot", async () => {
const snapshot = Effect.runSync(Deferred.make<Snapshot.ID | undefined>())
const { published, publisher } = capture("anthropic", { pendingSnapshot: Deferred.await(snapshot) })
const started = Effect.runFork(publisher.publish(LLMEvent.stepStart({ index: 0 })))
await Effect.runPromise(Effect.yieldNow)
expect(published).toEqual([])
Effect.runSync(Deferred.succeed(snapshot, Snapshot.ID.make("tree-start")))
await Effect.runPromise(Fiber.join(started))
expect(published.map((event) => [event.type, (event.data as { snapshot?: string }).snapshot])).toEqual([
["session.step.started.1", "tree-start"],
])
})
test("step finish records settlement without publishing step ended", async () => {
const { published, publisher } = capture()
await Effect.runPromise(publisher.publish(LLMEvent.stepStart({ index: 0 })))
+191
View File
@@ -127,6 +127,180 @@ describe("Snapshot", () => {
),
)
testEffect(Layer.empty).live("recovers from a corrupt index and ignores a stale index lock", () =>
Effect.acquireUseRelease(
Effect.promise(() => tmpdir()),
(tmp) =>
Effect.gen(function* () {
const project = path.join(tmp.path, "project")
yield* Effect.promise(async () => {
await fs.mkdir(project)
await fs.writeFile(path.join(project, "tracked.txt"), "one\n")
await initGit(project, true)
})
yield* Effect.gen(function* () {
const snapshot = yield* Snapshot.Service
const before = yield* snapshot.capture()
const storage = yield* snapshotDirectory(tmp.path)
yield* Effect.promise(async () => {
// A process killed mid-write in older releases left a zeroed index and a lock behind.
await fs.writeFile(path.join(storage, "index"), new Uint8Array(512))
await fs.writeFile(path.join(storage, "index.lock"), "")
await fs.writeFile(path.join(project, "tracked.txt"), "two\n")
})
const after = yield* snapshot.capture()
expect(after).toBeDefined()
if (!before || !after) return
expect(yield* snapshot.files({ from: before, to: after })).toEqual([RelativePath.make("tracked.txt")])
}).pipe(Effect.provide(snapshotLayer(tmp.path, project)))
}),
(tmp) => Effect.promise(() => tmp[Symbol.asyncDispose]()),
),
)
testEffect(Layer.empty).live("captures concurrently from independent processes", () =>
Effect.acquireUseRelease(
Effect.promise(() => tmpdir()),
(tmp) =>
Effect.gen(function* () {
const project = path.join(tmp.path, "project")
yield* Effect.promise(async () => {
await fs.mkdir(project)
await Promise.all(
Array.from({ length: 200 }, (_, index) =>
fs.writeFile(path.join(project, `f${index}.txt`), `${index}\n`),
),
)
await initGit(project, true)
})
// Each layer owns its own Git service, so their in-process locks do not coordinate.
const writer = (id: number) =>
Effect.gen(function* () {
const snapshot = yield* Snapshot.Service
return yield* Effect.forEach(
Array.from({ length: 8 }, (_, index) => index),
(index) =>
Effect.promise(() => fs.writeFile(path.join(project, `writer-${id}.txt`), `${index}\n`)).pipe(
Effect.andThen(snapshot.capture()),
),
)
}).pipe(Effect.provide(snapshotLayer(tmp.path, project)))
const results = yield* Effect.all([writer(0), writer(1), writer(2)], { concurrency: "unbounded" })
expect(results.flat().every((tree) => tree !== undefined)).toBe(true)
}),
(tmp) => Effect.promise(() => tmp[Symbol.asyncDispose]()),
),
)
testEffect(Layer.empty).live("captures a Location in a directory whose name starts with two dots", () =>
Effect.acquireUseRelease(
Effect.promise(() => tmpdir()),
(tmp) =>
Effect.gen(function* () {
const project = path.join(tmp.path, "project")
const location = path.join(project, "..scope")
yield* Effect.promise(async () => {
await fs.mkdir(location, { recursive: true })
await fs.writeFile(path.join(location, "tracked.txt"), "one\n")
await initGit(project)
})
yield* Effect.gen(function* () {
const snapshot = yield* Snapshot.Service
const before = yield* snapshot.capture()
yield* Effect.promise(() => fs.writeFile(path.join(location, "tracked.txt"), "two\n"))
const after = yield* snapshot.capture()
expect(before).toBeDefined()
expect(after).toBeDefined()
if (!before || !after) return
expect(yield* snapshot.files({ from: before, to: after })).toEqual([
RelativePath.make("..scope/tracked.txt"),
])
}).pipe(Effect.provide(snapshotLayer(tmp.path, location)))
}),
(tmp) => Effect.promise(() => tmp[Symbol.asyncDispose]()),
),
)
testEffect(Layer.empty).live("restores many files from several trees and removes paths absent from them", () =>
Effect.acquireUseRelease(
Effect.promise(() => tmpdir()),
(tmp) =>
Effect.gen(function* () {
const project = path.join(tmp.path, "project")
yield* Effect.promise(async () => {
await fs.mkdir(project)
await fs.writeFile(path.join(project, "a.txt"), "a1\n")
await fs.writeFile(path.join(project, "b[1].txt"), "b1\n")
await initGit(project, true)
})
yield* Effect.gen(function* () {
const snapshot = yield* Snapshot.Service
const first = yield* snapshot.capture()
yield* Effect.promise(async () => {
await fs.writeFile(path.join(project, "a.txt"), "a2\n")
await fs.writeFile(path.join(project, "c.txt"), "c2\n")
})
const second = yield* snapshot.capture()
yield* Effect.promise(async () => {
await fs.writeFile(path.join(project, "a.txt"), "a3\n")
await fs.writeFile(path.join(project, "b[1].txt"), "b3\n")
await fs.writeFile(path.join(project, "c.txt"), "c3\n")
await fs.writeFile(path.join(project, "d.txt"), "d3\n")
})
if (!first || !second) throw new globalThis.Error("capture failed")
yield* snapshot.restore({
files: new Map([
[RelativePath.make("a.txt"), second],
[RelativePath.make("b[1].txt"), first],
[RelativePath.make("c.txt"), first],
[RelativePath.make("d.txt"), second],
]),
})
expect(yield* read(path.join(project, "a.txt"))).toBe("a2\n")
expect(yield* read(path.join(project, "b[1].txt"))).toBe("b1\n")
expect(yield* exists(path.join(project, "c.txt"))).toBe(false)
expect(yield* exists(path.join(project, "d.txt"))).toBe(false)
}).pipe(Effect.provide(snapshotLayer(tmp.path, project)))
}),
(tmp) => Effect.promise(() => tmp[Symbol.asyncDispose]()),
),
)
testEffect(Layer.empty).live("restores the other files when removing one path fails", () =>
Effect.acquireUseRelease(
Effect.promise(() => tmpdir()),
(tmp) =>
Effect.gen(function* () {
const project = path.join(tmp.path, "project")
yield* Effect.promise(async () => {
await fs.mkdir(project)
await fs.writeFile(path.join(project, "c.txt"), "old\n")
await initGit(project, true)
})
yield* Effect.gen(function* () {
const snapshot = yield* Snapshot.Service
const first = yield* snapshot.capture()
if (!first) throw new globalThis.Error("capture failed")
// Removing `a/b` fails on POSIX because `a` is now a file.
yield* Effect.promise(async () => {
await fs.writeFile(path.join(project, "c.txt"), "changed\n")
await fs.writeFile(path.join(project, "a"), "file\n")
})
yield* snapshot
.restore({
files: new Map([
[RelativePath.make("a/b"), first],
[RelativePath.make("c.txt"), first],
]),
})
.pipe(Effect.exit)
expect(yield* read(path.join(project, "c.txt"))).toBe("old\n")
}).pipe(Effect.provide(snapshotLayer(tmp.path, project)))
}),
(tmp) => Effect.promise(() => tmp[Symbol.asyncDispose]()),
),
)
testEffect(Layer.empty).live("applies availability transforms", () =>
Effect.acquireUseRelease(
Effect.promise(() => tmpdir()),
@@ -219,6 +393,23 @@ function snapshotLayer(data: string, directory: string) {
])
}
function snapshotDirectory(data: string) {
return Effect.promise(async () => {
const projects = await fs.readdir(path.join(data, "snapshot"))
const project = path.join(data, "snapshot", projects[0]!)
return path.join(project, (await fs.readdir(project))[0]!)
})
}
function exists(file: string) {
return Effect.promise(() =>
fs.stat(file).then(
() => true,
() => false,
),
)
}
function read(file: string) {
return Effect.promise(() => fs.readFile(file, "utf8")).pipe(Effect.map((content) => content.replaceAll("\r\n", "\n")))
}