"""Project State Hub workplans/tasks into the graph-explorer payload. This is a *view* of State Hub coordination state, not Fabric topology. State Hub remains the authoring surface for workplans; Fabric only renders. """ from __future__ import annotations import json import os import re import urllib.request from datetime import datetime, timezone from typing import Any from urllib.parse import urlencode OPEN_WORKPLAN = frozenset({"proposed", "ready", "active", "blocked", "backlog"}) OPEN_TASK = frozenset({"todo", "progress", "wait"}) RESIDUAL_FLAVOR = "residual" WP_ID_RE = re.compile(r"\b([A-Z]{2,12}(?:-[A-Z]+)?-WP-[A-Z0-9-]+)\b") FLAVOR_COLORS = { "planning": "#6366f1", "implementation": "#1d4ed8", "refactoring": "#b45309", "extension": "#0f766e", "residual": "#64748b", } # CUST-WP-0074 qualified waits: external commitment (Case A) versus human gate (Case B). WAIT_KIND_COLORS = { "external": "#0891b2", "human": "#be123c", } WAIT_KINDS = frozenset({"external", "human", "both", "unqualified"}) BLOCKED_KINDS = frozenset({"external", "human", "none"}) HUMAN_WAIT_KINDS = frozenset({"human", "both"}) def is_residual_flavor(value: Any) -> bool: """True only for the closed flavor token. Titles are not consulted.""" return str(value or "").strip().lower() == RESIDUAL_FLAVOR def _flavor(record: dict[str, Any]) -> str | None: text = str(record.get("flavor") or "").strip().lower() return text or None def _indexed_depends_targets(workplan: dict[str, Any]) -> list[str]: """Hub-indexed workplan ids from depends_on stubs or id lists.""" raw = workplan.get("depends_on") if raw is None: raw = workplan.get("depends_on_workplans") if raw is None: return [] if not isinstance(raw, list): raw = [raw] targets: list[str] = [] seen: set[str] = set() for item in raw: if isinstance(item, dict): target = ( item.get("workplan_id") or item.get("workstream_id") or item.get("task_id") or item.get("id") ) else: target = item if not target: continue key = str(target) if key in seen: continue seen.add(key) targets.append(key) return targets def _task_depends_targets(task: dict[str, Any]) -> list[str]: """Task-level depends_on ids (workplan or task ids), when the payload carries them.""" return _indexed_depends_targets({"depends_on": task.get("depends_on")}) def _cited_tokens(task: dict[str, Any]) -> set[str]: blob = " ".join( str(task.get(key) or "") for key in ("description", "blocking_reason", "intervention_note", "title") ) return set(WP_ID_RE.findall(blob.upper())) def _hub_kind(value: Any, allowed: frozenset[str]) -> str | None: """Accept the hub field as-is or with a `blocked-`/`waiting-` prefix.""" text = str(value or "").strip().lower() for prefix in ("blocked-", "waiting-"): if text.startswith(prefix): text = text[len(prefix):] return text if text in allowed else None def task_wait_kind( task: dict[str, Any], *, has_external: bool, ) -> str | None: """`external` | `human` | `both` | `unqualified` for a wait task, else None. The hub's `wait_kind` wins when present; otherwise derive from `needs_human` (Case B) and dependency edges (Case A). """ if str(task.get("status") or "") != "wait": return None hub = _hub_kind(task.get("wait_kind"), WAIT_KINDS) if hub: return hub human = bool(task.get("needs_human")) if human and has_external: return "both" if human: return "human" if has_external: return "external" return "unqualified" def workplan_blocked_kind(workplan: dict[str, Any], wait_kinds: list[str | None]) -> str | None: """`human` when any wait task is a human gate, else `external`, else `none`.""" if str(workplan.get("status") or "") != "blocked": return None hub = _hub_kind(workplan.get("blocked_kind"), BLOCKED_KINDS) if hub: return hub kinds = {kind for kind in wait_kinds if kind} if kinds & HUMAN_WAIT_KINDS: return "human" if "external" in kinds: return "external" return "none" def _index_workplans( open_wps: list[dict[str, Any]], ) -> tuple[dict[str, dict[str, Any]], dict[str, str], dict[str, str]]: wp_by_id = {str(wp.get("id")): wp for wp in open_wps if wp.get("id")} slug_to_id = { str(wp.get("slug") or "").lower(): str(wp["id"]) for wp in open_wps if wp.get("slug") and wp.get("id") } record_to_id: dict[str, str] = {} for wp in open_wps: slug = str(wp.get("slug") or "") title = str(wp.get("title") or "") for token in WP_ID_RE.findall(slug.upper() + " " + title.upper()): record_to_id[token] = str(wp["id"]) if slug: record_to_id[slug.upper().replace("_", "-")] = str(wp["id"]) return wp_by_id, slug_to_id, record_to_id TASK_NODE_SIZE = 26 WORKPLAN_BASE_SIZE = 44 WORKPLAN_SIZE_STEP = 10 WORKPLAN_SIZE_CAP = 8 def chokepoint_size(in_degree: int) -> int: """Grow workplan nodes with the number of open items waiting on them.""" return WORKPLAN_BASE_SIZE + WORKPLAN_SIZE_STEP * min(max(in_degree, 0), WORKPLAN_SIZE_CAP) def coordination_graph_payload( workplans: list[dict[str, Any]], tasks: list[dict[str, Any]], *, include_residuals: bool = False, needs_human_only: bool = False, dependencies: list[dict[str, Any]] | None = None, ) -> dict[str, Any]: """Build a GraphExplorerPayload of open workplans and their open tasks. `dependencies` are raw `/workplans/{id}/dependencies/` rows (`from_workplan_id`, `to_workplan_id`, `to_task_id`); they qualify a wait task as external when the task cites the row's target. `needs_human_only` keeps only `blocked_kind: human` workplans and their human-wait tasks. """ open_all = [wp for wp in workplans if str(wp.get("status") or "") in OPEN_WORKPLAN] residual_open = sum(1 for wp in open_all if is_residual_flavor(_flavor(wp))) if include_residuals: open_wps = open_all else: open_wps = [wp for wp in open_all if not is_residual_flavor(_flavor(wp))] wp_by_id, slug_to_id, record_to_id = _index_workplans(open_wps) open_tasks = [ task for task in tasks if str(task.get("status") or "") in OPEN_TASK and str(task.get("workplan_id") or task.get("workstream_id") or "") in wp_by_id and (include_residuals or not is_residual_flavor(_flavor(task))) ] # Dependency targets per workplan: indexed stubs plus raw dependency rows. dep_targets: dict[str, set[str]] = {} for wp in open_wps: dep_targets[str(wp["id"])] = set(_indexed_depends_targets(wp)) for row in dependencies or []: src = str(row.get("from_workplan_id") or row.get("from_workstream_id") or "") dst = row.get("to_workplan_id") or row.get("to_workstream_id") or row.get("to_task_id") if src in dep_targets and dst: dep_targets[src].add(str(dst)) def _resolve(raw: str) -> str | None: if raw in wp_by_id: return raw return record_to_id.get(raw.upper()) or slug_to_id.get(raw.lower()) wait_kind_by_task: dict[str, str | None] = {} for task in open_tasks: wp_id = str(task.get("workplan_id") or task.get("workstream_id")) external = bool(_task_depends_targets(task)) if not external and str(task.get("status") or "") == "wait": resolved = {_resolve(dst) or dst for dst in dep_targets.get(wp_id, set())} cited = {_resolve(token) or token for token in _cited_tokens(task)} external = bool(resolved & cited) wait_kind_by_task[str(task["id"])] = task_wait_kind(task, has_external=external) blocked_kind_by_wp: dict[str, str | None] = {} for wp in open_wps: wp_id = str(wp["id"]) kinds = [ wait_kind_by_task[str(task["id"])] for task in open_tasks if str(task.get("workplan_id") or task.get("workstream_id")) == wp_id ] blocked_kind_by_wp[wp_id] = workplan_blocked_kind(wp, kinds) if needs_human_only: open_wps = [wp for wp in open_wps if blocked_kind_by_wp[str(wp["id"])] == "human"] wp_by_id, slug_to_id, record_to_id = _index_workplans(open_wps) open_tasks = [ task for task in open_tasks if str(task.get("workplan_id") or task.get("workstream_id")) in wp_by_id and wait_kind_by_task[str(task["id"])] in HUMAN_WAIT_KINDS ] task_ids = {str(task["id"]) for task in open_tasks} elements: list[dict[str, Any]] = [] for wp in open_wps: wp_id = str(wp["id"]) status = str(wp.get("status") or "unknown") flavor = _flavor(wp) elements.append( { "data": { "id": f"workplan:{wp_id}", "stableKey": f"workplan:{wp_id}", "kind": "Workplan", "layer": "workplan", "label": str(wp.get("slug") or wp.get("title") or wp_id)[:48], "name": wp.get("title"), "lifecycle": status, "status": status, "flavor": flavor, "nodeClass": flavor or "unspecified", "repo": wp.get("repo") or wp.get("slug"), "displayState": "show" if status != "blocked" else "highlight", "unresolved": status in {"blocked", "proposed"}, "blocked_kind": blocked_kind_by_wp[wp_id], "chokepoint": 0, "color": FLAVOR_COLORS.get(flavor or "", "#1d4ed8"), } } ) belongs = 0 wait_edges = 0 wait_kind_counts = {kind: 0 for kind in sorted(WAIT_KINDS)} for task in open_tasks: task_id = str(task["id"]) wp_id = str(task.get("workplan_id") or task.get("workstream_id")) status = str(task.get("status") or "todo") human = bool(task.get("needs_human")) flavor = _flavor(task) or _flavor(wp_by_id.get(wp_id) or {}) wait_kind = wait_kind_by_task[task_id] human_gate = wait_kind in HUMAN_WAIT_KINDS if wait_kind: wait_kind_counts[wait_kind] += 1 elements.append( { "data": { "id": f"task:{task_id}", "stableKey": f"task:{task_id}", "kind": "Task", "layer": "task", "label": str(task.get("title") or task_id)[:56], "name": task.get("title"), "lifecycle": status, "status": status, "flavor": flavor, "nodeClass": flavor or "unspecified", "needsHuman": human, "wait_kind": wait_kind, # Human gates have no edge target; the marker lives on the node. "humanGate": human_gate, "displayState": "highlight" if status == "wait" or human else "show", "unresolved": status == "wait" or human, "visualSize": TASK_NODE_SIZE, } } ) if human_gate: elements[-1]["data"]["color"] = WAIT_KIND_COLORS["human"] elif wait_kind == "external": elements[-1]["data"]["color"] = WAIT_KIND_COLORS["external"] elements.append( { "data": { "id": f"edge:belongs:{task_id}", "stableKey": f"edge:belongs:{task_id}", "kind": "Edge", "layer": "dependency", "displayState": "show", "source": f"task:{task_id}", "target": f"workplan:{wp_id}", "edgeType": "belongs_to", "edgeSource": "membership", "strength": "weak", "sourceLayer": "task", "targetLayer": "workplan", } } ) belongs += 1 if status == "wait" or human: wait_edges += 1 indexed_pairs: set[tuple[str, str]] = set() depends_count = 0 chokepoint: dict[str, int] = {str(wp["id"]): 0 for wp in open_wps} for wp in open_wps: src = str(wp["id"]) for dst_raw in _indexed_depends_targets(wp): dst = _resolve(dst_raw) if not dst or dst == src or dst not in wp_by_id: continue pair = (src, dst) if pair in indexed_pairs: continue indexed_pairs.add(pair) elements.append( { "data": { "id": f"edge:depends:{src}:{dst}", "stableKey": f"edge:depends:{src}:{dst}", "kind": "Edge", "layer": "dependency", "displayState": "show", "source": f"workplan:{src}", "target": f"workplan:{dst}", "edgeType": "depends_on", "edgeSource": "indexed", "edge_kind": "commitment", "strength": "strong", "sourceLayer": "workplan", "targetLayer": "workplan", } } ) depends_count += 1 chokepoint[dst] = chokepoint.get(dst, 0) + 1 # Task-level depends_on (CUST-WP-0074 Case A): task -> workplan or task. task_pairs: set[tuple[str, str]] = set() task_depends_count = 0 for task in open_tasks: task_id = str(task["id"]) src_wp = str(task.get("workplan_id") or task.get("workstream_id")) for dst_raw in _task_depends_targets(task): if dst_raw in task_ids and dst_raw != task_id: target, target_layer = f"task:{dst_raw}", "task" else: dst = _resolve(dst_raw) if not dst or dst == src_wp or dst not in wp_by_id: continue target, target_layer = f"workplan:{dst}", "workplan" chokepoint[dst] = chokepoint.get(dst, 0) + 1 if (task_id, target) in task_pairs: continue task_pairs.add((task_id, target)) elements.append( { "data": { "id": f"edge:task-depends:{task_id}:{target}", "stableKey": f"edge:task-depends:{task_id}:{target}", "kind": "Edge", "layer": "dependency", "displayState": "show", "source": f"task:{task_id}", "target": target, "edgeType": "depends_on", "edgeSource": "task", "edge_kind": "commitment", "strength": "strong", "sourceLayer": "task", "targetLayer": target_layer, } } ) task_depends_count += 1 cited = 0 for task in open_tasks: src_wp = str(task.get("workplan_id") or task.get("workstream_id")) for token in _cited_tokens(task): dst = record_to_id.get(token) or slug_to_id.get(token.lower()) if not dst or dst == src_wp or dst not in wp_by_id: continue if (src_wp, dst) in indexed_pairs: continue elements.append( { "data": { "id": f"edge:cites:{task['id']}:{dst}", "stableKey": f"edge:cites:{task['id']}:{dst}", "kind": "Edge", "layer": "dependency", "displayState": "show", "source": f"workplan:{src_wp}", "target": f"workplan:{dst}", "edgeType": "waits_on", "edgeSource": "citation", "edge_kind": "citation", "strength": "strong", "sourceLayer": "workplan", "targetLayer": "workplan", } } ) cited += 1 chokepoint[dst] = chokepoint.get(dst, 0) + 1 indexed_pairs.add((src_wp, dst)) for el in elements: data = el["data"] if data.get("kind") == "Edge": data.setdefault("edgeWidth", 3 if data.get("edgeType") in {"depends_on", "waits_on"} else 1.5) continue node_id = data["id"] if not node_id.startswith("workplan:"): continue wp_id = node_id.split(":", 1)[1] data["chokepoint"] = int(chokepoint.get(wp_id) or 0) data["visualSize"] = chokepoint_size(data["chokepoint"]) generated = datetime.now(timezone.utc).isoformat() return { "apiVersion": "railiance.fabric/v1alpha1", "kind": "GraphExplorerPayload", "manifest_id": "railiance-fabric.coordination-map", "generated_at": generated, "mode": "coordination", "metrics": { "open_workplans": len(open_wps), "open_tasks": len(open_tasks), "belongs_to_edges": belongs, "depends_on_edges": depends_count, "task_depends_on_edges": task_depends_count, "citation_edges": cited, "wait_or_human_tasks": wait_edges, **{f"wait_{kind}_tasks": count for kind, count in wait_kind_counts.items()}, "residual_open_workplans": residual_open, "include_residuals": include_residuals, "needs_human_only": needs_human_only, }, "elements": elements, "hidden_elements": [], } def fetch_hub_lists(api_base: str) -> tuple[list[dict[str, Any]], list[dict[str, Any]]]: """Open workplans and all tasks; `/state/deps` stubs carry the dependency rows.""" base = api_base.rstrip("/") workplans: list[dict[str, Any]] = [] for status in sorted(OPEN_WORKPLAN): workplans.extend(_get_json(f"{base}/workplans/?{urlencode({'status': status})}")) tasks = _get_json(f"{base}/tasks/") dep_rows = _get_json_optional(f"{base}/state/deps?include_residuals=true") by_id = { str(row.get("id")): row for row in dep_rows if isinstance(row, dict) and row.get("id") } for wp in workplans: extra = by_id.get(str(wp.get("id"))) if not extra: continue if extra.get("depends_on") and not wp.get("depends_on"): wp["depends_on"] = extra["depends_on"] if extra.get("flavor") and not wp.get("flavor"): wp["flavor"] = extra["flavor"] return workplans, tasks def _get_json(url: str) -> list[dict[str, Any]]: with urllib.request.urlopen(url, timeout=60) as response: payload = json.loads(response.read().decode()) if isinstance(payload, list): return payload raise RuntimeError(f"unexpected JSON from {url}") def _get_json_optional(url: str) -> list[dict[str, Any]]: try: return _get_json(url) except Exception: return [] def default_hub() -> str: return os.environ.get("STATE_HUB_URL") or os.environ.get("API_BASE") or "http://127.0.0.1:8000"