feat: expose read-only admin surface
This commit is contained in:
@@ -106,6 +106,21 @@ implementation state.
|
||||
broker-owned.
|
||||
- Completed: read-only source inventory is exposed through JSON-RPC HTTP and
|
||||
`wf source list` / `wf source inspect`.
|
||||
- Completed: read-only admin/config sibling surface covers connection
|
||||
inventory, connection status, and broker/server events. Keep it separate
|
||||
from `WorkflowApiSurface`; this is platform management, not workflow
|
||||
lifecycle.
|
||||
- Completed: read-only admin/config now has a neutral `WorkflowAdminApi` /
|
||||
`WorkflowAdminSurface`, is exposed through JSON-RPC HTTP, and is available
|
||||
through `wf admin connections`, `wf admin statuses`, and
|
||||
`wf admin events`.
|
||||
- Defer mutating source/connection config commands until the store-backed
|
||||
source registry is designed. Config can bootstrap sources, but server-owned
|
||||
dynamic source changes need persistence, validation, and auth rules before
|
||||
they are safe.
|
||||
- Longer term: make the MCP frontend an adapter over these neutral workflow,
|
||||
source-admin, and config-admin surfaces so the old `wf_mcp` server entry
|
||||
point can shrink or retire.
|
||||
|
||||
5. **CLI/API alignment**
|
||||
- Completed for the basic lifecycle: selected `wf` commands can target local
|
||||
|
||||
@@ -83,6 +83,8 @@ surface, or plain local CLI utilities.
|
||||
1. **Store-backed source registry**
|
||||
- Read-only source/admin operations are now available through JSON-RPC HTTP
|
||||
and `wf source list` / `wf source inspect`.
|
||||
- Read-only admin/config operations are now available through JSON-RPC HTTP
|
||||
and `wf admin connections`, `wf admin statuses`, and `wf admin events`.
|
||||
- Next source work is persistence for server-owned dynamic source changes.
|
||||
- Keep mutation out until the store-backed source registry is designed.
|
||||
|
||||
@@ -101,8 +103,9 @@ surface, or plain local CLI utilities.
|
||||
|
||||
## Open Questions
|
||||
|
||||
- Should source/admin operations live in `wf_api` as a sibling surface, or in a
|
||||
new package that composes `wf_api` plus server management services?
|
||||
- Source/admin and admin/config read-only operations currently live in `wf_api`
|
||||
as sibling surfaces. If mutation grows into a larger management domain, split
|
||||
that later instead of overloading `WorkflowApiSurface`.
|
||||
- Should `wf schema` describe local CLI command payloads only, or query a remote
|
||||
server for supported method schemas?
|
||||
- Should `wf docs` read packaged local docs, remote server docs, or both?
|
||||
|
||||
@@ -1,6 +1,11 @@
|
||||
from __future__ import annotations
|
||||
|
||||
from .listing import matches_query, paged_list_payload
|
||||
from .admin import (
|
||||
WorkflowAdminApi,
|
||||
WorkflowAdminConnectionProvider,
|
||||
WorkflowAdminEventProvider,
|
||||
)
|
||||
from .artifacts import WorkflowArtifactApi
|
||||
from .local_sources import builtin_sources, get_qualified_spec, qualify_spec
|
||||
from .models import RawWorkflowPlan, TraceRange
|
||||
@@ -20,6 +25,7 @@ from .runs import WorkflowRunApi
|
||||
from .service import WorkflowApi
|
||||
from .source_admin import WorkflowSourceAdminApi
|
||||
from .surface import (
|
||||
WorkflowAdminSurface,
|
||||
WorkflowApiSurface,
|
||||
WorkflowArtifactSurface,
|
||||
WorkflowCapabilitySurface,
|
||||
@@ -73,6 +79,10 @@ __all__ = [
|
||||
"RawWorkflowPlan",
|
||||
"RuntimeDependencies",
|
||||
"TraceRange",
|
||||
"WorkflowAdminApi",
|
||||
"WorkflowAdminConnectionProvider",
|
||||
"WorkflowAdminEventProvider",
|
||||
"WorkflowAdminSurface",
|
||||
"WorkflowApi",
|
||||
"WorkflowApiSurface",
|
||||
"WorkflowArtifactApi",
|
||||
|
||||
@@ -0,0 +1,68 @@
|
||||
from __future__ import annotations
|
||||
|
||||
from collections.abc import Mapping, Sequence
|
||||
from dataclasses import asdict, is_dataclass
|
||||
from typing import Any, Protocol
|
||||
|
||||
|
||||
class WorkflowAdminConnectionProvider(Protocol):
|
||||
"""Provides read-only connection inventory for admin frontends."""
|
||||
|
||||
def list_connections(self) -> Sequence[Mapping[str, Any] | object]: ...
|
||||
|
||||
def get_connection_statuses(self) -> Sequence[Mapping[str, Any] | object]: ...
|
||||
|
||||
|
||||
class WorkflowAdminEventProvider(Protocol):
|
||||
"""Provides read-only event history for admin frontends."""
|
||||
|
||||
def list_events(self) -> Sequence[Mapping[str, Any] | object]: ...
|
||||
|
||||
|
||||
class WorkflowAdminApi:
|
||||
"""Protocol-neutral read-only broker/server admin operations.
|
||||
|
||||
This surface is intentionally not part of WorkflowApiSurface. Connections
|
||||
and event history are platform management data, not workflow lifecycle data.
|
||||
"""
|
||||
|
||||
def __init__(
|
||||
self,
|
||||
*,
|
||||
connections: WorkflowAdminConnectionProvider,
|
||||
events: WorkflowAdminEventProvider,
|
||||
) -> None:
|
||||
self.connections = connections
|
||||
self.events = events
|
||||
|
||||
async def list_connections(self) -> dict[str, Any]:
|
||||
connections = sorted(
|
||||
(_payload(item) for item in self.connections.list_connections()),
|
||||
key=lambda item: str(item.get("id", item.get("connection_id", ""))),
|
||||
)
|
||||
return {"connections": connections, "total": len(connections)}
|
||||
|
||||
async def get_connection_statuses(self) -> dict[str, Any]:
|
||||
statuses = sorted(
|
||||
(_payload(item) for item in self.connections.get_connection_statuses()),
|
||||
key=lambda item: str(item.get("connection_id", item.get("id", ""))),
|
||||
)
|
||||
return {"statuses": statuses, "total": len(statuses)}
|
||||
|
||||
async def list_events(self) -> dict[str, Any]:
|
||||
events = [_payload(event) for event in self.events.list_events()]
|
||||
return {"events": events, "total": len(events)}
|
||||
|
||||
|
||||
def _payload(value: Mapping[str, Any] | object) -> dict[str, Any]:
|
||||
"""Normalize provider objects without depending on MCP event/config types."""
|
||||
if isinstance(value, Mapping):
|
||||
return dict(value)
|
||||
if is_dataclass(value) and not isinstance(value, type):
|
||||
return asdict(value)
|
||||
model_dump = getattr(value, "model_dump", None)
|
||||
if callable(model_dump):
|
||||
result = model_dump(mode="json")
|
||||
if isinstance(result, dict):
|
||||
return result
|
||||
raise TypeError(f"admin payload object is not serializable: {type(value)!r}")
|
||||
@@ -216,7 +216,18 @@ class WorkflowSourceAdminSurface(Protocol):
|
||||
) -> dict[str, Any]: ...
|
||||
|
||||
|
||||
class WorkflowAdminSurface(Protocol):
|
||||
"""Read-only connection/config admin methods exposed by platform frontends."""
|
||||
|
||||
async def list_connections(self) -> dict[str, Any]: ...
|
||||
|
||||
async def get_connection_statuses(self) -> dict[str, Any]: ...
|
||||
|
||||
async def list_events(self) -> dict[str, Any]: ...
|
||||
|
||||
|
||||
__all__ = [
|
||||
"WorkflowAdminSurface",
|
||||
"WorkflowApiSurface",
|
||||
"WorkflowArtifactSurface",
|
||||
"WorkflowCapabilitySurface",
|
||||
|
||||
@@ -5,6 +5,7 @@ from typing import Annotated
|
||||
import typer
|
||||
|
||||
from .commands import (
|
||||
admin,
|
||||
artifacts,
|
||||
caps,
|
||||
deployments,
|
||||
@@ -62,6 +63,7 @@ app.add_typer(artifacts.app, name="artifact")
|
||||
app.add_typer(deployments.app, name="deploy")
|
||||
app.add_typer(runs.app, name="run")
|
||||
app.add_typer(sources.app, name="source")
|
||||
app.add_typer(admin.app, name="admin")
|
||||
app.add_typer(docs.app, name="docs")
|
||||
app.add_typer(schema.app, name="schema")
|
||||
app.command("explain")(explain.explain_command)
|
||||
|
||||
@@ -0,0 +1,72 @@
|
||||
from __future__ import annotations
|
||||
|
||||
import asyncio
|
||||
from typing import Annotated
|
||||
|
||||
import typer
|
||||
|
||||
from wf_cli.context import load_cli_context_from_typer
|
||||
from wf_cli.formats import ListOutputFormat, emit_list_payload
|
||||
|
||||
app = typer.Typer(
|
||||
name="admin",
|
||||
help="Read workflow server admin and config state.",
|
||||
no_args_is_help=True,
|
||||
)
|
||||
|
||||
|
||||
@app.command("connections")
|
||||
def list_connections(
|
||||
ctx: typer.Context,
|
||||
output_format: Annotated[
|
||||
ListOutputFormat, typer.Option("--format", help="Output format.")
|
||||
] = ListOutputFormat.JSON,
|
||||
) -> None:
|
||||
"""List configured upstream connections known to the target."""
|
||||
context = load_cli_context_from_typer(ctx)
|
||||
payload = asyncio.run(context.admin.list_connections())
|
||||
emit_list_payload(
|
||||
payload,
|
||||
collection_key="connections",
|
||||
output_format=output_format,
|
||||
id_field="id",
|
||||
summary_fields=("server", "account", "enabled"),
|
||||
)
|
||||
|
||||
|
||||
@app.command("statuses")
|
||||
def get_connection_statuses(
|
||||
ctx: typer.Context,
|
||||
output_format: Annotated[
|
||||
ListOutputFormat, typer.Option("--format", help="Output format.")
|
||||
] = ListOutputFormat.JSON,
|
||||
) -> None:
|
||||
"""List connection catalog/status summaries."""
|
||||
context = load_cli_context_from_typer(ctx)
|
||||
payload = asyncio.run(context.admin.get_connection_statuses())
|
||||
emit_list_payload(
|
||||
payload,
|
||||
collection_key="statuses",
|
||||
output_format=output_format,
|
||||
id_field="connection_id",
|
||||
summary_fields=("server", "account", "enabled", "has_snapshot"),
|
||||
)
|
||||
|
||||
|
||||
@app.command("events")
|
||||
def list_events(
|
||||
ctx: typer.Context,
|
||||
output_format: Annotated[
|
||||
ListOutputFormat, typer.Option("--format", help="Output format.")
|
||||
] = ListOutputFormat.JSON,
|
||||
) -> None:
|
||||
"""List recorded workflow platform events."""
|
||||
context = load_cli_context_from_typer(ctx)
|
||||
payload = asyncio.run(context.admin.list_events())
|
||||
emit_list_payload(
|
||||
payload,
|
||||
collection_key="events",
|
||||
output_format=output_format,
|
||||
id_field="kind",
|
||||
summary_fields=("connection_id", "capability_id", "workflow_name"),
|
||||
)
|
||||
@@ -10,6 +10,8 @@ from pydantic import ValidationError
|
||||
|
||||
from wf_api import (
|
||||
WorkflowApi,
|
||||
WorkflowAdminApi,
|
||||
WorkflowAdminSurface,
|
||||
WorkflowApiSurface,
|
||||
WorkflowSourceAdminApi,
|
||||
WorkflowSourceAdminSurface,
|
||||
@@ -35,6 +37,7 @@ class CliContext:
|
||||
service: WfMcpService | None
|
||||
handlers: WorkflowApiSurface
|
||||
source_admin: WorkflowSourceAdminSurface
|
||||
admin: WorkflowAdminSurface
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
@@ -116,6 +119,7 @@ def load_cli_context(
|
||||
service=None,
|
||||
handlers=client,
|
||||
source_admin=client,
|
||||
admin=client,
|
||||
)
|
||||
|
||||
if _is_legacy_mcp_config(resolved_config_path):
|
||||
@@ -126,6 +130,10 @@ def load_cli_context(
|
||||
service=service,
|
||||
handlers=WorkflowApi(context_from_service(service)),
|
||||
source_admin=WorkflowSourceAdminApi(context_from_service(service)),
|
||||
admin=WorkflowAdminApi(
|
||||
connections=service.connection_service,
|
||||
events=service.events,
|
||||
),
|
||||
)
|
||||
|
||||
config = load_workflow_config(resolved_config_path)
|
||||
@@ -140,6 +148,7 @@ def load_cli_context(
|
||||
service=None,
|
||||
handlers=server.api,
|
||||
source_admin=server.source_admin,
|
||||
admin=server.admin,
|
||||
)
|
||||
if isinstance(target, RpcHttpTargetConfig):
|
||||
client = RpcWorkflowApiClient(
|
||||
@@ -155,6 +164,7 @@ def load_cli_context(
|
||||
service=None,
|
||||
handlers=client,
|
||||
source_admin=client,
|
||||
admin=client,
|
||||
)
|
||||
raise ValueError(f"unsupported workflow target {target!r}")
|
||||
|
||||
|
||||
@@ -1,7 +1,11 @@
|
||||
from dataclasses import asdict
|
||||
from typing import Any
|
||||
|
||||
from wf_api import WorkflowSourceAdminApi, WorkflowSourceAdminSurface
|
||||
from wf_api import (
|
||||
WorkflowAdminApi,
|
||||
WorkflowAdminSurface,
|
||||
WorkflowSourceAdminApi,
|
||||
WorkflowSourceAdminSurface,
|
||||
)
|
||||
from wf_mcp.broker.service import WfMcpService
|
||||
from wf_mcp.broker.service.workflow_operation_context import context_from_service
|
||||
from wf_mcp.shared.errors import error_payload
|
||||
@@ -15,18 +19,18 @@ class BrokerAdminHandlers:
|
||||
self.sources: WorkflowSourceAdminSurface = WorkflowSourceAdminApi(
|
||||
context_from_service(service)
|
||||
)
|
||||
self.admin: WorkflowAdminSurface = WorkflowAdminApi(
|
||||
connections=service.connection_service,
|
||||
events=service.events,
|
||||
)
|
||||
|
||||
def list_connections(self) -> list[dict[str, Any]]:
|
||||
return [
|
||||
asdict(connection)
|
||||
for connection in sorted(
|
||||
self.service.connections.list_all(),
|
||||
key=lambda connection: connection.id,
|
||||
)
|
||||
]
|
||||
async def list_connections(self) -> list[dict[str, Any]]:
|
||||
payload = await self.admin.list_connections()
|
||||
return payload["connections"]
|
||||
|
||||
def get_connection_statuses(self) -> list[dict[str, Any]]:
|
||||
return self.service.connection_statuses()
|
||||
async def get_connection_statuses(self) -> list[dict[str, Any]]:
|
||||
payload = await self.admin.get_connection_statuses()
|
||||
return payload["statuses"]
|
||||
|
||||
async def refresh_connection_catalog(self, connection_id: str) -> dict[str, Any]:
|
||||
try:
|
||||
@@ -93,5 +97,6 @@ class BrokerAdminHandlers:
|
||||
**error_payload(exc),
|
||||
}
|
||||
|
||||
def get_broker_events(self) -> list[dict[str, Any]]:
|
||||
return [asdict(event) for event in self.service.list_events()]
|
||||
async def get_broker_events(self) -> list[dict[str, Any]]:
|
||||
payload = await self.admin.list_events()
|
||||
return payload["events"]
|
||||
|
||||
@@ -46,7 +46,7 @@ def register_service_admin_tools(
|
||||
description="List configured MCP connections known to this server.",
|
||||
)
|
||||
async def list_connections() -> list[dict[str, Any]]:
|
||||
return handlers.list_connections()
|
||||
return await handlers.list_connections()
|
||||
|
||||
@server.tool(
|
||||
name=name("get_connection_statuses"),
|
||||
@@ -54,7 +54,7 @@ def register_service_admin_tools(
|
||||
description="Show configured MCP connection status and catalog counts.",
|
||||
)
|
||||
async def get_connection_statuses() -> list[dict[str, Any]]:
|
||||
return handlers.get_connection_statuses()
|
||||
return await handlers.get_connection_statuses()
|
||||
|
||||
@server.tool(
|
||||
name=name("refresh_connection_catalog"),
|
||||
@@ -177,4 +177,4 @@ def register_service_admin_tools(
|
||||
description="Return locally recorded broker/platform events.",
|
||||
)
|
||||
async def get_events() -> list[dict[str, Any]]:
|
||||
return handlers.get_broker_events()
|
||||
return await handlers.get_broker_events()
|
||||
|
||||
@@ -32,9 +32,17 @@ class ConnectionService:
|
||||
def list_all(self) -> list[ConnectionConfig]:
|
||||
return self.connections.list_all()
|
||||
|
||||
def list_connections(self) -> list[ConnectionConfig]:
|
||||
"""Return connection inventory for protocol-neutral admin surfaces."""
|
||||
return self.list_all()
|
||||
|
||||
def list_enabled(self) -> list[ConnectionConfig]:
|
||||
return self.connections.list_enabled()
|
||||
|
||||
def get_connection_statuses(self) -> list[dict[str, object]]:
|
||||
"""Return catalog-backed connection statuses for admin surfaces."""
|
||||
return self._source_catalog().connection_statuses()
|
||||
|
||||
def register_connection(self, connection: ConnectionConfig) -> None:
|
||||
self._validate_connection_id(connection.id)
|
||||
self.connections.register(connection)
|
||||
|
||||
@@ -5,7 +5,12 @@ from dataclasses import dataclass, field
|
||||
from pathlib import Path
|
||||
from typing import Any
|
||||
|
||||
from wf_api import WorkflowApi, WorkflowSourceAdminApi, durable_workflow_api
|
||||
from wf_api import (
|
||||
WorkflowAdminApi,
|
||||
WorkflowApi,
|
||||
WorkflowSourceAdminApi,
|
||||
durable_workflow_api,
|
||||
)
|
||||
from wf_api.local_sources import builtin_sources, get_qualified_spec
|
||||
from wf_api.models import RawWorkflowPlan, TraceRange
|
||||
from wf_api.operation_context import (
|
||||
@@ -64,6 +69,21 @@ class InMemoryWorkflowEventRecorder(WorkflowEventRecorder):
|
||||
}
|
||||
)
|
||||
|
||||
def list_events(self) -> list[dict[str, Any]]:
|
||||
"""Expose local server events for the read-only admin surface."""
|
||||
return list(self.events)
|
||||
|
||||
|
||||
@dataclass(frozen=True, slots=True)
|
||||
class EmptyWorkflowConnectionProvider:
|
||||
"""Read-only admin provider for local/static servers without upstream sources."""
|
||||
|
||||
def list_connections(self) -> list[dict[str, Any]]:
|
||||
return []
|
||||
|
||||
def get_connection_statuses(self) -> list[dict[str, Any]]:
|
||||
return []
|
||||
|
||||
|
||||
@dataclass(frozen=True, slots=True)
|
||||
class StaticWorkflowSpecProvider(WorkflowSpecProvider):
|
||||
@@ -238,6 +258,7 @@ class WorkflowServer:
|
||||
context: WorkflowOperationContext
|
||||
api: WorkflowApi
|
||||
source_admin: WorkflowSourceAdminApi
|
||||
admin: WorkflowAdminApi
|
||||
events: InMemoryWorkflowEventRecorder
|
||||
|
||||
@staticmethod
|
||||
@@ -266,11 +287,16 @@ def build_local_static_workflow_server(root: str | Path) -> WorkflowServer:
|
||||
)
|
||||
api = durable_workflow_api(context)
|
||||
source_admin = WorkflowSourceAdminApi(context)
|
||||
admin = WorkflowAdminApi(
|
||||
connections=EmptyWorkflowConnectionProvider(),
|
||||
events=events,
|
||||
)
|
||||
return WorkflowServer(
|
||||
config=config,
|
||||
stores=stores,
|
||||
context=context,
|
||||
api=api,
|
||||
source_admin=source_admin,
|
||||
admin=admin,
|
||||
events=events,
|
||||
)
|
||||
|
||||
@@ -4,6 +4,7 @@ from .app import create_rpc_app
|
||||
from .client import RpcWorkflowApiClient
|
||||
from .errors import WorkflowRpcError
|
||||
from .models import (
|
||||
AdminEmptyParams,
|
||||
CreateArtifactFromWorkspaceParams,
|
||||
CreateDraftFromCapabilityParams,
|
||||
CreateWrapperFromWorkspaceParams,
|
||||
@@ -35,6 +36,7 @@ from .models import (
|
||||
|
||||
__all__ = [
|
||||
"CreateArtifactFromWorkspaceParams",
|
||||
"AdminEmptyParams",
|
||||
"CreateDraftFromCapabilityParams",
|
||||
"CreateWrapperFromWorkspaceParams",
|
||||
"DeleteDeploymentParams",
|
||||
|
||||
@@ -7,6 +7,7 @@ import fastapi_jsonrpc as jsonrpc
|
||||
from wf_server import WorkflowServer
|
||||
|
||||
from .errors import WorkflowRpcError
|
||||
from .methods_admin import register_methods as register_admin_methods
|
||||
from .methods_artifacts import register_methods as register_artifact_methods
|
||||
from .methods_capabilities import register_methods as register_capability_methods
|
||||
from .methods_deployments import register_methods as register_deployment_methods
|
||||
@@ -45,6 +46,7 @@ def create_rpc_app(server: WorkflowServer, *, rpc_path: str = "/rpc") -> jsonrpc
|
||||
register_deployment_methods(entrypoint, server)
|
||||
register_run_methods(entrypoint, server)
|
||||
register_source_methods(entrypoint, server)
|
||||
register_admin_methods(entrypoint, server)
|
||||
|
||||
app.bind_entrypoint(entrypoint)
|
||||
return app
|
||||
|
||||
@@ -5,6 +5,7 @@ from dataclasses import dataclass
|
||||
import httpx # noqa: F401 # Backcompat for tests patching client.httpx.AsyncClient.
|
||||
|
||||
from .client_artifacts import RpcArtifactClientMixin
|
||||
from .client_admin import RpcAdminClientMixin
|
||||
from .client_base import RpcClientTransport
|
||||
from .client_capabilities import RpcCapabilityClientMixin
|
||||
from .client_deployments import RpcDeploymentClientMixin
|
||||
@@ -22,6 +23,7 @@ class RpcWorkflowApiClient(
|
||||
RpcDeploymentClientMixin,
|
||||
RpcRunClientMixin,
|
||||
RpcSourceAdminClientMixin,
|
||||
RpcAdminClientMixin,
|
||||
):
|
||||
"""WorkflowApiSurface implementation backed by JSON-RPC HTTP calls.
|
||||
|
||||
|
||||
@@ -0,0 +1,18 @@
|
||||
from __future__ import annotations
|
||||
|
||||
from typing import Any
|
||||
|
||||
|
||||
class RpcAdminClientMixin:
|
||||
"""JSON-RPC implementation of read-only admin/config surface methods."""
|
||||
|
||||
async def _call(self, method: str, params: dict[str, Any]) -> dict[str, Any]: ...
|
||||
|
||||
async def list_connections(self) -> dict[str, Any]:
|
||||
return await self._call("workflow.admin.connections.list", {})
|
||||
|
||||
async def get_connection_statuses(self) -> dict[str, Any]:
|
||||
return await self._call("workflow.admin.connection_statuses.list", {})
|
||||
|
||||
async def list_events(self) -> dict[str, Any]:
|
||||
return await self._call("workflow.admin.events.list", {})
|
||||
@@ -0,0 +1,51 @@
|
||||
from __future__ import annotations
|
||||
|
||||
from typing import Any
|
||||
|
||||
from fastapi import Body
|
||||
import fastapi_jsonrpc as jsonrpc
|
||||
|
||||
from wf_server import WorkflowServer
|
||||
|
||||
from .errors import WorkflowRpcError, raise_workflow_rpc_error
|
||||
from .models import AdminEmptyParams
|
||||
|
||||
|
||||
def register_methods(
|
||||
entrypoint: jsonrpc.Entrypoint,
|
||||
server: WorkflowServer,
|
||||
) -> None:
|
||||
"""Register read-only admin/config JSON-RPC methods."""
|
||||
|
||||
@entrypoint.method(
|
||||
name="workflow.admin.connections.list",
|
||||
errors=[WorkflowRpcError],
|
||||
)
|
||||
async def workflow_admin_connections_list(
|
||||
params: AdminEmptyParams = Body(default_factory=AdminEmptyParams),
|
||||
) -> dict[str, Any]:
|
||||
try:
|
||||
return await server.admin.list_connections()
|
||||
except (ValueError, KeyError, LookupError, FileNotFoundError) as exc:
|
||||
raise_workflow_rpc_error(exc)
|
||||
|
||||
@entrypoint.method(
|
||||
name="workflow.admin.connection_statuses.list",
|
||||
errors=[WorkflowRpcError],
|
||||
)
|
||||
async def workflow_admin_connection_statuses_list(
|
||||
params: AdminEmptyParams = Body(default_factory=AdminEmptyParams),
|
||||
) -> dict[str, Any]:
|
||||
try:
|
||||
return await server.admin.get_connection_statuses()
|
||||
except (ValueError, KeyError, LookupError, FileNotFoundError) as exc:
|
||||
raise_workflow_rpc_error(exc)
|
||||
|
||||
@entrypoint.method(name="workflow.admin.events.list", errors=[WorkflowRpcError])
|
||||
async def workflow_admin_events_list(
|
||||
params: AdminEmptyParams = Body(default_factory=AdminEmptyParams),
|
||||
) -> dict[str, Any]:
|
||||
try:
|
||||
return await server.admin.list_events()
|
||||
except (ValueError, KeyError, LookupError, FileNotFoundError) as exc:
|
||||
raise_workflow_rpc_error(exc)
|
||||
@@ -30,6 +30,10 @@ class HealthParams(RpcParamsModel):
|
||||
pass
|
||||
|
||||
|
||||
class AdminEmptyParams(RpcParamsModel):
|
||||
pass
|
||||
|
||||
|
||||
class ListCapabilitiesParams(RpcParamsModel):
|
||||
query: str | None = Field(default=None)
|
||||
source_id: str | None = Field(default=None)
|
||||
|
||||
@@ -0,0 +1,96 @@
|
||||
from __future__ import annotations
|
||||
|
||||
import asyncio
|
||||
from dataclasses import dataclass, field
|
||||
from typing import Any
|
||||
|
||||
from wf_api import WorkflowAdminApi, WorkflowAdminSurface
|
||||
|
||||
|
||||
@dataclass(frozen=True, slots=True)
|
||||
class FakeConnection:
|
||||
id: str
|
||||
server: str
|
||||
account: str
|
||||
enabled: bool = True
|
||||
metadata: dict[str, Any] = field(default_factory=dict)
|
||||
|
||||
|
||||
@dataclass(frozen=True, slots=True)
|
||||
class FakeEvent:
|
||||
kind: str
|
||||
timestamp_epoch_ms: int
|
||||
connection_id: str | None = None
|
||||
payload: dict[str, Any] = field(default_factory=dict)
|
||||
|
||||
|
||||
class FakeAdminProvider:
|
||||
def __init__(self) -> None:
|
||||
self.connections = [
|
||||
FakeConnection(id="zeta.personal", server="zeta", account="personal"),
|
||||
FakeConnection(id="alpha.work", server="alpha", account="work"),
|
||||
]
|
||||
self.statuses = [
|
||||
{"connection_id": "zeta.personal", "enabled": True},
|
||||
{"connection_id": "alpha.work", "enabled": False},
|
||||
]
|
||||
self.events = [
|
||||
FakeEvent(
|
||||
kind="connection_registered",
|
||||
timestamp_epoch_ms=123,
|
||||
connection_id="alpha.work",
|
||||
)
|
||||
]
|
||||
|
||||
def list_connections(self) -> list[FakeConnection]:
|
||||
return self.connections
|
||||
|
||||
def get_connection_statuses(self) -> list[dict[str, Any]]:
|
||||
return self.statuses
|
||||
|
||||
def list_events(self) -> list[FakeEvent]:
|
||||
return self.events
|
||||
|
||||
|
||||
def test_admin_api_lists_connections_in_id_order() -> None:
|
||||
provider = FakeAdminProvider()
|
||||
api = WorkflowAdminApi(connections=provider, events=provider)
|
||||
|
||||
payload = asyncio.run(api.list_connections())
|
||||
|
||||
assert payload["total"] == 2
|
||||
assert [connection["id"] for connection in payload["connections"]] == [
|
||||
"alpha.work",
|
||||
"zeta.personal",
|
||||
]
|
||||
|
||||
|
||||
def test_admin_api_lists_connection_statuses_in_id_order() -> None:
|
||||
provider = FakeAdminProvider()
|
||||
api = WorkflowAdminApi(connections=provider, events=provider)
|
||||
|
||||
payload = asyncio.run(api.get_connection_statuses())
|
||||
|
||||
assert payload["total"] == 2
|
||||
assert [status["connection_id"] for status in payload["statuses"]] == [
|
||||
"alpha.work",
|
||||
"zeta.personal",
|
||||
]
|
||||
|
||||
|
||||
def test_admin_api_lists_events() -> None:
|
||||
provider = FakeAdminProvider()
|
||||
api = WorkflowAdminApi(connections=provider, events=provider)
|
||||
|
||||
payload = asyncio.run(api.list_events())
|
||||
|
||||
assert payload["total"] == 1
|
||||
assert payload["events"][0]["kind"] == "connection_registered"
|
||||
assert payload["events"][0]["connection_id"] == "alpha.work"
|
||||
|
||||
|
||||
def test_admin_api_satisfies_surface_protocol() -> None:
|
||||
provider = FakeAdminProvider()
|
||||
api: WorkflowAdminSurface = WorkflowAdminApi(connections=provider, events=provider)
|
||||
|
||||
assert api is not None
|
||||
@@ -18,6 +18,7 @@ def test_wf_help_lists_lifecycle_groups() -> None:
|
||||
assert "deploy" in result.output
|
||||
assert "run" in result.output
|
||||
assert "source" in result.output
|
||||
assert "admin" in result.output
|
||||
assert "docs" in result.output
|
||||
assert "schema" in result.output
|
||||
assert "explain" in result.output
|
||||
@@ -76,6 +77,13 @@ def test_wf_source_list_help_exists() -> None:
|
||||
assert "--limit" in result.output
|
||||
|
||||
|
||||
def test_wf_admin_connections_help_exists() -> None:
|
||||
result = runner.invoke(app, ["admin", "connections", "--help"])
|
||||
|
||||
assert result.exit_code == 0
|
||||
assert "--format" in result.output
|
||||
|
||||
|
||||
def test_wf_artifact_list_help_exists() -> None:
|
||||
result = runner.invoke(app, ["artifact", "list", "--help"])
|
||||
|
||||
|
||||
@@ -310,6 +310,38 @@ def test_wf_source_commands_use_rpc_url_override(monkeypatch, tmp_path) -> None:
|
||||
assert '"id": "wf.std"' in inspected.output
|
||||
|
||||
|
||||
def test_wf_admin_commands_use_rpc_url_override(monkeypatch, tmp_path) -> None:
|
||||
server = build_local_static_workflow_server(tmp_path / "store")
|
||||
server.events.record_workflow_event(
|
||||
"workflow_test_event",
|
||||
capability_id="workflow.demo.v1",
|
||||
payload={"ok": True},
|
||||
)
|
||||
original_client = httpx.AsyncClient
|
||||
monkeypatch.setattr(
|
||||
"wf_transport_rpc_http.client.httpx.AsyncClient",
|
||||
lambda *args, **kwargs: original_client(
|
||||
transport=httpx.ASGITransport(app=create_rpc_app(server)),
|
||||
base_url="http://test",
|
||||
),
|
||||
)
|
||||
config_path = tmp_path / "wf.json"
|
||||
config_path.write_text('{"version": 1}', encoding="utf-8")
|
||||
runner = CliRunner()
|
||||
base_args = ["--config", str(config_path), "--url", "http://test/rpc"]
|
||||
|
||||
connections = runner.invoke(app, [*base_args, "admin", "connections"])
|
||||
statuses = runner.invoke(app, [*base_args, "admin", "statuses"])
|
||||
events = runner.invoke(app, [*base_args, "admin", "events"])
|
||||
|
||||
assert connections.exit_code == 0, connections.output
|
||||
assert '"connections": []' in connections.output
|
||||
assert statuses.exit_code == 0, statuses.output
|
||||
assert '"statuses": []' in statuses.output
|
||||
assert events.exit_code == 0, events.output
|
||||
assert '"kind": "workflow_test_event"' in events.output
|
||||
|
||||
|
||||
def test_wf_remote_draft_artifact_deploy_lifecycle(monkeypatch, tmp_path) -> None:
|
||||
server = build_local_static_workflow_server(tmp_path / "store")
|
||||
original_client = httpx.AsyncClient
|
||||
|
||||
@@ -18,8 +18,8 @@ def test_broker_admin_handlers_list_connections_and_events() -> None:
|
||||
)
|
||||
handlers = BrokerAdminHandlers(service)
|
||||
|
||||
connections = handlers.list_connections()
|
||||
events = handlers.get_broker_events()
|
||||
connections = _run(handlers.list_connections())
|
||||
events = _run(handlers.get_broker_events())
|
||||
|
||||
assert connections[0]["id"] == "demo.personal"
|
||||
assert connections[0]["server"] == "demo"
|
||||
|
||||
@@ -107,6 +107,37 @@ def test_rpc_workflow_client_lists_and_inspects_sources(tmp_path) -> None:
|
||||
asyncio.run(scenario())
|
||||
|
||||
|
||||
def test_rpc_workflow_client_reads_admin_state(tmp_path) -> None:
|
||||
async def scenario() -> None:
|
||||
server = build_local_static_workflow_server(tmp_path / "store")
|
||||
server.events.record_workflow_event(
|
||||
"workflow_test_event",
|
||||
capability_id="workflow.demo.v1",
|
||||
payload={"ok": True},
|
||||
)
|
||||
app = create_rpc_app(server)
|
||||
transport = httpx.ASGITransport(app=app)
|
||||
async with httpx.AsyncClient(
|
||||
transport=transport,
|
||||
base_url="http://test",
|
||||
) as http_client:
|
||||
client = RpcWorkflowApiClient(
|
||||
url="http://test/rpc",
|
||||
timeout_seconds=5,
|
||||
http_client=http_client,
|
||||
)
|
||||
connections = await client.list_connections()
|
||||
statuses = await client.get_connection_statuses()
|
||||
events = await client.list_events()
|
||||
|
||||
assert connections == {"connections": [], "total": 0}
|
||||
assert statuses == {"statuses": [], "total": 0}
|
||||
assert events["total"] == 1
|
||||
assert events["events"][0]["kind"] == "workflow_test_event"
|
||||
|
||||
asyncio.run(scenario())
|
||||
|
||||
|
||||
def test_rpc_workflow_client_runs_and_reads_trace(tmp_path) -> None:
|
||||
async def scenario() -> None:
|
||||
server = build_local_static_workflow_server(tmp_path / "store")
|
||||
|
||||
Reference in New Issue
Block a user