Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 4 additions & 0 deletions plugins/hermes-workflows-approvals/dashboard/dist/index.js
Original file line number Diff line number Diff line change
Expand Up @@ -667,6 +667,10 @@
const feedbackAction = (actions || []).find(function (item) { return item.requires_feedback && item.value !== "approve"; });
if (feedbackAction) selectedAction = feedbackAction.value;
}
if (String(selectedAction || "").toLowerCase().replace(/-/g, "_") === "request_changes" && !feedbackText && !editedOutputText) {
setUi(Object.assign({}, ui, { error: "request_changes requires nonblank feedback or valid edited_output" }));
return;
}
const payload = { action: selectedAction, feedback: feedbackText || undefined };
const editable = (step.input_surface || {}).editable_output || {};
const editField = editable.field || "edited_output";
Expand Down
36 changes: 27 additions & 9 deletions plugins/hermes-workflows-approvals/dashboard/plugin_api.py
Original file line number Diff line number Diff line change
Expand Up @@ -13,14 +13,18 @@
from typing import Any

from hermes_workflows import ApprovalDecisionInput, WorkflowEngine
from hermes_workflows.approvals import strip_client_controlled_provenance, validate_revision_response
from hermes_workflows.artifacts import artifact_descriptor, workflow_source_preview
from hermes_workflows.hermes_plugin_approvals import (
_configured_dbs,
_next_step_for_receipt,
_redact,
_receipt_to_payload,
_revision_schema_for_response,
_source_for_normalized_revision_replay,
approval_view_to_dict,
)
from hermes_workflows.revision_validation import RevisionActionValidationError
from hermes_workflows.workflow_loading import load_workflow_ref

try: # FastAPI is provided by Hermes Agent's dashboard process.
Expand All @@ -30,8 +34,8 @@


class _FallbackHTTPException(Exception):
def __init__(self, status_code: int, detail: str):
super().__init__(detail)
def __init__(self, status_code: int, detail: Any):
super().__init__(str(detail))
self.status_code = status_code
self.detail = detail

Expand Down Expand Up @@ -2240,7 +2244,7 @@ async def respond_review_request(body: dict[str, Any]) -> dict[str, Any]:
raw_payload = body.get("payload") if isinstance(body.get("payload"), dict) else {}
if not raw_payload:
raise HTTPException(status_code=400, detail="review response payload is required")
payload = {key: value for key, value in raw_payload.items() if key not in {"by", "source"}}
payload = strip_client_controlled_provenance(raw_payload)
raw_idempotency_key = str(body.get("idempotency_key") or body.get("event_id") or "").strip()
message_id = raw_idempotency_key or f"dashboard:{uuid.uuid4()}"
if not message_id.startswith("dashboard:"):
Expand All @@ -2250,15 +2254,27 @@ def record_and_resume() -> tuple[Any, dict[str, Any]]:
_ensure_workflow_project_on_path(db_path)
engine = WorkflowEngine(db_path)
normalized_payload = _normalize_review_payload_for_dashboard_request(engine, workflow_id, key, payload)
effective_idempotency_key = message_id
source = {
"channel": "local-dashboard",
"message_id": message_id,
}
if _revision_schema_for_response(engine, workflow_id, key) is not None:
normalized_payload, validated = validate_revision_response(normalized_payload)
effective_idempotency_key = f"revision:{workflow_id}:{key}:{validated.idempotency_key}"
source = _source_for_normalized_revision_replay(
engine,
workflow_id,
key,
effective_idempotency_key,
source,
)
receipt = engine.submit_operator_response(
workflow_id=workflow_id,
key=key,
payload=normalized_payload,
source={
"channel": "hermes-dashboard",
"message_id": message_id,
},
idempotency_key=message_id,
source=source,
idempotency_key=effective_idempotency_key,
resume=True,
)
post_resume = _status_packet(
Expand All @@ -2274,6 +2290,8 @@ def record_and_resume() -> tuple[Any, dict[str, Any]]:

try:
receipt, post_resume = await asyncio.to_thread(record_and_resume)
except RevisionActionValidationError as exc:
raise HTTPException(status_code=400, detail=exc.to_dict()) from exc
except Exception as exc:
raise HTTPException(status_code=400, detail=f"review response/resume failed: {type(exc).__name__}: {exc}") from exc
receipt_payload = _receipt_to_payload(receipt, resume_requested=True)
Expand Down Expand Up @@ -2307,7 +2325,7 @@ async def decide_approval(body: dict[str, Any]) -> dict[str, Any]:
key=key,
action=action,
source={
"channel": "hermes-dashboard",
"channel": "local-dashboard",
"message_id": message_id,
},
note=body.get("note"),
Expand Down
75 changes: 75 additions & 0 deletions src/hermes_workflows/approvals.py
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,73 @@
from dataclasses import dataclass
from typing import Any, Iterator, Mapping

from .provenance import project_response_provenance
from .revision_validation import ValidatedRevisionActionV1, validate_revision_action


CLIENT_CONTROLLED_PROVENANCE_FIELDS = frozenset(
{
"actor",
"authenticated_principal",
"by",
"display_label",
"principal",
"provenance",
"response_provenance",
"source",
"user",
}
)


def strip_client_controlled_provenance(payload: Mapping[str, Any]) -> dict[str, Any]:
"""Remove identity/provenance assertions from an untrusted response body."""

if not isinstance(payload, Mapping):
raise TypeError("response payload must be an object")
return {key: value for key, value in payload.items() if key not in CLIENT_CONTROLLED_PROVENANCE_FIELDS}


def is_revision_action_schema(schema_descriptor: Any) -> bool:
"""Return whether a typed operator request uses the v1 revision action contract."""

if not isinstance(schema_descriptor, Mapping):
return False
fields = schema_descriptor.get("fields")
if not isinstance(fields, list):
return False
action = next((field for field in fields if isinstance(field, Mapping) and field.get("name") == "action"), None)
if not isinstance(action, Mapping):
return False
options = action.get("options") if isinstance(action.get("options"), list) else schema_descriptor.get("choices")
field_names = {str(field.get("name")) for field in fields if isinstance(field, Mapping)}
return (
isinstance(options, list)
and {"approve", "request_changes"}.issubset(options)
and bool({"feedback", "edited_output"} & field_names)
)


def validate_revision_response(payload: Mapping[str, Any]) -> tuple[dict[str, Any], ValidatedRevisionActionV1]:
"""Validate and materialize a revision response for durable adapter submission."""

validated = validate_revision_action(payload)
return _materialize_json(validated.normalized_payload), validated


def response_provenance_for(*, by: str | None, source: Mapping[str, Any] | None) -> dict[str, Any]:
"""Truthfully classify an adapter response without upgrading client labels."""

return project_response_provenance({"by": by, "source": source}).to_dict()


def _materialize_json(value: Any) -> Any:
if isinstance(value, Mapping):
return {str(key): _materialize_json(item) for key, item in value.items()}
if isinstance(value, tuple):
return [_materialize_json(item) for item in value]
return value


@dataclass(frozen=True)
class ApprovalView:
Expand Down Expand Up @@ -60,6 +127,10 @@ def needs_revision(self) -> bool:
def feedback(self) -> str | None:
return self.direct_feedback or self.reason or self.note or self.message or self.comment

@property
def response_provenance(self) -> dict[str, Any]:
return response_provenance_for(by=self.by, source=self.source)

def to_dict(self) -> dict[str, Any]:
data: dict[str, Any] = {"action": self.action}
if self.by is not None:
Expand Down Expand Up @@ -115,6 +186,10 @@ class ApprovalReceipt:
result_summary: dict[str, Any] | None
workflow_ref: str | None = None

@property
def response_provenance(self) -> dict[str, Any]:
return response_provenance_for(by=self.by, source=self.source)


# Neutral names for the general human/operator checkpoint substrate. Approval
# remains a policy preset over this surface for now; these aliases let runtime,
Expand Down
91 changes: 82 additions & 9 deletions src/hermes_workflows/hermes_plugin_approvals.py
Original file line number Diff line number Diff line change
Expand Up @@ -7,9 +7,15 @@
from pathlib import Path
from typing import Any, cast

from .approvals import ApprovalDecisionInput
from .approvals import (
ApprovalDecisionInput,
is_revision_action_schema,
strip_client_controlled_provenance,
validate_revision_response,
)
from .engine import WorkflowEngine
from .receipts import redact_secrets
from .revision_validation import RevisionActionValidationError

PLUGIN_NAME = "hermes-workflows-approvals"
TOOLSET = "hermes_workflows_approvals"
Expand Down Expand Up @@ -293,10 +299,55 @@ def _source_from_args(args: dict[str, Any]) -> dict[str, Any]:

def _receipt_to_payload(receipt: Any, *, resume_requested: bool) -> dict[str, Any]:
payload = _as_payload(receipt)
response_provenance = getattr(receipt, "response_provenance", None)
if isinstance(response_provenance, dict):
payload["response_provenance"] = response_provenance
payload["resume_requested"] = resume_requested
return payload


def _revision_schema_for_response(engine: WorkflowEngine, workflow_id: str, key: str) -> dict[str, Any] | None:
matching_request = next(
(
event.get("payload")
for event in reversed(engine.events(workflow_id))
if event.get("type") == "ApprovalRequested" and event.get("key") == f"approval:{key}"
),
None,
)
if not isinstance(matching_request, dict):
return None
descriptor = matching_request.get("schema_descriptor")
return descriptor if isinstance(descriptor, dict) and is_revision_action_schema(descriptor) else None


def _source_for_normalized_revision_replay(
engine: WorkflowEngine,
workflow_id: str,
key: str,
idempotency_key: str,
source: dict[str, Any],
) -> dict[str, Any]:
"""Keep the first durable transport provenance for an exact revision replay."""

signal_key = f"signal:operator.response:{key}"
matching_signal = next(
(
event
for event in reversed(engine.events(workflow_id))
if event.get("type") == "SignalReceived"
and event.get("key") == signal_key
and event.get("idempotency_key") == idempotency_key
),
None,
)
if not isinstance(matching_signal, dict):
return source
event_payload = matching_signal.get("payload")
persisted_source = event_payload.get("source") if isinstance(event_payload, dict) else None
return dict(persisted_source) if isinstance(persisted_source, dict) else source


def _next_step_for_receipt(receipt_payload: dict[str, Any]) -> str | None:
status = receipt_payload.get("status")
if receipt_payload.get("resume_requested"):
Expand Down Expand Up @@ -363,15 +414,37 @@ def _handle_workflow_review_respond(args: dict[str, Any], **kwargs: Any) -> str:
payload = json.loads(payload)
if not isinstance(payload, dict) or not payload:
raise ValueError("payload must be a non-empty object")
receipt = WorkflowEngine(db_path).submit_operator_response(
workflow_id=str(args.get("workflow_id") or "").strip(),
key=str(args.get("key") or "").strip(),
workflow_id = str(args.get("workflow_id") or "").strip()
key = str(args.get("key") or "").strip()
payload = strip_client_controlled_provenance(payload)
engine = WorkflowEngine(db_path)
source = _source_from_args(args)
if _revision_schema_for_response(engine, workflow_id, key) is not None:
try:
payload, validated = validate_revision_response(payload)
except RevisionActionValidationError as exc:
return _tool_error(str(exc), validation=exc.to_dict())
idempotency_key = f"revision:{workflow_id}:{key}:{validated.idempotency_key}"
source = _source_for_normalized_revision_replay(
engine,
workflow_id,
key,
idempotency_key,
source,
)
else:
idempotency_key = (
args.get("idempotency_key")
or args.get("message_url")
or args.get("message_id")
or args.get("event_id")
)
receipt = engine.submit_operator_response(
workflow_id=workflow_id,
key=key,
payload=payload,
source=_source_from_args(args),
idempotency_key=args.get("idempotency_key")
or args.get("message_url")
or args.get("message_id")
or args.get("event_id"),
source=source,
idempotency_key=idempotency_key,
resume=resume,
)
receipt_payload = _receipt_to_payload(receipt, resume_requested=resume)
Expand Down
46 changes: 46 additions & 0 deletions tests/test_approval_provenance.py
Original file line number Diff line number Diff line change
Expand Up @@ -111,6 +111,52 @@ def test_human_approval_accepts_human_source_and_returns_provenance(tmp_path):
"channel": "discord",
"message_url": "discord://thread/123/message/456",
}
assert decision.to_dict() == {
"action": "approve",
"by": "skylar",
"source": {
"channel": "discord",
"message_url": "discord://thread/123/message/456",
},
}
assert decision.response_provenance == {
"schema_version": 1,
"kind": "legacy_unverified",
"principal": None,
"display_label": "skylar",
"event": {
"channel": "discord",
"message_id": None,
"message_url": "discord://thread/123/message/456",
"event_id": None,
},
}


def test_local_dashboard_decision_is_truthfully_unattributed():
from hermes_workflows.approvals import ApprovalDecision

decision = ApprovalDecision(
action="approve",
source={"channel": "local-dashboard", "event_id": "dashboard:click-1"},
)

assert decision.to_dict() == {
"action": "approve",
"source": {"channel": "local-dashboard", "event_id": "dashboard:click-1"},
}
assert decision.response_provenance == {
"schema_version": 1,
"kind": "unattributed_local_operator",
"principal": None,
"display_label": None,
"event": {
"channel": "local-dashboard",
"message_id": None,
"message_url": None,
"event_id": "dashboard:click-1",
},
}


def test_human_approval_rejects_missing_source(tmp_path):
Expand Down
Loading