feat: type exceptional draft results

This commit is contained in:
lda
2026-07-31 13:01:05 +07:00 Verified
parent e60c391610
commit c96fee4e7d
11 changed files with 391 additions and 122 deletions
+6 -6
View File
@@ -72,12 +72,12 @@
complete OpenRPC document for all 70 registered methods. Request payloads complete OpenRPC document for all 70 registered methods. Request payloads
retain useful Pydantic schemas, so OpenRPC is a viable transport input. retain useful Pydantic schemas, so OpenRPC is a viable transport input.
- Typed-result slices now give `workflow.health`, all artifact, deployment, - Typed-result slices now give `workflow.health`, all artifact, deployment,
and run operations, and 26 uniform persisted draft-workspace operations and run operations, every persisted draft-workspace operation, and the
named transport-neutral result schemas: 42 of 70 methods. The remaining 28 capability-bootstrap alias named transport-neutral result schemas: 47 of
success results still collapse to generic objects, including exceptional 70 methods. The remaining 23 success results still collapse to generic
draft compile/save shapes and capability/source/admin operations. Continue objects across capability discovery/call, stateless draft patch/validate,
introducing operation result DTOs before adopting generated TypeScript source discovery, and source-registry/admin operations. Continue introducing
contracts. operation result DTOs before adopting generated TypeScript contracts.
- The stock `@open-rpc/generator` TypeScript client is not suitable here. It - The stock `@open-rpc/generator` TypeScript client is not suitable here. It
exhausted a 4 GB Node heap on the full contract and emitted invalid dotted exhausted a 4 GB Node heap on the full contract and emitted invalid dotted
class members plus `any` results for a minimal `workflow.health` contract. class members plus `any` results for a minimal `workflow.health` contract.
+37 -27
View File
@@ -25,11 +25,14 @@ from .capability_requirements import observed_node_specs
from .drafts import WorkflowDraftApi from .drafts import WorkflowDraftApi
from .listing import matches_query, paged_list_payload from .listing import matches_query, paged_list_payload
from .models import ( from .models import (
CreateArtifactFromWorkspaceResult,
DeleteArtifactResult, DeleteArtifactResult,
JsonProjector, JsonProjector,
ListArtifactsResult, ListArtifactsResult,
RawWorkflowPlan, RawWorkflowPlan,
SaveArtifactResult, SaveArtifactResult,
SavedDraftArtifactResult,
UnsavedDraftArtifactResult,
WorkflowArtifactPayload, WorkflowArtifactPayload,
) )
from .operation_context import WorkflowOperationContext from .operation_context import WorkflowOperationContext
@@ -38,6 +41,8 @@ _PROJECT_ARTIFACT = JsonProjector(WorkflowArtifactPayload)
_PROJECT_ARTIFACT_LIST = JsonProjector(ListArtifactsResult) _PROJECT_ARTIFACT_LIST = JsonProjector(ListArtifactsResult)
_PROJECT_ARTIFACT_SAVE = JsonProjector(SaveArtifactResult) _PROJECT_ARTIFACT_SAVE = JsonProjector(SaveArtifactResult)
_PROJECT_ARTIFACT_DELETE = JsonProjector(DeleteArtifactResult) _PROJECT_ARTIFACT_DELETE = JsonProjector(DeleteArtifactResult)
_PROJECT_UNSAVED_DRAFT_ARTIFACT = JsonProjector(UnsavedDraftArtifactResult)
_PROJECT_SAVED_DRAFT_ARTIFACT = JsonProjector(SavedDraftArtifactResult)
class WorkflowArtifactApi: class WorkflowArtifactApi:
@@ -95,9 +100,7 @@ class WorkflowArtifactApi:
paged_list_payload("nodes", entries, cursor=cursor, limit=limit) paged_list_payload("nodes", entries, cursor=cursor, limit=limit)
) )
async def save_artifact( async def save_artifact(self, artifact: dict[str, Any]) -> SaveArtifactResult:
self, artifact: dict[str, Any]
) -> SaveArtifactResult:
workflow_artifact = WorkflowArtifact.model_validate(artifact) workflow_artifact = WorkflowArtifact.model_validate(artifact)
self._artifact_store().save_artifact(workflow_artifact) self._artifact_store().save_artifact(workflow_artifact)
self.context.events.record_workflow_event( self.context.events.record_workflow_event(
@@ -182,7 +185,7 @@ class WorkflowArtifactApi:
required_capabilities: dict[str, dict[str, Any]] | None = None, required_capabilities: dict[str, dict[str, Any]] | None = None,
source_bindings: dict[str, str] | None = None, source_bindings: dict[str, str] | None = None,
created_from_catalog_version: str | None = None, created_from_catalog_version: str | None = None,
) -> dict[str, Any]: ) -> SavedDraftArtifactResult:
from wf_artifacts import compile_workflow_draft from wf_artifacts import compile_workflow_draft
plan = compile_workflow_draft(draft) plan = compile_workflow_draft(draft)
@@ -202,6 +205,24 @@ class WorkflowArtifactApi:
observed_node_specs=observed_node_specs(self.context), observed_node_specs=observed_node_specs(self.context),
created_from_catalog_version=created_from_catalog_version, created_from_catalog_version=created_from_catalog_version,
) )
required_sources = _binding_required_sources(
workflow_artifact.required_capability_map(),
self.context.specs.capability_sources,
)
result = _PROJECT_SAVED_DRAFT_ARTIFACT(
{
"artifact_id": workflow_artifact.id,
"version": workflow_artifact.version,
"saved": True,
"required_logical_sources": required_sources,
"suggested_bindings": _suggested_self_bindings(
required_sources,
self.context.specs.capability_sources,
),
}
)
# Validate the public result before persistence so a projection bug
# cannot report failure after the artifact has already been saved.
self._artifact_store().save_artifact(workflow_artifact) self._artifact_store().save_artifact(workflow_artifact)
self.context.events.record_workflow_event( self.context.events.record_workflow_event(
"workflow_artifact_saved", "workflow_artifact_saved",
@@ -212,20 +233,7 @@ class WorkflowArtifactApi:
"created_from_draft": True, "created_from_draft": True,
}, },
) )
required_sources = _binding_required_sources( return result
workflow_artifact.required_capability_map(),
self.context.specs.capability_sources,
)
return {
"artifact_id": workflow_artifact.id,
"version": workflow_artifact.version,
"saved": True,
"required_logical_sources": required_sources,
"suggested_bindings": _suggested_self_bindings(
required_sources,
self.context.specs.capability_sources,
),
}
async def create_artifact_from_workspace( async def create_artifact_from_workspace(
self, self,
@@ -240,20 +248,22 @@ class WorkflowArtifactApi:
required_capabilities: dict[str, dict[str, Any]] | None = None, required_capabilities: dict[str, dict[str, Any]] | None = None,
source_bindings: dict[str, str] | None = None, source_bindings: dict[str, str] | None = None,
created_from_catalog_version: str | None = None, created_from_catalog_version: str | None = None,
) -> dict[str, Any]: ) -> CreateArtifactFromWorkspaceResult:
store = self.context.draft_workspace_store store = self.context.draft_workspace_store
if store is None: if store is None:
raise KeyError("draft workspace store is not configured") raise KeyError("draft workspace store is not configured")
workspace = store.get_workspace(workspace_id) workspace = store.get_workspace(workspace_id)
validation = await self.drafts.validate_draft(draft=workspace.draft) validation = await self.drafts.validate_draft(draft=workspace.draft)
if validation["status"] != "valid": if validation["status"] != "valid":
return { return _PROJECT_UNSAVED_DRAFT_ARTIFACT(
"saved": False, {
"workspace_id": workspace_id, "saved": False,
"revision": workspace.revision, "workspace_id": workspace_id,
"status": validation["status"], "revision": workspace.revision,
"diagnostics": validation["diagnostics"], "status": validation["status"],
} "diagnostics": validation["diagnostics"],
}
)
return await self.create_artifact_from_draft( return await self.create_artifact_from_draft(
artifact_id=artifact_id, artifact_id=artifact_id,
version=version, version=version,
@@ -279,7 +289,7 @@ class WorkflowArtifactApi:
required_capabilities: dict[str, dict[str, Any]] | None = None, required_capabilities: dict[str, dict[str, Any]] | None = None,
source_bindings: dict[str, str] | None = None, source_bindings: dict[str, str] | None = None,
created_from_catalog_version: str | None = None, created_from_catalog_version: str | None = None,
) -> dict[str, Any]: ) -> CreateArtifactFromWorkspaceResult:
"""Save the current draft workspace as a callable wrapper artifact.""" """Save the current draft workspace as a callable wrapper artifact."""
return await self.create_artifact_from_workspace( return await self.create_artifact_from_workspace(
workspace_id=workspace_id, workspace_id=workspace_id,
+28 -9
View File
@@ -1,7 +1,7 @@
from __future__ import annotations from __future__ import annotations
from collections.abc import Mapping, Sequence from collections.abc import Mapping, Sequence
from typing import Any from typing import Any, cast
from wf_artifacts import ( from wf_artifacts import (
DependencyDiagnostic, DependencyDiagnostic,
@@ -21,6 +21,12 @@ from .capability_requirements import required_capability_payloads
from .draft_authoring import WorkflowDraftAuthoringApi from .draft_authoring import WorkflowDraftAuthoringApi
from .drafts import WorkflowDraftApi from .drafts import WorkflowDraftApi
from .listing import matches_query, paged_list_payload from .listing import matches_query, paged_list_payload
from .models import (
CreateDraftWorkspaceFromCapabilityResult,
JsonProjector,
NextActionsPayload,
WrapperAuthoringHintsPayload,
)
from .next_actions import NextActions from .next_actions import NextActions
from .operation_context import WorkflowOperationContext from .operation_context import WorkflowOperationContext
from .refs import parse_workflow_surface_capability_id from .refs import parse_workflow_surface_capability_id
@@ -30,6 +36,8 @@ from .wrapper_hints import (
wrapper_hints_for_capability, wrapper_hints_for_capability,
) )
_PROJECT_WRAPPER_HINTS = JsonProjector(WrapperAuthoringHintsPayload)
def _schema_field_names(schema: dict[str, Any]) -> list[str]: def _schema_field_names(schema: dict[str, Any]) -> list[str]:
"""Return top-level JSON object property names for compact discovery rows.""" """Return top-level JSON object property names for compact discovery rows."""
@@ -351,10 +359,12 @@ class WorkflowCapabilityApi:
input_map: dict[str, str] | None = None, input_map: dict[str, str] | None = None,
output_map: dict[str, str] | None = None, output_map: dict[str, str] | None = None,
error_message_source: str | GraphSourcePath | None = None, error_message_source: str | GraphSourcePath | None = None,
) -> dict[str, Any]: ) -> CreateDraftWorkspaceFromCapabilityResult:
"""Create a patchable draft workspace from inspect_capability hints.""" """Create a patchable draft workspace from inspect_capability hints."""
capability = await self.inspect_capability(qualified_name=capability_name) capability = await self.inspect_capability(qualified_name=capability_name)
hints = capability["wrapper_hints"] # Validate capability-derived guidance before workspace creation. The
# workspace result is already projected by the draft-workspace API.
hints = _PROJECT_WRAPPER_HINTS(capability["wrapper_hints"])
result = await self.draft_authoring.create_minimal_draft_workspace( result = await self.draft_authoring.create_minimal_draft_workspace(
workspace_id=workspace_id, workspace_id=workspace_id,
name=name or _draft_name_from_capability(capability_name), name=name or _draft_name_from_capability(capability_name),
@@ -371,12 +381,21 @@ class WorkflowCapabilityApi:
error_message_source=error_message_source, error_message_source=error_message_source,
title=title, title=title,
) )
return { next_actions = cast(
**result, NextActionsPayload,
"wrapper_hints": hints, NextActions.from_wrapper_hints(
"next_actions": NextActions.from_wrapper_hints(
workspace_id=workspace_id, workspace_id=workspace_id,
revision=int(result["revision"]), revision=int(result["revision"]),
hints=hints, hints=dict(hints),
).model_dump(mode="json"), ).model_dump(mode="json"),
} )
# Both merged extensions are independently validated or model-dumped,
# while the workspace payload was projected by the mutating API.
return cast(
CreateDraftWorkspaceFromCapabilityResult,
{
**result,
"wrapper_hints": hints,
"next_actions": next_actions,
},
)
+24 -13
View File
@@ -53,8 +53,11 @@ from .draft_payloads import (
output_bindings_payload as _draft_output_bindings_payload, output_bindings_payload as _draft_output_bindings_payload,
) )
from .models import ( from .models import (
CompileDraftWorkspaceResult,
CompileDraftWorkspaceSuccess,
DeleteDraftWorkspaceResult, DeleteDraftWorkspaceResult,
DraftWorkspaceResult, DraftWorkspaceResult,
InvalidDraftResult,
JsonProjector, JsonProjector,
ListDraftWorkspacesResult, ListDraftWorkspacesResult,
) )
@@ -64,6 +67,8 @@ from .schema_projection import project_property_to_schema_path, schema_path_exis
_PROJECT_DRAFT_WORKSPACE = JsonProjector(DraftWorkspaceResult) _PROJECT_DRAFT_WORKSPACE = JsonProjector(DraftWorkspaceResult)
_PROJECT_DRAFT_WORKSPACE_LIST = JsonProjector(ListDraftWorkspacesResult) _PROJECT_DRAFT_WORKSPACE_LIST = JsonProjector(ListDraftWorkspacesResult)
_PROJECT_DRAFT_WORKSPACE_DELETE = JsonProjector(DeleteDraftWorkspaceResult) _PROJECT_DRAFT_WORKSPACE_DELETE = JsonProjector(DeleteDraftWorkspaceResult)
_PROJECT_DRAFT_COMPILE = JsonProjector(CompileDraftWorkspaceSuccess)
_PROJECT_INVALID_DRAFT = JsonProjector(InvalidDraftResult)
def _empty_object_schema() -> dict[str, Any]: def _empty_object_schema() -> dict[str, Any]:
@@ -176,18 +181,22 @@ class WorkflowDraftApi:
node_defs=self._node_defs_for_draft(draft), node_defs=self._node_defs_for_draft(draft),
) )
async def compile_draft(self, *, draft: dict[str, Any]) -> dict[str, Any]: async def compile_draft(
self, *, draft: dict[str, Any]
) -> CompileDraftWorkspaceSuccess:
plan = compile_workflow_draft(draft) plan = compile_workflow_draft(draft)
return { return _PROJECT_DRAFT_COMPILE(
"compiled_plan": plan, {
"required_capabilities": required_capability_payloads( "compiled_plan": plan,
required_capabilities_for_plan( "required_capabilities": required_capability_payloads(
plan, required_capabilities_for_plan(
source_bindings=None, plan,
context=self.context, source_bindings=None,
) context=self.context,
), )
} ),
}
)
async def patch_draft( async def patch_draft(
self, self,
@@ -331,12 +340,14 @@ class WorkflowDraftApi:
get_draft_workspace_record(store, workspace_id=workspace_id) get_draft_workspace_record(store, workspace_id=workspace_id)
) )
async def compile_draft_workspace(self, *, workspace_id: str) -> dict[str, Any]: async def compile_draft_workspace(
self, *, workspace_id: str
) -> CompileDraftWorkspaceResult:
"""Compile a stored draft workspace without mutating it.""" """Compile a stored draft workspace without mutating it."""
workspace = self._draft_store().get_workspace(workspace_id) workspace = self._draft_store().get_workspace(workspace_id)
validation = await self.validate_draft(draft=workspace.draft) validation = await self.validate_draft(draft=workspace.draft)
if validation["status"] != "valid": if validation["status"] != "valid":
return validation return _PROJECT_INVALID_DRAFT(validation)
return await self.compile_draft(draft=workspace.draft) return await self.compile_draft(draft=workspace.draft)
async def patch_draft_workspace( async def patch_draft_workspace(
+20
View File
@@ -35,11 +35,21 @@ from .deployments import (
WorkflowDeploymentPayload, WorkflowDeploymentPayload,
) )
from .drafts import ( from .drafts import (
CompileDraftWorkspaceResult,
CompileDraftWorkspaceSuccess,
CreateArtifactFromWorkspaceResult,
CreateDraftWorkspaceFromCapabilityResult,
DeleteDraftWorkspaceResult, DeleteDraftWorkspaceResult,
DraftDiagnosticPayload, DraftDiagnosticPayload,
DraftWorkspaceResult, DraftWorkspaceResult,
DraftWorkspaceSummary, DraftWorkspaceSummary,
InvalidDraftResult,
ListDraftWorkspacesResult, ListDraftWorkspacesResult,
SavedDraftArtifactResult,
UnsavedDraftArtifactResult,
WrapperAuthoringHintsPayload,
WrapperMissingDecisionPayload,
WrapperOutcomeCandidatePayload,
) )
from .runs import ( from .runs import (
InterruptPayload, InterruptPayload,
@@ -60,6 +70,10 @@ __all__ = [
"ArtifactKindPayload", "ArtifactKindPayload",
"CapabilityKindPayload", "CapabilityKindPayload",
"CapabilityRefPayload", "CapabilityRefPayload",
"CompileDraftWorkspaceResult",
"CompileDraftWorkspaceSuccess",
"CreateArtifactFromWorkspaceResult",
"CreateDraftWorkspaceFromCapabilityResult",
"DeleteArtifactResult", "DeleteArtifactResult",
"DeleteDeploymentResult", "DeleteDeploymentResult",
"DeleteDraftWorkspaceResult", "DeleteDraftWorkspaceResult",
@@ -72,6 +86,7 @@ __all__ = [
"GuidedResultPayload", "GuidedResultPayload",
"InterruptPayload", "InterruptPayload",
"InterruptRoutePayload", "InterruptRoutePayload",
"InvalidDraftResult",
"JsonObject", "JsonObject",
"JsonProjector", "JsonProjector",
"JsonSchema", "JsonSchema",
@@ -90,12 +105,17 @@ __all__ = [
"RunTraceResult", "RunTraceResult",
"RequiredCapabilityPayload", "RequiredCapabilityPayload",
"SaveArtifactResult", "SaveArtifactResult",
"SavedDraftArtifactResult",
"SaveDeploymentResult", "SaveDeploymentResult",
"SourceBindingPayload", "SourceBindingPayload",
"TraceRange", "TraceRange",
"TraceEntryPayload", "TraceEntryPayload",
"UnsavedDraftArtifactResult",
"ValidateDeploymentResult", "ValidateDeploymentResult",
"WorkflowDeploymentPayload", "WorkflowDeploymentPayload",
"WorkflowArtifactPayload", "WorkflowArtifactPayload",
"WorkflowRefPayload", "WorkflowRefPayload",
"WrapperAuthoringHintsPayload",
"WrapperMissingDecisionPayload",
"WrapperOutcomeCandidatePayload",
] ]
+95 -1
View File
@@ -2,7 +2,8 @@
from typing import Any, Literal, NotRequired, TypedDict from typing import Any, Literal, NotRequired, TypedDict
from .common import JsonObject from .artifacts import RequiredCapabilityPayload
from .common import JsonObject, NextActionsPayload
class DraftDiagnosticPayload(TypedDict): class DraftDiagnosticPayload(TypedDict):
@@ -51,3 +52,96 @@ class DeleteDraftWorkspaceResult(TypedDict):
workspace_id: str workspace_id: str
deleted: bool deleted: bool
status: Literal["deleted", "not_found"] status: Literal["deleted", "not_found"]
class InvalidDraftResult(TypedDict):
"""Validation-only result returned when a draft cannot be compiled."""
status: Literal["invalid"]
diagnostics: list[DraftDiagnosticPayload]
class CompileDraftWorkspaceSuccess(TypedDict):
"""Compiled raw plan and dependencies for one valid draft workspace."""
compiled_plan: JsonObject
required_capabilities: dict[str, RequiredCapabilityPayload]
type CompileDraftWorkspaceResult = CompileDraftWorkspaceSuccess | InvalidDraftResult
class WrapperOutcomeCandidatePayload(TypedDict):
"""Possible wrapper outcome mapping that still requires caller judgment."""
kind: Literal["boolean_control_field"]
source: str
candidate_outcomes: list[str]
confidence: Literal["high", "medium", "low"]
reason: str
automatic: bool
class WrapperMissingDecisionPayload(TypedDict):
"""One unresolved wrapper-authoring decision."""
kind: Literal[
"choose_output_fields",
"review_nested_output",
"confirm_boolean_outcomes",
"choose_error_mapping",
]
message: str
class WrapperAuthoringHintsPayload(TypedDict):
"""Transport-neutral projection of conservative wrapper scaffolding hints."""
capability_name: str
confidence: Literal["high", "medium", "low"]
declared_outcomes: list[str]
suggested_wrapper_outcomes: list[str]
outcome_policy: Literal[
"preserve_declared",
"manual_mapping_required",
]
input_schema: JsonObject
state_schema: JsonObject
output_schema: JsonObject
input_map: dict[str, str]
output_map: dict[str, str]
outcome_candidates: list[WrapperOutcomeCandidatePayload]
missing_decisions: list[WrapperMissingDecisionPayload]
notes: list[str]
class CreateDraftWorkspaceFromCapabilityResult(DraftWorkspaceResult):
"""Bootstrapped workspace plus the hints and next actions that shaped it."""
wrapper_hints: WrapperAuthoringHintsPayload
next_actions: NextActionsPayload
class UnsavedDraftArtifactResult(TypedDict):
"""Validation result when an invalid workspace is not saved as an artifact."""
saved: Literal[False]
workspace_id: str
revision: int
status: Literal["invalid"]
diagnostics: list[DraftDiagnosticPayload]
class SavedDraftArtifactResult(TypedDict):
"""Saved artifact result with source-binding guidance derived from a draft."""
artifact_id: str
version: int
saved: Literal[True]
required_logical_sources: list[str]
suggested_bindings: dict[str, str]
type CreateArtifactFromWorkspaceResult = (
SavedDraftArtifactResult | UnsavedDraftArtifactResult
)
+11 -6
View File
@@ -14,6 +14,10 @@ from .draft_authoring import RouteSource, WorkflowDraftAuthoringApi
from .draft_updates import CapabilityStepUpdate from .draft_updates import CapabilityStepUpdate
from .drafts import WorkflowDraftApi from .drafts import WorkflowDraftApi
from .models import ( from .models import (
CompileDraftWorkspaceResult,
CompileDraftWorkspaceSuccess,
CreateArtifactFromWorkspaceResult,
CreateDraftWorkspaceFromCapabilityResult,
DeleteArtifactResult, DeleteArtifactResult,
DeleteDeploymentResult, DeleteDeploymentResult,
DeleteDraftWorkspaceResult, DeleteDraftWorkspaceResult,
@@ -26,6 +30,7 @@ from .models import (
RunResult, RunResult,
RunTraceResult, RunTraceResult,
SaveArtifactResult, SaveArtifactResult,
SavedDraftArtifactResult,
SaveDeploymentResult, SaveDeploymentResult,
ValidateDeploymentResult, ValidateDeploymentResult,
WorkflowArtifactPayload, WorkflowArtifactPayload,
@@ -174,7 +179,7 @@ class WorkflowApi:
required_capabilities: dict[str, dict[str, Any]] | None = None, required_capabilities: dict[str, dict[str, Any]] | None = None,
source_bindings: dict[str, str] | None = None, source_bindings: dict[str, str] | None = None,
created_from_catalog_version: str | None = None, created_from_catalog_version: str | None = None,
) -> dict[str, Any]: ) -> SavedDraftArtifactResult:
return await self.artifacts.create_artifact_from_draft( return await self.artifacts.create_artifact_from_draft(
artifact_id=artifact_id, artifact_id=artifact_id,
version=version, version=version,
@@ -201,7 +206,7 @@ class WorkflowApi:
required_capabilities: dict[str, dict[str, Any]] | None = None, required_capabilities: dict[str, dict[str, Any]] | None = None,
source_bindings: dict[str, str] | None = None, source_bindings: dict[str, str] | None = None,
created_from_catalog_version: str | None = None, created_from_catalog_version: str | None = None,
) -> dict[str, Any]: ) -> CreateArtifactFromWorkspaceResult:
return await self.artifacts.create_artifact_from_workspace( return await self.artifacts.create_artifact_from_workspace(
workspace_id=workspace_id, workspace_id=workspace_id,
artifact_id=artifact_id, artifact_id=artifact_id,
@@ -227,7 +232,7 @@ class WorkflowApi:
required_capabilities: dict[str, dict[str, Any]] | None = None, required_capabilities: dict[str, dict[str, Any]] | None = None,
source_bindings: dict[str, str] | None = None, source_bindings: dict[str, str] | None = None,
created_from_catalog_version: str | None = None, created_from_catalog_version: str | None = None,
) -> dict[str, Any]: ) -> CreateArtifactFromWorkspaceResult:
return await self.artifacts.create_wrapper_from_workspace( return await self.artifacts.create_wrapper_from_workspace(
workspace_id=workspace_id, workspace_id=workspace_id,
artifact_id=artifact_id, artifact_id=artifact_id,
@@ -253,7 +258,7 @@ class WorkflowApi:
self, self,
*, *,
draft: dict[str, Any], draft: dict[str, Any],
) -> dict[str, Any]: ) -> CompileDraftWorkspaceSuccess:
return await self.drafts.compile_draft(draft=draft) return await self.drafts.compile_draft(draft=draft)
async def patch_draft( async def patch_draft(
@@ -332,7 +337,7 @@ class WorkflowApi:
self, self,
*, *,
workspace_id: str, workspace_id: str,
) -> dict[str, Any]: ) -> CompileDraftWorkspaceResult:
return await self.drafts.compile_draft_workspace(workspace_id=workspace_id) return await self.drafts.compile_draft_workspace(workspace_id=workspace_id)
async def patch_draft_workspace( async def patch_draft_workspace(
@@ -725,7 +730,7 @@ class WorkflowApi:
input_map: dict[str, str] | None = None, input_map: dict[str, str] | None = None,
output_map: dict[str, str] | None = None, output_map: dict[str, str] | None = None,
error_message_source: Any | None = None, error_message_source: Any | None = None,
) -> dict[str, Any]: ) -> CreateDraftWorkspaceFromCapabilityResult:
return await self.capabilities.create_draft_workspace_from_capability( return await self.capabilities.create_draft_workspace_from_capability(
workspace_id=workspace_id, workspace_id=workspace_id,
capability_name=capability_name, capability_name=capability_name,
+7 -4
View File
@@ -10,6 +10,9 @@ from wf_core.models.steps import InputBinding, OutputBinding
from .draft_authoring import RouteSource from .draft_authoring import RouteSource
from .draft_updates import CapabilityStepUpdate from .draft_updates import CapabilityStepUpdate
from .models import ( from .models import (
CompileDraftWorkspaceResult,
CreateArtifactFromWorkspaceResult,
CreateDraftWorkspaceFromCapabilityResult,
DeleteArtifactResult, DeleteArtifactResult,
DeleteDeploymentResult, DeleteDeploymentResult,
DeleteDraftWorkspaceResult, DeleteDraftWorkspaceResult,
@@ -87,7 +90,7 @@ class WorkflowDraftSurface(Protocol):
input_map: dict[str, str] | None = None, input_map: dict[str, str] | None = None,
output_map: dict[str, str] | None = None, output_map: dict[str, str] | None = None,
error_message_source: Any | None = None, error_message_source: Any | None = None,
) -> dict[str, Any]: ... ) -> CreateDraftWorkspaceFromCapabilityResult: ...
async def create_empty_draft_workspace( async def create_empty_draft_workspace(
self, self,
@@ -314,7 +317,7 @@ class WorkflowDraftSurface(Protocol):
self, self,
*, *,
workspace_id: str, workspace_id: str,
) -> dict[str, Any]: ... ) -> CompileDraftWorkspaceResult: ...
async def delete_draft_workspace( async def delete_draft_workspace(
self, self,
@@ -335,7 +338,7 @@ class WorkflowDraftSurface(Protocol):
required_capabilities: dict[str, dict[str, Any]] | None = None, required_capabilities: dict[str, dict[str, Any]] | None = None,
source_bindings: dict[str, str] | None = None, source_bindings: dict[str, str] | None = None,
created_from_catalog_version: str | None = None, created_from_catalog_version: str | None = None,
) -> dict[str, Any]: ... ) -> CreateArtifactFromWorkspaceResult: ...
async def create_wrapper_from_workspace( async def create_wrapper_from_workspace(
self, self,
@@ -349,7 +352,7 @@ class WorkflowDraftSurface(Protocol):
required_capabilities: dict[str, dict[str, Any]] | None = None, required_capabilities: dict[str, dict[str, Any]] | None = None,
source_bindings: dict[str, str] | None = None, source_bindings: dict[str, str] | None = None,
created_from_catalog_version: str | None = None, created_from_catalog_version: str | None = None,
) -> dict[str, Any]: ... ) -> CreateArtifactFromWorkspaceResult: ...
class WorkflowArtifactSurface(Protocol): class WorkflowArtifactSurface(Protocol):
+65 -50
View File
@@ -5,6 +5,9 @@ from typing import Any, Literal, cast
from wf_api import CapabilityStepUpdate from wf_api import CapabilityStepUpdate
from wf_api.models import ( from wf_api.models import (
CompileDraftWorkspaceResult,
CreateArtifactFromWorkspaceResult,
CreateDraftWorkspaceFromCapabilityResult,
DeleteDraftWorkspaceResult, DeleteDraftWorkspaceResult,
DraftWorkspaceResult, DraftWorkspaceResult,
ListDraftWorkspacesResult, ListDraftWorkspacesResult,
@@ -62,23 +65,26 @@ class RpcDraftClientMixin:
input_map: dict[str, str] | None = None, input_map: dict[str, str] | None = None,
output_map: dict[str, str] | None = None, output_map: dict[str, str] | None = None,
error_message_source: Any | None = None, error_message_source: Any | None = None,
) -> dict[str, Any]: ) -> CreateDraftWorkspaceFromCapabilityResult:
return await self._call( return cast(
"workflow.draft_workspaces.create_from_capability", CreateDraftWorkspaceFromCapabilityResult,
{ await self._call(
"workspace_id": workspace_id, "workflow.draft_workspaces.create_from_capability",
"capability_name": capability_name, {
"name": name, "workspace_id": workspace_id,
"title": title, "capability_name": capability_name,
"input_schema": input_schema, "name": name,
"state_schema": state_schema, "title": title,
"output_schema": output_schema, "input_schema": input_schema,
"input": input, "state_schema": state_schema,
"output": output, "output_schema": output_schema,
"input_map": input_map, "input": input,
"output_map": output_map, "output": output,
"error_message_source": error_message_source, "input_map": input_map,
}, "output_map": output_map,
"error_message_source": error_message_source,
},
),
) )
async def create_empty_draft_workspace( async def create_empty_draft_workspace(
@@ -546,10 +552,13 @@ class RpcDraftClientMixin:
self: RpcCaller, self: RpcCaller,
*, *,
workspace_id: str, workspace_id: str,
) -> dict[str, Any]: ) -> CompileDraftWorkspaceResult:
return await self._call( return cast(
"workflow.draft_workspaces.compile", CompileDraftWorkspaceResult,
{"workspace_id": workspace_id}, await self._call(
"workflow.draft_workspaces.compile",
{"workspace_id": workspace_id},
),
) )
async def delete_draft_workspace( async def delete_draft_workspace(
@@ -578,21 +587,24 @@ class RpcDraftClientMixin:
required_capabilities: dict[str, dict[str, Any]] | None = None, required_capabilities: dict[str, dict[str, Any]] | None = None,
source_bindings: dict[str, str] | None = None, source_bindings: dict[str, str] | None = None,
created_from_catalog_version: str | None = None, created_from_catalog_version: str | None = None,
) -> dict[str, Any]: ) -> CreateArtifactFromWorkspaceResult:
return await self._call( return cast(
"workflow.draft_workspaces.create_artifact", CreateArtifactFromWorkspaceResult,
{ await self._call(
"workspace_id": workspace_id, "workflow.draft_workspaces.create_artifact",
"artifact_id": artifact_id, {
"version": version, "workspace_id": workspace_id,
"title": title, "artifact_id": artifact_id,
"outcomes": list(outcomes), "version": version,
"kind": kind, "title": title,
"description": description, "outcomes": list(outcomes),
"required_capabilities": required_capabilities, "kind": kind,
"source_bindings": source_bindings, "description": description,
"created_from_catalog_version": created_from_catalog_version, "required_capabilities": required_capabilities,
}, "source_bindings": source_bindings,
"created_from_catalog_version": created_from_catalog_version,
},
),
) )
async def create_wrapper_from_workspace( async def create_wrapper_from_workspace(
@@ -607,18 +619,21 @@ class RpcDraftClientMixin:
required_capabilities: dict[str, dict[str, Any]] | None = None, required_capabilities: dict[str, dict[str, Any]] | None = None,
source_bindings: dict[str, str] | None = None, source_bindings: dict[str, str] | None = None,
created_from_catalog_version: str | None = None, created_from_catalog_version: str | None = None,
) -> dict[str, Any]: ) -> CreateArtifactFromWorkspaceResult:
return await self._call( return cast(
"workflow.draft_workspaces.create_wrapper", CreateArtifactFromWorkspaceResult,
{ await self._call(
"workspace_id": workspace_id, "workflow.draft_workspaces.create_wrapper",
"artifact_id": artifact_id, {
"version": version, "workspace_id": workspace_id,
"title": title, "artifact_id": artifact_id,
"outcomes": list(outcomes), "version": version,
"description": description, "title": title,
"required_capabilities": required_capabilities, "outcomes": list(outcomes),
"source_bindings": source_bindings, "description": description,
"created_from_catalog_version": created_from_catalog_version, "required_capabilities": required_capabilities,
}, "source_bindings": source_bindings,
"created_from_catalog_version": created_from_catalog_version,
},
),
) )
+9 -6
View File
@@ -9,6 +9,9 @@ from typing import Any
import fastapi_jsonrpc as jsonrpc import fastapi_jsonrpc as jsonrpc
from wf_api.models import ( from wf_api.models import (
CompileDraftWorkspaceResult,
CreateArtifactFromWorkspaceResult,
CreateDraftWorkspaceFromCapabilityResult,
DeleteDraftWorkspaceResult, DeleteDraftWorkspaceResult,
DraftWorkspaceResult, DraftWorkspaceResult,
ListDraftWorkspacesResult, ListDraftWorkspacesResult,
@@ -65,7 +68,7 @@ def register_methods(
) )
async def workflow_drafts_create_from_capability( async def workflow_drafts_create_from_capability(
params: CreateDraftFromCapabilityParams = RpcParams(), params: CreateDraftFromCapabilityParams = RpcParams(),
) -> dict[str, Any]: ) -> CreateDraftWorkspaceFromCapabilityResult:
try: try:
return await _create_from_capability(server, params) return await _create_from_capability(server, params)
except (ValueError, KeyError, LookupError, FileNotFoundError) as exc: except (ValueError, KeyError, LookupError, FileNotFoundError) as exc:
@@ -116,7 +119,7 @@ def register_methods(
) )
async def workflow_draft_workspaces_create_from_capability( async def workflow_draft_workspaces_create_from_capability(
params: CreateDraftFromCapabilityParams = RpcParams(), params: CreateDraftFromCapabilityParams = RpcParams(),
) -> dict[str, Any]: ) -> CreateDraftWorkspaceFromCapabilityResult:
try: try:
return await _create_from_capability(server, params) return await _create_from_capability(server, params)
except (ValueError, KeyError, LookupError, FileNotFoundError) as exc: except (ValueError, KeyError, LookupError, FileNotFoundError) as exc:
@@ -482,7 +485,7 @@ def register_methods(
) )
async def workflow_draft_workspaces_compile( async def workflow_draft_workspaces_compile(
params: CompileDraftWorkspaceParams = RpcParams(), params: CompileDraftWorkspaceParams = RpcParams(),
) -> dict[str, Any]: ) -> CompileDraftWorkspaceResult:
try: try:
return await server.api.compile_draft_workspace( return await server.api.compile_draft_workspace(
workspace_id=params.workspace_id, workspace_id=params.workspace_id,
@@ -556,7 +559,7 @@ def register_methods(
) )
async def workflow_draft_workspaces_create_artifact( async def workflow_draft_workspaces_create_artifact(
params: CreateArtifactFromWorkspaceParams = RpcParams(), params: CreateArtifactFromWorkspaceParams = RpcParams(),
) -> dict[str, Any]: ) -> CreateArtifactFromWorkspaceResult:
try: try:
return await server.api.create_artifact_from_workspace( return await server.api.create_artifact_from_workspace(
workspace_id=params.workspace_id, workspace_id=params.workspace_id,
@@ -578,7 +581,7 @@ def register_methods(
) )
async def workflow_draft_workspaces_create_wrapper( async def workflow_draft_workspaces_create_wrapper(
params: CreateWrapperFromWorkspaceParams = RpcParams(), params: CreateWrapperFromWorkspaceParams = RpcParams(),
) -> dict[str, Any]: ) -> CreateArtifactFromWorkspaceResult:
try: try:
return await server.api.create_wrapper_from_workspace( return await server.api.create_wrapper_from_workspace(
workspace_id=params.workspace_id, workspace_id=params.workspace_id,
@@ -598,7 +601,7 @@ def register_methods(
async def _create_from_capability( async def _create_from_capability(
server: WorkflowServer, server: WorkflowServer,
params: CreateDraftFromCapabilityParams, params: CreateDraftFromCapabilityParams,
) -> dict[str, Any]: ) -> CreateDraftWorkspaceFromCapabilityResult:
return await server.api.create_draft_workspace_from_capability( return await server.api.create_draft_workspace_from_capability(
workspace_id=params.workspace_id, workspace_id=params.workspace_id,
capability_name=params.capability_name, capability_name=params.capability_name,
@@ -202,6 +202,95 @@ def test_openrpc_pins_nested_draft_workspace_contract(
assert deleted["properties"]["status"]["enum"] == ["deleted", "not_found"] assert deleted["properties"]["status"]["enum"] == ["deleted", "not_found"]
def test_openrpc_exposes_typed_compile_draft_workspace_result(
openrpc_document: dict[str, Any],
) -> None:
method = _method_by_name(
openrpc_document,
"workflow.draft_workspaces.compile",
)
assert method["result"]["schema"] == {
"$ref": "#/components/schemas/CompileDraftWorkspaceResult"
}
assert openrpc_document["components"]["schemas"]["CompileDraftWorkspaceResult"] == {
"anyOf": [
{"$ref": ("#/components/schemas/CompileDraftWorkspaceSuccess")},
{"$ref": "#/components/schemas/InvalidDraftResult"},
]
}
schemas = openrpc_document["components"]["schemas"]
assert schemas["CompileDraftWorkspaceSuccess"]["required"] == [
"compiled_plan",
"required_capabilities",
]
assert schemas["InvalidDraftResult"]["properties"]["status"]["const"] == "invalid"
def test_openrpc_exposes_typed_capability_bootstrap_result(
openrpc_document: dict[str, Any],
) -> None:
for method_name in (
"workflow.drafts.create_from_capability",
"workflow.draft_workspaces.create_from_capability",
):
_assert_result_component(
openrpc_document,
method_name=method_name,
component_name="CreateDraftWorkspaceFromCapabilityResult",
properties={
"workspace_id",
"revision",
"wrapper_hints",
"next_actions",
},
)
result = openrpc_document["components"]["schemas"][
"CreateDraftWorkspaceFromCapabilityResult"
]
assert result["properties"]["wrapper_hints"] == {
"$ref": "#/components/schemas/WrapperAuthoringHintsPayload"
}
assert result["properties"]["next_actions"] == {
"$ref": "#/components/schemas/NextActionsPayload"
}
@pytest.mark.parametrize(
"method_name",
[
"workflow.draft_workspaces.create_artifact",
"workflow.draft_workspaces.create_wrapper",
],
)
def test_openrpc_exposes_typed_workspace_artifact_save_result(
openrpc_document: dict[str, Any],
method_name: str,
) -> None:
method = _method_by_name(openrpc_document, method_name)
assert method["result"]["schema"] == {
"$ref": "#/components/schemas/CreateArtifactFromWorkspaceResult"
}
assert openrpc_document["components"]["schemas"][
"CreateArtifactFromWorkspaceResult"
] == {
"anyOf": [
{"$ref": "#/components/schemas/SavedDraftArtifactResult"},
{"$ref": "#/components/schemas/UnsavedDraftArtifactResult"},
]
}
schemas = openrpc_document["components"]["schemas"]
assert schemas["SavedDraftArtifactResult"]["properties"]["saved"]["const"] is True
assert {
"required_logical_sources",
"suggested_bindings",
} <= schemas["SavedDraftArtifactResult"]["properties"].keys()
assert (
schemas["UnsavedDraftArtifactResult"]["properties"]["saved"]["const"] is False
)
@pytest.mark.parametrize( @pytest.mark.parametrize(
("method_name", "component_name", "properties"), ("method_name", "component_name", "properties"),
[ [