feat(skills): implement state loader DRY (OP-5) and Job-centric directory structure

- Extract load_state_json centralized helper inside lib.sh to unify state querying
- Refactor status, resume, stop, and monitor scripts to fetch state via MAM_STATE_JSON env var to avoid stdin pipeline collisions
- Restructure delegate-job to provision .mam/jobs/<job_id>/brief.md and direct agents to it, minimizing token size and preventing TUI paste freezes
- Harden send_keys_safe submission loop with was_popup state capture and edge case guards, preventing timing spin false-positives
- Passed cross-verification approved PASS from Planner Claude session
This commit is contained in:
2026-07-12 11:50:38 +09:00
parent 76ec0dd929
commit 35af8e33a2
7 changed files with 138 additions and 252 deletions
+62 -97
View File
@@ -119,52 +119,50 @@ tmux() {
# Fallback to TMUX_SERVER_NAME or 'default' if not registered or field is missing.
# Prints the resolved server name on stdout.
# ---------------------------------------------------------------------------
resolve_tmux_server() {
local session_name="$1"
SESSION_NAME="$session_name" env_python "$AGENT_SESSIONS_YAML" <<'PYEOF'
load_state_json() {
env_python "$AGENT_SESSIONS_YAML" <<'PYEOF'
import os, sys, sqlite3, json, yaml
name = os.environ['SESSION_NAME']
yaml_path = os.environ['YAML_PATH']
db_path = os.path.splitext(yaml_path)[0] + '.db'
d = {}
db_sessions = []
try:
if os.path.exists(db_path):
conn = sqlite3.connect(db_path, timeout=60.0)
try:
row = conn.execute('SELECT data FROM sessions WHERE name=?', (name,)).fetchone()
if row:
s = json.loads(row[0])
server = s.get('tmux_server')
if server:
print(server)
sys.exit(0)
except sqlite3.OperationalError:
pass
row = conn.execute('SELECT data FROM state WHERE id=1').fetchone()
if row:
d = json.loads(row[0])
for s in d.get('tmux_sessions', []):
if s.get('name') == name:
server = s.get('tmux_server')
if server:
print(server)
sys.exit(0)
try:
cursor = conn.execute('SELECT data FROM sessions')
for r in cursor.fetchall():
db_sessions.append(json.loads(r[0]))
d['tmux_sessions'] = db_sessions
except sqlite3.OperationalError:
pass
conn.close()
elif os.path.exists(yaml_path):
with open(yaml_path) as f:
d = yaml.safe_load(f) or {}
for s in d.get('tmux_sessions', []):
if s.get('name') == name:
server = s.get('tmux_server')
if server:
print(server)
sys.exit(0)
except Exception:
pass
# Fallback
print(os.environ.get('TMUX_SERVER_NAME', 'default'))
print(json.dumps(d, ensure_ascii=False))
PYEOF
}
resolve_tmux_server() {
local session_name="$1"
MAM_STATE_JSON="$(load_state_json)" SESSION_NAME="$session_name" python3 -c "
import sys, os, json
name = os.environ['SESSION_NAME']
d = json.loads(os.environ.get('MAM_STATE_JSON', '{}'))
for s in d.get('tmux_sessions', []):
if s.get('name') == name:
print(s.get('tmux_server', 'default'))
sys.exit(0)
print(os.environ.get('TMUX_SERVER_NAME', 'default'))
"
}
# ---------------------------------------------------------------------------
# derive_session_name <workspace> <agent>
#
@@ -495,15 +493,12 @@ PYEOF
find_workspace_uuid() {
local workspace="$1" agent="$2" session_name="${3:-}"
local abs; abs="$(cd "$workspace" 2>/dev/null && pwd)" || abs="$workspace"
WS_ABS="$abs" AGENT="$agent" TARGET_SESSION="$session_name" env_python "$AGENT_SESSIONS_YAML" <<'PYEOF'
import os, json, glob, sqlite3
import yaml
MAM_STATE_JSON="$(load_state_json)" WS_ABS="$abs" AGENT="$agent" TARGET_SESSION="$session_name" env_python "$AGENT_SESSIONS_YAML" <<'PYEOF'
import os, sys, json, glob, sqlite3
ws = os.environ['WS_ABS']
agent = os.environ['AGENT']
home = os.environ['HOME_DIR']
yaml_path = os.environ['YAML_PATH']
db_path = os.path.splitext(yaml_path)[0] + '.db'
claude_project_dir = os.environ.get('CLAUDE_PROJECT_DIR', f"{home}/.claude/projects")
target = os.environ.get('TARGET_SESSION', '')
@@ -512,13 +507,10 @@ OWN_KEY = {'claude': 'claude_session_id_own', 'agy': 'agy_conversation_id_own',
def iso_root_of(s):
# T5: per-row isolation root (all-L2). None => legacy global state dirs.
iso = s.get('isolation')
return iso.get('root') if isinstance(iso, dict) else None
# exists checks take an optional isolation root; path templates per lever
# (Phase 0 measured — note cline's isolated layout is <root>/sessions/, NOT data/sessions/).
def jsonl_exists(uuid, iso=None):
base = f"{iso}/projects" if iso else claude_project_dir
key = ws.replace('/', '-').replace('_', '-')
@@ -553,37 +545,21 @@ def own_exists(a, uuid, iso=None):
'hermes': hermes_exists, 'cline': cline_exists}[a](uuid, iso)
running_ids = set()
# Ingest the merged state dictionary from stdin pipeline
try:
if os.path.exists(db_path):
conn = sqlite3.connect(db_path, timeout=60.0)
try:
cursor = conn.execute("SELECT data FROM sessions WHERE status='running'")
for r_row in cursor.fetchall():
s_data = json.loads(r_row[0])
# the target row's own ids are legitimately "claimed by itself" — don't filter them
if target and s_data.get('name') == target:
continue
for k in ['claude_session_id_own', 'agy_conversation_id_own', 'hermes_conversation_id_own', 'cline_conversation_id_own']:
val = s_data.get(k)
if val:
running_ids.add(val)
except sqlite3.OperationalError:
pass
conn.close()
if not running_ids and os.path.exists(yaml_path):
with open(yaml_path) as f:
d_state = yaml.safe_load(f) or {}
for s_item in d_state.get('tmux_sessions', []):
if s_item.get('status') == 'running':
if target and s_item.get('name') == target:
continue
for k in ['claude_session_id_own', 'agy_conversation_id_own', 'hermes_conversation_id_own', 'cline_conversation_id_own']:
val = s_item.get(k)
if val:
running_ids.add(val)
d = json.loads(os.environ.get('MAM_STATE_JSON', '{}'))
except Exception:
pass
d = {}
running_ids = set()
for s_item in d.get('tmux_sessions', []):
if s_item.get('status') == 'running':
if target and s_item.get('name') == target:
continue
for k in ['claude_session_id_own', 'agy_conversation_id_own', 'hermes_conversation_id_own', 'cline_conversation_id_own']:
val = s_item.get(k)
if val:
running_ids.add(val)
def emit(u):
@@ -593,35 +569,12 @@ def emit(u):
raise SystemExit(0)
# 1) per-row own id for THIS workspace (optimized with direct sqlite query if db exists)
# Fetch all sessions matching this workspace
sessions = []
try:
if os.path.exists(db_path):
conn = sqlite3.connect(db_path, timeout=60.0)
has_sessions_table = False
try:
cursor = conn.execute('SELECT data FROM sessions WHERE pane_cwd=?', (ws,))
for row in cursor.fetchall():
sessions.append(json.loads(row[0]))
has_sessions_table = True
except sqlite3.OperationalError:
pass
if not has_sessions_table or not sessions:
row = conn.execute('SELECT data FROM state WHERE id=1').fetchone()
if row:
d = json.loads(row[0])
for s in d.get('tmux_sessions', []):
if isinstance(s, dict) and (s.get('pane') or {}).get('cwd') == ws:
sessions.append(s)
conn.close()
elif os.path.exists(yaml_path):
with open(yaml_path) as f:
d = yaml.safe_load(f) or {}
for s in d.get('tmux_sessions', []):
if isinstance(s, dict) and (s.get('pane') or {}).get('cwd') == ws:
sessions.append(s)
except Exception:
pass
for s in d.get('tmux_sessions', []):
if isinstance(s, dict) and (s.get('pane') or {}).get('cwd') == ws:
sessions.append(s)
# T5: target-session mode — resolve strictly within the named row's scope.
# An isolated row's conversation lives ONLY in its isolation root (which holds
@@ -1191,15 +1144,27 @@ send_keys_safe() {
_sks_tmux paste-buffer -b "sks_$job_id" -t "$sess"
_sks_tmux delete-buffer -b "sks_$job_id" 2>/dev/null || true
sleep 0.5
_pane_capture "$sess" | grep -Fq "$marker" || { echo "send_keys_safe: paste not visible ($sess)" >&2; return 3; }
local pane_content was_popup=0
pane_content=$(_pane_capture "$sess")
if printf '%s\n' "$pane_content" | grep -Eq "Pasted text|paste again to expand"; then
was_popup=1
elif ! printf '%s\n' "$pane_content" | grep -Fq "$marker"; then
echo "send_keys_safe: paste not visible ($sess)" >&2
return 3
fi
for try in 1 2 3; do
pre_submit=$(_pane_capture "$sess")
_sks_tmux send-keys -t "$sess" C-m
sleep "$try"
# Submitted = input area released the text AND rendering changed after Enter.
if ! _pane_tail "$sess" 3 | grep -Fq "$marker" \
&& [ "$(_pane_capture "$sess")" != "$pre_submit" ]; then
local cur_content
cur_content=$(_pane_capture "$sess")
# Hardened Submission Checks
if [ "$was_popup" = "0" ] && ! _pane_tail "$sess" 3 | grep -Fq "$marker" && [ "$cur_content" != "$pre_submit" ]; then
return 0
elif [ "$was_popup" = "1" ] && \
! printf '%s\n' "$cur_content" | grep -Eq "Pasted text|paste again to expand" && \
[ "$cur_content" != "$pre_submit" ]; then
return 0
fi
done
@@ -107,6 +107,25 @@ cmd_submit() {
--max-iterations "$MAX_ITERATIONS")"
echo "registered job: $JOB_ID"
# 1-1) Provision job directory and write direct brief.md (MAM Job Restructuring)
local job_dir="$REGISTRY_DIR/$JOB_ID"
if [[ "$DRY_RUN" != "1" ]]; then
mkdir -p "$job_dir"
cat <<EOF > "$job_dir/brief.md"
# 📋 Brief: Job $JOB_ID Delegation
- **Job ID**: $JOB_ID
- **Target Agent**: $AGENT (session: $AGENT_SESSION)
- **Role**: Worker
- **Timeout**: $TIMEOUT s (Idle: $IDLE_TIMEOUT s)
- **Output Report Path**: .mam/jobs/$JOB_ID/$AGENT-reports/report-final.md
## 🔎 Task Description
$PROMPT
EOF
echo "provisioned job directory: $job_dir"
fi
if [[ "$TYPE" == "direct" ]]; then
# 2) START THE SUBSCRIBER FIRST (ordering dependency — MQTT does not queue
# non-retained messages for absent subscribers).
@@ -140,17 +159,7 @@ cmd_submit() {
# an id from an earlier session is the #1 reason a delegated job sits idle and
# times out (see SKILL.md "Wrong job_id propagated to the agent"). We make the
# freshness explicit in the instruction header.
local instructions="Your job_id is \"$JOB_ID\" (the one just registered for THIS delegation — read it from the registry record, do NOT reuse any job_id you saw in earlier runs).
On start run: $pub --event started.
On permission/tool prompt run: $pub --event permission_required --detail '<tool>:<what>'.
On progress (optional): $pub --event progress --detail '<short status>'.
On success run: $pub --event completed --detail '<one-line summary>'.
On failure run: $pub --event error --detail '<one-line reason>'.
The subscriber for this job_id is already running; your completed/error event ends the job. Exit codes: 0 completed, 1 error, 2 publish failure.
Task: $PROMPT"
local instructions="Your job_id is \"$JOB_ID\". Detailed task requirements, instructions, and target output paths are documented in the task brief file at: .mam/jobs/$JOB_ID/brief.md. Please READ and follow .mam/jobs/$JOB_ID/brief.md to complete your work. Commands: start='$pub --event started', success='$pub --event completed --detail <summary>', error='$pub --event error --detail <reason>'."
run_agent "$JOB_ID" "$instructions"
@@ -203,6 +212,25 @@ Task: $PROMPT"
echo "Session: $current_session"
echo "=================================================="
# 1-1) Provision job directory and write iteration brief.md (MAM Job Restructuring)
local job_dir="$REGISTRY_DIR/$JOB_ID"
if [[ "$DRY_RUN" != "1" ]]; then
mkdir -p "$job_dir"
cat <<EOF > "$job_dir/brief.md"
# 📋 Brief: Job $JOB_ID Delegation (Iteration $iteration)
- **Job ID**: $JOB_ID
- **Target Agent/Session**: $current_session
- **Role**: $current_role
- **Iteration**: $iteration
- **Output Report Path**: .mam/jobs/$JOB_ID/${current_session}-reports/report-final.md
## 🔎 Task Description
$current_prompt
EOF
echo "provisioned iteration brief: $job_dir/brief.md"
fi
# Update job details in registry
"$PY" "$SCRIPT_DIR/scripts/registry.py" --registry-dir "$REGISTRY_DIR" update \
--job "$JOB_ID" \
@@ -238,17 +266,7 @@ Task: $PROMPT"
# Format instruction block
local pub="$PY $SCRIPT_DIR/scripts/publish_event.py --registry-dir $REGISTRY_DIR --job $JOB_ID"
local instructions="Your job_id is \"$JOB_ID\" (the one just registered for THIS delegation — read it from the registry record, do NOT reuse any job_id you saw in earlier runs).
On start run: $pub --event started.
On permission/tool prompt run: $pub --event permission_required --detail '<tool>:<what>'.
On progress (optional): $pub --event progress --detail '<short status>'.
On success run: $pub --event completed --detail '<one-line summary>'.
On failure run: $pub --event error --detail '<one-line reason>'.
The subscriber for this job_id is already running; your completed/error event ends the job. Exit codes: 0 completed, 1 error, 2 publish failure.
Task: $current_prompt"
local instructions="Your job_id is \"$JOB_ID\". Detailed task requirements, instructions, and target output paths for iteration $iteration are documented in the task brief file at: .mam/jobs/$JOB_ID/brief.md. Please READ and follow .mam/jobs/$JOB_ID/brief.md to complete your work. Commands: start='$pub --event started', success='$pub --event completed --detail <summary>', error='$pub --event error --detail <reason>'."
# Trigger agent
run_agent "$JOB_ID" "$instructions" "$current_session"
@@ -389,8 +407,8 @@ run_agent() {
# Before launching the agent, set up error trap to publish error event
if [ -n "${job_id:-}" ] && [ -n "${PY:-}" ]; then
local pub_script="$SCRIPT_DIR/scripts/publish_event.py"
trap 'rc=$?; if [ $rc -ne 0 ]; then "$PY" "$pub_script" --job "$job_id" --event error --detail "agent bootstrap failed (exit $rc)"; fi' EXIT
pub_script="$SCRIPT_DIR/scripts/publish_event.py"
trap "rc=\$?; if [ \$rc -ne 0 ]; then \"$PY\" \"$pub_script\" --job '$job_id' --event error --detail 'agent bootstrap failed (exit '\$rc')'; fi" EXIT
fi
echo "살아있는 에이전트 세션 '$sess'에 작업을 위임합니다..."
@@ -15,6 +15,7 @@
set -euo pipefail
source "$(cd "$(dirname "${BASH_SOURCE[0]}")/../.." && pwd)/lib.sh"
export WORKSPACE_ROOT
STATE_DIR="${AGENT_SESSIONS_STATE_DIR:-$WORKSPACE_ROOT/.cache/multi-agent-mux-monitor}"
@@ -304,31 +305,18 @@ claude_project_dir = os.environ.get('CLAUDE_PROJECT_DIR', f"{home}/.claude/proje
now_iso = datetime.now(timezone.utc).strftime('%Y-%m-%dT%H:%M:%SZ')
# atomic 래퍼에서는 d 가 이미 로드돼 있음. env_python(dry-run)에서는 여기서 로드.
try:
d
except NameError:
import sqlite3
db_path = os.path.splitext(yaml_path)[0] + '.db'
import subprocess
d = {}
try:
if os.path.exists(db_path):
conn = sqlite3.connect(db_path, timeout=60.0)
row = conn.execute('SELECT data FROM state WHERE id=1').fetchone()
if row: d = json.loads(row[0])
try:
db_sessions = []
cursor = conn.execute('SELECT data FROM sessions')
for s_row in cursor.fetchall():
db_sessions.append(json.loads(s_row[0]))
d['tmux_sessions'] = db_sessions
except sqlite3.OperationalError:
pass
conn.close()
elif os.path.exists(yaml_path):
with open(yaml_path) as f:
d = yaml.safe_load(f) or {}
ws_root = os.environ.get('WORKSPACE_ROOT')
if not ws_root:
ws_root = os.path.abspath(os.path.join(os.path.dirname(__file__), '../../../..'))
script = f"source '{ws_root}/.agents/skills/lib.sh' && load_state_json"
out = subprocess.check_output(['bash', '-c', script], stderr=subprocess.DEVNULL)
d = json.loads(out.decode('utf-8'))
except Exception:
pass
@@ -57,40 +57,15 @@ if { [ "$AGENT" = "agy" ] || [ "$AGENT" = "hermes" ] || [ "$AGENT" = "cline" ];
CHILD_PID="${CHILD_PID:-0}"
fi
DELEGATE_JOB_ID=$(env_python "$AGENT_SESSIONS_YAML" SESSION_NAME="$SESSION_NAME" <<'PYEOF'
import os, sys, sqlite3, json, yaml
DELEGATE_JOB_ID=$(MAM_STATE_JSON="$(load_state_json)" SESSION_NAME="$SESSION_NAME" python3 -c "
import sys, os, json
name = os.environ['SESSION_NAME']
yaml_path = os.environ['YAML_PATH']
db_path = os.path.splitext(yaml_path)[0] + '.db'
d = {}
try:
if os.path.exists(db_path):
conn = sqlite3.connect(db_path, timeout=60.0)
try:
row = conn.execute('SELECT data FROM sessions WHERE name=?', (name,)).fetchone()
if row:
s = json.loads(row[0])
print(s.get('delegate_job_id', '') or '')
raise SystemExit(0)
except sqlite3.OperationalError:
pass
row = conn.execute('SELECT data FROM state WHERE id=1').fetchone()
if row:
d = json.loads(row[0])
conn.close()
elif os.path.exists(yaml_path):
with open(yaml_path) as f:
d = yaml.safe_load(f) or {}
except Exception:
pass
d = json.loads(os.environ.get('MAM_STATE_JSON', '{}'))
for s in d.get('tmux_sessions', []):
if s.get('name') == name:
print(s.get('delegate_job_id', '') or '')
raise SystemExit(0)
raise SystemExit(0)
PYEOF
)
sys.exit(0)
")
atomic_dump_yaml "$AGENT_SESSIONS_YAML" \
SESSION_NAME="$SESSION_NAME" UUID="$UUID" AGENT="$AGENT" NOW_ISO="$NOW_ISO" \
@@ -28,38 +28,17 @@ fi
# Resolved relative to this script — no hardcoded absolute path (review item 6).
PROJECT_ROOT="$(cd "$(dirname "${BASH_SOURCE[0]}")/../../../../" && pwd)"
DRIFT_JSON="$DRIFT_JSON" env_python "$AGENT_SESSIONS_YAML" PROJECT_ROOT="$PROJECT_ROOT" <<'PYEOF'
import os, json, glob
import yaml
MAM_STATE_JSON="$(load_state_json)" DRIFT_JSON="$DRIFT_JSON" env_python "$AGENT_SESSIONS_YAML" PROJECT_ROOT="$PROJECT_ROOT" <<'PYEOF'
import os, sys, json, glob
yaml_path = os.environ['YAML_PATH']
home = os.environ['HOME_DIR']
claude_project_dir = os.environ.get('CLAUDE_PROJECT_DIR', f"{home}/.claude/projects")
drift = json.loads(os.environ['DRIFT_JSON'])
db_path = os.path.splitext(yaml_path)[0] + '.db'
d = {}
import sqlite3
try:
if os.path.exists(db_path):
conn = sqlite3.connect(db_path, timeout=60.0)
row = conn.execute('SELECT data FROM state WHERE id=1').fetchone()
if row: d = json.loads(row[0])
try:
db_sessions = []
cursor = conn.execute('SELECT data FROM sessions')
for s_row in cursor.fetchall():
db_sessions.append(json.loads(s_row[0]))
d['tmux_sessions'] = db_sessions
except sqlite3.OperationalError:
pass
conn.close()
elif os.path.exists(yaml_path):
with open(yaml_path) as f:
d = yaml.safe_load(f) or {}
d = json.loads(os.environ.get('MAM_STATE_JSON', '{}'))
except Exception:
pass
d = {}
alive = set(drift.get('tmux_sessions_alive', []))
drift_by_name = {}
@@ -83,44 +83,18 @@ fi
# 세션이 YAML 에 있는지 + 해당 row 의 워크스페이스 cwd 및 delegate_job_id 추출.
# JSON 으로 emit — cwd 에 '|' 가 들어가도 안전 (review item 7; 기존 cwd|jid 파서 대체).
MAPPED_DATA=$(env_python "$AGENT_SESSIONS_YAML" SESSION_NAME="$SESSION_NAME" <<'PYEOF'
import os, sys, json, yaml, sqlite3
MAPPED_DATA=$(MAM_STATE_JSON="$(load_state_json)" SESSION_NAME="$SESSION_NAME" python3 -c "
import sys, os, json
name = os.environ['SESSION_NAME']
yaml_path = os.environ['YAML_PATH']
db_path = os.path.splitext(yaml_path)[0] + '.db'
d = {}
try:
if os.path.exists(db_path):
conn = sqlite3.connect(db_path, timeout=60.0)
try:
row = conn.execute('SELECT data FROM sessions WHERE name=?', (name,)).fetchone()
if row:
s = json.loads(row[0])
cwd = (s.get('pane') or {}).get('cwd', '')
jid = s.get('delegate_job_id', '') or ''
print(json.dumps({"cwd": cwd, "job_id": jid}))
raise SystemExit(0)
except sqlite3.OperationalError:
pass
row = conn.execute('SELECT data FROM state WHERE id=1').fetchone()
if row:
d = json.loads(row[0])
conn.close()
elif os.path.exists(yaml_path):
with open(yaml_path) as f:
d = yaml.safe_load(f) or {}
except Exception:
pass
d = json.loads(os.environ.get('MAM_STATE_JSON', '{}'))
for s in d.get('tmux_sessions', []):
if s.get('name') == name:
cwd = (s.get('pane') or {}).get('cwd', '')
jid = s.get('delegate_job_id', '') or ''
print(json.dumps({"cwd": cwd, "job_id": jid}))
raise SystemExit(0)
raise SystemExit(7)
PYEOF
) || {
print(json.dumps({'cwd': cwd, 'job_id': jid}))
sys.exit(0)
sys.exit(7)
") || {
echo "ERROR: session '$SESSION_NAME' not in $AGENT_SESSIONS_YAML" >&2
exit 1
}