deerflow-code/offline-backend-20260512/backend/packages/harness/deerflow/persistence/position_roundtable/delivery.py
2026-09-07 18:24:55 +08:00

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