refactor(kanban_db_workspace): unify git subprocess/anchored-worktree/active-children helpers, flatten resolve_workspace (609->567 LOC)
This commit is contained in:
@@ -9,6 +9,7 @@ from __future__ import annotations
|
||||
|
||||
import logging
|
||||
import os
|
||||
import shutil
|
||||
import sqlite3
|
||||
import subprocess
|
||||
from pathlib import Path
|
||||
@@ -22,6 +23,33 @@ if TYPE_CHECKING:
|
||||
# Log-record parity with the origin module.
|
||||
_log = logging.getLogger("hermes_cli.kanban_db")
|
||||
|
||||
_REMOVABLE_KINDS = ("scratch", "worktree")
|
||||
|
||||
# Statuses after which a child no longer needs its parent's workspace artifacts.
|
||||
_ACTIVE_CHILDREN_SQL = (
|
||||
"SELECT 1 FROM task_links l "
|
||||
"JOIN tasks t ON t.id = l.child_id "
|
||||
"WHERE l.parent_id = ? AND t.status NOT IN ('done', 'archived', 'failed', 'cancelled') "
|
||||
"LIMIT 1"
|
||||
)
|
||||
|
||||
_WORKSPACE_ROW_SQL = "SELECT workspace_kind, workspace_path, branch_name FROM tasks WHERE id = ?"
|
||||
|
||||
|
||||
def _git(repo_root: Path, *args: str, timeout: int) -> subprocess.CompletedProcess:
|
||||
"""``git -C repo_root args``; never raises on a non-zero exit."""
|
||||
return subprocess.run(
|
||||
["git", "-C", str(repo_root), *args],
|
||||
capture_output=True,
|
||||
text=True, encoding='utf-8', errors='replace',
|
||||
timeout=timeout,
|
||||
check=False,
|
||||
)
|
||||
|
||||
|
||||
def _has_active_children(conn: sqlite3.Connection, task_id: str) -> bool:
|
||||
return conn.execute(_ACTIVE_CHILDREN_SQL, (task_id,)).fetchone() is not None
|
||||
|
||||
|
||||
def _managed_scratch_path_info(p: Path) -> tuple[bool, Optional[str]]:
|
||||
"""Return whether *p* is managed scratch storage and the matching board."""
|
||||
@@ -41,25 +69,13 @@ def _managed_scratch_path_info(p: Path) -> tuple[bool, Optional[str]]:
|
||||
if home is not None:
|
||||
with contextlib.suppress(OSError):
|
||||
roots.append(((home / "kanban" / "workspaces").resolve(strict=False), _kb.DEFAULT_BOARD))
|
||||
try:
|
||||
boards_parent = (home / "kanban" / "boards").resolve(strict=False)
|
||||
except OSError:
|
||||
boards_parent = None
|
||||
if boards_parent is not None:
|
||||
try:
|
||||
entries = list(boards_parent.iterdir())
|
||||
except OSError:
|
||||
entries = []
|
||||
for entry in entries:
|
||||
try:
|
||||
if not entry.is_dir():
|
||||
continue
|
||||
except OSError:
|
||||
continue
|
||||
try:
|
||||
entries: list[Path] = []
|
||||
with contextlib.suppress(OSError):
|
||||
entries = list((home / "kanban" / "boards").resolve(strict=False).iterdir())
|
||||
for entry in entries:
|
||||
with contextlib.suppress(OSError):
|
||||
if entry.is_dir():
|
||||
roots.append(((entry / "workspaces").resolve(strict=False), entry.name))
|
||||
except OSError:
|
||||
continue
|
||||
for root, board in roots:
|
||||
if p_abs == root:
|
||||
continue
|
||||
@@ -97,8 +113,7 @@ def _is_managed_scratch_path(p: Path) -> bool:
|
||||
``default_workdir`` pointing at a real source tree can otherwise pair with
|
||||
``workspace_kind='scratch'`` and make task completion delete user data.
|
||||
"""
|
||||
is_managed, _board = _managed_scratch_path_info(p)
|
||||
return is_managed
|
||||
return _managed_scratch_path_info(p)[0]
|
||||
|
||||
|
||||
def _cleanup_workspace(conn: sqlite3.Connection, task_id: str) -> None:
|
||||
@@ -108,43 +123,32 @@ def _cleanup_workspace(conn: sqlite3.Connection, task_id: str) -> None:
|
||||
only when provably free of work (clean tree, every commit reachable from a
|
||||
remote-tracking ref); ``dir`` is intentionally preserved."""
|
||||
try:
|
||||
row = conn.execute(
|
||||
"SELECT workspace_kind, workspace_path, branch_name FROM tasks WHERE id = ?",
|
||||
(task_id,),
|
||||
).fetchone()
|
||||
row = conn.execute(_WORKSPACE_ROW_SQL, (task_id,)).fetchone()
|
||||
if not row:
|
||||
return
|
||||
kind: Optional[str] = row["workspace_kind"]
|
||||
path: Optional[str] = row["workspace_path"]
|
||||
if kind not in ("scratch", "worktree") or not path:
|
||||
if kind not in _REMOVABLE_KINDS or not path:
|
||||
# Not removable itself, but completing may still unblock a deferred
|
||||
# parent scratch cleanup (e.g. a 'dir' child of a scratch parent).
|
||||
_try_cleanup_parent_workspaces(conn, task_id)
|
||||
return
|
||||
# Defer while any child is not yet terminal so it can still read
|
||||
# handoff artifacts from this workspace.
|
||||
_active_children = conn.execute(
|
||||
"SELECT 1 FROM task_links l "
|
||||
"JOIN tasks t ON t.id = l.child_id "
|
||||
"WHERE l.parent_id = ? AND t.status NOT IN ('done', 'archived', 'failed', 'cancelled') "
|
||||
"LIMIT 1",
|
||||
(task_id,),
|
||||
).fetchone()
|
||||
if _active_children:
|
||||
if _has_active_children(conn, task_id):
|
||||
_kb._log.debug(
|
||||
"Deferring %s workspace cleanup for task %s: "
|
||||
"active children still need workspace at %s",
|
||||
kind, task_id, path,
|
||||
)
|
||||
return
|
||||
# Kill the (dead) tmux worker session BEFORE removing a worktree so a
|
||||
# lingering worker never has its cwd deleted from under it.
|
||||
if kind == "worktree":
|
||||
# Kill the (dead) tmux worker session BEFORE removing the worktree
|
||||
# so a lingering worker never has its cwd deleted from under it.
|
||||
_cleanup_worker_tmux(conn, task_id)
|
||||
_cleanup_worktree_workspace(task_id, path, row["branch_name"])
|
||||
_try_cleanup_parent_workspaces(conn, task_id)
|
||||
return
|
||||
import shutil
|
||||
wp = Path(path)
|
||||
if wp.is_dir():
|
||||
# Containment guard: a board's ``default_workdir`` can pair
|
||||
@@ -199,13 +203,7 @@ def _cleanup_worktree_workspace(
|
||||
return
|
||||
# No --force: git's own dirty guard re-verifies at removal time, so if
|
||||
# the tree became dirty since our check (TOCTOU) removal fails safe.
|
||||
result = subprocess.run(
|
||||
["git", "-C", str(repo_root), "worktree", "remove", str(wp)],
|
||||
capture_output=True,
|
||||
text=True, encoding='utf-8', errors='replace',
|
||||
timeout=60,
|
||||
check=False,
|
||||
)
|
||||
result = _git(repo_root, "worktree", "remove", str(wp), timeout=60)
|
||||
if result.returncode != 0:
|
||||
_kb._log.warning(
|
||||
"git worktree remove failed for task %s at %s: %s",
|
||||
@@ -215,13 +213,7 @@ def _cleanup_worktree_workspace(
|
||||
_kb._log.debug("Removed worktree workspace: %s", wp)
|
||||
branch = (branch_name or "").strip() or f"wt/{task_id}"
|
||||
if branch.startswith("wt/"):
|
||||
subprocess.run(
|
||||
["git", "-C", str(repo_root), "branch", "-D", branch],
|
||||
capture_output=True,
|
||||
text=True, encoding='utf-8', errors='replace',
|
||||
timeout=30,
|
||||
check=False,
|
||||
)
|
||||
_git(repo_root, "branch", "-D", branch, timeout=30)
|
||||
except Exception:
|
||||
pass # best-effort — never block completion
|
||||
|
||||
@@ -236,31 +228,19 @@ def _try_cleanup_parent_workspaces(conn: sqlite3.Connection, task_id: str) -> No
|
||||
(task_id,),
|
||||
).fetchall()
|
||||
for (parent_id,) in parents:
|
||||
row = conn.execute(
|
||||
"SELECT workspace_kind, workspace_path, branch_name FROM tasks WHERE id = ?",
|
||||
(parent_id,),
|
||||
).fetchone()
|
||||
row = conn.execute(_WORKSPACE_ROW_SQL, (parent_id,)).fetchone()
|
||||
if (
|
||||
not row
|
||||
or row["workspace_kind"] not in ("scratch", "worktree")
|
||||
or row["workspace_kind"] not in _REMOVABLE_KINDS
|
||||
or not row["workspace_path"]
|
||||
or _has_active_children(conn, parent_id)
|
||||
):
|
||||
continue
|
||||
active = conn.execute(
|
||||
"SELECT 1 FROM task_links l "
|
||||
"JOIN tasks t ON t.id = l.child_id "
|
||||
"WHERE l.parent_id = ? AND t.status NOT IN ('done', 'archived', 'failed', 'cancelled') "
|
||||
"LIMIT 1",
|
||||
(parent_id,),
|
||||
).fetchone()
|
||||
if active:
|
||||
continue # still has active children
|
||||
if row["workspace_kind"] == "worktree":
|
||||
_cleanup_worktree_workspace(
|
||||
parent_id, row["workspace_path"], row["branch_name"]
|
||||
)
|
||||
continue
|
||||
import shutil
|
||||
wp = Path(row["workspace_path"])
|
||||
if wp.is_dir() and _is_managed_scratch_path(wp):
|
||||
shutil.rmtree(wp, ignore_errors=True)
|
||||
@@ -277,9 +257,8 @@ def _cleanup_worker_tmux(conn: sqlite3.Connection, task_id: str) -> None:
|
||||
).fetchone()
|
||||
if not row or not row["assignee"]:
|
||||
return
|
||||
assignee: str = row["assignee"]
|
||||
# Workers named swarm1-12 use tmux sessions named swarm-swarm1 etc.
|
||||
session = f"swarm-{assignee}"
|
||||
session = f"swarm-{row['assignee']}"
|
||||
out = subprocess.run(
|
||||
["tmux", "list-panes", "-t", session, "-F", "#{pane_dead}"],
|
||||
capture_output=True, text=True, encoding='utf-8', errors='replace', timeout=5,
|
||||
@@ -337,9 +316,7 @@ def _maybe_emit_scratch_tip(
|
||||
"""Emit the first-use scratch-workspace tip once per install, right after a
|
||||
scratch workspace is materialized. No-op for ``worktree``/``dir`` (preserved
|
||||
by design) and once the sentinel exists."""
|
||||
if (workspace_kind or "scratch") != "scratch":
|
||||
return
|
||||
if _scratch_tip_shown():
|
||||
if (workspace_kind or "scratch") != "scratch" or _scratch_tip_shown():
|
||||
return
|
||||
try:
|
||||
_kb._log.warning("kanban: %s (task %s)", _SCRATCH_TIP_MESSAGE, task_id)
|
||||
@@ -372,26 +349,23 @@ def _git_toplevel(path: Path) -> Optional[Path]:
|
||||
|
||||
def _git_branch_exists(repo_root: Path, branch_name: str) -> bool:
|
||||
try:
|
||||
result = subprocess.run(
|
||||
["git", "-C", str(repo_root), "show-ref", "--verify", f"refs/heads/{branch_name}"],
|
||||
capture_output=True,
|
||||
text=True, encoding='utf-8', errors='replace',
|
||||
timeout=30,
|
||||
check=False,
|
||||
)
|
||||
result = _git(repo_root, "show-ref", "--verify", f"refs/heads/{branch_name}", timeout=30)
|
||||
except Exception:
|
||||
return False
|
||||
return result.returncode == 0
|
||||
|
||||
|
||||
def _git_common_dir(path: Path) -> Optional[Path]:
|
||||
out = _kb._git_out(path, "rev-parse", "--path-format=absolute", "--git-common-dir")
|
||||
def _git_abs_path(path: Path, flag: str) -> Optional[Path]:
|
||||
out = _kb._git_out(path, "rev-parse", "--path-format=absolute", flag)
|
||||
return Path(out).expanduser().resolve(strict=False) if out else None
|
||||
|
||||
|
||||
def _git_common_dir(path: Path) -> Optional[Path]:
|
||||
return _git_abs_path(path, "--git-common-dir")
|
||||
|
||||
|
||||
def _git_dir(path: Path) -> Optional[Path]:
|
||||
out = _kb._git_out(path, "rev-parse", "--path-format=absolute", "--git-dir")
|
||||
return Path(out).expanduser().resolve(strict=False) if out else None
|
||||
return _git_abs_path(path, "--git-dir")
|
||||
|
||||
|
||||
def _git_current_branch(path: Path) -> Optional[str]:
|
||||
@@ -401,9 +375,7 @@ def _git_current_branch(path: Path) -> Optional[str]:
|
||||
def _is_linked_worktree_checkout(path: Path) -> bool:
|
||||
git_dir = _git_dir(path)
|
||||
common_dir = _git_common_dir(path)
|
||||
if git_dir is None or common_dir is None:
|
||||
return False
|
||||
return git_dir != common_dir
|
||||
return git_dir is not None and common_dir is not None and git_dir != common_dir
|
||||
|
||||
|
||||
def _nearest_existing_path(path: Path) -> Path:
|
||||
@@ -428,25 +400,14 @@ def _ensure_git_worktree(repo_root: Path, target: Path, branch_name: str) -> Non
|
||||
"""Materialize ``target`` as a linked git worktree under ``repo_root``."""
|
||||
target = target.expanduser()
|
||||
repo_common = _git_common_dir(repo_root)
|
||||
if target.exists() and repo_common is not None:
|
||||
target_common = _git_common_dir(target)
|
||||
if target_common == repo_common:
|
||||
return
|
||||
if target.exists() and repo_common is not None and _git_common_dir(target) == repo_common:
|
||||
return
|
||||
target.parent.mkdir(parents=True, exist_ok=True)
|
||||
if _git_branch_exists(repo_root, branch_name):
|
||||
cmd = ["git", "-C", str(repo_root), "worktree", "add", str(target), branch_name]
|
||||
args = ["worktree", "add", str(target), branch_name]
|
||||
else:
|
||||
cmd = [
|
||||
"git", "-C", str(repo_root), "worktree", "add", "-b", branch_name,
|
||||
str(target), "HEAD",
|
||||
]
|
||||
result = subprocess.run(
|
||||
cmd,
|
||||
capture_output=True,
|
||||
text=True, encoding='utf-8', errors='replace',
|
||||
timeout=60,
|
||||
check=False,
|
||||
)
|
||||
args = ["worktree", "add", "-b", branch_name, str(target), "HEAD"]
|
||||
result = _git(repo_root, *args, timeout=60)
|
||||
if result.returncode != 0:
|
||||
stderr = (result.stderr or result.stdout or "").strip()
|
||||
raise RuntimeError(
|
||||
@@ -454,6 +415,13 @@ def _ensure_git_worktree(repo_root: Path, target: Path, branch_name: str) -> Non
|
||||
)
|
||||
|
||||
|
||||
def _anchored_worktree(repo_root: Path, task_id: str, branch_name: str) -> tuple[Path, str]:
|
||||
"""Materialize the canonical ``<repo>/.worktrees/<task-id>`` worktree."""
|
||||
target = repo_root / ".worktrees" / task_id
|
||||
_ensure_git_worktree(repo_root, target, branch_name)
|
||||
return target, branch_name
|
||||
|
||||
|
||||
def _resolve_worktree_workspace(
|
||||
task: Task, *, board: Optional[str] = None
|
||||
) -> tuple[Path, str]:
|
||||
@@ -485,9 +453,7 @@ def _resolve_worktree_workspace(
|
||||
f"task {task.id} has workspace_kind=worktree but board "
|
||||
f"{board_slug!r} default_workdir {board_default!r} is not inside a git repo"
|
||||
)
|
||||
target = repo_root / ".worktrees" / task.id
|
||||
_ensure_git_worktree(repo_root, target, branch_name)
|
||||
return target, branch_name
|
||||
return _anchored_worktree(repo_root, task.id, branch_name)
|
||||
|
||||
requested = Path(task.workspace_path).expanduser()
|
||||
if not requested.is_absolute():
|
||||
@@ -518,9 +484,7 @@ def _resolve_worktree_workspace(
|
||||
|
||||
repo_root = _git_toplevel(requested)
|
||||
if repo_root is not None and requested_resolved == repo_root:
|
||||
target = repo_root / ".worktrees" / task.id
|
||||
_ensure_git_worktree(repo_root, target, branch_name)
|
||||
return target, branch_name
|
||||
return _anchored_worktree(repo_root, task.id, branch_name)
|
||||
|
||||
repo_root = _repo_root_for_worktree_target(requested.parent)
|
||||
if repo_root is None:
|
||||
@@ -550,21 +514,20 @@ def resolve_workspace(task: Task, *, board: Optional[str] = None) -> Path:
|
||||
Persist the resolved path via ``set_workspace_path`` so later runs reuse it.
|
||||
"""
|
||||
kind = task.workspace_kind or "scratch"
|
||||
if kind == "scratch":
|
||||
if task.workspace_path:
|
||||
# Legacy explicit-path scratch tasks get the same absolute-path
|
||||
# guard as dir: — same threat model.
|
||||
p = Path(task.workspace_path).expanduser()
|
||||
if not p.is_absolute():
|
||||
raise ValueError(
|
||||
f"task {task.id} has non-absolute workspace_path "
|
||||
f"{task.workspace_path!r}; workspace paths must be absolute"
|
||||
)
|
||||
else:
|
||||
p = _kb.workspaces_root(board=board) / task.id
|
||||
p.mkdir(parents=True, exist_ok=True)
|
||||
return p
|
||||
if kind == "dir":
|
||||
if kind == "worktree":
|
||||
return _resolve_worktree_workspace(task, board=board)[0]
|
||||
if kind == "scratch" and not task.workspace_path:
|
||||
p = _kb.workspaces_root(board=board) / task.id
|
||||
elif kind == "scratch":
|
||||
# Legacy explicit-path scratch tasks get the same absolute-path guard
|
||||
# as dir: — same threat model.
|
||||
p = Path(task.workspace_path).expanduser()
|
||||
if not p.is_absolute():
|
||||
raise ValueError(
|
||||
f"task {task.id} has non-absolute workspace_path "
|
||||
f"{task.workspace_path!r}; workspace paths must be absolute"
|
||||
)
|
||||
elif kind == "dir":
|
||||
if not task.workspace_path:
|
||||
raise ValueError(
|
||||
f"task {task.id} has workspace_kind=dir but no workspace_path"
|
||||
@@ -576,34 +539,29 @@ def resolve_workspace(task: Task, *, board: Optional[str] = None) -> Path:
|
||||
f"{task.workspace_path!r}; use an absolute path "
|
||||
f"(relative paths are ambiguous against the dispatcher's CWD)"
|
||||
)
|
||||
p.mkdir(parents=True, exist_ok=True)
|
||||
return p
|
||||
if kind == "worktree":
|
||||
p, _branch_name = _resolve_worktree_workspace(task, board=board)
|
||||
return p
|
||||
raise ValueError(f"unknown workspace_kind: {kind}")
|
||||
else:
|
||||
raise ValueError(f"unknown workspace_kind: {kind}")
|
||||
p.mkdir(parents=True, exist_ok=True)
|
||||
return p
|
||||
|
||||
|
||||
def _set_task_column(conn: sqlite3.Connection, task_id: str, column: str, value: str) -> None:
|
||||
with _kb.write_txn(conn):
|
||||
conn.execute(f"UPDATE tasks SET {column} = ? WHERE id = ?", (value, task_id))
|
||||
|
||||
|
||||
def set_workspace_path(
|
||||
conn: sqlite3.Connection, task_id: str, path: Path | str
|
||||
) -> None:
|
||||
with _kb.write_txn(conn):
|
||||
conn.execute(
|
||||
"UPDATE tasks SET workspace_path = ? WHERE id = ?",
|
||||
(str(path), task_id),
|
||||
)
|
||||
_set_task_column(conn, task_id, "workspace_path", str(path))
|
||||
|
||||
|
||||
def set_branch_name(
|
||||
conn: sqlite3.Connection, task_id: str, branch_name: str
|
||||
) -> None:
|
||||
with _kb.write_txn(conn):
|
||||
conn.execute(
|
||||
"UPDATE tasks SET branch_name = ? WHERE id = ?",
|
||||
(str(branch_name), task_id),
|
||||
)
|
||||
_set_task_column(conn, task_id, "branch_name", str(branch_name))
|
||||
|
||||
|
||||
# Late-bound origin namespace (see module docstring). Imported LAST so this
|
||||
# Late-bound origin namespace (see module docstring); imported LAST so this
|
||||
# module is fully populated before ``kanban_db`` re-exports from it.
|
||||
from hermes_cli import kanban_db as _kb # noqa: E402
|
||||
|
||||
Reference in New Issue
Block a user