refactor(semantics): close dashboard testing protocol audit
This commit is contained in:
@@ -1,5 +1,5 @@
|
||||
#!/usr/bin/env python3
|
||||
#region FullFlow.Environment [C:3] [TYPE Module] [SEMANTICS docker,secrets,isolated-fixture]
|
||||
# #region FullFlow.Environment [C:3] [TYPE Module] [SEMANTICS docker,secrets,isolated-fixture]
|
||||
# @PURPOSE Generate isolated fixture credentials without reading production settings.
|
||||
# @PRE Target does not exist; parent directory is writable.
|
||||
# @POST A mode-0600 env file contains fresh secrets; stdout contains its path only.
|
||||
@@ -23,4 +23,4 @@ values = {
|
||||
with os.fdopen(os.open(path, os.O_WRONLY | os.O_CREAT | os.O_EXCL, 0o600), 'w') as file:
|
||||
file.write(''.join(f'{key}={value}\n' for key, value in values.items()))
|
||||
print(path)
|
||||
#endregion FullFlow.Environment
|
||||
# #endregion FullFlow.Environment
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
#!/usr/bin/env python3
|
||||
#region FullFlow.Readiness [C:4] [TYPE Module] [SEMANTICS docker,readiness,authentication,evidence]
|
||||
# #region FullFlow.Readiness [C:4] [TYPE Module] [SEMANTICS docker,readiness,authentication,evidence]
|
||||
# @PURPOSE Probe real API readiness and password authentication without printing tokens.
|
||||
# @PRE The isolated stack is running; credentials come only from the named fixture env file.
|
||||
# @POST A nonzero exit identifies unavailable or unauthenticated services; successful probes list IDs.
|
||||
@@ -37,4 +37,4 @@ except Exception as error:
|
||||
failures.append('ss-tools')
|
||||
print(f'ss-tools: probe failed ({type(error).__name__}: {getattr(error, "reason", getattr(error, "code", "unexpected response"))})')
|
||||
sys.exit(bool(failures))
|
||||
#endregion FullFlow.Readiness
|
||||
# #endregion FullFlow.Readiness
|
||||
|
||||
@@ -2,6 +2,7 @@
|
||||
# @BRIEF Exercise fresh MCP profiles using only server-issued questions and CAS identities.
|
||||
# @INVARIANT All results originate from the real Docker backend API.
|
||||
import json
|
||||
from functools import partial
|
||||
from common import Blocked
|
||||
|
||||
|
||||
@@ -12,46 +13,12 @@ class McpProfileMixin:
|
||||
# @BRIEF Resolve only questions issued by a fresh external MCP profile.
|
||||
def profile(self):
|
||||
headers = {"Accept": "application/json, text/event-stream", "Content-Type": "application/json"}
|
||||
# #region Tooling.FullFlowMcpProfileMixin.rpc [C:3] [TYPE Function]
|
||||
# @BRIEF Validate the actual external JSON-RPC and MCP tool response.
|
||||
def rpc(name, payload):
|
||||
response = self.client.post("/mcp/", json=payload, headers=headers)
|
||||
if response.headers.get("mcp-session-id"):
|
||||
headers["Mcp-Session-Id"] = response.headers["mcp-session-id"]
|
||||
if not response.content:
|
||||
value = {}
|
||||
elif "text/event-stream" in response.headers.get("content-type", ""):
|
||||
data = [line[6:] for line in response.text.splitlines() if line.startswith("data: ")]
|
||||
value = json.loads(data[-1]) if data else {}
|
||||
else:
|
||||
value = response.json()
|
||||
tool_error = value.get("result", {}).get("isError") is True
|
||||
self.record(name, "PASS" if response.status_code in (200,202) and "error" not in value and not tool_error else "BLOCKED",
|
||||
http_status=response.status_code, response=value)
|
||||
if response.status_code not in (200,202) or "error" in value or tool_error:
|
||||
raise Blocked(name)
|
||||
return value
|
||||
|
||||
# #endregion Tooling.FullFlowMcpProfileMixin.rpc
|
||||
rpc = partial(self._profile_rpc, headers)
|
||||
rpc("mcp-initialize", {"jsonrpc":"2.0", "id":1, "method":"initialize", "params":{
|
||||
"protocolVersion":"2025-03-26", "capabilities":{}, "clientInfo":{"name":"docker-full-flow-verifier","version":"1"}}})
|
||||
headers["MCP-Protocol-Version"] = "2025-03-26"
|
||||
rpc("mcp-initialized", {"jsonrpc":"2.0", "method":"notifications/initialized"})
|
||||
# #region Tooling.FullFlowMcpProfileMixin.unwrap [C:3] [TYPE Function]
|
||||
# @BRIEF Decode structured MCP domain results without hiding tool failures.
|
||||
def unwrap(response):
|
||||
result = response.get("result", {})
|
||||
structured = result.get("structuredContent")
|
||||
if not structured:
|
||||
content = next((c["text"] for c in result.get("content",[]) if c.get("type") == "text"), "{}")
|
||||
try:
|
||||
structured = json.loads(content)
|
||||
except json.JSONDecodeError:
|
||||
self.record("mcp-tool-domain-response", "BLOCKED", error_text=content)
|
||||
raise Blocked("MCP tool returned unstructured error text") from None
|
||||
return structured
|
||||
|
||||
# #endregion Tooling.FullFlowMcpProfileMixin.unwrap
|
||||
unwrap = self._profile_unwrap
|
||||
intent = {"environment_id":"full-flow-preprod", "dashboard_id":self.config.get("dashboard_id",1),
|
||||
"objective":"Verify filtered sales revenue against approved baseline", "selected_case_ids":["C05"]}
|
||||
structured = unwrap(rpc("mcp-propose-real-profile", {"jsonrpc":"2.0", "id":2, "method":"tools/call", "params":{
|
||||
@@ -90,6 +57,44 @@ class McpProfileMixin:
|
||||
|
||||
# #endregion Tooling.FullFlowMcpProfileMixin.profile
|
||||
|
||||
# #region Tooling.FullFlowMcpProfileMixin.rpc [C:3] [TYPE Function]
|
||||
# @BRIEF Validate the actual external JSON-RPC and MCP tool response.
|
||||
def _profile_rpc(self, headers, name, payload):
|
||||
response = self.client.post("/mcp/", json=payload, headers=headers)
|
||||
if response.headers.get("mcp-session-id"):
|
||||
headers["Mcp-Session-Id"] = response.headers["mcp-session-id"]
|
||||
if not response.content:
|
||||
value = {}
|
||||
elif "text/event-stream" in response.headers.get("content-type", ""):
|
||||
data = [line[6:] for line in response.text.splitlines() if line.startswith("data: ")]
|
||||
value = json.loads(data[-1]) if data else {}
|
||||
else:
|
||||
value = response.json()
|
||||
tool_error = value.get("result", {}).get("isError") is True
|
||||
self.record(name, "PASS" if response.status_code in (200,202) and "error" not in value and not tool_error else "BLOCKED",
|
||||
http_status=response.status_code, response=value)
|
||||
if response.status_code not in (200,202) or "error" in value or tool_error:
|
||||
raise Blocked(name)
|
||||
return value
|
||||
|
||||
# #endregion Tooling.FullFlowMcpProfileMixin.rpc
|
||||
|
||||
# #region Tooling.FullFlowMcpProfileMixin.unwrap [C:3] [TYPE Function]
|
||||
# @BRIEF Decode structured MCP domain results without hiding tool failures.
|
||||
def _profile_unwrap(self, response):
|
||||
result = response.get("result", {})
|
||||
structured = result.get("structuredContent")
|
||||
if not structured:
|
||||
content = next((c["text"] for c in result.get("content",[]) if c.get("type") == "text"), "{}")
|
||||
try:
|
||||
structured = json.loads(content)
|
||||
except json.JSONDecodeError:
|
||||
self.record("mcp-tool-domain-response", "BLOCKED", error_text=content)
|
||||
raise Blocked("MCP tool returned unstructured error text") from None
|
||||
return structured
|
||||
|
||||
# #endregion Tooling.FullFlowMcpProfileMixin.unwrap
|
||||
|
||||
# #endregion Tooling.FullFlowMcpProfileMixin.McpProfileMixin
|
||||
|
||||
|
||||
|
||||
@@ -1,7 +1,6 @@
|
||||
# Real-time isolated Docker soak
|
||||
|
||||
## @{ Tooling.Stage6Soak.Protocol [C:4] [TYPE ADR] [SEMANTICS soak,protocol,durable,isolated]
|
||||
|
||||
@BRIEF Executable opt-in protocol for the persistent M01 collector and independent auditor.
|
||||
@RELATION DEPENDS_ON -> [Tooling.Stage6Soak.Collector]
|
||||
@RELATION DEPENDS_ON -> [Tooling.Stage6Soak.Controller]
|
||||
|
||||
@@ -51,16 +51,7 @@ def assess(manifest, journal, runs, byte_findings, *, evaluated_at):
|
||||
requested = {r["payload"]["index"]: timestamp(r["utc"]) for r in journal if r["event"] == "backend_restart_requested"}
|
||||
restored = {r["payload"]["index"]: timestamp(r["utc"]) for r in journal if r["event"] == "backend_restart_completed"}
|
||||
down_windows = [(begin, restored[index]) for index, begin in requested.items() if index in restored]
|
||||
if started:
|
||||
start, end = timestamp(started), timestamp(evaluated_at)
|
||||
due = start.replace(minute=start.minute - start.minute % 5, second=0, microsecond=0) + timedelta(minutes=5)
|
||||
expected = set()
|
||||
while due <= end:
|
||||
if not any(begin <= due <= finish for begin, finish in down_windows):
|
||||
expected.add(due.isoformat())
|
||||
due += timedelta(minutes=5)
|
||||
if expected - set(slots):
|
||||
issues.append("MISSING_LOGICAL_DUE_SLOTS")
|
||||
_missing_due_slots(started, evaluated_at, down_windows, slots, issues)
|
||||
observations = [timestamp(row["utc"]) for row in journal if row["event"] == "observation"]
|
||||
window = [timestamp(started)] + observations + [timestamp(evaluated_at)] if started else []
|
||||
gap = max(((b - a).total_seconds() for a, b in zip(window, window[1:])), default=0)
|
||||
@@ -72,15 +63,7 @@ def assess(manifest, journal, runs, byte_findings, *, evaluated_at):
|
||||
if any(row["event"] == "observation" and (not row["payload"].get("queue_slo_met")
|
||||
or not row["payload"].get("schedule_enabled")) for row in journal):
|
||||
issues.append("QUEUE_OR_SCHEDULE_SLO_VIOLATION")
|
||||
restarts = [row for row in journal if row["event"] == "backend_restart_completed"]
|
||||
if {row["payload"].get("index") for row in restarts} != {1, 2, 3} or len(restarts) != 3:
|
||||
issues.append("THREE_RESTARTS_NOT_PROVEN")
|
||||
elif any(not row["payload"].get("before", {}).get("started_at")
|
||||
or not row["payload"].get("after", {}).get("started_at")
|
||||
or row["payload"]["before"]["started_at"] == row["payload"]["after"]["started_at"]
|
||||
or any(row["payload"][side].get("project") != "ss-tools-full-flow"
|
||||
or row["payload"][side].get("service") != "backend" for side in ("before", "after")) for row in restarts):
|
||||
issues.append("REAL_RESTART_IDENTITY_NOT_PROVEN")
|
||||
_restart_issues(journal, issues)
|
||||
tabs = [row["payload"] for row in journal if row["event"] == "tab_canary"]
|
||||
if not all(any(row.get("tabs") == n and row.get("status") == "passed" for row in tabs) for n in (5, 15, 50)):
|
||||
issues.append("TAB_CANARIES_NOT_PROVEN")
|
||||
@@ -134,18 +117,7 @@ def retained_run(cursor, row, manifest, storage):
|
||||
raw = (Path(storage) / "drafts" / row["id"] / sha).read_bytes()
|
||||
if hashlib.sha256(raw).hexdigest() != sha:
|
||||
return issues + ["ACTUAL_RAW_HASH_MISMATCH"]
|
||||
response = json.loads(raw)["result"]
|
||||
metric = details["metric_coordinate"]["metric_name"]
|
||||
if len(response) != 1 or len(response[0]["data"]) != 1 or response[0]["colnames"].count(metric) != 1:
|
||||
return issues + ["ACTUAL_SCALAR_AMBIGUOUS"]
|
||||
value = response[0]["data"][0][metric]
|
||||
position = response[0]["colnames"].index(metric)
|
||||
column_type = response[0].get("coltypes", [])[position]
|
||||
if type(column_type) is not int or column_type != 0:
|
||||
issues.append("ACTUAL_NUMERIC_SCHEMA_NOT_PROVEN")
|
||||
if type(value) not in {int, float} or not Decimal(str(value)).is_finite() or Decimal(str(value)) != Decimal("16350"):
|
||||
issues.append("ACTUAL_ORACLE_MISMATCH")
|
||||
return issues
|
||||
return _scalar_issues(raw, details, issues)
|
||||
# #endregion Tooling.Stage6Soak.Audit.RetainedRun
|
||||
|
||||
|
||||
@@ -193,4 +165,53 @@ def execute(root, env, storage):
|
||||
atomic(root / "independent-audit.json", result)
|
||||
return result
|
||||
# #endregion Tooling.Stage6Soak.Audit.Execute
|
||||
|
||||
# #region Tooling.Stage6Soak.Audit.MissingSlots [C:3] [TYPE Function]
|
||||
# @BRIEF Preserve the extracted operation phase and its ordering.
|
||||
def _missing_due_slots(started, evaluated_at, down_windows, slots, issues):
|
||||
if started:
|
||||
start, end = timestamp(started), timestamp(evaluated_at)
|
||||
due = start.replace(minute=start.minute - start.minute % 5, second=0, microsecond=0) + timedelta(minutes=5)
|
||||
expected = set()
|
||||
while due <= end:
|
||||
if not any(begin <= due <= finish for begin, finish in down_windows):
|
||||
expected.add(due.isoformat())
|
||||
due += timedelta(minutes=5)
|
||||
if expected - set(slots):
|
||||
issues.append("MISSING_LOGICAL_DUE_SLOTS")
|
||||
# #endregion Tooling.Stage6Soak.Audit.MissingSlots
|
||||
|
||||
|
||||
# #region Tooling.Stage6Soak.Audit.RestartIssues [C:3] [TYPE Function]
|
||||
# @BRIEF Preserve the extracted operation phase and its ordering.
|
||||
def _restart_issues(journal, issues):
|
||||
restarts = [row for row in journal if row["event"] == "backend_restart_completed"]
|
||||
if {row["payload"].get("index") for row in restarts} != {1, 2, 3} or len(restarts) != 3:
|
||||
issues.append("THREE_RESTARTS_NOT_PROVEN")
|
||||
elif any(not row["payload"].get("before", {}).get("started_at")
|
||||
or not row["payload"].get("after", {}).get("started_at")
|
||||
or row["payload"]["before"]["started_at"] == row["payload"]["after"]["started_at"]
|
||||
or any(row["payload"][side].get("project") != "ss-tools-full-flow"
|
||||
or row["payload"][side].get("service") != "backend" for side in ("before", "after")) for row in restarts):
|
||||
issues.append("REAL_RESTART_IDENTITY_NOT_PROVEN")
|
||||
# #endregion Tooling.Stage6Soak.Audit.RestartIssues
|
||||
|
||||
|
||||
# #region Tooling.Stage6Soak.Audit.ScalarIssues [C:3] [TYPE Function]
|
||||
# @BRIEF Preserve the extracted operation phase and its ordering.
|
||||
def _scalar_issues(raw, details, issues):
|
||||
response = json.loads(raw)["result"]
|
||||
metric = details["metric_coordinate"]["metric_name"]
|
||||
if len(response) != 1 or len(response[0]["data"]) != 1 or response[0]["colnames"].count(metric) != 1:
|
||||
return issues + ["ACTUAL_SCALAR_AMBIGUOUS"]
|
||||
value = response[0]["data"][0][metric]
|
||||
position = response[0]["colnames"].index(metric)
|
||||
column_type = response[0].get("coltypes", [])[position]
|
||||
if type(column_type) is not int or column_type != 0:
|
||||
issues.append("ACTUAL_NUMERIC_SCHEMA_NOT_PROVEN")
|
||||
if type(value) not in {int, float} or not Decimal(str(value)).is_finite() or Decimal(str(value)) != Decimal("16350"):
|
||||
issues.append("ACTUAL_ORACLE_MISMATCH")
|
||||
return issues
|
||||
# #endregion Tooling.Stage6Soak.Audit.ScalarIssues
|
||||
|
||||
# #endregion Tooling.Stage6Soak.Audit
|
||||
|
||||
@@ -128,15 +128,7 @@ def observe(root, manifest, env, client):
|
||||
# @POST Signal interruption leaves the owned schedule paused and journal incomplete; normal completion requires real 72h.
|
||||
def collect(args, env, client):
|
||||
root, manifest = args.root, manifest_at(args.root)
|
||||
if manifest.get("stopped_at"):
|
||||
raise ValueError("SOAK_ALREADY_STOPPED")
|
||||
if args.command == "start":
|
||||
if manifest.get("started_at"):
|
||||
raise ValueError("SOAK_ALREADY_STARTED_USE_RESUME")
|
||||
manifest["started_at"] = now()
|
||||
atomic(root / "manifest.json", manifest)
|
||||
elif not manifest.get("started_at"):
|
||||
raise ValueError("SOAK_NOT_STARTED")
|
||||
_begin_collection(args, manifest, root)
|
||||
switch(client, manifest, True)
|
||||
append(root, "collector_boot", {"pid": __import__("os").getpid(), "monotonic": time.monotonic()})
|
||||
stopping, finished = False, False
|
||||
@@ -174,16 +166,7 @@ def collect(args, env, client):
|
||||
manifest["stopped_at"] = now()
|
||||
atomic(root / "manifest.json", manifest)
|
||||
append(root, "collector_stopped", {"elapsed_seconds": (timestamp(manifest["stopped_at"]) - timestamp(manifest["started_at"])).total_seconds()})
|
||||
deadline = time.monotonic() + manifest["max_queue_age_seconds"]
|
||||
while time.monotonic() < deadline:
|
||||
try:
|
||||
sample = observe(root, manifest, env, client)
|
||||
append(root, "drain_observation", sample)
|
||||
if sample["queued_or_running"] == 0:
|
||||
break
|
||||
except Exception as error:
|
||||
append(root, "drain_error", {"error_type": type(error).__name__})
|
||||
time.sleep(manifest["poll_seconds"])
|
||||
_drain_collection(manifest, root, env, client)
|
||||
return execute(root, env, args.storage)
|
||||
# #endregion Tooling.Stage6Soak.Collector.Run
|
||||
|
||||
@@ -242,4 +225,35 @@ if __name__ == "__main__":
|
||||
except Exception as failure:
|
||||
print(json.dumps({"status": "BLOCKED", "error_type": type(failure).__name__}))
|
||||
sys.exit(1)
|
||||
|
||||
# #region Tooling.Stage6Soak.Collector.Begin [C:3] [TYPE Function]
|
||||
# @BRIEF Preserve the extracted operation phase and its ordering.
|
||||
def _begin_collection(args, manifest, root):
|
||||
if manifest.get("stopped_at"):
|
||||
raise ValueError("SOAK_ALREADY_STOPPED")
|
||||
if args.command == "start":
|
||||
if manifest.get("started_at"):
|
||||
raise ValueError("SOAK_ALREADY_STARTED_USE_RESUME")
|
||||
manifest["started_at"] = now()
|
||||
atomic(root / "manifest.json", manifest)
|
||||
elif not manifest.get("started_at"):
|
||||
raise ValueError("SOAK_NOT_STARTED")
|
||||
# #endregion Tooling.Stage6Soak.Collector.Begin
|
||||
|
||||
|
||||
# #region Tooling.Stage6Soak.Collector.Drain [C:3] [TYPE Function]
|
||||
# @BRIEF Preserve the extracted operation phase and its ordering.
|
||||
def _drain_collection(manifest, root, env, client):
|
||||
deadline = time.monotonic() + manifest["max_queue_age_seconds"]
|
||||
while time.monotonic() < deadline:
|
||||
try:
|
||||
sample = observe(root, manifest, env, client)
|
||||
append(root, "drain_observation", sample)
|
||||
if sample["queued_or_running"] == 0:
|
||||
break
|
||||
except Exception as error:
|
||||
append(root, "drain_error", {"error_type": type(error).__name__})
|
||||
time.sleep(manifest["poll_seconds"])
|
||||
# #endregion Tooling.Stage6Soak.Collector.Drain
|
||||
|
||||
# #endregion Tooling.Stage6Soak.Collector
|
||||
|
||||
Reference in New Issue
Block a user