Repository navigation
feat(tools): guarded write CAS core with per-path FIFO chain (S4a, #1375) #1405
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
base: main
Are you sure you want to change the base?
Changes from all commits
a37dd24
3dd8700
ba1332e
6813013
588b95f
7a25fc0
56ce4bf
be894d9
58bb5ce
bb93d62
4d4d511
1ce0421
968078d
a2c364d
9f8a857
cb72ac7
8cb5920
2810c8f
006389d
2dd792c
51f0cea
8e5061e
bacb9d3
fc94f62
171a26e
17c736e
f1b2ae4
17c53e0
509208e
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|
|
|
@@ -976,6 +976,71 @@ describe("Cline", () => { | |||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| }) | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| }) | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
|
|
||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| describe("observation registry lifecycle (S4a, epic #1375)", () => { | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| it("gives each Task its own observation registry", () => { | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| const firstTask = new Task({ | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| provider: mockProvider, | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| apiConfiguration: mockApiConfig, | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| task: "first observation task", | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| startTask: false, | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| }) | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| const secondTask = new Task({ | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| provider: mockProvider, | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| apiConfiguration: mockApiConfig, | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| task: "second observation task", | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| startTask: false, | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| }) | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
|
|
||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| // The guarded-write contract assumes an observation in one task never validates | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| // a write issued by another task. | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| expect(firstTask.observationRegistry).not.toBe(secondTask.observationRegistry) | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| firstTask.observationRegistry.observe("/workspace/a.ts", "v-a") | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| expect(firstTask.observationRegistry.get("/workspace/a.ts")?.version).toBe("v-a") | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| expect(secondTask.observationRegistry.get("/workspace/a.ts")).toBeUndefined() | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| }) | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
|
|
||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| it("clears the observation registry when the task is disposed", async () => { | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| const task = new Task({ | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| provider: mockProvider, | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| apiConfiguration: mockApiConfig, | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| task: "disposed observation task", | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| startTask: false, | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| }) | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| task.observationRegistry.observe("/workspace/a.ts", "v-a") | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| task.observationRegistry.observe("/workspace/b.ts", "v-b") | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
|
|
||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| await task.dispose() | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
|
|
||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| // A disposed task cannot serve another guarded write, so its observed paths | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| // (version token + timestamp each) must not stay reachable for the host lifetime. | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| expect(task.observationRegistry.get("/workspace/a.ts")).toBeUndefined() | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| // A disposed task cannot serve another guarded write, so its observed paths | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| // (version token + timestamp each) must not stay reachable for the host lifetime. | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| expect(task.observationRegistry.get("/workspace/b.ts")).toBeUndefined() | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| // disposeOnce() closes the registry rather than only clearing the map. | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| }) | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
|
Comment on lines
+1002
to
+1021
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. 📐 Maintainability & Code Quality | 🔵 Trivial | 💤 Low value Remove the duplicated comment and assert the closed state. Lines 1014-1015 and 1017-1018 contain the same comment twice. Line 1020 says that Proposed fix expect(task.observationRegistry.get("/workspace/a.ts")).toBeUndefined()
- // A disposed task cannot serve another guarded write, so its observed paths
- // (version token + timestamp each) must not stay reachable for the host lifetime.
expect(task.observationRegistry.get("/workspace/b.ts")).toBeUndefined()
// disposeOnce() closes the registry rather than only clearing the map.
+ expect(task.observationRegistry.isClosed).toBe(true)📝 Committable suggestion
Suggested change
🤖 Prompt for AI AgentsSource: Path instructions |
||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
|
|
||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| it("refuses an observation recorded after the task was disposed", async () => { | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| const task = new Task({ | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| provider: mockProvider, | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| apiConfiguration: mockApiConfig, | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| task: "late observation task", | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| startTask: false, | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| }) | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
|
|
||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| await task.dispose() | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
|
|
||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| // A read that was already in flight when the task was disposed can finish late and | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| // call observe(). Recording then would repopulate a registry no task owns and hand a | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| // version token to a guarded write that will never happen. | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| task.observationRegistry.observe("/workspace/late.ts", "v-late") | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
|
|
||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| expect(task.observationRegistry.get("/workspace/late.ts")).toBeUndefined() | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| expect(task.observationRegistry.size).toBe(0) | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| expect(task.observationRegistry.isClosed).toBe(true) | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| }) | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| }) | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
|
|
||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| describe("constructor", () => { | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| it.each([{ apiConfigName: "parent-local-profile" }, { apiConfigName: undefined }])( | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| "uses an explicit delegated-child context without shared state or startup persistence", | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
|
|
||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,91 @@ | ||
| import { describe, it, expect, vi } from "vitest" | ||
|
|
||
| import { ObservationRegistry } from "../observationRegistry" | ||
|
|
||
| describe("ObservationRegistry", () => { | ||
| it("observe → get returns the recorded version and observedAt", () => { | ||
| const reg = new ObservationRegistry() | ||
| reg.observe("/a/b/c.ts", "1:2:300:4000000000:5000000000") | ||
|
|
||
| const obs = reg.get("/a/b/c.ts") | ||
| expect(obs).toBeDefined() | ||
| expect(obs!.version).toBe("1:2:300:4000000000:5000000000") | ||
| expect(typeof obs!.observedAt).toBe("number") | ||
| }) | ||
|
|
||
| it("re-observe replaces the entry with a fresh observedAt", () => { | ||
| vi.useFakeTimers() | ||
| const reg = new ObservationRegistry() | ||
| reg.observe("/a/b/c.ts", "v1") | ||
| const first = reg.get("/a/b/c.ts")! | ||
| expect(first.version).toBe("v1") | ||
|
|
||
| vi.advanceTimersByTime(50) | ||
| reg.observe("/a/b/c.ts", "v2") | ||
| const second = reg.get("/a/b/c.ts")! | ||
| expect(second.version).toBe("v2") | ||
| expect(second.observedAt).toBeGreaterThan(first.observedAt) | ||
|
|
||
| vi.useRealTimers() | ||
| }) | ||
|
|
||
| it("has returns true for observed paths, false otherwise", () => { | ||
| const reg = new ObservationRegistry() | ||
| reg.observe("/x.ts", "t1") | ||
| expect(reg.has("/x.ts")).toBe(true) | ||
| expect(reg.has("/y.ts")).toBe(false) | ||
| }) | ||
|
|
||
| it("size reflects the number of observed entries", () => { | ||
| const reg = new ObservationRegistry() | ||
| expect(reg.size).toBe(0) | ||
| reg.observe("/a.ts", "t1") | ||
| reg.observe("/b.ts", "t2") | ||
| expect(reg.size).toBe(2) | ||
| }) | ||
|
|
||
| it("clear removes all entries and resets size to 0", () => { | ||
| const reg = new ObservationRegistry() | ||
| reg.observe("/a.ts", "t1") | ||
| reg.observe("/b.ts", "t2") | ||
| reg.clear() | ||
| expect(reg.size).toBe(0) | ||
| expect(reg.get("/a.ts")).toBeUndefined() | ||
| expect(reg.has("/b.ts")).toBe(false) | ||
| }) | ||
|
|
||
| it("get on empty registry returns undefined", () => { | ||
| const reg = new ObservationRegistry() | ||
| expect(reg.get("/any.ts")).toBeUndefined() | ||
| }) | ||
|
|
||
| it("separate instances are independent — observing in one does not appear in the other", () => { | ||
| const regA = new ObservationRegistry() | ||
| const regB = new ObservationRegistry() | ||
| regA.observe("/shared.ts", "v1") | ||
| expect(regA.get("/shared.ts")).toBeDefined() | ||
| expect(regB.get("/shared.ts")).toBeUndefined() | ||
| regB.observe("/shared.ts", "v2") | ||
| expect(regA.get("/shared.ts")!.version).toBe("v1") | ||
| expect(regB.get("/shared.ts")!.version).toBe("v2") | ||
| }) | ||
| }) | ||
|
|
||
|
|
||
| describe("close() - disposal is terminal", () => { | ||
| it("drops every observation and refuses later ones", () => { | ||
| const registry = new ObservationRegistry() | ||
| registry.observe("/workspace/a.ts", "v-a") | ||
| expect(registry.size).toBe(1) | ||
|
|
||
| registry.close() | ||
|
|
||
| expect(registry.size).toBe(0) | ||
| expect(registry.isClosed).toBe(true) | ||
|
|
||
| // A read that was in flight when the task was disposed must not repopulate it. | ||
| registry.observe("/workspace/late.ts", "v-late") | ||
| expect(registry.get("/workspace/late.ts")).toBeUndefined() | ||
| expect(registry.size).toBe(0) | ||
| }) | ||
| }) |
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,75 @@ | ||
| /** | ||
| * Per-task file observation registry (upstream epic #1375, phase A2). | ||
| * | ||
| * Each Task owns its own instance so parent and subtask observations are | ||
| * independent. The S4 guarded-write will compare these versions against the | ||
| * token recomputed pre-write to detect stale reads or file replacement. | ||
| * | ||
| * Pure in-memory — zero I/O, no dependencies. The observations ARE consulted: | ||
| * guardedWrite reads this registry before publishing (src/core/tools/guardedWrite.ts) | ||
| * and compares the recorded version token against the token recomputed from disk, so a | ||
| * stale read or an out-of-band replacement that the check detects is rejected instead of | ||
| * published over. Detection is best effort against a non-cooperating process: the token is | ||
| * recomputed before the publish, so a replacement that lands after that check and before | ||
| * the rename is not observable from here and can still win. Closing that last window needs | ||
| * a cross-process lock or an atomic create, not a token comparison. | ||
| */ | ||
|
|
||
| export interface FileObservation { | ||
| /** Version token derived from on-disk fs.stat (bigint mode). */ | ||
| version: string | ||
| /** Millisecond timestamp when the observation was recorded. */ | ||
| observedAt: number | ||
| } | ||
|
|
||
| export class ObservationRegistry { | ||
| private readonly entries = new Map<string, FileObservation>() | ||
|
|
||
| /** Set by close(): after disposal the registry refuses further observations. */ | ||
| private closed = false | ||
|
|
||
| /** | ||
| * Record an observation for a file at its absolute path, unless the registry is closed. | ||
| * | ||
| * Re-observing replaces the entry with a fresh observedAt timestamp and | ||
| * the new version token. A read that was already in flight can finish after | ||
| * Task.disposeOnce() dropped the observations; recording then would hand a version token | ||
| * to a task that no longer serves any request, and a later guarded write could consult | ||
| * it. close() therefore makes this a no-op, so disposal is terminal at this layer. | ||
| */ | ||
| observe(absolutePath: string, version: string): void { | ||
| if (this.closed) { | ||
| return | ||
| } | ||
| this.entries.set(absolutePath, { version, observedAt: Date.now() }) | ||
| } | ||
|
|
||
| get(absolutePath: string): FileObservation | undefined { | ||
| return this.entries.get(absolutePath) | ||
| } | ||
|
|
||
| has(absolutePath: string): boolean { | ||
| return this.entries.has(absolutePath) | ||
| } | ||
|
|
||
| clear(): void { | ||
| this.entries.clear() | ||
| } | ||
|
|
||
| /** | ||
| * Drop every observation and refuse any later one. Task.disposeOnce() calls this so a | ||
| * disposed task's registry cannot be repopulated by a read that finishes late. | ||
| */ | ||
| close(): void { | ||
| this.closed = true | ||
| this.entries.clear() | ||
| } | ||
|
|
||
| get isClosed(): boolean { | ||
| return this.closed | ||
| } | ||
|
|
||
| get size(): number { | ||
| return this.entries.size | ||
| } | ||
| } | ||
|
coderabbitai[bot] marked this conversation as resolved.
|
||
Uh oh!
There was an error while loading. Please reload this page.