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
48 changes: 43 additions & 5 deletions src/core/worktree-service.js
Original file line number Diff line number Diff line change
Expand Up @@ -358,7 +358,7 @@ function classifyEntry(entry, repoRoot, activeWorktree, protectedBranches, merge
return { ...item, status: "review", reason: "not merged into base branch by git ancestry", selectable: true };
}

export function createWorktreeWorkflowService({ directory, git, stateStore }) {
export function createWorktreeWorkflowService({ directory, git, stateStore, logger = null }) {
async function computeCleanupPreview({ repoRoot, activeWorktree }) {
const config = await loadWorkflowConfig(repoRoot);
const { defaultBranch, baseBranch, baseRef } = await resolveBaseTarget(repoRoot, config);
Expand Down Expand Up @@ -443,6 +443,17 @@ export function createWorktreeWorkflowService({ directory, git, stateStore }) {
async function updateStateForPrepare(repoRoot, sessionID, prepared, createdBy = "manual", workspaceRole = "linear-flow") {
if (!sessionID || !stateStore) return;
const state = await stateStore.loadSessionState(repoRoot, sessionID);
const previous = stateStore.findTaskByID(state, prepared.branch) || stateStore.findTaskByWorktreePath(state, prepared.worktree_path);
const previousActiveTaskID = stateStore.getActiveTask(state);
const isMeaningfulBindingChange =
!previous ||
previousActiveTaskID !== prepared.branch ||
previous.task_id !== prepared.branch ||
previous.branch !== prepared.branch ||
previous.worktree_path !== prepared.worktree_path ||
(previous.title ?? null) !== (prepared.title ?? null) ||
previous.created_by !== createdBy ||
previous.workspace_role !== workspaceRole;
const next = stateStore.setActiveTask(
stateStore.upsertTask(state, {
task_id: prepared.branch,
Expand All @@ -456,6 +467,16 @@ export function createWorktreeWorkflowService({ directory, git, stateStore }) {
prepared.branch,
);
await stateStore.saveSessionState(repoRoot, sessionID, next);
if (isMeaningfulBindingChange) {
logger?.info(previous ? "session_binding_updated" : "session_binding_created", {
session_id: sessionID,
task_id: prepared.branch,
branch: prepared.branch,
worktree_path: prepared.worktree_path,
created_by: createdBy,
workspace_role: workspaceRole,
});
}
}
async function updateStateForCleanup(repoRoot, sessionID, removed) {
if (!sessionID || !stateStore || removed.length === 0) return;
Expand All @@ -474,12 +495,17 @@ export function createWorktreeWorkflowService({ directory, git, stateStore }) {
});
if (stateStore.getActiveTask(state) === taskID) {
state = stateStore.setActiveTask(state, null);
logger?.info("session_binding_cleared", {
session_id: sessionID,
task_id: taskID,
reason: "cleanup",
});
}
}
await stateStore.saveSessionState(repoRoot, sessionID, state);
}

async function prepare({ title, sessionID, createdBy = "manual" }) {
async function prepare({ title, sessionID, createdBy = "manual", workspaceRole = "linear-flow" }) {
const repoRoot = await getRepoRoot();
const config = await loadWorkflowConfig(repoRoot);
const { defaultBranch, baseBranch, baseRef } = await resolveBaseTarget(repoRoot, config);
Expand All @@ -495,7 +521,7 @@ export function createWorktreeWorkflowService({ directory, git, stateStore }) {
const branchCommit = (await git(["rev-parse", branchName], { cwd: repoRoot })).stdout;
if (branchCommit !== baseCommit) throw new Error(`New branch ${branchName} does not match ${baseBranch} at ${baseCommit}. Found ${branchCommit} instead.`);
const result = buildPrepareResult({ title, branch: branchName, worktreePath, defaultBranch, baseBranch, baseRef, baseCommit });
await updateStateForPrepare(repoRoot, sessionID, result, createdBy);
await updateStateForPrepare(repoRoot, sessionID, result, createdBy, workspaceRole);
return result;
}

Expand All @@ -521,12 +547,19 @@ export function createWorktreeWorkflowService({ directory, git, stateStore }) {
);
await stateStore.saveSessionState(repoRoot, sessionID, next);
const refreshed = stateStore.getActiveTaskRecord(next);
logger?.info("session_binding_updated", {
session_id: sessionID,
task_id: refreshed?.task_id,
branch: refreshed?.branch,
worktree_path: refreshed?.worktree_path,
created_by: refreshed?.created_by,
workspace_role: refreshed?.workspace_role,
});
return { repoRoot, task: refreshed };
}
return { repoRoot, task: activeTask };
}
const prepared = await prepare({ title, sessionID, createdBy: "harness" });
await updateStateForPrepare(repoRoot, sessionID, prepared, "harness", workspaceRole);
const prepared = await prepare({ title, sessionID, createdBy: "harness", workspaceRole });
return {
repoRoot,
task: {
Expand Down Expand Up @@ -569,6 +602,11 @@ export function createWorktreeWorkflowService({ directory, git, stateStore }) {
});
if (stateStore.getActiveTask(state) === current.task_id) {
state = stateStore.setActiveTask(state, null);
logger?.info("session_binding_cleared", {
session_id: sessionID,
task_id: current.task_id,
reason: nextStatus === "blocked" ? "blocked" : "completed",
});
}
await stateStore.saveSessionState(repoRoot, sessionID, state);
return { task_id: current.task_id, status: nextStatus };
Expand Down
60 changes: 53 additions & 7 deletions src/index.js
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@ import {
inferTaskLifecycleTransition,
rewriteRepoScopedPathIntoWorktree,
} from "./core/task-binding.js";
import { createDecisionLogger } from "./runtime/decision-logger.js";
import { createRuntimeStateStore } from "./runtime/state-store.js";

function publishStructuredResult(context, result) {
Expand Down Expand Up @@ -165,11 +166,13 @@ export const __internal = {

export const pluginID = "@sven1103/opencode-worktree-workflow";

export const WorktreeWorkflowPlugin = async ({ $, directory }) => {
export const WorktreeWorkflowPlugin = async ({ $, directory, logger: providedLogger = null }) => {
const logger = providedLogger || createDecisionLogger();
const service = createWorktreeWorkflowService({
directory,
git: createGitRunner($, directory),
stateStore: createRuntimeStateStore(),
logger,
});

async function onToolExecuteBefore(input, output) {
Expand All @@ -181,12 +184,15 @@ export const WorktreeWorkflowPlugin = async ({ $, directory }) => {

const sessionID = input?.sessionID;
let binding = null;
let repoRootForLogging = null;

if (classification.requiresIsolation) {
if (!sessionID) throw new Error(`Isolation required for ${toolName || "tool"} but sessionID is missing.`);
const repoRoot = await service.getRepoRoot();
repoRootForLogging = repoRoot;
for (const key of rewritePolicy.opaqueArgKeys) {
if (hasOpaqueRepoRootAbsoluteReference({ value: args[key], repoRoot })) {
logger.info("tool_path_rewrite_skipped", { tool: toolName, arg_key: key, reason: "opaque-repo-root-reference", session_id: sessionID });
throw new Error(`Blocked: ${toolName} ${key} includes repo-root absolute path that cannot be safely rewritten.`);
}
}
Expand All @@ -197,11 +203,23 @@ export const WorktreeWorkflowPlugin = async ({ $, directory }) => {
});
} else if (sessionID) {
const repoRoot = await service.getRepoRoot();
repoRootForLogging = repoRoot;
const { activeTask } = await service.getSessionBinding({ repoRoot, sessionID });
if (activeTask?.worktree_path) binding = { repoRoot, task: activeTask };
}

if (!binding) return;
if (!binding) {
for (const key of rewritePolicy.pathArgKeys) {
const hasArg = key in args;
logger.info("tool_path_rewrite_skipped", {
tool: toolName,
arg_key: key,
reason: hasArg ? "no-active-binding" : "arg-missing",
session_id: sessionID,
});
}
return;
}

if (toolName === "task") {
const handoffPath = resolveSafeHandoffPath({
Expand Down Expand Up @@ -230,13 +248,39 @@ export const WorktreeWorkflowPlugin = async ({ $, directory }) => {
nextArgs.path = binding.task.worktree_path;
}
for (const key of rewritePolicy.pathArgKeys) {
if (key in nextArgs) {
nextArgs[key] = rewriteRepoScopedPathIntoWorktree({
value: nextArgs[key],
repoRoot: binding.repoRoot,
worktreePath: binding.task.worktree_path,
if (!(key in nextArgs)) {
logger.info("tool_path_rewrite_skipped", { tool: toolName, arg_key: key, reason: "arg-missing", session_id: sessionID, task_id: binding.task.task_id });
continue;
}
const before = nextArgs[key];
const after = rewriteRepoScopedPathIntoWorktree({
value: nextArgs[key],
repoRoot: binding.repoRoot,
worktreePath: binding.task.worktree_path,
});
nextArgs[key] = after;
if (before !== after) {
logger.info("tool_path_rewrite_applied", { tool: toolName, arg_key: key, session_id: sessionID, task_id: binding.task.task_id });
logger.debug("tool_path_rewrite_applied", {
tool: toolName,
arg_key: key,
session_id: sessionID,
task_id: binding.task.task_id,
before_path: typeof before === "string" ? before : String(before),
after_path: typeof after === "string" ? after : String(after),
});
continue;
}
let reason = "outside-repo-root";
if (typeof before === "string" && before.trim()) {
const resolved = path.resolve(before);
const inWorktree = (() => {
const relative = path.relative(path.resolve(binding.task.worktree_path), resolved);
return relative === "" || (!relative.startsWith("..") && !path.isAbsolute(relative));
})();
if (inWorktree) reason = "already-in-worktree";
}
logger.info("tool_path_rewrite_skipped", { tool: toolName, arg_key: key, reason, session_id: sessionID, task_id: binding.task.task_id, repo_root_known: Boolean(repoRootForLogging) });
}

output.args = nextArgs;
Expand Down Expand Up @@ -282,11 +326,13 @@ export const WorktreeWorkflowPlugin = async ({ $, directory }) => {
};
return;
} catch {
logger.info("nonfatal_plugin_error", { stage: "task_advisory_cleanup_preview", session_id: sessionID, message: "Cleanup advisory preview failed." });
// Advisory preview is non-fatal.
}
}
}
} catch {
logger.info("nonfatal_plugin_error", { stage: "task_lifecycle_inference", session_id: sessionID, message: "Task lifecycle correlation failed." });
// Artifact correlation/lifecycle inference is non-fatal.
}
}
Expand Down
49 changes: 49 additions & 0 deletions src/runtime/decision-logger.js
Original file line number Diff line number Diff line change
@@ -0,0 +1,49 @@
function normalizeLevel(level) {
const raw = typeof level === "string" ? level.trim().toLowerCase() : "";
if (raw === "silent") return "silent";
if (raw === "debug") return "debug";
return "info";
}

function sanitizeValue(value) {
if (value == null) return value;
if (typeof value === "string") return value.length > 500 ? `${value.slice(0, 500)}…` : value;
if (typeof value === "number" || typeof value === "boolean") return value;
if (Array.isArray(value)) {
return value.slice(0, 20).map((entry) => sanitizeValue(entry)).filter((entry) => entry !== undefined);
}
return undefined;
}

function sanitizeFields(fields) {
if (!fields || typeof fields !== "object") return {};
const blocked = /(secret|token|password|prompt|content|payload|patch)/i;
const result = {};
for (const [key, value] of Object.entries(fields)) {
if (blocked.test(key)) continue;
const sanitized = sanitizeValue(value);
if (sanitized !== undefined) result[key] = sanitized;
}
return result;
}

export function createDecisionLogger({ level = process.env.OPENCODE_WORKTREE_LOG_LEVEL, write } = {}) {
const resolvedLevel = normalizeLevel(level);
const sink = typeof write === "function" ? write : (line) => process.stderr.write(`${line}\n`);

function emit(logLevel, event, fields) {
if (resolvedLevel === "silent") return;
if (resolvedLevel === "info" && logLevel === "debug") return;
if (typeof event !== "string" || !event) return;
sink(JSON.stringify({ ts: new Date().toISOString(), level: logLevel, event, ...sanitizeFields(fields) }));
}

return {
info(event, fields = {}) {
emit("info", event, fields);
},
debug(event, fields = {}) {
emit("debug", event, fields);
},
};
}
22 changes: 21 additions & 1 deletion test-support/helpers.js
Original file line number Diff line number Diff line change
Expand Up @@ -109,10 +109,30 @@ async function createRemoteRepo() {
};
}

async function createPlugin(repoPath) {
function createMemoryDecisionLogger(captureLogs, logLevel = "debug") {
const rank = logLevel === "debug" ? 2 : logLevel === "info" ? 1 : 0;
function emit(level, event, fields = {}) {
if (rank === 0) return;
if (rank === 1 && level === "debug") return;
captureLogs.push({ level, event, ...fields });
}
return {
info(event, fields = {}) {
emit("info", event, fields);
},
debug(event, fields = {}) {
emit("debug", event, fields);
},
};
}

async function createPlugin(repoPath, options = {}) {
const { logLevel, captureLogs } = options;
const logger = Array.isArray(captureLogs) ? createMemoryDecisionLogger(captureLogs, logLevel || "debug") : null;
return pluginModule.server({
$: createShell(repoPath),
directory: repoPath,
...(logger ? { logger } : {}),
});
}

Expand Down
12 changes: 9 additions & 3 deletions test/advisory-cleanup.test.js
Original file line number Diff line number Diff line change
Expand Up @@ -16,10 +16,11 @@ import {

test("task completion emits advisory cleanup preview and persists terminal lifecycle", async () => {
const fixture = await createRemoteRepo();
const logs = [];
const previous = process.env.OPENCODE_WORKTREE_STATE_DIR;
process.env.OPENCODE_WORKTREE_STATE_DIR = fixture.stateDir;
try {
const plugin = await createPlugin(fixture.repoPath);
const plugin = await createPlugin(fixture.repoPath, { captureLogs: logs });
const handoffPath = await createHandoffArtifact(fixture.repoPath, "session-advisory-1", "handoff-1");
const delegated = await runTaskDelegationHook(plugin, {
sessionID: "session-advisory-1",
Expand All @@ -43,6 +44,7 @@ test("task completion emits advisory cleanup preview and persists terminal lifec
const state = await store.loadSessionState(repoRoot, "session-advisory-1");
assert.equal(state.active_task_id, null);
assert.equal(state.tasks[0].status, "completed");
assert.equal(logs.some((entry) => entry.event === "session_binding_cleared" && entry.reason === "completed"), true);
} finally {
process.env.OPENCODE_WORKTREE_STATE_DIR = previous;
await fixture.cleanup();
Expand Down Expand Up @@ -135,10 +137,11 @@ test("advisory preview marks unknown provenance for unmanaged candidates", async

test("advisory preview failures are non-fatal", async () => {
const fixture = await createRemoteRepo();
const logs = [];
const previous = process.env.OPENCODE_WORKTREE_STATE_DIR;
process.env.OPENCODE_WORKTREE_STATE_DIR = fixture.stateDir;
try {
const plugin = await createPlugin(fixture.repoPath);
const plugin = await createPlugin(fixture.repoPath, { captureLogs: logs });
const handoffPath = await createHandoffArtifact(fixture.repoPath, "session-advisory-5", "handoff-5");
const delegated = await runTaskDelegationHook(plugin, {
sessionID: "session-advisory-5",
Expand All @@ -156,6 +159,7 @@ test("advisory preview failures are non-fatal", async () => {

assert.equal(output.output?.toolName, "task");
assert.equal(output.advisoryMetadata, null);
assert.equal(logs.some((entry) => entry.event === "nonfatal_plugin_error" && entry.stage === "task_advisory_cleanup_preview"), true);
} finally {
process.env.OPENCODE_WORKTREE_STATE_DIR = previous;
await fixture.cleanup();
Expand All @@ -164,10 +168,11 @@ test("advisory preview failures are non-fatal", async () => {

test("malformed result artifact is non-fatal in task after-hook", async () => {
const fixture = await createRemoteRepo();
const logs = [];
const previous = process.env.OPENCODE_WORKTREE_STATE_DIR;
process.env.OPENCODE_WORKTREE_STATE_DIR = fixture.stateDir;
try {
const plugin = await createPlugin(fixture.repoPath);
const plugin = await createPlugin(fixture.repoPath, { captureLogs: logs });
const handoffPath = await createHandoffArtifact(fixture.repoPath, "session-advisory-6", "handoff-6");
const delegated = await runTaskDelegationHook(plugin, {
sessionID: "session-advisory-6",
Expand All @@ -186,6 +191,7 @@ test("malformed result artifact is non-fatal in task after-hook", async () => {
assert.equal(output.output?.toolName, "task");
assert.equal(output.advisoryMetadata, null);
assert.equal(output.advisoryTextParts.length, 0);
assert.equal(logs.some((entry) => entry.event === "nonfatal_plugin_error" && entry.stage === "task_lifecycle_inference"), true);
} finally {
process.env.OPENCODE_WORKTREE_STATE_DIR = previous;
await fixture.cleanup();
Expand Down
Loading
Loading