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
5 changes: 5 additions & 0 deletions .changeset/handoff-store-shutdown.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,5 @@
---
"@paleo/alignfirst-developer-openclaw-plugin": patch
---

Fixed a leaked handoff database connection when the gateway stops during a takeover turn.
5 changes: 5 additions & 0 deletions .changeset/workspace-untracked-branches.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,5 @@
---
"@paleo/workspace": patch
---

Prevented new workspace branches from inheriting an upstream from `--from`.
2 changes: 1 addition & 1 deletion docs/alignfirst-developer/openclaw-plugin.md
Original file line number Diff line number Diff line change
Expand Up @@ -78,4 +78,4 @@ Host entry points checked and refused for an external plugin, so the plugin owns
- With `claude-sonnet-5`, the bot ran ninety-six `alcode` executions in one test day in the foreground with a 60-second timeout, ignoring the background rule. Terra ran none. The 30-second mock run hides the harm a real coding agent would suffer.
- Terra chains one background run per completion step (code, log review, tests, push), so the user sees several intermediate acknowledgements. A product question, not a defect.
- A `⚠️ Message blocked` host notice appeared twice in Slack threads after a silent takeover turn. Its text is in neither the OpenClaw sources nor the build.
- Before the takeover acknowledgement was added, the first visible thread post landed 1.5 to 2.5 minutes after the starter, at the end of the setup turn. The reaction acknowledgement now provides earlier surface feedback.
- The takeover reaction lands immediately after the history read. The first visible thread post still arrives 1.5 to 2.5 minutes after the starter, at the end of the setup turn.
Original file line number Diff line number Diff line change
Expand Up @@ -60,6 +60,8 @@ export function registerThreadHandoff(api: OpenClawPluginApi): void {
);
},
async stop() {
// `service.stop()` guarantees no detached completion calls `getStore()` afterwards; such a
// call would reopen the database behind the close below.
await service.stop();
store?.close();
store = undefined;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -30,15 +30,15 @@ export function createHandoffService(params: HandoffServiceParams): HandoffServi
const now = params.now ?? Date.now;
const scanIntervalMs = params.scanIntervalMs ?? SCAN_INTERVAL_MS;
const attemptSpacingMs = params.attemptSpacingMs ?? ATTEMPT_SPACING_MS;
const dispatchResource = new AsyncResource("alignfirst.thread-handoff.dispatch", {
requireManualDestroy: true,
});
const dispatchResource = new AsyncResource("alignfirst.thread-handoff.dispatch");
const targetWork = new Map<string, Promise<unknown>>();
const inFlight = new Map<string, Promise<void>>();
let timer: ReturnType<typeof setInterval> | undefined;
let scan: Promise<void> | undefined;
let dispatchResourceDestroyed = false;
let stopped = true;
// The host closes the store once `stop()` resolves. A dispatch that settles afterwards must not
// call `getStore()`, which would reopen the database and leak the connection.
let storeReleased = false;

const service: HandoffService = {
async startTurn(record) {
Expand All @@ -60,7 +60,7 @@ export function createHandoffService(params: HandoffServiceParams): HandoffServi
: Promise.reject(
new Error(`Channel ${updated.channelId} is not configured for handoff.`),
);
const completion = finishAttempt(params, updated, turn, now)
const completion = finishAttempt(params, updated, turn, now, () => storeReleased)
.finally(() => {
if (inFlight.get(updated.handoffId) === completion) {
inFlight.delete(updated.handoffId);
Expand All @@ -79,6 +79,7 @@ export function createHandoffService(params: HandoffServiceParams): HandoffServi
async start() {
if (!stopped) return;
stopped = false;
storeReleased = false;
params.getStore();
await recoverPending(service, params, inFlight, now(), attemptSpacingMs);
timer = setInterval(() => {
Expand All @@ -95,16 +96,10 @@ export function createHandoffService(params: HandoffServiceParams): HandoffServi
},
async stop() {
stopped = true;
storeReleased = true;
if (timer) clearInterval(timer);
timer = undefined;
try {
await scan;
} finally {
if (!dispatchResourceDestroyed) {
dispatchResource.emitDestroy();
dispatchResourceDestroyed = true;
}
}
await scan;
},
};
return service;
Expand Down Expand Up @@ -159,6 +154,7 @@ async function finishAttempt(
record: HandoffRecord,
turn: Promise<void>,
now: () => number,
storeReleased: () => boolean,
): Promise<void> {
let failure: unknown;
try {
Expand All @@ -167,7 +163,6 @@ async function finishAttempt(
failure = error;
}
try {
const current = params.getStore().recordAttemptEnd(record.routeKey, now());
if (failure === undefined) {
params.logger.debug?.(
`thread-handoff ${record.handoffId} start attempt ${record.attemptCount} completed`,
Expand All @@ -177,6 +172,9 @@ async function finishAttempt(
`thread-handoff ${record.handoffId} start attempt ${record.attemptCount} failed: ${errorMessage(failure)}`,
);
}
// A stale `lastAttemptedAt` only makes the record eligible for recovery sooner after a restart.
if (storeReleased()) return;
const current = params.getStore().recordAttemptEnd(record.routeKey, now());
if (current?.state === "pending" && current.attemptCount >= MAX_ATTEMPTS) {
params.logger.warn(
`thread-handoff ${record.handoffId} stays pending after ${current.attemptCount} start attempts; it remains claimable by the next human message in the thread; inspect it with: openclaw thread-handoff list`,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -233,6 +233,26 @@ describe("handoff turn start and recovery", () => {
await fixture.service.stop();
fixture.store.close();
});

it("leaves the store alone when a turn settles after stop", async () => {
const fixture = serviceFixture();
const record = handoff();
fixture.store.insertHandoff(record);
const turn = createDeferred<void>();
dispatchTurn.mockReturnValue(turn.promise);

await fixture.service.startTurn(record);
await fixture.service.stop();
const callsBeforeSettling = fixture.getStore.mock.calls.length;
turn.resolve();
await vi.waitFor(() =>
expect(fixture.logger.debug).toHaveBeenCalledWith(expect.stringContaining("completed")),
);

expect(fixture.getStore).toHaveBeenCalledTimes(callsBeforeSettling);
expect(fixture.logger.error).not.toHaveBeenCalled();
fixture.store.close();
});
});

function createDeferred<T>(): {
Expand All @@ -252,6 +272,7 @@ async function nextEventLoopTurn(): Promise<void> {

function serviceFixture(options: { now?: () => number } = {}) {
const store = createHandoffStore(temporaryStateDir());
const getStore = vi.fn(() => store);
const runtime = {
config: { current: () => ({}) },
} as unknown as OpenClawPluginApi["runtime"];
Expand All @@ -263,12 +284,13 @@ function serviceFixture(options: { now?: () => number } = {}) {
};
return {
store,
getStore,
runtime,
logger,
service: createHandoffService({
runtime,
configuration: { channelSurfaces: { slack: "slack", discord: "discord" } },
getStore: () => store,
getStore,
logger: logger as unknown as PluginLogger,
now: options.now ?? (() => 100_000),
}),
Expand Down
9 changes: 7 additions & 2 deletions packages/openclaw-test/test/docker-compose.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -4,8 +4,13 @@ import { describe, expect, it } from "vitest";
describe("gateway Compose service", () => {
it("declares the externally managed runtime config read-only", () => {
const compose = readFileSync(new URL("../docker-compose.yml", import.meta.url), "utf8");
const gateway = compose.slice(compose.indexOf(" gateway:"), compose.indexOf(" runner:"));
const gatewayStart = compose.indexOf("\n gateway:");
const gatewayEnd = compose.indexOf("\n runner:");

expect(gateway).toContain('OPENCLAW_CONFIG_READONLY: "1"');
// Both markers must be found and ordered, otherwise the slice would silently widen to the
// rest of the file and the assertion would accept the variable on any other service.
expect(gatewayStart).toBeGreaterThan(-1);
expect(gatewayEnd).toBeGreaterThan(gatewayStart);
expect(compose.slice(gatewayStart, gatewayEnd)).toContain('OPENCLAW_CONFIG_READONLY: "1"');
});
});
10 changes: 9 additions & 1 deletion packages/workspace/src/worktree.ts
Original file line number Diff line number Diff line change
Expand Up @@ -74,7 +74,15 @@ export function createBranch(
const worktreePath = dedupeWorktreePath(
computeWorktreePath(ctx.mainWorktree, finalBranch, dirNameFn),
);
const addArgs = ["worktree", "add", "-b", finalBranch, "--end-of-options", worktreePath];
const addArgs = [
"worktree",
"add",
"--no-track",
"-b",
finalBranch,
"--end-of-options",
worktreePath,
];
if (from !== undefined) addArgs.push(from);
execFileSync("git", addArgs, { stdio: stdioFor(run) });
return { ...ctx, currentWorktree: worktreePath, isMainWorktree: false };
Expand Down
2 changes: 1 addition & 1 deletion packages/workspace/templates/guide.md
Original file line number Diff line number Diff line change
Expand Up @@ -20,7 +20,7 @@ A **workspace** is a git worktree (with its branch) plus its own dev setup: syml
{{COMMANDS:setup}}
```

With `-c`, the new branch starts at the current worktree's HEAD (like `git switch -c`); `--from <ref>` accepts any commit-ish as the base. When the branch name is already taken, `setup -c` errors; add `--dedupe` to append `-2`, `-3`… instead.
With `-c`, the new branch starts at the current worktree's HEAD (like `git switch -c`); `--from <ref>` accepts any commit-ish as the base. New branches have no upstream; set one when you first push. When the branch name is already taken, `setup -c` errors; add `--dedupe` to append `-2`, `-3`… instead.

{{#PORTS}}
`setup` creates the worktree (branch, ports, symlinks, config files), then runs the project's finalize step: install dependencies, build, provision the database.
Expand Down
22 changes: 22 additions & 0 deletions packages/workspace/test/worktree-git.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,9 @@ beforeEach(() => {
git(["init", "-b", "main"]);
git(["config", "user.email", "test@example.com"]);
git(["config", "user.name", "Test"]);
// `createBranch` passes `--no-track`; pin the key it neutralises so the assertion cannot pass
// because of the developer's own `branch.autoSetupMerge`.
git(["config", "branch.autoSetupMerge", "always"]);
git(["commit", "--allow-empty", "-m", "init"]);
process.chdir(repo);
ctx = { currentWorktree: repo, mainWorktree: repo, isMainWorktree: true };
Expand Down Expand Up @@ -76,6 +79,21 @@ describe("createBranch conflicts", () => {
});
});

describe("createBranch start point", () => {
it("does not track a remote branch passed with --from", () => {
git(["remote", "add", "origin", repo]);
git(["update-ref", "refs/remotes/origin/main", "HEAD"]);
git(["commit", "--allow-empty", "-m", "local commit"]);

createBranch("feature", ctx, run, { from: "origin/main" });

expect(gitOutput(["rev-parse", "feature"])).toBe(gitOutput(["rev-parse", "origin/main"]));
expect(gitOutput(["for-each-ref", "--format=%(upstream:short)", "refs/heads/feature"])).toBe(
"",
);
});
});

describe("worktree error paths", () => {
it("throws when --from does not resolve", () => {
expect(() => createBranch("feat", ctx, run, { from: "does-not-exist" })).toThrow(
Expand All @@ -92,6 +110,10 @@ function git(args: string[]): void {
execFileSync("git", args, { cwd: repo, stdio: "pipe" });
}

function gitOutput(args: string[]): string {
return execFileSync("git", args, { cwd: repo, encoding: "utf-8" }).trim();
}

function branchExists(branch: string): boolean {
const out = execFileSync("git", ["branch", "--list", branch], {
cwd: repo,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -99,7 +99,7 @@ The `{ask}` is one sentence, and it reflects the first unresolved requirement:
- No TASK → ask what needs to be done.
- A resource URL that may provide the project or ticket → ask for neither; state that the working session will inspect the URL.
- A request explicitly spanning several projects, or work independent of any project → ask for no main project; state `Ready for the work session.` in the user's language.
- Nothing else needs an answer → state `Ready for the work session.` in the user's language (for example, "Ready for the work session.").
- Nothing else needs an answer → state `Ready for the work session.` in the user's language.

Whichever case applies, the `{ask}` never says that this channel session handles the work, and never says that work has begun.

Expand Down