Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
30 changes: 24 additions & 6 deletions .agents/skills/verify-mcode/scripts/provider-completeness.mjs
Original file line number Diff line number Diff line change
Expand Up @@ -348,31 +348,47 @@ export async function runWorkspaceInvalidationJourney({ repoRoot, workspace, run
await observer.rpc("file.refresh", { workspaceId });
recordOwnedFile(run.run, ownerFile);
recordOwnedFile(run.run, observerFile);
const ownerPath = expectedChangePath(ownerFile);
const observerPath = expectedChangePath(observerFile);
await io.writeFile(ownerFile, "WATCHER_OWNER_MARKER\n", "utf8");
await owner.rpc("file.refresh", { workspaceId });
await Promise.all([
waitForExactWorkspaceInvalidation(ownerEvents, workspaceId, NodePath.basename(ownerFile), timeoutMs),
waitForExactWorkspaceInvalidation(observerEvents, workspaceId, NodePath.basename(ownerFile), timeoutMs),
waitForExactWorkspaceInvalidation(ownerEvents, workspaceId, ownerPath, timeoutMs),
waitForExactWorkspaceInvalidation(observerEvents, workspaceId, ownerPath, timeoutMs),
]);
await owner.close();
const ownerEventCount = ownerEvents.length;
const observerStart = observerEvents.length;
await io.appendFile(observerFile, "WATCHER_OBSERVER_MARKER\n", "utf8");
await observer.rpc("file.refresh", { workspaceId });
await waitForExactWorkspaceInvalidation(observerEvents, workspaceId, NodePath.basename(observerFile), timeoutMs, observerStart);
await waitForExactWorkspaceInvalidation(observerEvents, workspaceId, observerPath, timeoutMs, observerStart);
if (ownerEvents.length !== ownerEventCount) throw new Error("Condition: disconnected client received a later files.changed push.");
return {
kind: "live-rpc-proof",
control: "public file.refresh RPC and files.changed push",
workspaceId,
owner: { closed: true, changes: [NodePath.basename(ownerFile)] },
observer: { active: true, changes: [NodePath.basename(ownerFile), NodePath.basename(observerFile)] },
owner: { closed: true, changes: [ownerPath] },
observer: { active: true, changes: [ownerPath, observerPath] },
};
} finally {
await Promise.all([owner.close(), observer.close()]);
}
}

/** `files.changed` reports git-root-relative paths; a fixture inside a parent repo carries its directory prefix. */
function expectedChangePath(file) {
try {
const prefix = NodeChildProcess.execFileSync(
"git",
["-C", NodePath.dirname(file), "rev-parse", "--show-prefix"],
{ encoding: "utf8", timeout: 10_000 },
).trim();
return prefix + NodePath.basename(file);
} catch {
return NodePath.basename(file);
}
}

function watcherFixtureFiles(run) {
if (!run?.run || typeof run.fixtureDirectory !== "string") throw new Error("Condition: watcher proof requires one owned workspace and fixture directory.");
const ownerFile = NodePath.join(run.fixtureDirectory, "watch-owner-sentinel.txt");
Expand Down Expand Up @@ -405,8 +421,10 @@ function recordOwnedFile(run, file) {
async function waitForExactWorkspaceInvalidation(events, workspaceId, path, timeoutMs, start = 0) {
const deadline = Date.now() + timeoutMs;
while (Date.now() < deadline) {
// files.changed is a global broadcast; unrelated workspaces and scopes
// legitimately push during the wait. Only the absence of the exact owned
// push (timeout) proves a defect.
const changes = events.slice(start).filter(isFilesChangedPush);
if (changes.some((event) => !isExactWorkspaceInvalidation(event, workspaceId, path))) throw new Error(`Condition: files.changed push did not match the owned ${path} watcher scope.`);
if (changes.some((event) => isExactWorkspaceInvalidation(event, workspaceId, path))) return;
await delay(25);
}
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,145 @@
import "reflect-metadata";
import { describe, expect, it, vi } from "vitest";
import { OpenCodeProvider } from "../opencode-provider.js";
import { OpenCodeServerPool } from "../opencode-server-pool.js";
import type { ProviderTurnDiffUpdate, TurnRequest } from "@mcode/contracts";

function testProvider(http: never, pool: OpenCodeServerPool) {
const settingsService = { get: () => ({ provider: { cli: { opencode: "opencode" } } }) };
const envService = { getEnv: () => ({}) };
const host = {
events: { submit: async () => ({ commit: {}, delivery: { ingress: "queued" } }) },
processes: { attach: () => {}, terminateTree: async () => {} },
runtime: { platform: "win32" },
environment: { snapshot: () => ({}) },
browser: {},
threadControl: {},
grants: {},
};
const provider = new OpenCodeProvider(settingsService as never, envService as never, host as never);
provider.configureTestSeams({
pool,
http: http as never,
probeCli: async () => ({ binaryPath: "opencode", version: "test" }),
idleConfirm: { intervalMs: 5, requiredPolls: 2, timeoutMs: 500, maxPollErrors: 2 },
});
return provider;
}

function testPool(): OpenCodeServerPool {
return new OpenCodeServerPool({
spawn: () => ({ pid: 1, on: () => {}, off: () => {}, kill: () => true }) as never,
waitForHealth: async () => {},
terminateTree: async () => {},
findFreePort: async () => 4096,
now: () => Date.now(),
env: () => ({}),
});
}

function turnRequest(): TurnRequest<"opencode"> {
return {
turnId: "turn-1",
turnExecutionId: "11111111-1111-4111-8111-111111111111",
sessionId: "mcode-thread-1",
workspaceId: "ws-1",
threadId: "thread-1",
message: "hello",
cwd: "/w/a",
model: "anthropic/claude-sonnet-4-6",
permissionMode: "full",
interactionMode: "build",
providerOptions: {},
} as TurnRequest<"opencode">;
}

/** Full-context upstream patch for one changed line. */
function diffEvent(file: string, before: string, after: string) {
return {
type: "session.diff",
properties: {
sessionID: "ses_1",
diff: [{
file,
patch: [
`Index: ${file}`,
"===================================================================",
`--- ${file}`,
`+++ ${file}`,
"@@ -1,1 +1,1 @@",
`-${before}`,
`+${after}`,
"",
].join("\n"),
additions: 1,
deletions: 1,
status: "modified",
}],
},
};
}

function fakeHttp(envelopes: unknown[]) {
return {
createSession: vi.fn(async () => ({ id: "ses_1" })),
promptAsync: vi.fn(async () => {}),
abortSession: vi.fn(async () => {}),
listModels: vi.fn(async () => []),
listSessionMessages: vi.fn(async () => []),
getSessionStatus: vi.fn(async () => new Proxy({}, { get: () => ({ type: "idle" }) })),
subscribeEvents: vi.fn(async (_url: string, _signal: AbortSignal, onEnvelope: (e: unknown) => void) => {
for (const envelope of envelopes) onEnvelope(envelope);
onEnvelope({ type: "session.idle", properties: { sessionID: "ses_1" } });
}),
};
}

describe("OpenCodeProvider native turn diff", () => {
it("pushes a complete native snapshot keyed by the dispatched turn identity", async () => {
const provider = testProvider(fakeHttp([diffEvent("notes.txt", "agent marker: old", "agent marker: new")]) as never, testPool());
const updates: ProviderTurnDiffUpdate[] = [];
provider.onTurnDiff((event) => updates.push(event));
await provider.sendTurn(turnRequest());
expect(updates).toHaveLength(1);
expect(updates[0]).toMatchObject({
turnId: "turn-1",
turnExecutionId: "11111111-1111-4111-8111-111111111111",
deliveryAttempt: 1,
revision: 1,
state: "snapshot",
nativeFidelity: "agent",
});
const patch = (updates[0] as { patch: string }).patch;
expect(patch).toContain("diff --git a/notes.txt b/notes.txt");
expect(patch).toContain("-agent marker: old");
expect(patch).toContain("+agent marker: new");
provider.shutdown();
});

it("ignores a diff event owned by another upstream session", async () => {
const foreign = diffEvent("notes.txt", "a", "b");
foreign.properties.sessionID = "ses_other";
const provider = testProvider(fakeHttp([foreign]) as never, testPool());
const updates: ProviderTurnDiffUpdate[] = [];
provider.onTurnDiff((event) => updates.push(event));
await provider.sendTurn(turnRequest());
expect(updates).toHaveLength(0);
provider.shutdown();
});

it("pushes rejected evidence whole when one entry is unusable", async () => {
const bad = { type: "session.diff", properties: { sessionID: "ses_1", diff: [{ file: "blob.bin", additions: 0, deletions: 0 }] } };
const provider = testProvider(fakeHttp([diffEvent("notes.txt", "a", "b"), bad]) as never, testPool());
const updates: ProviderTurnDiffUpdate[] = [];
provider.onTurnDiff((event) => updates.push(event));
await provider.sendTurn(turnRequest());
expect(updates.map((update) => update.state)).toEqual(["snapshot", "rejected"]);
provider.shutdown();
});

it("declares the turn-diff capability", () => {
const provider = testProvider(fakeHttp([]) as never, testPool());
expect(provider.descriptor.capabilities).toContainEqual({ name: "turn-diff", support: "supported" });
provider.shutdown();
});
});
Original file line number Diff line number Diff line change
Expand Up @@ -12,11 +12,13 @@ import type {
ProviderId,
ProviderModelInfo,
ProviderIdentity,
ProviderTurnDiffUpdate,
SessionForker,
TurnRequest,
} from "@mcode/contracts";
import { AgentEventType, providerRuntimeEvent } from "@mcode/contracts";
import type { ProviderHostPorts } from "@mcode/providers";
import { OpenCodeNativeTurnDiff } from "@mcode/providers";
import { SettingsService } from "../../../settings/settings-service.js";
import { EnvService } from "../../../../runtime/environment/env-service.js";
import { CleanForker } from "../../../handoff/index.js";
Expand All @@ -40,7 +42,7 @@ import {
import { formatOpenCodeResumeCursor, parseOpenCodeResumeCursor } from "./opencode-resume-cursor.js";
import { probeOpenCodeCli } from "./opencode-cli.js";

const OPENCODE_SUPPORTED_CAPABILITIES = ["build", "plan", "permissions", "session-eviction"] as const;
const OPENCODE_SUPPORTED_CAPABILITIES = ["build", "plan", "permissions", "session-eviction", "turn-diff"] as const;

/** Visible notice when a missing upstream session forces a fresh start. */
const OPENCODE_SESSION_INVALIDATED_SUBTYPE = "sdk_session_invalidated";
Expand Down Expand Up @@ -113,6 +115,9 @@ interface OpenCodeTurnState {
messageRoles: Map<string, string>;
/** Text already forwarded per message, for exactly-once streaming. */
forwardedText: Map<string, string>;
/** Native turn-diff accumulator and its monotonic push cursor for this turn. */
nativeDiff: OpenCodeNativeTurnDiff;
nativeDiffRevision: number;
}

function nestedSessionId(holder: unknown): string | undefined {
Expand Down Expand Up @@ -259,6 +264,7 @@ export class OpenCodeProvider extends NodeEvents.EventEmitter implements IAgentP
private readonly canonicalEventPublisher: CanonicalLiveEventPublisher | undefined;
private readonly pendingPermissions = new Map<string, OpenCodePendingAsk>();
private readonly seenNotices = new Map<string, Set<string>>();
private readonly turnDiffListeners = new Set<(event: ProviderTurnDiffUpdate) => void>();
private pool: OpenCodeServerPool;
private http: OpenCodeHttpClient;
private probeCli: (cliPath: string, platform: string) => Promise<{ binaryPath: string; version: string }>;
Expand Down Expand Up @@ -547,6 +553,8 @@ export class OpenCodeProvider extends NodeEvents.EventEmitter implements IAgentP
idleConfirmPromise: null,
messageRoles: new Map<string, string>(),
forwardedText: new Map<string, string>(),
nativeDiff: new OpenCodeNativeTurnDiff(),
nativeDiffRevision: 0,
};
this.turns.set(sessionId, created);
return created;
Expand Down Expand Up @@ -594,6 +602,8 @@ export class OpenCodeProvider extends NodeEvents.EventEmitter implements IAgentP
state.abortController = new AbortController();
state.messageRoles.clear();
state.forwardedText.clear();
state.nativeDiff = new OpenCodeNativeTurnDiff();
state.nativeDiffRevision = 0;
emit({ type: AgentEventType.TurnStarted, threadId } satisfies AgentEvent);
const entry = await this.acquireTurnEntry(req, routing, state, emit);
if (!entry) {
Expand Down Expand Up @@ -744,6 +754,37 @@ export class OpenCodeProvider extends NodeEvents.EventEmitter implements IAgentP
emit({ type: AgentEventType.System, threadId, subtype: `sdk_session_id:${created.id}` } satisfies AgentEvent);
}

/** Subscribe to complete native turn diffs without routing their bytes through renderer events. */
onTurnDiff(handler: (event: ProviderTurnDiffUpdate) => void): () => void {
this.turnDiffListeners.add(handler);
return () => this.turnDiffListeners.delete(handler);
}

/**
* Fold one upstream `session.diff` event into the turn's native aggregate
* and push the complete result. Upstream diffs are computed between its own
* step snapshots, so they carry agent-attributed evidence; the event mapper
* keeps classifying these envelopes as noise for the narrative timeline.
*/
private publishNativeTurnDiff(
req: TurnRequest<"opencode">,
state: OpenCodeTurnState,
properties: Record<string, unknown>,
): void {
const result = state.nativeDiff.observe(Array.isArray(properties.diff) ? properties.diff : []);
if (result === null) return;
const identity = {
turnId: req.turnId,
turnExecutionId: req.turnExecutionId,
deliveryAttempt: req.deliveryAttempt ?? 1,
revision: ++state.nativeDiffRevision,
};
const event: ProviderTurnDiffUpdate = result.state === "snapshot"
? { ...identity, ...result, nativeFidelity: "agent" }
: { ...identity, ...result };
for (const listener of this.turnDiffListeners) listener(event);
}

private handleTurnEnvelope(
envelope: unknown,
req: TurnRequest<"opencode">,
Expand All @@ -759,6 +800,9 @@ export class OpenCodeProvider extends NodeEvents.EventEmitter implements IAgentP
this.trackMessageRole(state, normalized);
const owner = sessionIdOfNormalized(normalized);
if (owner && owner !== upstreamId) return;
if (normalized.type === "session.diff") {
this.publishNativeTurnDiff(req, state, normalized.properties);
}
}
const mapped = mapOpenCodeEnvelope(envelope, {
threadId,
Expand Down
28 changes: 28 additions & 0 deletions apps/web/src/components/diff/LastTurnView.tsx
Original file line number Diff line number Diff line change
@@ -1,5 +1,6 @@
import type { ReviewComparison } from "@mcode/contracts";
import { FileList } from "./FileList";
import { Badge } from "@/components/ui/badge";

/** Props for LastTurnView. */
interface LastTurnViewProps {
Expand Down Expand Up @@ -35,6 +36,15 @@ export function LastTurnView({ threadId, comparison, cacheVersion, refreshing, o

return (
<div data-testid="review-last-turn" className="flex h-full min-h-0 flex-col">
<div className="flex items-center gap-2 px-3 py-2 border-b border-border/15">
<span className="font-mono text-[11px] tabular-nums text-foreground/70">
{comparison.files.length}
</span>
<span className="font-mono text-[10.5px] uppercase tracking-[0.14em] text-muted-foreground">
{turnLabel(comparison)}
</span>
<TurnDiffSource comparison={comparison} />
</div>
<FileList
files={comparison.files}
source="turn-diff"
Expand All @@ -50,3 +60,21 @@ export function LastTurnView({ threadId, comparison, cacheVersion, refreshing, o
</div>
);
}

function turnLabel(comparison: ReviewComparison): string {
const files = comparison.files.length === 1 ? "file" : "files";
const phase = comparison.turnDiff?.phase === "live" ? "Live" : "last turn";
return `${files} · ${phase}`;
}

function TurnDiffSource({ comparison }: { comparison: ReviewComparison }) {
if (!comparison.turnDiff) return null;
return <Badge
variant="secondary"
data-testid="review-turn-source"
data-review-source={comparison.turnDiff.source}
data-review-fidelity={comparison.turnDiff.fidelity}
>
{comparison.turnDiff.source === "git" ? "Git fallback: same-file edits may appear" : comparison.turnDiff.source === "tracked" ? "Tracked file evidence" : "Agent changes"}
</Badge>;
}
Loading
Loading