feat: persist owner-scoped profile resolutions
This commit is contained in:
53
backend/alembic/versions/0028_durable_test_pack_profiles.py
Normal file
53
backend/alembic/versions/0028_durable_test_pack_profiles.py
Normal file
@@ -0,0 +1,53 @@
|
||||
# #region Migrations.DurableTestPackProfiles [C:2] [TYPE Module] [SEMANTICS migration,profile,cas,idempotency]
|
||||
# @BRIEF Persist owner-scoped typed profile sessions and mutation receipts.
|
||||
"""Persist durable test-pack profile resolution state."""
|
||||
|
||||
from alembic import op
|
||||
import sqlalchemy as sa
|
||||
|
||||
revision = "0028_durable_test_pack_profiles"
|
||||
down_revision = "0027_schedule_baseline"
|
||||
branch_labels = None
|
||||
depends_on = None
|
||||
|
||||
|
||||
def upgrade() -> None:
|
||||
op.create_table(
|
||||
"scenario_test_pack_profiles",
|
||||
sa.Column("profile_handle_id", sa.String(36), primary_key=True),
|
||||
sa.Column("owner_principal", sa.String(128), nullable=False),
|
||||
sa.Column("environment_id", sa.String(128), nullable=False),
|
||||
sa.Column("dashboard_id", sa.Integer(), nullable=False),
|
||||
sa.Column("objective", sa.String(2000), nullable=False),
|
||||
sa.Column("selected_case_ids", sa.JSON(), nullable=False),
|
||||
sa.Column("profile_digest", sa.String(64), nullable=False),
|
||||
sa.Column("context_fingerprint", sa.String(64), nullable=False, server_default=""),
|
||||
sa.Column("resolutions", sa.JSON(), nullable=False),
|
||||
sa.Column("profile_snapshot", sa.JSON(), nullable=False),
|
||||
sa.Column("cas_version", sa.Integer(), nullable=False, server_default="0"),
|
||||
sa.Column("created_at", sa.DateTime(), nullable=False),
|
||||
sa.Column("updated_at", sa.DateTime(), nullable=False),
|
||||
)
|
||||
op.create_index("ix_scenario_test_pack_profiles_owner_principal", "scenario_test_pack_profiles", ["owner_principal"])
|
||||
op.create_index("ix_scenario_test_pack_profiles_dashboard_id", "scenario_test_pack_profiles", ["dashboard_id"])
|
||||
op.create_index("ix_scenario_test_pack_profiles_profile_digest", "scenario_test_pack_profiles", ["profile_digest"])
|
||||
op.create_table(
|
||||
"scenario_test_pack_profile_receipts",
|
||||
sa.Column("receipt_id", sa.String(36), primary_key=True),
|
||||
sa.Column("profile_handle_id", sa.String(36), nullable=False),
|
||||
sa.Column("owner_principal", sa.String(128), nullable=False),
|
||||
sa.Column("idempotency_key", sa.String(255), nullable=False),
|
||||
sa.Column("request_hash", sa.String(64), nullable=False),
|
||||
sa.Column("response", sa.JSON(), nullable=False),
|
||||
sa.Column("created_at", sa.DateTime(), nullable=False),
|
||||
sa.UniqueConstraint("profile_handle_id", "owner_principal", "idempotency_key", name="uq_profile_receipt_key"),
|
||||
)
|
||||
|
||||
|
||||
def downgrade() -> None:
|
||||
op.drop_table("scenario_test_pack_profile_receipts")
|
||||
op.drop_index("ix_scenario_test_pack_profiles_profile_digest", table_name="scenario_test_pack_profiles")
|
||||
op.drop_index("ix_scenario_test_pack_profiles_dashboard_id", table_name="scenario_test_pack_profiles")
|
||||
op.drop_index("ix_scenario_test_pack_profiles_owner_principal", table_name="scenario_test_pack_profiles")
|
||||
op.drop_table("scenario_test_pack_profiles")
|
||||
# #endregion Migrations.DurableTestPackProfiles
|
||||
@@ -386,6 +386,9 @@ class ResolveTestPackProfileInput(BaseModel):
|
||||
objective: StrictStr = Field(min_length=1, max_length=2000)
|
||||
selected_case_ids: list[StrictStr] = Field(min_length=1, max_length=19)
|
||||
expected_profile_digest: StrictStr = Field(pattern=r"^[a-f0-9]{64}$")
|
||||
profile_handle_id: StrictStr | None = Field(default=None, min_length=1, max_length=36)
|
||||
idempotency_key: StrictStr | None = Field(default=None, min_length=1, max_length=255)
|
||||
expected_cas_version: StrictInt = Field(default=0, ge=0)
|
||||
resolutions: list[TestPackProfileResolution] = Field(min_length=1, max_length=100)
|
||||
|
||||
@field_validator("selected_case_ids")
|
||||
|
||||
@@ -18,7 +18,10 @@
|
||||
from __future__ import annotations
|
||||
|
||||
import re
|
||||
import hashlib
|
||||
import json
|
||||
from typing import Any
|
||||
from uuid import uuid4
|
||||
|
||||
from mcp.server.fastmcp import Context
|
||||
|
||||
@@ -34,6 +37,7 @@ from src.mcp_server.scenario_inputs import (
|
||||
InspectContextInput,
|
||||
RegisterDraftPackInput,
|
||||
ResolveTestPackProfileInput,
|
||||
TestPackProfileResolution,
|
||||
ScenarioCompileInput,
|
||||
ScenarioResolveInput,
|
||||
ScenarioStartInput,
|
||||
@@ -41,6 +45,7 @@ from src.mcp_server.scenario_inputs import (
|
||||
)
|
||||
from src.models.maintenance import MaintenanceEvent
|
||||
from src.models.agent_run import AgentRun
|
||||
from src.models.scenario_handles import TestPackProfileReceipt, TestPackProfileSession
|
||||
from src.core.utils.client_registry import get_superset_client
|
||||
from src.services.dashboard_testing.query_model import inspect_dashboard_query_model
|
||||
from src.services.dashboard_testing.scenario.compiler import CompileScenarioRequest, compile_scenario
|
||||
@@ -73,6 +78,28 @@ from src.services.llm_provider import LLMProviderService
|
||||
from src.services.mcp_approvals import decide_mcp_approval, list_pending_mcp_approvals
|
||||
|
||||
|
||||
# #region McpServer.TestPackProfile.PersistPreview [C:3] [TYPE Function] [SEMANTICS mcp,profile,durable,owner]
|
||||
# @ingroup McpServer.TestPackProfile.ResolutionHelpers
|
||||
# @BRIEF Store one server-derived profile snapshot under the authenticated principal.
|
||||
def _persist_profile_preview(profile, request, _scenario) -> str:
|
||||
access = _access_token_context.get()
|
||||
if access is None or not access.subject:
|
||||
raise ValueError("PROFILE_OWNER_REQUIRED")
|
||||
profile_handle_id = str(uuid4())
|
||||
with SessionLocal() as db:
|
||||
db.add(TestPackProfileSession(
|
||||
profile_handle_id=profile_handle_id, owner_principal=str(access.subject),
|
||||
environment_id=request.environment_id, dashboard_id=request.dashboard_id,
|
||||
objective=request.objective, selected_case_ids=list(request.selected_case_ids),
|
||||
profile_digest=profile.profile_digest,
|
||||
context_fingerprint=profile.query_model_fingerprint,
|
||||
resolutions=[], profile_snapshot=profile.model_dump(mode="json"), cas_version=0,
|
||||
))
|
||||
db.commit()
|
||||
return profile_handle_id
|
||||
# #endregion McpServer.TestPackProfile.PersistPreview
|
||||
|
||||
|
||||
# #region McpServer.TestPackProfile.ResolutionHelpers [C:3] [TYPE Module] [SEMANTICS mcp,profile,resolution,validation]
|
||||
# @ingroup McpServer
|
||||
# @BRIEF Validate typed profile resolutions against current unresolved profile items and build server compiler inputs.
|
||||
@@ -391,6 +418,8 @@ def register_scenario_tools(server) -> None:
|
||||
return {
|
||||
"status": profile.status,
|
||||
"profile": profile.model_dump(mode="json"),
|
||||
"profile_handle_id": _persist_profile_preview(profile, request, scenario),
|
||||
"cas_version": 0,
|
||||
"preview": {
|
||||
"step_count": len(scenario.steps),
|
||||
"artifacts": pack.get("artifacts", []),
|
||||
@@ -403,14 +432,42 @@ def register_scenario_tools(server) -> None:
|
||||
|
||||
# #region McpServer.ScenarioTools.ResolveTestPackProfile [C:4] [TYPE Function] [SEMANTICS mcp,profile,resolve,cas]
|
||||
# @ingroup McpServer
|
||||
# @BRIEF Re-inspect context, CAS-check the profile digest, and apply only reviewed selector hints.
|
||||
# @PRE Every resolution ID names a current needs_selector blocker in the server-rebuilt profile.
|
||||
# @POST Returns a fresh deterministic preview; no handles or graph bytes are returned.
|
||||
# @SIDE_EFFECT Performs bounded Superset reads and compiler logging; it writes no database state.
|
||||
# @BRIEF Re-inspect context, CAS-check a durable owner profile, and apply reviewed typed answers.
|
||||
# @PRE Every resolution ID names a current unresolved item; the owner and CAS match.
|
||||
# @POST Accumulated answers and replay receipt commit atomically after fresh-context validation.
|
||||
# @SIDE_EFFECT Performs bounded Superset reads and one profile/receipt database transaction.
|
||||
# @INVARIANT Unsupported targets, stale profile digests and caller graph/value claims fail closed.
|
||||
# #region McpServer.ScenarioTools.ResolveTestPackProfile.Call [C:4] [TYPE Function] [SEMANTICS mcp,profile,resolve,inspection]
|
||||
@server.tool(name="resolve_test_pack_profile", structured_output=True)
|
||||
async def resolve_test_pack_profile(request: ResolveTestPackProfileInput) -> dict[str, Any]:
|
||||
access = _access_token_context.get()
|
||||
if access is None or not access.subject or not request.profile_handle_id or not request.idempotency_key:
|
||||
return {"status": "blocked", "error": "PROFILE_OWNER_OR_IDEMPOTENCY_REQUIRED"}
|
||||
owner = str(access.subject)
|
||||
request_body = request.model_dump(mode="json", exclude={"idempotency_key"})
|
||||
request_hash = hashlib.sha256(json.dumps(request_body, sort_keys=True, separators=(",", ":")).encode()).hexdigest()
|
||||
with SessionLocal() as db:
|
||||
session = db.query(TestPackProfileSession).filter(
|
||||
TestPackProfileSession.profile_handle_id == request.profile_handle_id,
|
||||
TestPackProfileSession.owner_principal == owner,
|
||||
).with_for_update().first()
|
||||
if session is None:
|
||||
return {"status": "blocked", "error": "PROFILE_ACCESS_DENIED"}
|
||||
receipt = db.query(TestPackProfileReceipt).filter_by(
|
||||
profile_handle_id=session.profile_handle_id, owner_principal=owner,
|
||||
idempotency_key=request.idempotency_key,
|
||||
).first()
|
||||
if receipt is not None:
|
||||
if receipt.request_hash != request_hash:
|
||||
return {"status": "conflict", "error": "IDEMPOTENCY_CONFLICT"}
|
||||
return {**receipt.response, "replayed": True}
|
||||
if (session.cas_version != request.expected_cas_version
|
||||
or session.environment_id != request.environment_id
|
||||
or session.dashboard_id != request.dashboard_id
|
||||
or session.objective != request.objective
|
||||
or session.selected_case_ids != request.selected_case_ids):
|
||||
return {"status": "conflict", "error": "PROFILE_CAS_CONFLICT",
|
||||
"current_profile_digest": session.profile_digest, "cas_version": session.cas_version}
|
||||
environment = get_config_manager().get_environment(request.environment_id)
|
||||
if environment is None:
|
||||
return {"status": "blocked", "error": "ENV_NOT_FOUND"}
|
||||
@@ -426,9 +483,13 @@ def register_scenario_tools(server) -> None:
|
||||
selected_case_ids=request.selected_case_ids,
|
||||
browser_available=resolve_browser_availability(),
|
||||
)
|
||||
if (baseline.query_model_fingerprint != session.context_fingerprint
|
||||
and session.cas_version == 0):
|
||||
return {"status": "conflict", "error": "PROFILE_STALE_CONTEXT",
|
||||
"current_profile_digest": baseline.profile_digest, "cas_version": session.cas_version}
|
||||
if baseline.profile_digest != request.expected_profile_digest:
|
||||
return {"status": "conflict", "error": "PROFILE_STALE",
|
||||
"current_profile_digest": baseline.profile_digest}
|
||||
return {"status": "conflict", "error": "PROFILE_STALE_CONTEXT",
|
||||
"current_profile_digest": baseline.profile_digest, "cas_version": session.cas_version}
|
||||
selector_resolutions = [item for item in request.resolutions if item.selector_hint is not None]
|
||||
coordinate_resolutions = [item for item in request.resolutions if item.coordinate_id is not None]
|
||||
has_selector_items = any(item.kind == "needs_selector" for item in baseline.unresolved)
|
||||
@@ -441,14 +502,29 @@ def register_scenario_tools(server) -> None:
|
||||
)
|
||||
if invalid_resolution or (selector_resolutions and not selector_changes) or (coordinate_resolutions and not coordinate_choices):
|
||||
return {"status": "blocked", "error": "PROFILE_RESOLUTION_INVALID"}
|
||||
parameters = _selector_parameters(selector_changes or [])
|
||||
accumulated = list(session.resolutions or [])
|
||||
accumulated.extend(item.model_dump(mode="json") for item in request.resolutions)
|
||||
all_selector_changes = _selector_profile_changes(baseline, [
|
||||
TestPackProfileResolution.model_validate(item) for item in accumulated
|
||||
if item.get("selector_hint") is not None
|
||||
], scenario)
|
||||
if accumulated and all_selector_changes is None:
|
||||
return {"status": "conflict", "error": "PROFILE_STALE_CONTEXT",
|
||||
"current_profile_digest": baseline.profile_digest, "cas_version": session.cas_version}
|
||||
parameters = _selector_parameters(all_selector_changes or [])
|
||||
profile, updated, pack = build_test_pack_profile(
|
||||
query_model=query_model, objective=request.objective,
|
||||
selected_case_ids=request.selected_case_ids, parameters=parameters,
|
||||
browser_available=resolve_browser_availability(),
|
||||
)
|
||||
if coordinate_choices:
|
||||
selected_coordinates = {item["unresolved_id"]: item["coordinate_id"] for item in coordinate_choices}
|
||||
all_coordinates = [TestPackProfileResolution.model_validate(item) for item in accumulated
|
||||
if item.get("coordinate_id") is not None]
|
||||
all_coordinate_choices = _coordinate_profile_choices(baseline, all_coordinates) if all_coordinates else []
|
||||
if all_coordinate_choices is None:
|
||||
return {"status": "conflict", "error": "PROFILE_STALE_CONTEXT",
|
||||
"current_profile_digest": baseline.profile_digest, "cas_version": session.cas_version}
|
||||
if all_coordinate_choices:
|
||||
selected_coordinates = {item["unresolved_id"]: item["coordinate_id"] for item in all_coordinate_choices}
|
||||
profile = apply_coordinate_choices(profile, selected_coordinates)
|
||||
except (KeyError, ValueError) as exc:
|
||||
logger.explore("Test-pack profile resolution rejected", src="McpServer.ScenarioTools.ResolveTestPackProfile",
|
||||
@@ -458,11 +534,31 @@ def register_scenario_tools(server) -> None:
|
||||
logger.explore("Test-pack profile resolution inspection failed", src="McpServer.ScenarioTools.ResolveTestPackProfile",
|
||||
error_code="INSPECTION_FAILED", payload={"dashboard_id": request.dashboard_id}, error=str(exc)[:300])
|
||||
return {"status": "blocked", "error": "INSPECTION_FAILED"}
|
||||
return {
|
||||
response = {
|
||||
"status": profile.status, "profile": profile.model_dump(mode="json"),
|
||||
"profile_handle_id": session.profile_handle_id,
|
||||
"cas_version": session.cas_version + 1,
|
||||
"preview": {"step_count": len(updated.steps), "artifacts": pack.get("artifacts", []),
|
||||
"validation": pack.get("validation_summary", {})},
|
||||
}
|
||||
with SessionLocal() as db:
|
||||
current = db.query(TestPackProfileSession).filter(
|
||||
TestPackProfileSession.profile_handle_id == session.profile_handle_id,
|
||||
TestPackProfileSession.owner_principal == owner,
|
||||
TestPackProfileSession.cas_version == request.expected_cas_version,
|
||||
).with_for_update().first()
|
||||
if current is None:
|
||||
return {"status": "conflict", "error": "PROFILE_CAS_CONFLICT"}
|
||||
current.resolutions = accumulated
|
||||
current.profile_digest = profile.profile_digest
|
||||
current.profile_snapshot = profile.model_dump(mode="json")
|
||||
current.cas_version += 1
|
||||
db.add(TestPackProfileReceipt(
|
||||
profile_handle_id=current.profile_handle_id, owner_principal=owner,
|
||||
idempotency_key=request.idempotency_key, request_hash=request_hash, response=response,
|
||||
))
|
||||
db.commit()
|
||||
return response
|
||||
# #endregion McpServer.ScenarioTools.ResolveTestPackProfile.Call
|
||||
# #endregion McpServer.ScenarioTools.ResolveTestPackProfile
|
||||
|
||||
|
||||
@@ -21,7 +21,7 @@ from __future__ import annotations
|
||||
from datetime import UTC, datetime
|
||||
import uuid
|
||||
|
||||
from sqlalchemy import JSON, Boolean, Column, DateTime, Integer, String
|
||||
from sqlalchemy import JSON, Boolean, Column, DateTime, Integer, String, UniqueConstraint
|
||||
|
||||
from .mapping import Base
|
||||
|
||||
@@ -92,4 +92,41 @@ class DraftPackHandle(Base):
|
||||
consumed_at = Column(DateTime, nullable=True)
|
||||
created_at = Column(DateTime, nullable=False, default=_now)
|
||||
# #endregion Models.ScenarioHandles.DraftPack
|
||||
|
||||
|
||||
# #region Models.ScenarioHandles.ProfileSession [C:3] [TYPE Class] [SEMANTICS scenario,profile,durable,cas]
|
||||
# @ingroup Models
|
||||
# @BRIEF Durable owner-scoped typed profile snapshot and accumulated resolutions.
|
||||
class TestPackProfileSession(Base):
|
||||
__tablename__ = "scenario_test_pack_profiles"
|
||||
profile_handle_id = Column(String(36), primary_key=True, default=_id)
|
||||
owner_principal = Column(String(128), nullable=False, index=True)
|
||||
environment_id = Column(String(128), nullable=False)
|
||||
dashboard_id = Column(Integer, nullable=False, index=True)
|
||||
objective = Column(String(2000), nullable=False)
|
||||
selected_case_ids = Column(JSON, nullable=False, default=list)
|
||||
profile_digest = Column(String(64), nullable=False, index=True)
|
||||
context_fingerprint = Column(String(64), nullable=False, default="")
|
||||
resolutions = Column(JSON, nullable=False, default=list)
|
||||
profile_snapshot = Column(JSON, nullable=False, default=dict)
|
||||
cas_version = Column(Integer, nullable=False, default=0)
|
||||
created_at = Column(DateTime, nullable=False, default=_now)
|
||||
updated_at = Column(DateTime, nullable=False, default=_now)
|
||||
|
||||
|
||||
# #region Models.ScenarioHandles.ProfileReceipt [C:2] [TYPE Class] [SEMANTICS scenario,profile,idempotency]
|
||||
# @ingroup Models
|
||||
# @BRIEF Append-only owner-scoped replay receipt for one profile mutation.
|
||||
class TestPackProfileReceipt(Base):
|
||||
__tablename__ = "scenario_test_pack_profile_receipts"
|
||||
receipt_id = Column(String(36), primary_key=True, default=_id)
|
||||
profile_handle_id = Column(String(36), nullable=False, index=True)
|
||||
owner_principal = Column(String(128), nullable=False, index=True)
|
||||
idempotency_key = Column(String(255), nullable=False)
|
||||
request_hash = Column(String(64), nullable=False)
|
||||
response = Column(JSON, nullable=False, default=dict)
|
||||
created_at = Column(DateTime, nullable=False, default=_now)
|
||||
__table_args__ = (UniqueConstraint("profile_handle_id", "owner_principal", "idempotency_key", name="uq_profile_receipt_key"),)
|
||||
# #endregion Models.ScenarioHandles.ProfileReceipt
|
||||
# #endregion Models.ScenarioHandles.ProfileSession
|
||||
# #endregion Models.ScenarioHandles
|
||||
|
||||
@@ -351,6 +351,9 @@ async def test_propose_test_pack_profile_returns_unresolved_preview_only(monkeyp
|
||||
baseline_item = next(item for item in profile["profile"]["unresolved"] if item["kind"] == "needs_metric")
|
||||
coordinate_request = {
|
||||
**request,
|
||||
"profile_handle_id": profile["profile_handle_id"],
|
||||
"idempotency_key": "profile-coordinate-1",
|
||||
"expected_cas_version": 0,
|
||||
"expected_profile_digest": profile["profile"]["profile_digest"],
|
||||
"resolutions": [{"unresolved_id": baseline_item["id"], "coordinate_id": coordinates[0]["coordinate_id"],
|
||||
"reason": "Analyst selected this metric."}],
|
||||
@@ -362,6 +365,8 @@ async def test_propose_test_pack_profile_returns_unresolved_preview_only(monkeyp
|
||||
assert updated_metric["coordinate_ids"] == [coordinates[0]["coordinate_id"]]
|
||||
assert updated_metric["kind"] == "needs_baseline"
|
||||
assert coordinate_result["profile"]["eligible"] is False
|
||||
replay = _unwrap(await server.call_tool("resolve_test_pack_profile", {"request": coordinate_request}))
|
||||
assert replay["replayed"] is True
|
||||
forged_coordinate = _unwrap(await server.call_tool("resolve_test_pack_profile", {
|
||||
"request": {**coordinate_request, "resolutions": [{**coordinate_request["resolutions"][0],
|
||||
"coordinate_id": "0" * 32}]}
|
||||
@@ -372,6 +377,9 @@ async def test_propose_test_pack_profile_returns_unresolved_preview_only(monkeyp
|
||||
selected_step = unresolved["step_id"]
|
||||
resolved_request = {
|
||||
**request,
|
||||
"profile_handle_id": profile["profile_handle_id"],
|
||||
"idempotency_key": "profile-selector-1",
|
||||
"expected_cas_version": 0,
|
||||
"expected_profile_digest": profile["profile"]["profile_digest"],
|
||||
"resolutions": [{"unresolved_id": unresolved_id, "step_id": selected_step,
|
||||
"selector_hint": "#sales-filter", "reason": "Analyst verified this control."}],
|
||||
@@ -383,10 +391,16 @@ async def test_propose_test_pack_profile_returns_unresolved_preview_only(monkeyp
|
||||
}
|
||||
assert "scenario" not in resolved and "draft_pack" not in resolved
|
||||
stale = _unwrap(await server.call_tool("resolve_test_pack_profile", {
|
||||
"request": {**coordinate_request, "expected_profile_digest": "0" * 64}
|
||||
"request": {**coordinate_request, "expected_profile_digest": "0" * 64,
|
||||
"idempotency_key": "profile-stale-1"}
|
||||
}))
|
||||
assert stale == {"status": "conflict", "error": "PROFILE_STALE",
|
||||
"current_profile_digest": profile["profile"]["profile_digest"]}
|
||||
assert stale["status"] == "conflict"
|
||||
assert stale["error"] in {"PROFILE_CAS_CONFLICT", "PROFILE_STALE_CONTEXT"}
|
||||
conflict = _unwrap(await server.call_tool("resolve_test_pack_profile", {
|
||||
"request": {**coordinate_request, "idempotency_key": "profile-coordinate-1",
|
||||
"resolutions": [{**coordinate_request["resolutions"][0], "reason": "changed"}]}
|
||||
}))
|
||||
assert conflict == {"status": "conflict", "error": "IDEMPOTENCY_CONFLICT"}
|
||||
forged = _unwrap(await server.call_tool("resolve_test_pack_profile", {
|
||||
"request": {**resolved_request, "resolutions": [{**resolved_request["resolutions"][0],
|
||||
"step_id": "phase-2-B02-apply_native_filter"}]}
|
||||
|
||||
Reference in New Issue
Block a user