Compare commits

...
13 changed files with 2428 additions and 120 deletions

No files matched your search

@@ -0,0 +1,64 @@
/**
* Copy existing snapshot stores and run the real compaction 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-compact.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-compact")
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.compact(
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 compaction.
if (!("objects" in git)) return
yield* git.objects.compact(
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)} compaction ${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
+706
View File
@@ -0,0 +1,706 @@
/**
* 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))
+678 -104
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 privateIndexPrefix = "index.opencode-"
const excludeMarker = "# opencode: mirrored from the source repository; edits below are replaced\n"
// Like `git gc --auto`, pack once loose objects accumulate, and merge packs before lookups slow down.
const compactLooseLimit = 2048
const compactPackLimit = 16
export interface CaptureInput {
readonly repository: Repository
readonly scopes: readonly RelativePath[]
/** Source repository whose ignore rules decide which paths are recorded. */
readonly ignores?: Repository
/** Repository whose index rebuilds a corrupt snapshot index without rehashing tracked files. */
readonly seed?: 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",
"compact",
]),
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 merge 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 compact: (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,148 @@ const layer = Layer.effect(
})
})
/**
* A new store is initialized in a private staging 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* initialize({ ...input, directory: input.gitDirectory })
return repository
}
const staging = AbsolutePath.make(
`${input.gitDirectory}.init-${process.pid}-${Math.random().toString(36).slice(2)}`,
)
yield* Effect.acquireUseRelease(
Effect.succeed(staging),
(directory) =>
Effect.gen(function* () {
yield* initialize({ ...input, directory })
const renamed = yield* fs.rename(directory, 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),
})
}),
(directory) => fs.remove(directory, { recursive: true, force: true }).pipe(Effect.ignore),
)
return repository
})
const initialize = 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.gitDirectory, snapshotConfigFile), snapshotConfig)
const config = path.join(input.gitDirectory, "config")
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
})
/**
* Report modified, deleted, and non-ignored untracked paths against the index.
* Both commands only read the index, so they run against the shared index
* without taking `index.lock`. Two parallel processes beat one combined
* `ls-files -m -o`, whose lstat pass is not threaded like diff-files'.
*/
const scan = Effect.fnUntraced(function* (repository: Repository, scope: RelativePath) {
const list = (args: string[]) =>
repositoryOperation("refresh", repository, args).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", "--", literal(scope)]),
list(["ls-files", "--others", "--exclude-standard", "-z", "--", literal(scope)]),
],
{ concurrency: 2 },
)
return { tracked, untracked }
})
const stage = Effect.fnUntraced(function* (input: {
repository: Repository
changes: { tracked: readonly RelativePath[]; untracked: readonly RelativePath[] }
ignores?: Repository
maximumUntrackedFileBytes?: number
index?: string
}) {
const candidates = [...input.changes.tracked, ...input.changes.untracked]
if (!candidates.length) return { skipped: [] }
const 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 staged = candidates.filter((item) => !excluded.has(item) && !skip.has(item))
const removed = [...excluded, ...skipped]
// 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,
})
if (staged.length)
yield* repositoryOperation(
"refresh",
input.repository,
["update-index", "--add", "--remove", "--replace", "-z", "--stdin"],
{ stdin: staged.join("\0") + "\0", index: input.index },
)
return { skipped }
})
const refresh = Effect.fn("Git.index.refresh")(function* (input: {
@@ -387,52 +514,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* stage({ ...input, changes: yield* scan(input.repository, input.scope) })
})
const ignored = Effect.fn("Git.index.ignored")(function* (input: {
@@ -472,26 +554,240 @@ 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 sharedIndex = (repository: Repository) => path.join(repository.gitDirectory, "index")
/** Identity of an index file; Git only replaces an index by renaming a new file into place. */
const stamp = (file: string) =>
fs.stat(file).pipe(
Effect.map((info) =>
[Option.getOrUndefined(info.ino), info.size, Option.getOrUndefined(info.mtime)?.getTime()].join(":"),
),
Effect.orElseSucceed(() => undefined),
)
/**
* Run index writes against a private index, then publish it with an atomic
* rename. Processes never contend on `index.lock`, an interrupted or killed
* writer cannot leave a partial shared index, and a stale lock is irrelevant.
* The private index starts as a hard link: Git never writes an index in place,
* it writes a lock file and renames it over the private name, so the shared
* file is never modified and no bytes are copied. The shared index is only a
* stat cache over the object store; when writers race, the last publish wins
* and the next capture reconciles against the worktree.
*/
const privateIndex = <A, E, R>(repository: Repository, use: (index: string) => Effect.Effect<A, E, R>) =>
Effect.acquireUseRelease(
Effect.gen(function* () {
const shared = sharedIndex(repository)
const index = path.join(
repository.gitDirectory,
`${privateIndexPrefix}${process.pid}-${Math.random().toString(36).slice(2)}`,
)
if (!(yield* fs.existsSafe(shared))) return index
const linked = yield* fs.link(shared, index).pipe(
Effect.as(true),
Effect.orElseSucceed(() => false),
)
if (linked) return index
const info = yield* fs.stat(shared).pipe(Effect.option)
yield* fs.copyFile(shared, index).pipe(
Effect.mapError(
(cause) =>
new OperationError({
operation: "refresh",
directory: repository.gitDirectory,
message: "Failed to prepare a private index",
cause,
}),
),
)
// A copy must keep the timestamp Git uses to detect racily clean entries.
const mtime = Option.isSome(info) ? Option.getOrUndefined(info.value.mtime) : undefined
if (mtime && Option.isSome(info))
yield* fs
.utimes(
index,
Option.getOrElse(info.value.atime, () => mtime),
mtime,
)
.pipe(Effect.ignore)
return index
}),
(index) =>
Effect.gen(function* () {
const value = yield* use(index)
const published = yield* Effect.uninterruptible(
Effect.gen(function* () {
// Stamp the private file before it takes the shared name, so a concurrent publish can never
// pair another writer's index with this tree.
const identity = yield* stamp(index)
const renamed = yield* fs.rename(index, sharedIndex(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 ? identity : undefined
}),
)
return { value, published }
}),
(index) => fs.remove(index, { force: true }).pipe(Effect.ignore),
)
// Clean captures return the tree last written from the exact shared index they scanned.
const trees = new Map<string, { readonly stamp: string; readonly tree: TreeID }>()
const swept = new Set<string>()
const captureOnce = Effect.fnUntraced(function* (input: CaptureInput) {
yield* syncExcludes(input.repository, input.ignores)
const changes = yield* Effect.forEach(input.scopes, (scope) => scan(input.repository, scope), {
concurrency: "unbounded",
})
const current = yield* stamp(sharedIndex(input.repository))
const cached = trees.get(input.repository.gitDirectory)
if (
current &&
cached?.stamp === current &&
changes.every((change) => !change.tracked.length && !change.untracked.length)
)
return cached.tree
const result = yield* privateIndex(input.repository, (index) =>
Effect.gen(function* () {
yield* Effect.forEach(changes, (change) => stage({ ...input, changes: change, index }), { discard: true })
return yield* writeTree(input.repository, index)
}),
)
if (result.published) trees.set(input.repository.gitDirectory, { stamp: result.published, tree: result.value })
if (!result.published) trees.delete(input.repository.gitDirectory)
return result.value
})
const captureTree = Effect.fn("Git.tree.capture")((input: CaptureInput) =>
locked(
input.repository,
Effect.gen(function* () {
if (!swept.has(input.repository.gitDirectory)) {
swept.add(input.repository.gitDirectory)
yield* sweepPrivateIndexes(input.repository)
yield* upgradeConfig(input.repository)
}
return yield* captureOnce(input).pipe(
Effect.catch((error) =>
Effect.gen(function* () {
if (!(yield* indexCorrupt(input.repository))) return yield* error
yield* Effect.logWarning("rebuilding corrupt snapshot index", {
directory: input.repository.gitDirectory,
message: error.message,
})
yield* reseedIndex(input.repository, input.seed)
return yield* captureOnce(input)
}),
),
)
}),
),
)
/**
* Mirror the source repository's info/exclude into a marked block of the
* store's own file, so one ls-files pass applies both, exactly like listing
* with the store's rules and then filtering with the source's. Content above
* the marker, such as patterns from an init template, is kept.
*/
const syncExcludes = Effect.fnUntraced(function* (repository: Repository, source?: Repository) {
if (!source || source.gitDirectory === repository.gitDirectory) return
const target = path.join(repository.commonDirectory, "info", "exclude")
const [wanted, current] = yield* Effect.all(
[fs.readFileStringSafe(path.join(source.commonDirectory, "info", "exclude")), fs.readFileStringSafe(target)],
{ concurrency: 2 },
).pipe(Effect.orElseSucceed(() => [undefined, undefined] as const))
const own = (current ?? "").split(excludeMarker)[0]!
const next = wanted ? `${own}${own && !own.endsWith("\n") ? "\n" : ""}${excludeMarker}${wanted}` : own
if (next === (current ?? "")) return
// Concurrent Git processes must never read a partially written file.
const staged = `${target}.${process.pid}-${Math.random().toString(36).slice(2)}`
yield* fs
.writeWithDirs(staged, next)
.pipe(
Effect.andThen(fs.rename(staged, target)),
Effect.ignore,
Effect.ensuring(fs.remove(staged, { force: true }).pipe(Effect.ignore)),
)
})
/**
* 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, surface unchanged.
*/
const indexCorrupt = Effect.fnUntraced(function* (repository: Repository) {
if (!(yield* fs.existsSafe(sharedIndex(repository)))) return false
return yield* repositoryOperation("refresh", repository, [
"ls-files",
"-z",
"--",
literal(".opencode-probe"),
]).pipe(
Effect.as(false),
Effect.orElseSucceed(() => true),
)
})
const reseedIndex = Effect.fnUntraced(function* (repository: Repository, seed?: Repository) {
yield* fs.remove(sharedIndex(repository), { force: true }).pipe(Effect.ignore)
trees.delete(repository.gitDirectory)
if (!seed) return
// Seeding keeps the source's stat cache and cache-tree, so recovery does not rehash every tracked file.
yield* privateIndex(repository, (index) =>
fs.copyFile(path.join(seed.gitDirectory, "index"), index).pipe(Effect.ignore),
)
})
// Stores created by earlier releases keep the settings they were created with; only OpenCode's own file is rewritten.
const upgradeConfig = 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 staged = `${file}.${process.pid}-${Math.random().toString(36).slice(2)}`
yield* fs
.writeFileString(staged, snapshotConfig)
.pipe(
Effect.andThen(fs.rename(staged, file)),
Effect.ignore,
Effect.ensuring(fs.remove(staged, { force: true }).pipe(Effect.ignore)),
)
})
// Private indexes left by killed processes; live captures finish long before the cutoff.
const sweepPrivateIndexes = 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(privateIndexPrefix)),
(entry) =>
fs.stat(path.join(repository.gitDirectory, entry)).pipe(
Effect.flatMap((info) =>
(Option.getOrUndefined(info.mtime)?.getTime() ?? 0) < cutoff
? fs.remove(path.join(repository.gitDirectory, entry), { force: true })
: Effect.void,
),
Effect.ignore,
),
{ discard: true },
)
})
const treeFiles = Effect.fn("Git.tree.files")(function* (input: {
repository: Repository
from: TreeID
@@ -523,7 +819,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(literal)]
// Patch headers have no -z form: unquoted paths keep chunksByFile matching non-ASCII names.
const [names, numbers, patch] = yield* Effect.all(
[
@@ -581,7 +877,7 @@ const layer = Layer.effect(
"-z",
tree,
"--",
file,
literal(file),
])).text.replace(/\0$/, "")
if (!text) return false
if (!/^\d+\s+\w+\s+[0-9a-f]+\t/.test(text))
@@ -593,35 +889,288 @@ 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,
}),
),
)
/** Paths that have an entry in `tree`, listed in chunks to stay below argument limits. */
const entries = 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(literal)]).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())
})
/**
* Restores path by path in map order, like the original implementation, unless
* every path is independent. Operations on independent paths commute, so they
* are batched into 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,
Effect.gen(function* () {
if (!input.files.size) return
if (!independent([...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* privateIndex(input.repository, (index) =>
repositoryOperation("restore", input.repository, ["checkout", tree, "--", literal(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]) =>
entries(input.repository, tree, files).pipe(
Effect.map((present) => ({
tree,
present: files.filter((file) => present.has(file)),
absent: files.filter((file) => !present.has(file)),
})),
),
)
yield* Effect.forEach(
plan.flatMap((item) => item.absent),
(file) => removePath(input.repository, file),
{ concurrency: 16, discard: true },
)
const checkouts = plan.filter((item) => item.present.length)
if (!checkouts.length) return
yield* privateIndex(input.repository, (index) =>
Effect.forEach(
checkouts,
(item) =>
repositoryOperation(
"restore",
input.repository,
["checkout", item.tree, "--pathspec-from-file=-", "--pathspec-file-nul"],
{ stdin: item.present.map(literal).join("\0") + "\0", index },
),
)
}),
{ discard: true },
),
{ discard: true },
),
)
}),
),
)
/**
* Snapshot trees have no refs, so Git's own repack and prune would treat every
* snapshot as garbage. Compaction instead hands pack-objects an explicit list of
* every loose object (and, when merging, every object in the old packs), checks
* the new index lists all of them, and only then deletes the loose copies and
* the merged packs. Nothing is pruned, and objects are never delegated to the
* source repository through alternates.
*/
const compact = Effect.fn("Git.objects.compact")(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 merge = packs.length >= compactPackLimit
if (loose.length < compactLooseLimit && !merge) return
yield* compactionLock(
repository,
Effect.gen(function* () {
yield* sweepTemporaryPacks(packDirectory)
const merged = merge
? yield* Effect.forEach(packs, (pack) =>
packIndex(repository, path.join(packDirectory, `${pack}.idx`)).pipe(
Effect.map((oids) => ({ pack, oids })),
),
)
: []
const wanted = [...new Set([...loose.map((item) => item.oid), ...merged.flatMap((item) => item.oids)])]
if (!wanted.length) return
const written = (yield* repositoryOperation(
"compact",
repository,
[
"-c",
"pack.threads=2",
"-c",
"pack.windowMemory=64m",
"pack-objects",
"-q",
"--non-empty",
path.join(packDirectory, "pack"),
],
{ stdin: wanted.join("\n") + "\n" },
)).text
.split("\n")
.map((line) => line.trim())
.filter(Boolean)
const packed = new Set(
(yield* Effect.forEach(written, (name) =>
packIndex(repository, path.join(packDirectory, `pack-${name}.idx`)),
)).flat(),
)
const missing = wanted.filter((oid) => !packed.has(oid))
if (missing.length)
return yield* new OperationError({
operation: "compact",
directory: repository.gitDirectory,
message: `Packed ${packed.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 kept = new Set(written.map((name) => `pack-${name}`))
yield* Effect.forEach(
merged.filter((item) => !kept.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: "compact",
directory: repository.gitDirectory,
message: `Failed to read ${file}`,
cause,
}),
),
)
const result = yield* repositoryOperation("compact", 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 compactions would stay lossless but duplicate work and objects.
const compactionLock = <A, E, R>(repository: Repository, effect: Effect.Effect<A, E, R>) => {
const file = path.join(repository.gitDirectory, "opencode-compact.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: { compact },
})
}),
)
@@ -749,6 +1299,30 @@ 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 literal(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 independent(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 resolvePath(cwd: string, value: string) {
const trimmed = value.replace(/[\r\n]+$/, "")
if (!trimmed) return cwd
@@ -19,7 +19,8 @@ type Input = {
readonly agent: Agent.ID
readonly model: Model.Ref
readonly providerMetadataKey: string
readonly snapshot?: Snapshot.ID
/** Awaited before `Step.Started`, so the capture can overlap the provider request. */
readonly snapshot?: Effect.Effect<Snapshot.ID | undefined>
readonly started: number
readonly assistantMessageID: SessionMessage.ID
}
@@ -97,6 +98,9 @@ export const createLLMEventPublisher = (bus: Pick<Bus.Interface, "publish">, inp
let stepSettlement: StepRecord["finish"]
const startAssistant = Effect.fnUntraced(function* () {
if (stepStarted) return assistantMessageID
// Await before claiming the start: an interruption while waiting must leave the step unstarted.
const snapshot = input.snapshot ? yield* input.snapshot : undefined
if (stepStarted) return assistantMessageID
stepStarted = true
yield* bus.publish(SessionEvent.Step.Started, {
@@ -104,7 +108,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
+5 -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 preimage only has to precede local tool execution, which cannot begin before Step.Started,
// so the capture overlaps the provider's time to first event instead of delaying the request.
const startCapture = 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,
snapshot: Fiber.join(startCapture),
started: yield* Clock.currentTimeMillis,
})
const toolRuns: Array<{
@@ -224,6 +226,7 @@ export const make = Effect.gen(function* () {
const record = publisher.record()
if (record.finish || record.failure) {
const startSnapshot = yield* Fiber.join(startCapture)
const snapshot = yield* snapshots.capture()
const files =
startSnapshot && snapshot
+30 -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,44 @@ 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))
// A directory named like `..scope` is inside the project; only a `..` segment escapes it.
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
// Background and rate-limited: a capture never waits for compaction, and the check itself is a few directory reads.
let maintenance = { checked: Number.NEGATIVE_INFINITY, running: false }
const maintain = Effect.fnUntraced(function* (repository: Git.Repository) {
const now = yield* Clock.currentTimeMillis
if (maintenance.running || now - maintenance.checked < 10 * 60 * 1000) return
maintenance = { checked: now, running: true }
yield* git.objects.compact(repository).pipe(
Effect.catch((cause) => Effect.logWarning("failed to compact snapshot objects", { cause })),
Effect.ensuring(
Effect.sync(() => {
maintenance = { ...maintenance, running: false }
}),
),
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,
seed: repo.source,
maximumUntrackedFileBytes: 2 * 1024 * 1024,
})
yield* maintain(repo.snapshotRepository)
return ID.make(tree)
}).pipe(
Effect.catch((cause) => Effect.logWarning("failed to capture snapshot", { cause }).pipe(Effect.as(undefined))),
)
+134
View File
@@ -296,3 +296,137 @@ describe("Git trees", () => {
}),
)
})
describe("Git objects", () => {
it.live(
"compacts loose objects and merges 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, then several small batches as separate packs.
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(async () =>
(await $`find ${objects} -path '*/objects/??/*' -type f`.text()).split("\n").filter(Boolean),
)
const before = yield* everything()
expect((yield* loose()).length).toBeGreaterThan(2048)
yield* git.objects.compact(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 oids = await $`find ${objects} -path '*/objects/??/*' -type f`.text()
const list = oids
.split("\n")
.filter(Boolean)
.map((file) => file.split("/").slice(-2).join(""))
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.compact(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",
)
}),
)
})
+156
View File
@@ -127,6 +127,145 @@ 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("applies availability transforms", () =>
Effect.acquireUseRelease(
Effect.promise(() => tmpdir()),
@@ -219,6 +358,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")))
}