mirror of
https://github.com/anomalyco/opencode.git
synced 2026-10-03 05:56:18 +00:00
Compare commits
2
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
ef062be96c | ||
|
|
2e3b4408cf |
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)
|
||||
@@ -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")
|
||||
Executable
+20
@@ -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
|
||||
@@ -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
@@ -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
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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))),
|
||||
)
|
||||
|
||||
@@ -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",
|
||||
)
|
||||
}),
|
||||
)
|
||||
})
|
||||
@@ -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")))
|
||||
}
|
||||
|
||||
Reference in new issue
Block a user