Files
pulse-bridge/bridge.py
T
Nirav PatelandClaude Fable 5 aaedf758d9 suppress outbound mirror comments for intake-managed items (double-comment fix)
intake_track() already narrates status on the user's post; mirroring the
same item's ledger update/close ops added near-duplicate comments (seen in
the E2E smoke: 'stage → build' next to '▶️ running').

Co-Authored-By: Claude Fable 5 <[email protected]>
Claude-Session: https://claude.ai/code/session_01YECpkAwQUwgu7NVy91R8fW
2026-08-22 15:06:19 -04:00

578 lines
24 KiB
Python
Raw Blame History

This file contains invisible Unicode characters
This file contains invisible Unicode characters that are indistinguishable to humans but may be processed differently by a computer. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
#!/usr/bin/env python3
"""pulse-bridge v0 — mirror the estate's shared work ledger into Pulse posts.
Source: ~/.local/bin/shre-items (JSONL ledger, `feed --since` op stream).
Target: local mib007 (:5520) comms API — the backing store of the Pulse /
Activity feed. One ledger item = one post in the "activity" channel; stage
changes / failures / escalations = threaded comments; close = final comment.
Stdlib only. State (checkpoint + ledger-id -> post-id map) lives in
~/.shre/pulse-bridge/state.json, written atomically (tmp + rename).
"""
import json
import os
import re
import subprocess
import sys
import time
import urllib.error
import urllib.request
from datetime import datetime, timezone
from pathlib import Path
# ── Config (env file, then environment, then defaults) ───────────────────────
HOME = Path.home()
ENV_FILE = Path(os.environ.get("PULSE_BRIDGE_ENV",
Path(__file__).resolve().parent / "bridge.env"))
def _load_env_file(path: Path) -> dict:
conf = {}
if path.exists():
for line in path.read_text(encoding="utf-8").splitlines():
line = line.strip()
if not line or line.startswith("#") or "=" not in line:
continue
k, v = line.split("=", 1)
conf[k.strip()] = v.strip().strip('"').strip("'")
return conf
_FILE = _load_env_file(ENV_FILE)
def cfg(key: str, default: str) -> str:
return os.environ.get(key) or _FILE.get(key) or default
POLL_SEC = int(cfg("PULSE_BRIDGE_POLL_SEC", "30"))
MIB_BASE = cfg("PULSE_BRIDGE_MIB_BASE", "http://127.0.0.1:5520").rstrip("/")
WORKSPACE_ID = cfg("PULSE_BRIDGE_WORKSPACE_ID",
"293d29db-4978-4c54-a92e-b24d4c0a7115") # Nirlab Command Center
CHANNEL_NAME = cfg("PULSE_BRIDGE_CHANNEL", "activity")
MENTION = cfg("PULSE_BRIDGE_MENTION", "nir") # user handle to @mention
POSTER_NAME = cfg("PULSE_BRIDGE_POSTER_NAME", "Ledger")
SHRE_ITEMS = cfg("PULSE_BRIDGE_SHRE_ITEMS", str(HOME / ".local/bin/shre-items"))
STATE_FILE = Path(cfg("PULSE_BRIDGE_STATE",
str(HOME / ".shre/pulse-bridge/state.json")))
TOKEN_FILE = Path(cfg("PULSE_BRIDGE_TOKEN_FILE",
str(HOME / ".shre/service-tokens.json")))
TOKEN_KEY = cfg("PULSE_BRIDGE_TOKEN_KEY", "mib007")
STAGES = ["queued", "build", "review", "merge", "deploy", "verify"]
# ── Intake (reverse direction: Pulse post -> ledger item + kanban task) ──────
INTAKE_ENABLED = cfg("PULSE_BRIDGE_INTAKE", "1") == "1"
INTAKE_MAX_CONCURRENT = int(cfg("PULSE_BRIDGE_INTAKE_MAX_CONCURRENT", "3"))
HERMES_DIR = cfg("PULSE_BRIDGE_HERMES_DIR", str(HOME / ".hermes/hermes-agent"))
HERMES_HOME_DIR = cfg("PULSE_BRIDGE_HERMES_HOME", str(HOME / ".hermes"))
# user_names never ingested (service accounts / the bridge itself)
INTAKE_IGNORE_USERS = {u.strip().lower() for u in
cfg("PULSE_BRIDGE_INTAKE_IGNORE",
f"{POSTER_NAME},Ellie,AROS").split(",") if u.strip()}
def log(msg: str) -> None:
print(f"{datetime.now(timezone.utc).isoformat(timespec='seconds')} {msg}",
flush=True)
# ── mib007 HTTP ──────────────────────────────────────────────────────────────
def _token() -> str:
return json.loads(TOKEN_FILE.read_text(encoding="utf-8"))[TOKEN_KEY]
def api(method: str, path: str, body=None, timeout=15):
req = urllib.request.Request(
f"{MIB_BASE}{path}",
data=json.dumps(body).encode() if body is not None else None,
method=method,
headers={
"Authorization": f"Bearer {_token()}",
"Content-Type": "application/json",
},
)
with urllib.request.urlopen(req, timeout=timeout) as resp:
raw = resp.read()
return json.loads(raw) if raw else None
# ── State ────────────────────────────────────────────────────────────────────
def load_state() -> dict:
if STATE_FILE.exists():
try:
return json.loads(STATE_FILE.read_text(encoding="utf-8"))
except json.JSONDecodeError:
# Starting fresh here would re-seed every open item -> a duplicate
# post per item. Refuse to run until a human inspects the file.
log(f"FATAL state file corrupt: {STATE_FILE} — fix or remove it "
"deliberately; refusing to reseed automatically")
sys.exit(1)
return {"checkpoint": "", "seen": [], "map": {}, "seeded": False,
"channel_id": "", "intake_checkpoint": 0, "intake": {}}
def save_state(state: dict) -> None:
STATE_FILE.parent.mkdir(parents=True, exist_ok=True)
tmp = STATE_FILE.with_suffix(".json.tmp")
with open(tmp, "w", encoding="utf-8") as f:
json.dump(state, f, indent=1)
f.flush()
os.fsync(f.fileno())
os.replace(tmp, STATE_FILE)
# ── Content helpers ──────────────────────────────────────────────────────────
def neutralize(text: str) -> str:
"""Insert a zero-width space after '@' in ledger-derived text.
mib007's comms route treats the FIRST '@handle' in a message as an agent
mention: a matching agent fires an AI reply into the thread, and on this
instance ANY mention 500s after the insert (the agent lookup references a
nonexistent url_key column). So no bridge-emitted text may ever contain a
bare '@handle' — including the bridge's own needs-you escalation, which is
emitted zwsp-neutralised and renders identically in the UI.
"""
return (text or "").replace("@", "@")
def tag_line(item: dict) -> str:
bits = [f"#{item.get('kind', 'item')}"]
if item.get("stage"):
bits.append(f"#{item['stage']}")
for t in item.get("tags", []) or []:
bits.append(f"#{t}")
if item.get("surface"):
bits.append(f"surface:{item['surface']}")
if item.get("project"):
bits.append(f"project:{item['project']}")
bits.append(f"ledger:{item['id']}")
return " · ".join(bits)
def post_body(item: dict) -> str:
head = {"needs-you": "🙋", "failed": "⛔", "pipeline": "🔧",
"recommended": "💡", "gap": "🕳"}.get(item.get("kind", ""), "📌")
parts = [f"{head} {neutralize(item.get('title', '(untitled)'))}"]
if item.get("detail"):
parts.append(neutralize(item["detail"]))
parts.append(tag_line(item))
return "\n\n".join(parts)
def update_body(rec: dict) -> str:
lines = []
kind = rec.get("kind")
stage = rec.get("stage")
note = neutralize(rec.get("note", ""))
if kind == "needs-you":
lines.append(neutralize(f"@{MENTION}") + " NEEDS YOU: this item is now"
" waiting on a human." + (f" — {note}" if note else ""))
note = "" # already included
elif kind == "failed":
where = f" at {stage}" if stage else ""
attempt = f", attempt {rec['attempt']}" if rec.get("attempt") else ""
lines.append(f"⛔ failed{where}{attempt}" + (f" — {note}" if note else ""))
note = ""
else:
if kind:
lines.append(f"kind → {kind}")
if stage:
lines.append(f"stage → {stage}")
if rec.get("attempt") and not kind:
lines.append(f"attempt {rec['attempt']}")
if rec.get("detail"):
lines.append(f"detail updated: {neutralize(rec['detail'])}")
if note:
lines.append(note)
return "\n".join(lines) or "(updated)"
def close_body(rec: dict) -> str:
why = neutralize(rec.get("why", ""))
if rec.get("status") == "dropped":
return "🗑 dropped" + (f" — {why}" if why else "")
return "✅ done" + (f" — {why}" if why else "")
# ── Ledger access ────────────────────────────────────────────────────────────
def shre_items(*args: str) -> str:
out = subprocess.run([SHRE_ITEMS, *args], capture_output=True, text=True,
timeout=30)
if out.returncode != 0:
raise RuntimeError(f"shre-items {' '.join(args)} rc={out.returncode}: "
f"{out.stderr.strip()[:300]}")
return out.stdout
def feed_since(checkpoint: str) -> list:
args = ["feed"] + (["--since", checkpoint] if checkpoint else [])
records = []
for line in shre_items(*args).splitlines():
line = line.strip()
if not line:
continue
try:
records.append(json.loads(line))
except json.JSONDecodeError:
continue
records.sort(key=lambda r: r.get("at", ""))
return records
def open_items() -> list:
return json.loads(shre_items("list", "--json") or "[]")
# ── Pulse actions ────────────────────────────────────────────────────────────
def ensure_channel(state: dict) -> str:
if state.get("channel_id"):
return state["channel_id"]
channels = api("GET", f"/api/workspaces/{WORKSPACE_ID}/comms/channels") or []
for ch in channels:
if (ch.get("name") or "").lower() == CHANNEL_NAME.lower():
state["channel_id"] = ch["id"]
return ch["id"]
ch = api("POST", f"/api/workspaces/{WORKSPACE_ID}/comms/channels",
{"name": CHANNEL_NAME, "description": "Workspace activity feed",
"type": "public"})
state["channel_id"] = ch["id"]
return ch["id"]
def send_message(channel_id: str, content: str, thread_id=None) -> dict:
body = {"content": content, "type": "text", "userName": POSTER_NAME}
if thread_id:
body["threadId"] = thread_id
return api("POST",
f"/api/workspaces/{WORKSPACE_ID}/comms/channels/{channel_id}/messages",
body)
def find_existing_post(state: dict, ledger_id: str):
"""Recovery lookup: the target API can insert a row and then fail, so a
retried create must first check whether a post tagged ledger:<id> already
exists. Bounded to the most recent 200 channel messages."""
try:
msgs = api("GET", f"/api/workspaces/{WORKSPACE_ID}/comms/channels/"
f"{ensure_channel(state)}/messages?limit=200") or []
except Exception:
return None
needle = f"ledger:{ledger_id}"
for m in msgs:
if not m.get("thread_id") and needle in (m.get("content") or ""):
return m.get("id")
return None
def create_post(state: dict, item: dict) -> None:
if item["id"] in state["map"]:
return # already mirrored (seed/feed overlap)
try:
msg = send_message(ensure_channel(state), post_body(item))
except urllib.error.HTTPError:
# The insert may have landed before the error (proven insert-then-500
# behaviour). Adopt an existing post instead of duplicating on retry.
existing = find_existing_post(state, item["id"])
if existing:
state["map"][item["id"]] = existing
log(f"post {item['id'][:8]} -> adopted existing {existing}")
return
raise # genuinely not created; retry next poll
state["map"][item["id"]] = msg["id"]
log(f"post {item['id'][:8]} -> {msg['id']} ({item.get('title', '')[:60]})")
def comment(state: dict, ledger_id: str, content: str, what: str) -> None:
post_id = state["map"].get(ledger_id)
if not post_id:
log(f"SKIP {what} {ledger_id[:8]}: no mapped post (predates bridge)")
return
try:
msg = send_message(ensure_channel(state), content, thread_id=post_id)
except urllib.error.HTTPError as e:
# 500: mib007 comms can fail AFTER the insert (proven mention-path
# crash) — retrying risks duplicates, so drop. 404/410: the post is
# gone — retrying can never succeed. Everything else (401/403/429/
# 502/503) is safe to retry next poll.
if e.code in (500, 404, 410):
log(f"DROP {what} {ledger_id[:8]}: HTTP {e.code} "
"(comment not retried)")
return
raise
log(f"{what:6} {ledger_id[:8]} -> comment {msg['id']}")
# ── Intake: Pulse posts -> ledger + kanban ───────────────────────────────────
def kanban(*args: str, timeout: int = 60) -> str:
out = subprocess.run(
[f"{HERMES_DIR}/venv/bin/python", "-m", "shiva_cli.main",
"kanban", *args],
capture_output=True, text=True, timeout=timeout, cwd=HERMES_DIR,
env={**os.environ, "HERMES_HOME": HERMES_HOME_DIR})
if out.returncode != 0:
raise RuntimeError(f"kanban {args[0]} rc={out.returncode}: "
f"{out.stderr.strip()[:300]}")
return out.stdout
def ledger_add_for_post(post: dict) -> str:
"""Create (or adopt) a pipeline ledger item for a Pulse post; return its
FULL id. Idempotent: a crash after `shre-items add` but before the state
save must not create a second item on retry, so existing open items are
searched for this post's marker first."""
marker = f"post:{post['id']}"
for item in open_items():
if marker in (item.get("detail") or ""):
return item["id"]
content = (post.get("content") or "").strip()
title = " ".join(content.split())[:80] or "(untitled pulse post)"
detail = f"{content}\n\n[pulse-intake] {marker} " \
f"author:{post.get('user_name') or post.get('user_id')}"
out = shre_items("add", "--kind", "pipeline", "--stage", "queued",
"--title", title, "--detail", detail,
"--surface", "pulse", "--tag", "pulse-intake")
short = out.split()[0]
for item in open_items():
if item["id"].startswith(short):
return item["id"]
raise RuntimeError(f"cannot resolve full ledger id for {short}")
def intake_active_count(state: dict) -> int:
return sum(1 for e in state["intake"].values()
if e.get("task_id") and not e.get("closed"))
def intake_start_task(state: dict, post_id: str, entry: dict) -> None:
"""Create the kanban task for an intake entry and acknowledge on the post."""
content = entry.get("content", "")
title = " ".join(content.split())[:80] or "(pulse task)"
# Crash recovery leans on the CLI's documented idempotency contract:
# "If a non-archived task with this key exists, its id is returned
# instead of creating a duplicate" — so re-calling after a crash between
# create and save_state adopts the existing task rather than forking one.
raw = kanban("create", title,
"--body", f"{content}\n\n(from Pulse post {post_id}, "
f"ledger {entry['ledger_id'][:8]})",
"--created-by", "pulse-bridge", "--assignee", "default",
"--idempotency-key", f"pulse-{post_id}", "--json")
task = json.loads(raw)
task_id = task["id"] if isinstance(task, dict) else task[0]["id"]
entry["task_id"] = task_id
entry["last_status"] = task.get("status", "todo") if isinstance(task, dict) else "todo"
entry.pop("queued", None)
save_state(state)
log(f"intake {post_id[:8]} -> kanban {task_id}")
try:
send_message(ensure_channel(state),
f"🤖 picked up — ledger {entry['ledger_id'][:8]}, "
f"kanban {task_id}", thread_id=post_id)
except urllib.error.HTTPError as e:
log(f"DROP ack {post_id[:8]}: HTTP {e.code}")
def intake_ingest(state: dict) -> None:
"""Turn new non-bridge top-level activity posts into ledger items + tasks."""
msgs = api("GET", f"/api/workspaces/{WORKSPACE_ID}/comms/channels/"
f"{ensure_channel(state)}/messages?limit=100") or []
cp = int(state.get("intake_checkpoint") or 0)
if cp == 0:
# First intake run: start strictly after the newest existing post;
# never ingest history. (+1 because the window below is INCLUSIVE.)
state["intake_checkpoint"] = 1 + max(
[int(m.get("created_at") or 0) for m in msgs] or [0])
save_state(state)
return
# Inclusive window (ts >= cp) so two posts sharing one millisecond can't
# lose the second to a crash between them; the intake map (keyed by post
# id) is the dedupe, mirroring the outbound seen-set discipline.
candidates = []
for m in msgs:
ts = int(m.get("created_at") or 0)
content = (m.get("content") or "").strip()
if (ts < cp or m.get("thread_id") or m["id"] in state["intake"]
or (m.get("user_name") or "").lower() in INTAKE_IGNORE_USERS
or (m.get("user_id") or "").startswith("agent:")
or content.startswith("🤖") or "ledger:" in content):
continue
candidates.append(m)
for post in sorted(candidates, key=lambda m: int(m["created_at"])):
ledger_id = ledger_add_for_post(post)
# Dedupe with the outbound mirror: map the ledger item to the USER'S
# post before its op:add reaches the feed, so create_post() skips it.
state["map"][ledger_id] = post["id"]
entry = {"ledger_id": ledger_id, "content": post.get("content") or "",
"author": post.get("user_name") or post.get("user_id")}
state["intake"][post["id"]] = entry
state["intake_checkpoint"] = max(int(state["intake_checkpoint"]),
int(post["created_at"]))
save_state(state)
log(f"intake {post['id'][:8]} -> ledger {ledger_id[:8]}")
backlog = any(e.get("queued") for e in state["intake"].values())
if not backlog and intake_active_count(state) < INTAKE_MAX_CONCURRENT:
intake_start_task(state, post["id"], entry)
else:
entry["queued"] = True
save_state(state)
n = intake_active_count(state)
log(f"intake {post['id'][:8]} queued behind {n} tasks")
try:
send_message(ensure_channel(state),
f"⏳ queued behind {n} tasks",
thread_id=post["id"])
except urllib.error.HTTPError as e:
log(f"DROP queue-note {post['id'][:8]}: HTTP {e.code}")
STATUS_COMMENT = {"running": "▶️ running", "blocked": "⛔ blocked",
"todo": "📋 waiting for a worker"}
def intake_track(state: dict) -> None:
"""Advance ledger + post comments as kanban task statuses change."""
open_entries = {pid: e for pid, e in state["intake"].items()
if e.get("task_id") and not e.get("closed")}
if not open_entries:
return
tasks = {t["id"]: t for t in json.loads(kanban("list", "--json"))}
for post_id, entry in open_entries.items():
task = tasks.get(entry["task_id"])
if task is None:
continue # archived/missing; leave for manual attention
status = task.get("status")
if status == entry.get("last_status"):
continue
entry["last_status"] = status
lid = entry["ledger_id"]
if status == "done":
note = (task.get("result") or "").strip()
body = "✅ completed" + (f" — {neutralize(note[:400])}" if note
else "")
entry["closed"] = True
try:
shre_items("done", lid[:8], "--why",
f"kanban {entry['task_id']} completed")
except RuntimeError as e:
log(f"WARN close {lid[:8]}: {e}")
else:
body = STATUS_COMMENT.get(status, f"status → {status}")
if status == "running":
try:
shre_items("update", lid[:8], "--stage", "build",
"--note", f"kanban {entry['task_id']} running")
except RuntimeError as e:
log(f"WARN stage {lid[:8]}: {e}")
save_state(state)
log(f"track {entry['task_id']} {status} (post {post_id[:8]})")
try:
send_message(ensure_channel(state), body, thread_id=post_id)
except urllib.error.HTTPError as e:
log(f"DROP track {post_id[:8]}: HTTP {e.code}")
def intake_promote(state: dict) -> None:
"""Start queued intake entries when concurrency slots free up."""
for post_id, entry in state["intake"].items():
if not entry.get("queued"):
continue
if intake_active_count(state) >= INTAKE_MAX_CONCURRENT:
return
intake_start_task(state, post_id, entry)
# ── Main loop ────────────────────────────────────────────────────────────────
def seed(state: dict) -> None:
"""First run: mirror currently-OPEN items only, then start the op stream
from 'now'. Historic adds/closes are deliberately not replayed."""
checkpoint = datetime.now(timezone.utc).isoformat(timespec="seconds")
items = open_items()
log(f"seeding {len(items)} open items")
for item in items:
try:
create_post(state, item)
except Exception as e: # keep seeding the rest
log(f"ERROR seeding {item.get('id', '?')[:8]}: {e}")
state["seeded"] = True
state["checkpoint"] = checkpoint
save_state(state)
def process(state: dict, rec: dict) -> None:
op, lid = rec.get("op"), rec.get("id")
if not lid or op not in ("add", "update", "close"):
return
if op == "add":
create_post(state, rec)
return
# Intake-managed items already get status comments from intake_track();
# mirroring their ledger ops too would double-comment the user's post.
if any(e.get("ledger_id") == lid for e in state.get("intake", {}).values()):
return
if op == "update":
comment(state, lid, update_body(rec), "update")
else:
comment(state, lid, close_body(rec), "close")
def poll(state: dict) -> None:
records = feed_since(state["checkpoint"])
seen = {tuple(s) for s in state.get("seen", [])}
new = [r for r in records
if (r.get("id"), r.get("op"), r.get("at")) not in seen]
for rec in new:
try:
process(state, rec)
except Exception as e:
log(f"ERROR {rec.get('op')} {str(rec.get('id'))[:8]}: {e}")
# don't advance past a failed record's timestamp; retry next poll
break
seen.add((rec.get("id"), rec.get("op"), rec.get("at")))
state["checkpoint"] = max(state["checkpoint"], rec.get("at", ""))
# keep only boundary-second entries for the inclusive-window dedupe
state["seen"] = [list(t) for t in seen if t[2] >= state["checkpoint"]]
save_state(state)
def main() -> int:
log(f"pulse-bridge v0 starting: mib={MIB_BASE} ws={WORKSPACE_ID} "
f"channel={CHANNEL_NAME} poll={POLL_SEC}s state={STATE_FILE}")
state = load_state()
if not state.get("seeded"):
try:
seed(state)
except Exception as e:
log(f"FATAL seed failed: {e}")
time.sleep(POLL_SEC)
return 1 # launchd KeepAlive restarts us
# Older state files predate intake; give them the new keys.
state.setdefault("intake_checkpoint", 0)
state.setdefault("intake", {})
while True:
try:
poll(state)
except Exception as e:
log(f"ERROR poll: {e}")
if INTAKE_ENABLED:
# track first (frees slots), then promote the queued backlog
# (FIFO), then ingest new posts — so newcomers can't jump ahead.
for step in (intake_track, intake_promote, intake_ingest):
try:
step(state)
except Exception as e:
log(f"ERROR {step.__name__}: {e}")
time.sleep(POLL_SEC)
if __name__ == "__main__":
sys.exit(main())