feat: replace complete draft documents safely
This commit is contained in:
@@ -24,6 +24,9 @@ from wf_artifacts import (
|
||||
from wf_artifacts import (
|
||||
patch_draft_workspace as patch_draft_workspace_record,
|
||||
)
|
||||
from wf_artifacts import (
|
||||
replace_draft_workspace_document as replace_draft_workspace_document_record,
|
||||
)
|
||||
from wf_core.models.schemas import NodeDef
|
||||
from wf_core.models.steps import (
|
||||
InputBinding,
|
||||
@@ -341,6 +344,22 @@ class WorkflowDraftApi:
|
||||
draft=draft,
|
||||
)
|
||||
|
||||
async def replace_draft_workspace_document(
|
||||
self,
|
||||
*,
|
||||
workspace_id: str,
|
||||
revision: int,
|
||||
draft: dict[str, Any],
|
||||
) -> dict[str, Any]:
|
||||
"""Replace and semantically revalidate one complete workspace draft."""
|
||||
return replace_draft_workspace_document_record(
|
||||
self._draft_store(),
|
||||
workspace_id=workspace_id,
|
||||
revision=revision,
|
||||
draft=draft,
|
||||
node_defs_for_draft=self._node_defs_for_draft,
|
||||
)
|
||||
|
||||
async def set_draft_name(
|
||||
self,
|
||||
*,
|
||||
|
||||
@@ -331,6 +331,20 @@ class WorkflowApi:
|
||||
patch=patch,
|
||||
)
|
||||
|
||||
async def replace_draft_workspace_document(
|
||||
self,
|
||||
*,
|
||||
workspace_id: str,
|
||||
revision: int,
|
||||
draft: dict[str, Any],
|
||||
) -> dict[str, Any]:
|
||||
"""Replace and semantically revalidate one complete workspace draft."""
|
||||
return await self.drafts.replace_draft_workspace_document(
|
||||
workspace_id=workspace_id,
|
||||
revision=revision,
|
||||
draft=draft,
|
||||
)
|
||||
|
||||
async def set_draft_name(
|
||||
self,
|
||||
*,
|
||||
|
||||
@@ -92,6 +92,16 @@ class WorkflowDraftSurface(Protocol):
|
||||
patch: list[dict[str, Any]],
|
||||
) -> dict[str, Any]: ...
|
||||
|
||||
async def replace_draft_workspace_document(
|
||||
self,
|
||||
*,
|
||||
workspace_id: str,
|
||||
revision: int,
|
||||
draft: dict[str, Any],
|
||||
) -> dict[str, Any]:
|
||||
"""Replace and semantically revalidate one complete workspace draft."""
|
||||
...
|
||||
|
||||
async def set_draft_name(
|
||||
self,
|
||||
*,
|
||||
|
||||
@@ -12,6 +12,7 @@ from .draft_workspaces import (
|
||||
ensure_workspace_id,
|
||||
get_draft_workspace,
|
||||
patch_draft_workspace,
|
||||
replace_draft_workspace_document,
|
||||
replace_validated_draft_document,
|
||||
summarize_draft_workspace,
|
||||
)
|
||||
@@ -92,6 +93,7 @@ __all__ = [
|
||||
"logical_ref_for_concrete_ref",
|
||||
"normalize_plan_node_refs",
|
||||
"patch_draft_workspace",
|
||||
"replace_draft_workspace_document",
|
||||
"replace_validated_draft_document",
|
||||
"patch_workflow_draft",
|
||||
"summarize_draft_workspace",
|
||||
|
||||
@@ -2,6 +2,7 @@ from .api import (
|
||||
create_draft_workspace,
|
||||
get_draft_workspace,
|
||||
patch_draft_workspace,
|
||||
replace_draft_workspace_document,
|
||||
replace_validated_draft_document,
|
||||
)
|
||||
from .models import (
|
||||
@@ -24,6 +25,7 @@ __all__ = [
|
||||
"ensure_workspace_id",
|
||||
"get_draft_workspace",
|
||||
"patch_draft_workspace",
|
||||
"replace_draft_workspace_document",
|
||||
"replace_validated_draft_document",
|
||||
"summarize_draft_workspace",
|
||||
]
|
||||
|
||||
@@ -103,6 +103,46 @@ def patch_draft_workspace(
|
||||
return summarize_draft_workspace(next_workspace)
|
||||
|
||||
|
||||
def replace_draft_workspace_document(
|
||||
store: DraftWorkspaceStore,
|
||||
*,
|
||||
workspace_id: str,
|
||||
revision: int,
|
||||
draft: JsonObject,
|
||||
node_defs_for_draft: NodeDefsForDraft,
|
||||
) -> JsonObject:
|
||||
"""Replace and semantically revalidate one complete draft document."""
|
||||
workspace = store.get_workspace(workspace_id)
|
||||
if workspace.revision != revision:
|
||||
return _revision_conflict_payload(workspace, revision)
|
||||
WorkflowDraft.model_validate(draft)
|
||||
if draft == workspace.draft:
|
||||
return summarize_draft_workspace(workspace)
|
||||
|
||||
stored_draft = deepcopy(draft)
|
||||
validation = validate_workflow_draft(
|
||||
stored_draft,
|
||||
node_defs=node_defs_for_draft(stored_draft),
|
||||
)
|
||||
next_workspace = workspace.model_copy(
|
||||
update={
|
||||
"revision": workspace.revision + 1,
|
||||
"draft": _canonical_draft_if_valid(
|
||||
stored_draft,
|
||||
validation_status=validation["status"],
|
||||
),
|
||||
"status": validation["status"],
|
||||
"diagnostics": validation["diagnostics"],
|
||||
"updated_at_epoch_ms": _now_ms(),
|
||||
}
|
||||
)
|
||||
try:
|
||||
store.replace_workspace(next_workspace, expected_revision=revision)
|
||||
except DraftWorkspaceConflictError as exc:
|
||||
return _revision_conflict_payload(exc.workspace, revision)
|
||||
return summarize_draft_workspace(next_workspace)
|
||||
|
||||
|
||||
def replace_validated_draft_document(
|
||||
store: DraftWorkspaceStore,
|
||||
*,
|
||||
|
||||
@@ -2,6 +2,9 @@ from __future__ import annotations
|
||||
|
||||
from typing import Any
|
||||
|
||||
import pytest
|
||||
from pydantic import ValidationError
|
||||
|
||||
from wf_artifacts import (
|
||||
DraftWorkspaceStore,
|
||||
FileDraftWorkspaceStore,
|
||||
@@ -9,6 +12,7 @@ from wf_artifacts import (
|
||||
create_draft_workspace,
|
||||
get_draft_workspace,
|
||||
patch_draft_workspace,
|
||||
replace_draft_workspace_document,
|
||||
replace_validated_draft_document,
|
||||
summarize_draft_workspace,
|
||||
)
|
||||
@@ -292,6 +296,114 @@ def test_replace_validated_draft_document_isolates_persisted_draft() -> None:
|
||||
assert store.get_workspace("echo_draft").draft["name"] == "replacement"
|
||||
|
||||
|
||||
def test_replace_draft_workspace_document_revalidates_and_increments_revision(
|
||||
tmp_path,
|
||||
) -> None:
|
||||
store = FileDraftWorkspaceStore(tmp_path)
|
||||
create_draft_workspace(store, workspace_id="echo_draft", draft=_draft())
|
||||
replacement = {**_draft(), "name": "replacement"}
|
||||
|
||||
result = replace_draft_workspace_document(
|
||||
store,
|
||||
workspace_id="echo_draft",
|
||||
revision=1,
|
||||
draft=replacement,
|
||||
node_defs_for_draft=lambda _draft: [],
|
||||
)
|
||||
|
||||
assert result["revision"] == 2
|
||||
assert store.get_workspace("echo_draft").draft["name"] == "replacement"
|
||||
|
||||
|
||||
def test_replace_draft_workspace_document_identical_draft_is_a_noop(
|
||||
tmp_path,
|
||||
) -> None:
|
||||
store = FileDraftWorkspaceStore(tmp_path)
|
||||
create_draft_workspace(store, workspace_id="echo_draft", draft=_draft())
|
||||
current = store.get_workspace("echo_draft")
|
||||
|
||||
result = replace_draft_workspace_document(
|
||||
store,
|
||||
workspace_id="echo_draft",
|
||||
revision=1,
|
||||
draft=current.draft,
|
||||
node_defs_for_draft=lambda _draft: [],
|
||||
)
|
||||
|
||||
assert result["revision"] == 1
|
||||
assert store.get_workspace("echo_draft") == current
|
||||
|
||||
|
||||
def test_replace_draft_workspace_document_rejects_stale_revision(
|
||||
tmp_path,
|
||||
) -> None:
|
||||
store = FileDraftWorkspaceStore(tmp_path)
|
||||
create_draft_workspace(store, workspace_id="echo_draft", draft=_draft())
|
||||
patch_draft_workspace(
|
||||
store,
|
||||
workspace_id="echo_draft",
|
||||
revision=1,
|
||||
patch=[{"op": "replace", "path": "/name", "value": "current"}],
|
||||
)
|
||||
before = store.get_workspace("echo_draft")
|
||||
|
||||
result = replace_draft_workspace_document(
|
||||
store,
|
||||
workspace_id="echo_draft",
|
||||
revision=1,
|
||||
draft={**_draft(), "name": "stale"},
|
||||
node_defs_for_draft=lambda _draft: [],
|
||||
)
|
||||
|
||||
assert result["status"] == "conflict"
|
||||
assert result["diagnostics"][0]["code"] == "revision_conflict"
|
||||
assert store.get_workspace("echo_draft") == before
|
||||
|
||||
|
||||
def test_replace_draft_workspace_document_rejects_structural_failure(
|
||||
tmp_path,
|
||||
) -> None:
|
||||
store = FileDraftWorkspaceStore(tmp_path)
|
||||
create_draft_workspace(store, workspace_id="echo_draft", draft=_draft())
|
||||
before = store.get_workspace("echo_draft")
|
||||
|
||||
with pytest.raises(ValidationError):
|
||||
replace_draft_workspace_document(
|
||||
store,
|
||||
workspace_id="echo_draft",
|
||||
revision=1,
|
||||
draft={"name": "incomplete"},
|
||||
node_defs_for_draft=lambda _draft: [],
|
||||
)
|
||||
|
||||
assert store.get_workspace("echo_draft") == before
|
||||
|
||||
|
||||
def test_replace_draft_workspace_document_persists_fresh_invalid_diagnostics(
|
||||
tmp_path,
|
||||
) -> None:
|
||||
store = FileDraftWorkspaceStore(tmp_path)
|
||||
create_draft_workspace(store, workspace_id="echo_draft", draft=_draft())
|
||||
replacement = _draft()
|
||||
replacement["routes"] = {"echo": {"ok": "missing_step"}}
|
||||
|
||||
result = replace_draft_workspace_document(
|
||||
store,
|
||||
workspace_id="echo_draft",
|
||||
revision=1,
|
||||
draft=replacement,
|
||||
node_defs_for_draft=lambda _draft: [],
|
||||
)
|
||||
stored = store.get_workspace("echo_draft")
|
||||
|
||||
assert result["revision"] == 2
|
||||
assert result["status"] == "invalid"
|
||||
assert result["diagnostics"]
|
||||
assert stored.status == "invalid"
|
||||
assert stored.diagnostics == result["diagnostics"]
|
||||
assert stored.draft["routes"] == replacement["routes"]
|
||||
|
||||
|
||||
def test_get_draft_workspace_includes_full_draft_only_when_requested(tmp_path) -> None:
|
||||
store = FileDraftWorkspaceStore(tmp_path)
|
||||
create_draft_workspace(store, workspace_id="echo_draft", draft=_draft())
|
||||
|
||||
@@ -1539,6 +1539,74 @@ async def test_validate_draft_workspace_refreshes_status(tmp_path: Path) -> None
|
||||
assert fetched["status"] == "invalid"
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_replace_draft_workspace_document_uses_current_capability_definitions(
|
||||
tmp_path: Path,
|
||||
) -> None:
|
||||
api, _service, _authoring = _draft_api(
|
||||
FileWorkflowArtifactStore(tmp_path / "replace_document_capability_defs"),
|
||||
register_echo=True,
|
||||
)
|
||||
await api.create_draft_workspace(
|
||||
workspace_id="report",
|
||||
draft=_echo_draft(),
|
||||
)
|
||||
replacement = _echo_draft()
|
||||
replacement["name"] = "replacement"
|
||||
replacement["steps"]["echo"]["output"] = [
|
||||
{"source": "missing", "target": "state.echoed"}
|
||||
]
|
||||
|
||||
result = await api.replace_draft_workspace_document(
|
||||
workspace_id="report",
|
||||
revision=1,
|
||||
draft=replacement,
|
||||
)
|
||||
|
||||
assert result["revision"] == 2
|
||||
assert result["status"] == "invalid"
|
||||
assert any(
|
||||
"missing" in str(diagnostic.get("message", ""))
|
||||
for diagnostic in result["diagnostics"]
|
||||
)
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_replace_draft_workspace_document_refreshes_semantic_diagnostics(
|
||||
tmp_path: Path,
|
||||
) -> None:
|
||||
api, _service, _authoring = _draft_api(
|
||||
FileWorkflowArtifactStore(tmp_path / "replace_document_diagnostics"),
|
||||
register_echo=True,
|
||||
)
|
||||
await api.create_draft_workspace(
|
||||
workspace_id="report",
|
||||
draft=_echo_draft(),
|
||||
)
|
||||
before = await api.get_draft_workspace(
|
||||
workspace_id="report",
|
||||
include_draft=True,
|
||||
)
|
||||
replacement = _echo_draft()
|
||||
replacement["routes"] = {"echo": {"ok": "missing_step"}}
|
||||
|
||||
result = await api.replace_draft_workspace_document(
|
||||
workspace_id="report",
|
||||
revision=1,
|
||||
draft=replacement,
|
||||
)
|
||||
after = await api.get_draft_workspace(
|
||||
workspace_id="report",
|
||||
include_draft=True,
|
||||
)
|
||||
|
||||
assert before["diagnostics"] == []
|
||||
assert result["status"] == "invalid"
|
||||
assert result["diagnostics"]
|
||||
assert after["diagnostics"] == result["diagnostics"]
|
||||
assert after["draft"]["routes"] == replacement["routes"]
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_validate_draft_workspace_suggests_bind(
|
||||
tmp_path: Path,
|
||||
|
||||
Reference in New Issue
Block a user