118 lines
4.4 KiB
Python
118 lines
4.4 KiB
Python
"""Pure delivery-decision helpers for position-roundtable nodes.
|
|
|
|
These functions decide whether a turn delivered a real Markdown report, whether
|
|
the delivery changed, and whether a node is the special intent position. They
|
|
are intentionally **pure** (no DB / no app imports) so the exact same logic can
|
|
run both:
|
|
|
|
- in the API router (for read-only projections and validation), and
|
|
- inside ``PositionRoundtableRepository.apply_node_command`` — *under the row
|
|
locks* — where every state-machine decision must be computed from the locked
|
|
rows, never from a stale pre-transaction snapshot.
|
|
|
|
Keeping one implementation avoids the classic drift where the router's decision
|
|
and the store's transition disagree under concurrency.
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
from typing import Any
|
|
|
|
_OUTPUTS_PREFIX = "mnt/user-data/outputs/"
|
|
|
|
|
|
def _artifact_path(artifact: dict[str, Any]) -> str:
|
|
return str(artifact.get("path") or artifact.get("artifact_url") or "").strip().lower()
|
|
|
|
|
|
def _artifact_name(artifact: dict[str, Any]) -> str:
|
|
return str(artifact.get("name") or artifact.get("filename") or "").strip().lower()
|
|
|
|
|
|
def _artifact_mime(artifact: dict[str, Any]) -> str:
|
|
return str(artifact.get("mime_type") or artifact.get("mimeType") or "").strip().lower()
|
|
|
|
|
|
def _in_outputs(path: str) -> bool:
|
|
"""Only files under the thread outputs dir count as deliverables."""
|
|
return not path or path.replace("\\", "/").lstrip("/").startswith(_OUTPUTS_PREFIX)
|
|
|
|
|
|
def is_markdown_artifact(artifact: Any) -> bool:
|
|
if not isinstance(artifact, dict):
|
|
return False
|
|
path = _artifact_path(artifact)
|
|
name = _artifact_name(artifact)
|
|
mime = _artifact_mime(artifact)
|
|
if not _in_outputs(path):
|
|
return False
|
|
return (
|
|
"markdown" in mime
|
|
or path.endswith(".md")
|
|
or path.endswith(".markdown")
|
|
or name.endswith(".md")
|
|
or name.endswith(".markdown")
|
|
)
|
|
|
|
|
|
def is_json_artifact(artifact: Any) -> bool:
|
|
if not isinstance(artifact, dict):
|
|
return False
|
|
path = _artifact_path(artifact)
|
|
name = _artifact_name(artifact)
|
|
mime = _artifact_mime(artifact)
|
|
if not _in_outputs(path):
|
|
return False
|
|
return "json" in mime or path.endswith(".json") or name.endswith(".json")
|
|
|
|
|
|
def artifact_signature(artifact: dict[str, Any]) -> tuple[str, str, str]:
|
|
"""Use server-returned path/mtime/size to detect an actually rewritten file."""
|
|
return (
|
|
str(artifact.get("path") or artifact.get("artifact_url") or ""),
|
|
str(artifact.get("modified_at") or artifact.get("modifiedAt") or ""),
|
|
str(artifact.get("size_bytes") or artifact.get("sizeBytes") or ""),
|
|
)
|
|
|
|
|
|
def markdown_delivery_changed(
|
|
previous_artifacts: Any,
|
|
incoming_artifacts: list[dict[str, Any]],
|
|
) -> bool:
|
|
"""Whether this turn rewrote a Markdown deliverable, not just chatted.
|
|
|
|
A completed node can receive normal follow-up Q&A. That must not bump its
|
|
effective delivery revision or invalidate every downstream stage unless a
|
|
Markdown report really changed. The standard artifact listing includes
|
|
``modified_at`` and ``size_bytes``, which together make a stable lightweight
|
|
version marker without reading a potentially large file again.
|
|
"""
|
|
old = previous_artifacts if isinstance(previous_artifacts, list) else []
|
|
old_signatures = {
|
|
artifact_signature(item) for item in old if is_markdown_artifact(item)
|
|
}
|
|
new_signatures = {
|
|
artifact_signature(item) for item in incoming_artifacts if is_markdown_artifact(item)
|
|
}
|
|
return bool(new_signatures - old_signatures)
|
|
|
|
|
|
def is_intent_position_node(
|
|
position_id: str | None,
|
|
intent_position_id: str | None = None,
|
|
) -> bool:
|
|
"""Whether a node is the special intent position (answers need no report)."""
|
|
position = str(position_id or "").strip()
|
|
if intent_position_id:
|
|
return position == intent_position_id
|
|
# Existing persisted sessions predate the configurable directory. Retain
|
|
# their historical interpretation until they are archived/recreated.
|
|
return position.lower() in {"intelligence", "intel", "intent"}
|
|
|
|
|
|
def session_intent_position_id(chain_snapshot: Any) -> str | None:
|
|
"""Read the frozen intent-position id from an activated chain snapshot."""
|
|
chain = chain_snapshot if isinstance(chain_snapshot, dict) else {}
|
|
value = chain.get("_intent_position_id")
|
|
return value.strip() if isinstance(value, str) and value.strip() else None
|