Compare commits
3
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
7ecf7bf2d6 | ||
|
|
0589ec8069 | ||
|
|
300e8acd13 |
+326
-1
@@ -29,7 +29,9 @@ from dataclasses import dataclass
|
||||
from datetime import datetime, timedelta, timezone
|
||||
from typing import Any, Iterator, Sequence
|
||||
|
||||
SCHEMA_VERSION = 3
|
||||
import dependency_graph
|
||||
|
||||
SCHEMA_VERSION = 4
|
||||
|
||||
# Assignable work kinds only — raw monitoring incidents are never work items.
|
||||
WORK_KINDS = frozenset({"issue", "pr"})
|
||||
@@ -147,7 +149,41 @@ CREATE TABLE IF NOT EXISTS incident_links (
|
||||
UNIQUE (provider, provider_base_url, provider_org, provider_project, provider_issue_id)
|
||||
);
|
||||
|
||||
-- Durable dependency graph (#784, umbrella #628 scope item 6). Dependencies
|
||||
-- were previously re-parsed per allocation run and discarded; each row here is
|
||||
-- one relationship with its conditions, current state, and evidence. Creating
|
||||
-- the table is itself the v3→v4 migration: additive, idempotent, and it never
|
||||
-- touches the pre-existing tables.
|
||||
CREATE TABLE IF NOT EXISTS dependency_edges (
|
||||
edge_id TEXT PRIMARY KEY,
|
||||
remote TEXT NOT NULL,
|
||||
org TEXT NOT NULL,
|
||||
repo TEXT NOT NULL,
|
||||
source_kind TEXT NOT NULL CHECK (source_kind IN ('issue', 'pr')),
|
||||
source_number INTEGER NOT NULL,
|
||||
target_kind TEXT NOT NULL CHECK (target_kind IN ('issue', 'pr')),
|
||||
target_number INTEGER NOT NULL,
|
||||
edge_type TEXT NOT NULL,
|
||||
blocking_condition TEXT NOT NULL DEFAULT '',
|
||||
completion_condition TEXT NOT NULL DEFAULT '',
|
||||
state TEXT NOT NULL,
|
||||
evidence TEXT NOT NULL DEFAULT '{}',
|
||||
created_at TEXT NOT NULL,
|
||||
updated_at TEXT NOT NULL,
|
||||
last_observed_at TEXT NOT NULL,
|
||||
UNIQUE (
|
||||
remote, org, repo, source_kind, source_number,
|
||||
target_kind, target_number, edge_type
|
||||
)
|
||||
);
|
||||
|
||||
CREATE INDEX IF NOT EXISTS idx_leases_work_status ON leases(work_item_id, status);
|
||||
-- Reverse lookup ("what waits on this target") is the query automatic
|
||||
-- resumption needs, so it gets its own index alongside the forward one.
|
||||
CREATE INDEX IF NOT EXISTS idx_dependency_edges_source
|
||||
ON dependency_edges(remote, org, repo, source_kind, source_number);
|
||||
CREATE INDEX IF NOT EXISTS idx_dependency_edges_target
|
||||
ON dependency_edges(remote, org, repo, target_kind, target_number);
|
||||
CREATE INDEX IF NOT EXISTS idx_assignments_session ON assignments(session_id, status);
|
||||
CREATE INDEX IF NOT EXISTS idx_incident_gitea ON incident_links(gitea_org, gitea_repo, gitea_issue_number);
|
||||
"""
|
||||
@@ -1848,3 +1884,292 @@ class ControlPlaneDB:
|
||||
f"transferred lease ownership from {owner} to {adopter_session_id}"
|
||||
],
|
||||
}
|
||||
|
||||
# --- Dependency graph (#784, umbrella #628 scope item 6) ----------------
|
||||
|
||||
@staticmethod
|
||||
def _dependency_edge_row(row: sqlite3.Row | None) -> dict[str, Any] | None:
|
||||
"""Return a stored edge as a plain dict with evidence decoded."""
|
||||
if row is None:
|
||||
return None
|
||||
edge = dict(row)
|
||||
raw = edge.get("evidence")
|
||||
try:
|
||||
edge["evidence"] = json.loads(raw) if raw else {}
|
||||
except (TypeError, ValueError):
|
||||
# A row written by an older/foreign writer must not break reads.
|
||||
edge["evidence"] = {"unparsed": str(raw)}
|
||||
return edge
|
||||
|
||||
def upsert_dependency_edge(
|
||||
self,
|
||||
*,
|
||||
remote: str,
|
||||
org: str,
|
||||
repo: str,
|
||||
source_kind: str,
|
||||
source_number: int,
|
||||
target_kind: str,
|
||||
target_number: int,
|
||||
edge_type: str,
|
||||
state: str,
|
||||
blocking_condition: str | None = None,
|
||||
completion_condition: str | None = None,
|
||||
evidence: Any = None,
|
||||
) -> dict[str, Any]:
|
||||
"""Insert or refresh one dependency edge, keyed by its relationship.
|
||||
|
||||
Uniqueness is (scope, source, target, edge_type), so re-observing the
|
||||
same relationship updates one row instead of appending history — the
|
||||
edge is current state, and transitions are recorded as ``events``.
|
||||
|
||||
Edge type, state, and both endpoint kinds are validated fail-closed;
|
||||
an unknown value writes nothing. Evidence is sanitized before storage.
|
||||
"""
|
||||
edge_type_norm = dependency_graph.normalize_edge_type(edge_type)
|
||||
state_norm = dependency_graph.normalize_edge_state(state)
|
||||
source_kind_norm = dependency_graph.normalize_work_kind(source_kind)
|
||||
target_kind_norm = dependency_graph.normalize_work_kind(target_kind)
|
||||
source_no = int(source_number)
|
||||
target_no = int(target_number)
|
||||
if blocking_condition is None or completion_condition is None:
|
||||
defaults = dependency_graph.default_conditions(edge_type_norm)
|
||||
blocking_condition = (
|
||||
defaults[0] if blocking_condition is None else blocking_condition
|
||||
)
|
||||
completion_condition = (
|
||||
defaults[1] if completion_condition is None else completion_condition
|
||||
)
|
||||
evidence_json = json.dumps(
|
||||
dependency_graph.sanitize_evidence(evidence if evidence is not None else {})
|
||||
)
|
||||
now_s = _ts()
|
||||
|
||||
with self._tx() as conn:
|
||||
existing = conn.execute(
|
||||
"""
|
||||
SELECT * FROM dependency_edges
|
||||
WHERE remote = ? AND org = ? AND repo = ?
|
||||
AND source_kind = ? AND source_number = ?
|
||||
AND target_kind = ? AND target_number = ? AND edge_type = ?
|
||||
""",
|
||||
(
|
||||
remote,
|
||||
org,
|
||||
repo,
|
||||
source_kind_norm,
|
||||
source_no,
|
||||
target_kind_norm,
|
||||
target_no,
|
||||
edge_type_norm,
|
||||
),
|
||||
).fetchone()
|
||||
|
||||
if existing is None:
|
||||
edge_id = uuid.uuid4().hex
|
||||
conn.execute(
|
||||
"""
|
||||
INSERT INTO dependency_edges(
|
||||
edge_id, remote, org, repo,
|
||||
source_kind, source_number, target_kind, target_number,
|
||||
edge_type, blocking_condition, completion_condition,
|
||||
state, evidence, created_at, updated_at, last_observed_at
|
||||
) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
|
||||
""",
|
||||
(
|
||||
edge_id,
|
||||
remote,
|
||||
org,
|
||||
repo,
|
||||
source_kind_norm,
|
||||
source_no,
|
||||
target_kind_norm,
|
||||
target_no,
|
||||
edge_type_norm,
|
||||
blocking_condition,
|
||||
completion_condition,
|
||||
state_norm,
|
||||
evidence_json,
|
||||
now_s,
|
||||
now_s,
|
||||
now_s,
|
||||
),
|
||||
)
|
||||
else:
|
||||
edge_id = str(existing["edge_id"])
|
||||
conn.execute(
|
||||
"""
|
||||
UPDATE dependency_edges
|
||||
SET blocking_condition = ?, completion_condition = ?,
|
||||
state = ?, evidence = ?, updated_at = ?,
|
||||
last_observed_at = ?
|
||||
WHERE edge_id = ?
|
||||
""",
|
||||
(
|
||||
blocking_condition,
|
||||
completion_condition,
|
||||
state_norm,
|
||||
evidence_json,
|
||||
now_s,
|
||||
now_s,
|
||||
edge_id,
|
||||
),
|
||||
)
|
||||
prior_state = str(existing["state"])
|
||||
if prior_state != state_norm:
|
||||
self._record_edge_transition_conn(
|
||||
conn,
|
||||
edge_id=edge_id,
|
||||
prior_state=prior_state,
|
||||
new_state=state_norm,
|
||||
detail="observed during upsert",
|
||||
now_s=now_s,
|
||||
)
|
||||
|
||||
row = conn.execute(
|
||||
"SELECT * FROM dependency_edges WHERE edge_id = ?", (edge_id,)
|
||||
).fetchone()
|
||||
return self._dependency_edge_row(row) or {}
|
||||
|
||||
@staticmethod
|
||||
def _record_edge_transition_conn(
|
||||
conn: sqlite3.Connection,
|
||||
*,
|
||||
edge_id: str,
|
||||
prior_state: str,
|
||||
new_state: str,
|
||||
detail: str,
|
||||
now_s: str,
|
||||
) -> None:
|
||||
"""Append a state transition to the shared ``events`` audit table.
|
||||
|
||||
``work_item_id`` stays NULL: an edge endpoint is a Gitea issue/PR that
|
||||
may never have been assigned, so it has no work_items row to reference.
|
||||
"""
|
||||
message = (
|
||||
f"dependency edge {edge_id} state {prior_state} -> {new_state}"
|
||||
f" ({detail})"
|
||||
)
|
||||
conn.execute(
|
||||
"""
|
||||
INSERT INTO events(work_item_id, event_type, message, created_at)
|
||||
VALUES (NULL, 'dependency_edge_state_change', ?, ?)
|
||||
""",
|
||||
(message, now_s),
|
||||
)
|
||||
|
||||
def list_dependency_edges(
|
||||
self,
|
||||
*,
|
||||
remote: str | None = None,
|
||||
org: str | None = None,
|
||||
repo: str | None = None,
|
||||
source_kind: str | None = None,
|
||||
source_number: int | None = None,
|
||||
target_kind: str | None = None,
|
||||
target_number: int | None = None,
|
||||
edge_type: str | None = None,
|
||||
state: str | None = None,
|
||||
limit: int = 500,
|
||||
) -> list[dict[str, Any]]:
|
||||
"""Return stored edges, filtered.
|
||||
|
||||
Filtering by *target* answers "what is waiting on this work unit",
|
||||
which is the query automatic resumption needs and which body-text
|
||||
parsing could never serve.
|
||||
"""
|
||||
clauses: list[str] = []
|
||||
params: list[Any] = []
|
||||
if remote:
|
||||
clauses.append("remote = ?")
|
||||
params.append(remote)
|
||||
if org:
|
||||
clauses.append("org = ?")
|
||||
params.append(org)
|
||||
if repo:
|
||||
clauses.append("repo = ?")
|
||||
params.append(repo)
|
||||
if source_kind:
|
||||
clauses.append("source_kind = ?")
|
||||
params.append(dependency_graph.normalize_work_kind(source_kind))
|
||||
if source_number is not None:
|
||||
clauses.append("source_number = ?")
|
||||
params.append(int(source_number))
|
||||
if target_kind:
|
||||
clauses.append("target_kind = ?")
|
||||
params.append(dependency_graph.normalize_work_kind(target_kind))
|
||||
if target_number is not None:
|
||||
clauses.append("target_number = ?")
|
||||
params.append(int(target_number))
|
||||
if edge_type:
|
||||
clauses.append("edge_type = ?")
|
||||
params.append(dependency_graph.normalize_edge_type(edge_type))
|
||||
if state:
|
||||
clauses.append("state = ?")
|
||||
params.append(dependency_graph.normalize_edge_state(state))
|
||||
|
||||
sql = "SELECT * FROM dependency_edges"
|
||||
if clauses:
|
||||
sql += " WHERE " + " AND ".join(clauses)
|
||||
sql += " ORDER BY source_number ASC, target_number ASC, edge_type ASC LIMIT ?"
|
||||
params.append(int(limit))
|
||||
|
||||
with self._tx(immediate=False) as conn:
|
||||
rows = conn.execute(sql, params).fetchall()
|
||||
return [edge for edge in (self._dependency_edge_row(r) for r in rows) if edge]
|
||||
|
||||
def record_dependency_edge_observation(
|
||||
self,
|
||||
edge_id: str,
|
||||
*,
|
||||
state: str,
|
||||
evidence: Any = None,
|
||||
detail: str = "observation recorded",
|
||||
) -> dict[str, Any]:
|
||||
"""Update an existing edge's state and evidence, auditing the change.
|
||||
|
||||
A transition writes an ``events`` row carrying both the prior and the
|
||||
new state, so a later blocked/resume decision can be reconstructed from
|
||||
durable state rather than from a recomputed reason string.
|
||||
"""
|
||||
state_norm = dependency_graph.normalize_edge_state(state)
|
||||
now_s = _ts()
|
||||
with self._tx() as conn:
|
||||
existing = conn.execute(
|
||||
"SELECT * FROM dependency_edges WHERE edge_id = ?", (edge_id,)
|
||||
).fetchone()
|
||||
if existing is None:
|
||||
raise ControlPlaneError(
|
||||
f"dependency edge '{edge_id}' does not exist (fail closed)"
|
||||
)
|
||||
prior_state = str(existing["state"])
|
||||
if evidence is None:
|
||||
evidence_json = str(existing["evidence"] or "{}")
|
||||
else:
|
||||
evidence_json = json.dumps(
|
||||
dependency_graph.sanitize_evidence(evidence)
|
||||
)
|
||||
conn.execute(
|
||||
"""
|
||||
UPDATE dependency_edges
|
||||
SET state = ?, evidence = ?, updated_at = ?, last_observed_at = ?
|
||||
WHERE edge_id = ?
|
||||
""",
|
||||
(state_norm, evidence_json, now_s, now_s, edge_id),
|
||||
)
|
||||
if prior_state != state_norm:
|
||||
self._record_edge_transition_conn(
|
||||
conn,
|
||||
edge_id=edge_id,
|
||||
prior_state=prior_state,
|
||||
new_state=state_norm,
|
||||
detail=detail,
|
||||
now_s=now_s,
|
||||
)
|
||||
row = conn.execute(
|
||||
"SELECT * FROM dependency_edges WHERE edge_id = ?", (edge_id,)
|
||||
).fetchone()
|
||||
edge = self._dependency_edge_row(row) or {}
|
||||
edge["prior_state"] = prior_state
|
||||
edge["state_changed"] = prior_state != state_norm
|
||||
return edge
|
||||
|
||||
@@ -0,0 +1,307 @@
|
||||
"""Durable dependency-edge vocabulary for the control plane (#784, umbrella #628).
|
||||
|
||||
Umbrella #628 scope item 6 requires dependencies to be durable structured state
|
||||
carrying source, target, type, blocking condition, completion condition, current
|
||||
state, and evidence. Before this module the only dependency knowledge in the
|
||||
system was the per-run parse performed by :mod:`allocator_dependencies`, which
|
||||
collapsed into two in-memory ``WorkCandidate`` fields and was then discarded.
|
||||
|
||||
This module owns the vocabulary half of that store:
|
||||
|
||||
* the seven relationship types #628 enumerates;
|
||||
* the three observation states, matching the outcome of
|
||||
:func:`allocator_dependencies.resolve_dependency_state`;
|
||||
* fail-closed normalization for both, plus for work kinds;
|
||||
* the default blocking/completion condition text for each type;
|
||||
* evidence sanitization, so no credential or endpoint ever reaches the store.
|
||||
|
||||
Persistence lives in :mod:`control_plane_db`; ingestion from a live allocation
|
||||
run is :func:`record_issue_dependency_edges`. Nothing here changes allocator
|
||||
selection — this slice records the graph, it does not act on it.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import re
|
||||
from typing import Any, Iterable, Mapping
|
||||
|
||||
# --- Work kinds -------------------------------------------------------------
|
||||
# Mirrors control_plane_db.WORK_KINDS. Declared locally so this module stays
|
||||
# import-light and usable from the DB layer without a circular import.
|
||||
WORK_KIND_ISSUE = "issue"
|
||||
WORK_KIND_PR = "pr"
|
||||
WORK_KINDS = frozenset({WORK_KIND_ISSUE, WORK_KIND_PR})
|
||||
|
||||
# --- Edge types (#628 scope item 6) -----------------------------------------
|
||||
EDGE_ISSUE_BLOCKED_BY_ISSUE = "issue_blocked_by_issue"
|
||||
EDGE_PR_WAITING_FOR_REQUESTED_CHANGES = "pr_waiting_for_requested_changes"
|
||||
EDGE_MERGE_WAITING_FOR_APPROVAL = "merge_waiting_for_approval"
|
||||
EDGE_RECONCILIATION_WAITING_FOR_MERGE = "reconciliation_waiting_for_merge"
|
||||
EDGE_DEPLOYMENT_WAITING_FOR_INFRASTRUCTURE = "deployment_waiting_for_infrastructure"
|
||||
EDGE_ACCEPTANCE_WAITING_FOR_VALIDATION = "acceptance_waiting_for_validation"
|
||||
EDGE_TASK_WAITING_FOR_DEFECT_FIX = "task_waiting_for_defect_fix"
|
||||
|
||||
EDGE_TYPES: frozenset[str] = frozenset(
|
||||
{
|
||||
EDGE_ISSUE_BLOCKED_BY_ISSUE,
|
||||
EDGE_PR_WAITING_FOR_REQUESTED_CHANGES,
|
||||
EDGE_MERGE_WAITING_FOR_APPROVAL,
|
||||
EDGE_RECONCILIATION_WAITING_FOR_MERGE,
|
||||
EDGE_DEPLOYMENT_WAITING_FOR_INFRASTRUCTURE,
|
||||
EDGE_ACCEPTANCE_WAITING_FOR_VALIDATION,
|
||||
EDGE_TASK_WAITING_FOR_DEFECT_FIX,
|
||||
}
|
||||
)
|
||||
|
||||
# --- Edge states ------------------------------------------------------------
|
||||
# Deliberately three-valued: unavailable evidence is never recorded as met,
|
||||
# matching resolve_dependency_state's fail-closed contract (#758 AC6/AC7).
|
||||
STATE_UNMET = "unmet"
|
||||
STATE_MET = "met"
|
||||
STATE_UNAVAILABLE = "unavailable"
|
||||
|
||||
EDGE_STATES: frozenset[str] = frozenset({STATE_UNMET, STATE_MET, STATE_UNAVAILABLE})
|
||||
|
||||
# Default condition text per edge type: (blocking_condition, completion_condition).
|
||||
DEFAULT_CONDITIONS: dict[str, tuple[str, str]] = {
|
||||
EDGE_ISSUE_BLOCKED_BY_ISSUE: (
|
||||
"target issue is not closed",
|
||||
"target issue is closed",
|
||||
),
|
||||
EDGE_PR_WAITING_FOR_REQUESTED_CHANGES: (
|
||||
"requested changes are outstanding at the current head",
|
||||
"requested changes are addressed at the current head",
|
||||
),
|
||||
EDGE_MERGE_WAITING_FOR_APPROVAL: (
|
||||
"no approval exists at the current head",
|
||||
"an approval exists at the current head",
|
||||
),
|
||||
EDGE_RECONCILIATION_WAITING_FOR_MERGE: (
|
||||
"target pull request is not merged",
|
||||
"target pull request is merged",
|
||||
),
|
||||
EDGE_DEPLOYMENT_WAITING_FOR_INFRASTRUCTURE: (
|
||||
"required infrastructure is unavailable",
|
||||
"required infrastructure is available",
|
||||
),
|
||||
EDGE_ACCEPTANCE_WAITING_FOR_VALIDATION: (
|
||||
"required validation evidence is missing",
|
||||
"required validation evidence is recorded",
|
||||
),
|
||||
EDGE_TASK_WAITING_FOR_DEFECT_FIX: (
|
||||
"blocking defect is unresolved or undeployed",
|
||||
"blocking defect is fixed and the runtime carries the fix",
|
||||
),
|
||||
}
|
||||
|
||||
|
||||
class DependencyGraphError(ValueError):
|
||||
"""Base error for dependency-edge vocabulary violations."""
|
||||
|
||||
|
||||
class InvalidEdgeTypeError(DependencyGraphError):
|
||||
"""Raised when an edge type outside :data:`EDGE_TYPES` is supplied."""
|
||||
|
||||
|
||||
class InvalidEdgeStateError(DependencyGraphError):
|
||||
"""Raised when a state outside :data:`EDGE_STATES` is supplied."""
|
||||
|
||||
|
||||
class InvalidEdgeEndpointError(DependencyGraphError):
|
||||
"""Raised when an edge endpoint is not an assignable work unit."""
|
||||
|
||||
|
||||
def normalize_edge_type(value: Any) -> str:
|
||||
"""Return the canonical edge type, or raise fail-closed.
|
||||
|
||||
Unknown values are never coerced to a default: an unrecognized relationship
|
||||
would be stored as an unqueryable free-text row and would silently break
|
||||
reverse lookup for whichever consumer expected the real type.
|
||||
"""
|
||||
text = str(value or "").strip().lower()
|
||||
if text not in EDGE_TYPES:
|
||||
raise InvalidEdgeTypeError(
|
||||
f"unknown dependency edge_type '{value}'; expected one of "
|
||||
f"{sorted(EDGE_TYPES)} (fail closed)"
|
||||
)
|
||||
return text
|
||||
|
||||
|
||||
def normalize_edge_state(value: Any) -> str:
|
||||
"""Return the canonical edge state, or raise fail-closed."""
|
||||
text = str(value or "").strip().lower()
|
||||
if text not in EDGE_STATES:
|
||||
raise InvalidEdgeStateError(
|
||||
f"unknown dependency edge state '{value}'; expected one of "
|
||||
f"{sorted(EDGE_STATES)} (fail closed)"
|
||||
)
|
||||
return text
|
||||
|
||||
|
||||
def normalize_work_kind(value: Any) -> str:
|
||||
"""Return the canonical work kind for an edge endpoint, or raise."""
|
||||
text = str(value or "").strip().lower()
|
||||
if text not in WORK_KINDS:
|
||||
raise InvalidEdgeEndpointError(
|
||||
f"dependency edge endpoint kind '{value}' is not assignable work; "
|
||||
f"expected one of {sorted(WORK_KINDS)} (never raw incidents)"
|
||||
)
|
||||
return text
|
||||
|
||||
|
||||
def default_conditions(edge_type: str) -> tuple[str, str]:
|
||||
"""Return ``(blocking_condition, completion_condition)`` for *edge_type*."""
|
||||
return DEFAULT_CONDITIONS[normalize_edge_type(edge_type)]
|
||||
|
||||
|
||||
# --- Evidence sanitization --------------------------------------------------
|
||||
|
||||
_SECRET_KEY_PATTERN = re.compile(
|
||||
r"token|secret|password|passwd|authorization|auth_header|credential|api_key"
|
||||
r"|apikey|private_key|cookie|session_token",
|
||||
re.IGNORECASE,
|
||||
)
|
||||
_URL_PATTERN = re.compile(r"\b[a-z][a-z0-9+.-]*://\S+", re.IGNORECASE)
|
||||
|
||||
REDACTED = "[redacted]"
|
||||
|
||||
# Evidence is a small observation record; a deep or huge payload is a sign the
|
||||
# caller is dumping API responses into the store.
|
||||
_MAX_EVIDENCE_DEPTH = 6
|
||||
_MAX_EVIDENCE_STRING = 2000
|
||||
|
||||
|
||||
def sanitize_evidence(payload: Any, *, _depth: int = 0) -> Any:
|
||||
"""Return *payload* with credentials and endpoint URLs removed.
|
||||
|
||||
Applies to every stored evidence record. Keys naming a secret are replaced
|
||||
wholesale; any value containing a URL has the URL replaced, so an endpoint
|
||||
can never be persisted or handed back through a read tool.
|
||||
"""
|
||||
if _depth > _MAX_EVIDENCE_DEPTH:
|
||||
return REDACTED
|
||||
if isinstance(payload, Mapping):
|
||||
clean: dict[str, Any] = {}
|
||||
for key, value in payload.items():
|
||||
name = str(key)
|
||||
if _SECRET_KEY_PATTERN.search(name):
|
||||
clean[name] = REDACTED
|
||||
else:
|
||||
clean[name] = sanitize_evidence(value, _depth=_depth + 1)
|
||||
return clean
|
||||
if isinstance(payload, (list, tuple)):
|
||||
return [sanitize_evidence(item, _depth=_depth + 1) for item in payload]
|
||||
if isinstance(payload, str):
|
||||
text = _URL_PATTERN.sub(REDACTED, payload)
|
||||
if len(text) > _MAX_EVIDENCE_STRING:
|
||||
text = text[:_MAX_EVIDENCE_STRING] + "…"
|
||||
return text
|
||||
if isinstance(payload, (int, float, bool)) or payload is None:
|
||||
return payload
|
||||
return sanitize_evidence(str(payload), _depth=_depth + 1)
|
||||
|
||||
|
||||
# --- Ingestion from a live allocation run -----------------------------------
|
||||
|
||||
# Observed live state as recorded in evidence. The exact Gitea state string is
|
||||
# not stored for the unmet case: resolve_dependency_state has already reduced
|
||||
# "any live value other than closed" to unmet, and re-deriving it here would
|
||||
# invent evidence the resolver never produced.
|
||||
OBSERVED_CLOSED = "closed"
|
||||
OBSERVED_NOT_CLOSED = "not_closed"
|
||||
OBSERVED_UNAVAILABLE = "unavailable"
|
||||
|
||||
OBSERVATION_SOURCE_ALLOCATOR = "allocator_live_issue_lookup"
|
||||
|
||||
_OBSERVED_STATE_BY_EDGE_STATE = {
|
||||
STATE_MET: OBSERVED_CLOSED,
|
||||
STATE_UNMET: OBSERVED_NOT_CLOSED,
|
||||
STATE_UNAVAILABLE: OBSERVED_UNAVAILABLE,
|
||||
}
|
||||
|
||||
|
||||
def _observation(state: str, *, observed_by: str | None, subject: str) -> dict[str, Any]:
|
||||
return {
|
||||
"observed_state": _OBSERVED_STATE_BY_EDGE_STATE[state],
|
||||
"observation_source": OBSERVATION_SOURCE_ALLOCATOR,
|
||||
"observed_by_session": observed_by,
|
||||
"declaration": "Depends declaration in issue body",
|
||||
"subject": subject,
|
||||
}
|
||||
|
||||
|
||||
def edges_from_dependency_resolution(
|
||||
resolution: Mapping[str, Any],
|
||||
*,
|
||||
source_number: int,
|
||||
observed_by: str | None = None,
|
||||
) -> list[dict[str, Any]]:
|
||||
"""Convert one resolver result into edge records ready for persistence.
|
||||
|
||||
*resolution* is the dict returned by
|
||||
:func:`allocator_dependencies.resolve_dependency_state`. Its ``met`` /
|
||||
``unmet`` / ``unavailable`` partitions map one-to-one onto the stored
|
||||
states, so no dependency is re-classified here.
|
||||
"""
|
||||
subject = f"issue#{int(source_number)}"
|
||||
blocking, completion = default_conditions(EDGE_ISSUE_BLOCKED_BY_ISSUE)
|
||||
records: list[dict[str, Any]] = []
|
||||
partitions: tuple[tuple[str, Iterable[Any]], ...] = (
|
||||
(STATE_MET, resolution.get("met") or ()),
|
||||
(STATE_UNMET, resolution.get("unmet") or ()),
|
||||
(STATE_UNAVAILABLE, resolution.get("unavailable") or ()),
|
||||
)
|
||||
for state, refs in partitions:
|
||||
for ref in refs:
|
||||
records.append(
|
||||
{
|
||||
"source_kind": WORK_KIND_ISSUE,
|
||||
"source_number": int(source_number),
|
||||
"target_kind": WORK_KIND_ISSUE,
|
||||
"target_number": int(ref),
|
||||
"edge_type": EDGE_ISSUE_BLOCKED_BY_ISSUE,
|
||||
"state": state,
|
||||
"blocking_condition": blocking,
|
||||
"completion_condition": completion,
|
||||
"evidence": _observation(
|
||||
state, observed_by=observed_by, subject=subject
|
||||
),
|
||||
}
|
||||
)
|
||||
return records
|
||||
|
||||
|
||||
def record_issue_dependency_edges(
|
||||
db: Any,
|
||||
*,
|
||||
remote: str,
|
||||
org: str,
|
||||
repo: str,
|
||||
source_number: int,
|
||||
resolution: Mapping[str, Any],
|
||||
observed_by: str | None = None,
|
||||
) -> list[str]:
|
||||
"""Persist the edges implied by one candidate's dependency resolution.
|
||||
|
||||
Best-effort by contract: allocation correctness must not depend on this
|
||||
store existing or being writable, so every failure is returned as a reason
|
||||
string and never raised. The caller keeps using the in-memory resolution it
|
||||
already holds.
|
||||
"""
|
||||
try:
|
||||
records = edges_from_dependency_resolution(
|
||||
resolution, source_number=source_number, observed_by=observed_by
|
||||
)
|
||||
except Exception as exc: # noqa: BLE001 — ingestion never breaks allocation
|
||||
return [f"dependency edge ingestion skipped for issue#{source_number}: {exc}"]
|
||||
|
||||
reasons: list[str] = []
|
||||
for record in records:
|
||||
try:
|
||||
db.upsert_dependency_edge(remote=remote, org=org, repo=repo, **record)
|
||||
except Exception as exc: # noqa: BLE001 — see docstring
|
||||
reasons.append(
|
||||
f"dependency edge not persisted for issue#{source_number} → "
|
||||
f"issue#{record['target_number']}: {exc}"
|
||||
)
|
||||
return reasons
|
||||
@@ -102,6 +102,7 @@ that gates each call, not which tools exist.
|
||||
- `gitea_heartbeat_reviewer_pr_lease`
|
||||
- `gitea_inspect_workflow_lease`
|
||||
- `gitea_issue_irrecoverable_provenance_authorization`
|
||||
- `gitea_list_dependency_edges`
|
||||
- `gitea_list_issue_comments`
|
||||
- `gitea_list_issues`
|
||||
- `gitea_list_labels`
|
||||
|
||||
+118
-3
@@ -1932,6 +1932,7 @@ import mcp_session_state # noqa: E402
|
||||
import stale_review_decision_lock # noqa: E402
|
||||
import allocator_service # noqa: E402
|
||||
import allocator_dependencies # noqa: E402
|
||||
import dependency_graph # noqa: E402 # #784 durable dependency edges
|
||||
import control_plane_db # noqa: E402
|
||||
import lease_lifecycle # noqa: E402
|
||||
import workflow_dashboard # noqa: E402 # #605 live queue/lease dashboard
|
||||
@@ -19004,6 +19005,8 @@ def _allocator_candidates_from_gitea(
|
||||
repo: str,
|
||||
include_issues: bool = True,
|
||||
include_prs: bool = True,
|
||||
dependency_store: Any = None,
|
||||
observed_by: str | None = None,
|
||||
) -> tuple[list[Any], list[str], bool]:
|
||||
"""Build allocator candidates from live Gitea open issues/PRs.
|
||||
|
||||
@@ -19013,6 +19016,12 @@ def _allocator_candidates_from_gitea(
|
||||
candidate construction. ``inventory_complete`` is False when a required
|
||||
listing failed, so the caller can fail closed instead of ranking a
|
||||
silently short candidate set.
|
||||
|
||||
When *dependency_store* is a control-plane DB, each resolved dependency is
|
||||
also persisted as a durable edge (#784). Persistence is best-effort: a
|
||||
write failure is appended to ``reasons`` and never changes the candidate
|
||||
set, so allocation keeps working exactly as it did before the store
|
||||
existed. Callers that only read (the dashboard) pass no store.
|
||||
"""
|
||||
reasons: list[str] = []
|
||||
candidates: list[Any] = []
|
||||
@@ -19163,6 +19172,27 @@ def _allocator_candidates_from_gitea(
|
||||
)
|
||||
dep_unmet = dep_result["dependency_unmet"]
|
||||
dep_reason = dep_result["reason"]
|
||||
# #784: record the same resolution as durable graph state. The
|
||||
# in-memory fields below stay the selection input; this write only
|
||||
# makes the observation queryable and auditable afterwards.
|
||||
if dependency_store is not None and dep_result["refs"]:
|
||||
try:
|
||||
reasons.extend(
|
||||
dependency_graph.record_issue_dependency_edges(
|
||||
dependency_store,
|
||||
remote=remote,
|
||||
org=o,
|
||||
repo=r,
|
||||
source_number=int(number),
|
||||
resolution=dep_result,
|
||||
observed_by=observed_by,
|
||||
)
|
||||
)
|
||||
except Exception as exc: # noqa: BLE001 — never fail allocation
|
||||
reasons.append(
|
||||
"dependency edge persistence failed for issue#"
|
||||
f"{number}: {_redact(str(exc))}"
|
||||
)
|
||||
try:
|
||||
candidates.append(
|
||||
allocator_service.WorkCandidate(
|
||||
@@ -19950,6 +19980,12 @@ def gitea_allocate_next_work(
|
||||
"comment_lease_only": False,
|
||||
}
|
||||
|
||||
# Resolved before inventory so dependency-edge evidence can name the
|
||||
# observing session (#784); the value is unchanged from its prior use.
|
||||
sid = (session_id or "").strip() or (
|
||||
f"{profile_name or 'session'}-{os.getpid()}-{uuid.uuid4().hex[:8]}"
|
||||
)
|
||||
|
||||
inv_reasons: list[str] = []
|
||||
candidates: list[Any] = []
|
||||
if candidates_json is not None and candidates_json != "":
|
||||
@@ -19977,6 +20013,8 @@ def gitea_allocate_next_work(
|
||||
repo=r,
|
||||
include_issues=include_issues,
|
||||
include_prs=include_prs,
|
||||
dependency_store=db,
|
||||
observed_by=sid,
|
||||
)
|
||||
# #758 AC3: ranking a silently short inventory could select the wrong
|
||||
# candidate, so an incomplete listing is a fail-closed stop.
|
||||
@@ -19997,9 +20035,6 @@ def gitea_allocate_next_work(
|
||||
"comment_lease_only": False,
|
||||
}
|
||||
|
||||
sid = (session_id or "").strip() or (
|
||||
f"{profile_name or 'session'}-{os.getpid()}-{uuid.uuid4().hex[:8]}"
|
||||
)
|
||||
try:
|
||||
result = allocator_service.allocate_next_work(
|
||||
db,
|
||||
@@ -20118,6 +20153,86 @@ def gitea_list_workflow_leases(
|
||||
return result
|
||||
|
||||
|
||||
@mcp.tool()
|
||||
def gitea_list_dependency_edges(
|
||||
remote: str = "dadeschools",
|
||||
host: str | None = None,
|
||||
org: str | None = None,
|
||||
repo: str | None = None,
|
||||
source_kind: str | None = None,
|
||||
source_number: int | None = None,
|
||||
target_kind: str | None = None,
|
||||
target_number: int | None = None,
|
||||
edge_type: str | None = None,
|
||||
state: str | None = None,
|
||||
limit: int = 200,
|
||||
) -> dict:
|
||||
"""List durable dependency edges from the control plane (#784).
|
||||
|
||||
Read-only. Edges are recorded by the allocator when it resolves declared
|
||||
dependencies against live state; this tool never creates or changes one.
|
||||
|
||||
Filtering by *target* answers "what is waiting on this work unit", which
|
||||
body-text parsing could not serve. Unknown ``edge_type``/``state`` values
|
||||
fail closed rather than returning an unfiltered set.
|
||||
"""
|
||||
read_block = _profile_operation_gate("gitea.read")
|
||||
if read_block:
|
||||
return {
|
||||
"success": False,
|
||||
"reasons": read_block,
|
||||
"edges": [],
|
||||
"permission_report": _permission_block_report("gitea.read"),
|
||||
"read_only": True,
|
||||
}
|
||||
try:
|
||||
h, o, r = _resolve(remote, host, org, repo)
|
||||
except ValueError as exc:
|
||||
return {"success": False, "reasons": [str(exc)], "edges": [], "read_only": True}
|
||||
db, errs = _control_plane_db_or_error()
|
||||
if db is None:
|
||||
return {"success": False, "reasons": errs, "edges": [], "read_only": True}
|
||||
try:
|
||||
edges = db.list_dependency_edges(
|
||||
remote=remote if remote in REMOTES else remote,
|
||||
org=o,
|
||||
repo=r,
|
||||
source_kind=source_kind,
|
||||
source_number=source_number,
|
||||
target_kind=target_kind,
|
||||
target_number=target_number,
|
||||
edge_type=edge_type,
|
||||
state=state,
|
||||
limit=max(1, int(limit)),
|
||||
)
|
||||
except dependency_graph.DependencyGraphError as exc:
|
||||
return {
|
||||
"success": False,
|
||||
"reasons": [f"{exc} (fail closed)"],
|
||||
"edges": [],
|
||||
"read_only": True,
|
||||
}
|
||||
except Exception as exc: # noqa: BLE001
|
||||
return {
|
||||
"success": False,
|
||||
"reasons": [f"dependency edge listing failed: {_redact(str(exc))}"],
|
||||
"edges": [],
|
||||
"read_only": True,
|
||||
}
|
||||
return {
|
||||
"success": True,
|
||||
"read_only": True,
|
||||
"edges": edges,
|
||||
"count": len(edges),
|
||||
"remote": remote,
|
||||
"org": o,
|
||||
"repo": r,
|
||||
"edge_types": sorted(dependency_graph.EDGE_TYPES),
|
||||
"edge_states": sorted(dependency_graph.EDGE_STATES),
|
||||
"reasons": [],
|
||||
}
|
||||
|
||||
|
||||
@mcp.tool()
|
||||
def gitea_workflow_dashboard(
|
||||
remote: str = "dadeschools",
|
||||
|
||||
@@ -36,7 +36,7 @@ class ControlPlaneDBTest(unittest.TestCase):
|
||||
rows = dict(conn.execute("SELECT key, value FROM schema_meta").fetchall())
|
||||
finally:
|
||||
conn.close()
|
||||
self.assertEqual(rows["schema_version"], "3")
|
||||
self.assertEqual(rows["schema_version"], "4")
|
||||
self.assertIn("DB coordinates", rows["architecture"])
|
||||
self.assertIn("bridge", rows["architecture"].lower())
|
||||
|
||||
|
||||
@@ -0,0 +1,722 @@
|
||||
"""Tests for durable dependency edges (#784, umbrella #628 scope item 6)."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
import os
|
||||
import sqlite3
|
||||
import tempfile
|
||||
import unittest
|
||||
from unittest.mock import patch
|
||||
|
||||
import allocator_dependencies
|
||||
import dependency_graph
|
||||
import gitea_mcp_server as srv
|
||||
import mcp_tool_inventory
|
||||
from control_plane_db import SCHEMA_VERSION, ControlPlaneDB, ControlPlaneError
|
||||
|
||||
ISSUE = dependency_graph.WORK_KIND_ISSUE
|
||||
PR = dependency_graph.WORK_KIND_PR
|
||||
EDGE_BLOCKED = dependency_graph.EDGE_ISSUE_BLOCKED_BY_ISSUE
|
||||
|
||||
# Schema as it stood before this change, used to prove a real v3 → v4 migration
|
||||
# rather than a fresh-database creation dressed up as one.
|
||||
_V3_SCHEMA = """
|
||||
CREATE TABLE schema_meta (key TEXT PRIMARY KEY, value TEXT NOT NULL);
|
||||
CREATE TABLE sessions (
|
||||
session_id TEXT PRIMARY KEY,
|
||||
role TEXT NOT NULL,
|
||||
profile TEXT,
|
||||
namespace TEXT,
|
||||
pid INTEGER,
|
||||
started_at TEXT NOT NULL,
|
||||
last_heartbeat_at TEXT NOT NULL,
|
||||
status TEXT NOT NULL DEFAULT 'active'
|
||||
);
|
||||
CREATE TABLE work_items (
|
||||
work_item_id INTEGER PRIMARY KEY AUTOINCREMENT,
|
||||
remote TEXT NOT NULL,
|
||||
org TEXT NOT NULL,
|
||||
repo TEXT NOT NULL,
|
||||
kind TEXT NOT NULL CHECK (kind IN ('issue', 'pr')),
|
||||
number INTEGER NOT NULL,
|
||||
state TEXT NOT NULL DEFAULT 'open',
|
||||
priority INTEGER NOT NULL DEFAULT 0,
|
||||
current_head_sha TEXT,
|
||||
updated_at TEXT NOT NULL,
|
||||
UNIQUE (remote, org, repo, kind, number)
|
||||
);
|
||||
CREATE TABLE leases (
|
||||
lease_id TEXT PRIMARY KEY,
|
||||
work_item_id INTEGER NOT NULL REFERENCES work_items(work_item_id),
|
||||
session_id TEXT NOT NULL REFERENCES sessions(session_id),
|
||||
role TEXT NOT NULL,
|
||||
phase TEXT NOT NULL DEFAULT 'claimed',
|
||||
expires_at TEXT NOT NULL,
|
||||
heartbeat_at TEXT NOT NULL,
|
||||
status TEXT NOT NULL DEFAULT 'active'
|
||||
);
|
||||
CREATE TABLE assignments (
|
||||
assignment_id TEXT PRIMARY KEY,
|
||||
work_item_id INTEGER NOT NULL REFERENCES work_items(work_item_id),
|
||||
session_id TEXT NOT NULL REFERENCES sessions(session_id),
|
||||
lease_id TEXT NOT NULL REFERENCES leases(lease_id),
|
||||
allowed_actions TEXT NOT NULL,
|
||||
forbidden_actions TEXT NOT NULL,
|
||||
expected_head_sha TEXT,
|
||||
role TEXT NOT NULL,
|
||||
status TEXT NOT NULL DEFAULT 'active',
|
||||
created_at TEXT NOT NULL
|
||||
);
|
||||
CREATE TABLE terminal_locks (
|
||||
terminal_lock_id INTEGER PRIMARY KEY AUTOINCREMENT,
|
||||
remote TEXT NOT NULL,
|
||||
org TEXT NOT NULL,
|
||||
repo TEXT NOT NULL,
|
||||
terminal_pr INTEGER NOT NULL,
|
||||
review_id TEXT,
|
||||
decision TEXT,
|
||||
status TEXT NOT NULL DEFAULT 'active',
|
||||
cleanup_state TEXT,
|
||||
created_at TEXT NOT NULL,
|
||||
UNIQUE (remote, org, repo, terminal_pr)
|
||||
);
|
||||
CREATE TABLE events (
|
||||
event_id INTEGER PRIMARY KEY AUTOINCREMENT,
|
||||
work_item_id INTEGER REFERENCES work_items(work_item_id),
|
||||
event_type TEXT NOT NULL,
|
||||
message TEXT NOT NULL,
|
||||
created_at TEXT NOT NULL
|
||||
);
|
||||
CREATE TABLE incident_links (
|
||||
link_id INTEGER PRIMARY KEY AUTOINCREMENT,
|
||||
provider TEXT NOT NULL,
|
||||
provider_base_url TEXT NOT NULL DEFAULT '',
|
||||
provider_org TEXT NOT NULL DEFAULT '',
|
||||
provider_project TEXT NOT NULL DEFAULT '',
|
||||
provider_issue_id TEXT NOT NULL,
|
||||
provider_short_id TEXT,
|
||||
provider_permalink TEXT,
|
||||
fingerprint TEXT,
|
||||
gitea_org TEXT NOT NULL,
|
||||
gitea_repo TEXT NOT NULL,
|
||||
gitea_issue_number INTEGER NOT NULL,
|
||||
linked_pr_numbers TEXT,
|
||||
first_seen TEXT,
|
||||
last_seen TEXT,
|
||||
event_count INTEGER,
|
||||
status TEXT NOT NULL DEFAULT 'open',
|
||||
release_resolved_at TEXT,
|
||||
last_sync_at TEXT,
|
||||
UNIQUE (provider, provider_base_url, provider_org, provider_project,
|
||||
provider_issue_id)
|
||||
);
|
||||
"""
|
||||
|
||||
|
||||
def _edge_kwargs(**overrides):
|
||||
base = {
|
||||
"remote": "prgs",
|
||||
"org": "Scaled-Tech-Consulting",
|
||||
"repo": "Gitea-Tools",
|
||||
"source_kind": ISSUE,
|
||||
"source_number": 784,
|
||||
"target_kind": ISSUE,
|
||||
"target_number": 628,
|
||||
"edge_type": EDGE_BLOCKED,
|
||||
"state": dependency_graph.STATE_UNMET,
|
||||
}
|
||||
base.update(overrides)
|
||||
return base
|
||||
|
||||
|
||||
class VocabularyTest(unittest.TestCase):
|
||||
"""AC4, AC5: the edge vocabulary is complete and fails closed."""
|
||||
|
||||
def test_all_seven_umbrella_relationship_types_exist(self) -> None:
|
||||
self.assertEqual(len(dependency_graph.EDGE_TYPES), 7)
|
||||
for edge_type in (
|
||||
dependency_graph.EDGE_ISSUE_BLOCKED_BY_ISSUE,
|
||||
dependency_graph.EDGE_PR_WAITING_FOR_REQUESTED_CHANGES,
|
||||
dependency_graph.EDGE_MERGE_WAITING_FOR_APPROVAL,
|
||||
dependency_graph.EDGE_RECONCILIATION_WAITING_FOR_MERGE,
|
||||
dependency_graph.EDGE_DEPLOYMENT_WAITING_FOR_INFRASTRUCTURE,
|
||||
dependency_graph.EDGE_ACCEPTANCE_WAITING_FOR_VALIDATION,
|
||||
dependency_graph.EDGE_TASK_WAITING_FOR_DEFECT_FIX,
|
||||
):
|
||||
self.assertIn(edge_type, dependency_graph.EDGE_TYPES)
|
||||
blocking, completion = dependency_graph.default_conditions(edge_type)
|
||||
self.assertTrue(blocking and completion)
|
||||
|
||||
def test_states_match_the_resolver_partitions(self) -> None:
|
||||
self.assertEqual(
|
||||
dependency_graph.EDGE_STATES,
|
||||
frozenset({"unmet", "met", "unavailable"}),
|
||||
)
|
||||
|
||||
def test_unknown_edge_type_is_rejected(self) -> None:
|
||||
with self.assertRaises(dependency_graph.InvalidEdgeTypeError):
|
||||
dependency_graph.normalize_edge_type("waits_for_vibes")
|
||||
|
||||
def test_unknown_state_is_rejected(self) -> None:
|
||||
with self.assertRaises(dependency_graph.InvalidEdgeStateError):
|
||||
dependency_graph.normalize_edge_state("probably_fine")
|
||||
|
||||
def test_non_work_endpoint_kind_is_rejected(self) -> None:
|
||||
with self.assertRaises(dependency_graph.InvalidEdgeEndpointError):
|
||||
dependency_graph.normalize_work_kind("incident")
|
||||
|
||||
def test_evidence_sanitization_strips_credentials_and_urls(self) -> None:
|
||||
clean = dependency_graph.sanitize_evidence(
|
||||
{
|
||||
"token": "abc123",
|
||||
"authorization": "Bearer xyz",
|
||||
"note": "fetched from https://gitea.example.invalid/api/v1/x",
|
||||
"nested": [{"api_key": "k"}, "plain"],
|
||||
"observed_state": "closed",
|
||||
}
|
||||
)
|
||||
self.assertEqual(clean["token"], dependency_graph.REDACTED)
|
||||
self.assertEqual(clean["authorization"], dependency_graph.REDACTED)
|
||||
self.assertNotIn("https://", clean["note"])
|
||||
self.assertEqual(clean["nested"][0]["api_key"], dependency_graph.REDACTED)
|
||||
self.assertEqual(clean["observed_state"], "closed")
|
||||
|
||||
|
||||
class SchemaTest(unittest.TestCase):
|
||||
"""AC1-AC3: schema creation, migration, and idempotence."""
|
||||
|
||||
def setUp(self) -> None:
|
||||
self._tmp = tempfile.TemporaryDirectory()
|
||||
self.db_path = os.path.join(self._tmp.name, "cp.sqlite3")
|
||||
|
||||
def tearDown(self) -> None:
|
||||
self._tmp.cleanup()
|
||||
|
||||
def _tables(self) -> set[str]:
|
||||
conn = sqlite3.connect(self.db_path)
|
||||
try:
|
||||
return {
|
||||
row[0]
|
||||
for row in conn.execute(
|
||||
"SELECT name FROM sqlite_master WHERE type = 'table'"
|
||||
).fetchall()
|
||||
}
|
||||
finally:
|
||||
conn.close()
|
||||
|
||||
def _schema_version(self) -> str:
|
||||
conn = sqlite3.connect(self.db_path)
|
||||
try:
|
||||
row = conn.execute(
|
||||
"SELECT value FROM schema_meta WHERE key = 'schema_version'"
|
||||
).fetchone()
|
||||
finally:
|
||||
conn.close()
|
||||
return str(row[0]) if row else ""
|
||||
|
||||
def test_fresh_database_is_v4_with_the_edge_table(self) -> None:
|
||||
ControlPlaneDB(self.db_path)
|
||||
self.assertEqual(SCHEMA_VERSION, 4)
|
||||
self.assertEqual(self._schema_version(), "4")
|
||||
self.assertIn("dependency_edges", self._tables())
|
||||
|
||||
def _seed_v3(self) -> None:
|
||||
conn = sqlite3.connect(self.db_path)
|
||||
try:
|
||||
conn.executescript(_V3_SCHEMA)
|
||||
conn.execute(
|
||||
"INSERT INTO schema_meta(key, value) VALUES ('schema_version', '3')"
|
||||
)
|
||||
conn.execute(
|
||||
"""
|
||||
INSERT INTO work_items(
|
||||
remote, org, repo, kind, number, state, priority, updated_at
|
||||
) VALUES ('prgs', 'O', 'R', 'issue', 601, 'open', 20,
|
||||
'2026-07-01T00:00:00Z')
|
||||
"""
|
||||
)
|
||||
conn.execute(
|
||||
"""
|
||||
INSERT INTO sessions(session_id, role, started_at, last_heartbeat_at)
|
||||
VALUES ('legacy-session', 'author', '2026-07-01T00:00:00Z',
|
||||
'2026-07-01T00:00:00Z')
|
||||
"""
|
||||
)
|
||||
conn.execute(
|
||||
"""
|
||||
INSERT INTO events(work_item_id, event_type, message, created_at)
|
||||
VALUES (1, 'legacy', 'kept', '2026-07-01T00:00:00Z')
|
||||
"""
|
||||
)
|
||||
conn.commit()
|
||||
finally:
|
||||
conn.close()
|
||||
|
||||
def test_v3_database_migrates_in_place_without_losing_rows(self) -> None:
|
||||
self._seed_v3()
|
||||
self.assertNotIn("dependency_edges", self._tables())
|
||||
|
||||
ControlPlaneDB(self.db_path)
|
||||
|
||||
self.assertEqual(self._schema_version(), "4")
|
||||
self.assertIn("dependency_edges", self._tables())
|
||||
conn = sqlite3.connect(self.db_path)
|
||||
try:
|
||||
self.assertEqual(
|
||||
conn.execute("SELECT COUNT(*) FROM work_items").fetchone()[0], 1
|
||||
)
|
||||
self.assertEqual(
|
||||
conn.execute("SELECT COUNT(*) FROM sessions").fetchone()[0], 1
|
||||
)
|
||||
self.assertEqual(
|
||||
conn.execute(
|
||||
"SELECT message FROM events WHERE event_type = 'legacy'"
|
||||
).fetchone()[0],
|
||||
"kept",
|
||||
)
|
||||
for table in ("leases", "assignments", "terminal_locks", "incident_links"):
|
||||
conn.execute(f"SELECT COUNT(*) FROM {table}").fetchone()
|
||||
finally:
|
||||
conn.close()
|
||||
|
||||
def test_migration_is_idempotent(self) -> None:
|
||||
self._seed_v3()
|
||||
ControlPlaneDB(self.db_path)
|
||||
db = ControlPlaneDB(self.db_path) # second open re-runs the migration
|
||||
ControlPlaneDB(self.db_path)
|
||||
|
||||
self.assertEqual(self._schema_version(), "4")
|
||||
conn = sqlite3.connect(self.db_path)
|
||||
try:
|
||||
tables = conn.execute(
|
||||
"SELECT name FROM sqlite_master WHERE type = 'table' "
|
||||
"AND name = 'dependency_edges'"
|
||||
).fetchall()
|
||||
finally:
|
||||
conn.close()
|
||||
self.assertEqual(len(tables), 1)
|
||||
self.assertEqual(db.list_dependency_edges(), [])
|
||||
|
||||
|
||||
class EdgePersistenceTest(unittest.TestCase):
|
||||
"""AC5-AC10: storage, uniqueness, lookup, scope, audit, redaction."""
|
||||
|
||||
def setUp(self) -> None:
|
||||
self._tmp = tempfile.TemporaryDirectory()
|
||||
self.db = ControlPlaneDB(os.path.join(self._tmp.name, "cp.sqlite3"))
|
||||
|
||||
def tearDown(self) -> None:
|
||||
self._tmp.cleanup()
|
||||
|
||||
def _events(self) -> list[tuple[str, str]]:
|
||||
conn = sqlite3.connect(self.db.db_path)
|
||||
try:
|
||||
return [
|
||||
(str(row[0]), str(row[1]))
|
||||
for row in conn.execute(
|
||||
"SELECT event_type, message FROM events"
|
||||
).fetchall()
|
||||
]
|
||||
finally:
|
||||
conn.close()
|
||||
|
||||
def test_invalid_values_write_nothing(self) -> None:
|
||||
with self.assertRaises(dependency_graph.InvalidEdgeTypeError):
|
||||
self.db.upsert_dependency_edge(**_edge_kwargs(edge_type="nonsense"))
|
||||
with self.assertRaises(dependency_graph.InvalidEdgeStateError):
|
||||
self.db.upsert_dependency_edge(**_edge_kwargs(state="maybe"))
|
||||
with self.assertRaises(dependency_graph.InvalidEdgeEndpointError):
|
||||
self.db.upsert_dependency_edge(**_edge_kwargs(target_kind="incident"))
|
||||
self.assertEqual(self.db.list_dependency_edges(), [])
|
||||
|
||||
def test_stored_edge_carries_the_full_contract(self) -> None:
|
||||
edge = self.db.upsert_dependency_edge(
|
||||
**_edge_kwargs(evidence={"observed_state": "not_closed"})
|
||||
)
|
||||
self.assertEqual(edge["source_number"], 784)
|
||||
self.assertEqual(edge["target_number"], 628)
|
||||
self.assertEqual(edge["edge_type"], EDGE_BLOCKED)
|
||||
self.assertEqual(edge["state"], "unmet")
|
||||
self.assertEqual(edge["blocking_condition"], "target issue is not closed")
|
||||
self.assertEqual(edge["completion_condition"], "target issue is closed")
|
||||
self.assertEqual(edge["evidence"], {"observed_state": "not_closed"})
|
||||
self.assertTrue(edge["created_at"])
|
||||
self.assertTrue(edge["last_observed_at"])
|
||||
|
||||
def test_repeated_upsert_updates_one_row(self) -> None:
|
||||
first = self.db.upsert_dependency_edge(**_edge_kwargs())
|
||||
second = self.db.upsert_dependency_edge(
|
||||
**_edge_kwargs(state="met", evidence={"observed_state": "closed"})
|
||||
)
|
||||
self.assertEqual(first["edge_id"], second["edge_id"])
|
||||
edges = self.db.list_dependency_edges()
|
||||
self.assertEqual(len(edges), 1)
|
||||
self.assertEqual(edges[0]["state"], "met")
|
||||
self.assertEqual(edges[0]["evidence"], {"observed_state": "closed"})
|
||||
|
||||
def test_upsert_state_change_is_audited(self) -> None:
|
||||
self.db.upsert_dependency_edge(**_edge_kwargs())
|
||||
self.db.upsert_dependency_edge(**_edge_kwargs()) # unchanged: no event
|
||||
self.assertEqual(self._events(), [])
|
||||
self.db.upsert_dependency_edge(**_edge_kwargs(state="met"))
|
||||
events = self._events()
|
||||
self.assertEqual(len(events), 1)
|
||||
self.assertEqual(events[0][0], "dependency_edge_state_change")
|
||||
self.assertIn("unmet -> met", events[0][1])
|
||||
|
||||
def test_reverse_lookup_finds_every_waiter(self) -> None:
|
||||
self.db.upsert_dependency_edge(**_edge_kwargs(source_number=784))
|
||||
self.db.upsert_dependency_edge(**_edge_kwargs(source_number=790))
|
||||
self.db.upsert_dependency_edge(
|
||||
**_edge_kwargs(
|
||||
source_kind=PR,
|
||||
source_number=791,
|
||||
edge_type=dependency_graph.EDGE_TASK_WAITING_FOR_DEFECT_FIX,
|
||||
)
|
||||
)
|
||||
self.db.upsert_dependency_edge(
|
||||
**_edge_kwargs(source_number=792, target_number=999)
|
||||
)
|
||||
|
||||
waiters = self.db.list_dependency_edges(target_kind=ISSUE, target_number=628)
|
||||
self.assertEqual(
|
||||
sorted(edge["source_number"] for edge in waiters), [784, 790, 791]
|
||||
)
|
||||
|
||||
def test_forward_lookup_and_state_filter(self) -> None:
|
||||
self.db.upsert_dependency_edge(**_edge_kwargs(target_number=628))
|
||||
self.db.upsert_dependency_edge(**_edge_kwargs(target_number=603, state="met"))
|
||||
blockers = self.db.list_dependency_edges(source_number=784, state="unmet")
|
||||
self.assertEqual([edge["target_number"] for edge in blockers], [628])
|
||||
|
||||
def test_scope_isolation(self) -> None:
|
||||
self.db.upsert_dependency_edge(**_edge_kwargs())
|
||||
self.db.upsert_dependency_edge(**_edge_kwargs(repo="Other-Repo"))
|
||||
self.assertEqual(
|
||||
len(self.db.list_dependency_edges(remote="prgs", repo="Gitea-Tools")), 1
|
||||
)
|
||||
self.assertEqual(
|
||||
len(self.db.list_dependency_edges(remote="prgs", repo="Other-Repo")), 1
|
||||
)
|
||||
self.assertEqual(len(self.db.list_dependency_edges(remote="dadeschools")), 0)
|
||||
|
||||
def test_observation_records_transition_with_prior_state(self) -> None:
|
||||
edge = self.db.upsert_dependency_edge(**_edge_kwargs())
|
||||
updated = self.db.record_dependency_edge_observation(
|
||||
edge["edge_id"],
|
||||
state="met",
|
||||
evidence={"observed_state": "closed"},
|
||||
detail="target closed by merge",
|
||||
)
|
||||
self.assertEqual(updated["prior_state"], "unmet")
|
||||
self.assertEqual(updated["state"], "met")
|
||||
self.assertTrue(updated["state_changed"])
|
||||
events = self._events()
|
||||
self.assertEqual(len(events), 1)
|
||||
self.assertIn("unmet -> met", events[0][1])
|
||||
self.assertIn("target closed by merge", events[0][1])
|
||||
|
||||
def test_observation_on_unknown_edge_fails_closed(self) -> None:
|
||||
with self.assertRaises(ControlPlaneError):
|
||||
self.db.record_dependency_edge_observation("no-such-edge", state="met")
|
||||
|
||||
def test_evidence_never_persists_a_credential_or_endpoint(self) -> None:
|
||||
self.db.upsert_dependency_edge(
|
||||
**_edge_kwargs(
|
||||
evidence={
|
||||
"token": "super-secret",
|
||||
"source": "GET https://gitea.example.invalid/api/v1/issues/628",
|
||||
}
|
||||
)
|
||||
)
|
||||
conn = sqlite3.connect(self.db.db_path)
|
||||
try:
|
||||
raw = conn.execute("SELECT evidence FROM dependency_edges").fetchone()[0]
|
||||
finally:
|
||||
conn.close()
|
||||
self.assertNotIn("super-secret", raw)
|
||||
self.assertNotIn("https://", raw)
|
||||
stored = json.loads(raw)
|
||||
self.assertEqual(stored["token"], dependency_graph.REDACTED)
|
||||
self.assertEqual(
|
||||
self.db.list_dependency_edges()[0]["evidence"]["token"],
|
||||
dependency_graph.REDACTED,
|
||||
)
|
||||
|
||||
|
||||
class ResolutionIngestionTest(unittest.TestCase):
|
||||
"""AC11, AC12: allocation-run ingestion and write-failure tolerance."""
|
||||
|
||||
def setUp(self) -> None:
|
||||
self._tmp = tempfile.TemporaryDirectory()
|
||||
self.db = ControlPlaneDB(os.path.join(self._tmp.name, "cp.sqlite3"))
|
||||
|
||||
def tearDown(self) -> None:
|
||||
self._tmp.cleanup()
|
||||
|
||||
def _resolution(self):
|
||||
# Same call the allocator makes: parse the body, resolve live state.
|
||||
body = "* Parent: #628 · Depends: #601, #603, #999 · Related: #613"
|
||||
refs = allocator_dependencies.parse_dependency_refs(body)
|
||||
live = {601: "closed", 603: "open", 999: None}
|
||||
return allocator_dependencies.resolve_dependency_state(
|
||||
refs, lambda n: live[n], subject="issue#784"
|
||||
)
|
||||
|
||||
def test_one_edge_per_reference_with_matching_state(self) -> None:
|
||||
resolution = self._resolution()
|
||||
reasons = dependency_graph.record_issue_dependency_edges(
|
||||
self.db,
|
||||
remote="prgs",
|
||||
org="Scaled-Tech-Consulting",
|
||||
repo="Gitea-Tools",
|
||||
source_number=784,
|
||||
resolution=resolution,
|
||||
observed_by="prgs-author-1234-abcd",
|
||||
)
|
||||
self.assertEqual(reasons, [])
|
||||
|
||||
edges = {
|
||||
edge["target_number"]: edge
|
||||
for edge in self.db.list_dependency_edges(source_number=784)
|
||||
}
|
||||
self.assertEqual(sorted(edges), [601, 603, 999])
|
||||
self.assertEqual(edges[601]["state"], "met")
|
||||
self.assertEqual(edges[603]["state"], "unmet")
|
||||
self.assertEqual(edges[999]["state"], "unavailable")
|
||||
self.assertEqual(edges[999]["evidence"]["observed_state"], "unavailable")
|
||||
self.assertEqual(
|
||||
edges[603]["evidence"]["observed_by_session"], "prgs-author-1234-abcd"
|
||||
)
|
||||
self.assertEqual(edges[601]["edge_type"], EDGE_BLOCKED)
|
||||
|
||||
def test_unavailable_evidence_is_never_recorded_as_met(self) -> None:
|
||||
resolution = self._resolution()
|
||||
dependency_graph.record_issue_dependency_edges(
|
||||
self.db,
|
||||
remote="prgs",
|
||||
org="O",
|
||||
repo="R",
|
||||
source_number=784,
|
||||
resolution=resolution,
|
||||
)
|
||||
met = self.db.list_dependency_edges(state="met")
|
||||
self.assertEqual([edge["target_number"] for edge in met], [601])
|
||||
|
||||
def test_store_write_failure_is_reported_not_raised(self) -> None:
|
||||
class BrokenStore:
|
||||
def upsert_dependency_edge(self, **_kwargs):
|
||||
raise RuntimeError("disk is on fire")
|
||||
|
||||
reasons = dependency_graph.record_issue_dependency_edges(
|
||||
BrokenStore(),
|
||||
remote="prgs",
|
||||
org="O",
|
||||
repo="R",
|
||||
source_number=784,
|
||||
resolution=self._resolution(),
|
||||
)
|
||||
self.assertEqual(len(reasons), 3)
|
||||
self.assertTrue(all("disk is on fire" in reason for reason in reasons))
|
||||
|
||||
def test_no_declared_dependencies_writes_nothing(self) -> None:
|
||||
resolution = allocator_dependencies.resolve_dependency_state(
|
||||
(), lambda n: "closed", subject="issue#784"
|
||||
)
|
||||
reasons = dependency_graph.record_issue_dependency_edges(
|
||||
self.db,
|
||||
remote="prgs",
|
||||
org="O",
|
||||
repo="R",
|
||||
source_number=784,
|
||||
resolution=resolution,
|
||||
)
|
||||
self.assertEqual(reasons, [])
|
||||
self.assertEqual(self.db.list_dependency_edges(), [])
|
||||
|
||||
|
||||
class AllocationRunIngestionTest(unittest.TestCase):
|
||||
"""AC11, AC12, AC14: the live allocator path writes edges without changing
|
||||
what it selects."""
|
||||
|
||||
def setUp(self) -> None:
|
||||
self._tmp = tempfile.TemporaryDirectory()
|
||||
self.db = ControlPlaneDB(os.path.join(self._tmp.name, "cp.sqlite3"))
|
||||
|
||||
def tearDown(self) -> None:
|
||||
self._tmp.cleanup()
|
||||
|
||||
@staticmethod
|
||||
def _issue(number: int, *, body: str = "") -> dict:
|
||||
return {
|
||||
"number": number,
|
||||
"title": f"issue {number}",
|
||||
"body": body,
|
||||
"labels": [{"name": "status:ready"}],
|
||||
"state": "open",
|
||||
}
|
||||
|
||||
def _fake_gitea(self, issues, *, closed=()):
|
||||
closed_set = set(closed)
|
||||
|
||||
def api_get_all(url, _auth, **_kw):
|
||||
if "/pulls" in url:
|
||||
return []
|
||||
return list(issues)
|
||||
|
||||
def api_request(_method, url, _auth, **_kw):
|
||||
number = int(url.rsplit("/", 1)[-1])
|
||||
state = "closed" if number in closed_set else "open"
|
||||
return {"number": number, "state": state}
|
||||
|
||||
return api_get_all, api_request
|
||||
|
||||
def _allocate(self, issues, *, closed=(), db, **kwargs):
|
||||
api_get_all, api_request = self._fake_gitea(issues, closed=closed)
|
||||
with patch(
|
||||
"gitea_mcp_server._profile_operation_gate", return_value=None
|
||||
), patch(
|
||||
"gitea_mcp_server._resolve", return_value=("h", "O", "R")
|
||||
), patch(
|
||||
"gitea_mcp_server._auth", return_value="token REDACTED"
|
||||
), patch(
|
||||
"gitea_mcp_server.get_profile",
|
||||
return_value={"profile_name": "prgs-author", "role": "author"},
|
||||
), patch(
|
||||
"gitea_mcp_server._authenticated_username", return_value="jcwalker3"
|
||||
), patch(
|
||||
"gitea_mcp_server._control_plane_db_or_error", return_value=(db, [])
|
||||
), patch(
|
||||
"gitea_mcp_server.api_get_all", side_effect=api_get_all
|
||||
), patch(
|
||||
"gitea_mcp_server.api_request", side_effect=api_request
|
||||
), patch(
|
||||
"gitea_mcp_server.sentry_observability.monitor_checkin", return_value=None
|
||||
):
|
||||
return srv.gitea_allocate_next_work(
|
||||
remote="prgs", org="O", repo="R", role="author", **kwargs
|
||||
)
|
||||
|
||||
def test_live_run_persists_one_edge_per_declared_reference(self) -> None:
|
||||
issues = [
|
||||
self._issue(600, body="* Parent: #900 · Depends: #601, #500"),
|
||||
self._issue(601),
|
||||
self._issue(602),
|
||||
]
|
||||
result = self._allocate(issues, closed={500}, db=self.db)
|
||||
self.assertTrue(result["success"])
|
||||
|
||||
edges = self.db.list_dependency_edges(remote="prgs", org="O", repo="R")
|
||||
by_target = {edge["target_number"]: edge for edge in edges}
|
||||
self.assertEqual(sorted(by_target), [500, 601])
|
||||
self.assertEqual(by_target[601]["state"], "unmet")
|
||||
self.assertEqual(by_target[500]["state"], "met")
|
||||
self.assertEqual(by_target[601]["source_number"], 600)
|
||||
self.assertEqual(
|
||||
by_target[601]["edge_type"], dependency_graph.EDGE_ISSUE_BLOCKED_BY_ISSUE
|
||||
)
|
||||
self.assertTrue(by_target[601]["evidence"]["observed_by_session"])
|
||||
|
||||
def test_selection_is_unchanged_by_the_store(self) -> None:
|
||||
issues = [
|
||||
self._issue(600, body="* Depends: #601"),
|
||||
self._issue(601),
|
||||
self._issue(602),
|
||||
]
|
||||
|
||||
class DeadStore:
|
||||
"""Stands in for a control-plane DB whose edge writes all fail."""
|
||||
|
||||
def __init__(self, real):
|
||||
self._real = real
|
||||
|
||||
def __getattr__(self, name):
|
||||
return getattr(self._real, name)
|
||||
|
||||
def upsert_dependency_edge(self, **_kwargs):
|
||||
raise RuntimeError("edge store unavailable")
|
||||
|
||||
healthy = self._allocate(issues, db=self.db)
|
||||
broken = self._allocate(issues, db=DeadStore(self.db))
|
||||
|
||||
self.assertEqual(
|
||||
healthy["selected"]["number"], broken["selected"]["number"]
|
||||
)
|
||||
self.assertEqual(
|
||||
{s["number"] for s in healthy["skipped"]},
|
||||
{s["number"] for s in broken["skipped"]},
|
||||
)
|
||||
self.assertEqual(healthy["candidate_count"], broken["candidate_count"])
|
||||
self.assertTrue(broken["success"])
|
||||
warnings = broken.get("inventory_warnings") or []
|
||||
self.assertTrue(
|
||||
any("edge store unavailable" in str(w) for w in warnings),
|
||||
f"write failure must surface in reasons, got {warnings}",
|
||||
)
|
||||
|
||||
def test_repeated_runs_do_not_duplicate_edges(self) -> None:
|
||||
issues = [self._issue(600, body="* Depends: #601"), self._issue(601)]
|
||||
self._allocate(issues, db=self.db)
|
||||
self._allocate(issues, db=self.db)
|
||||
self.assertEqual(len(self.db.list_dependency_edges()), 1)
|
||||
|
||||
|
||||
class ListDependencyEdgesToolTest(unittest.TestCase):
|
||||
"""AC13: the read-only tool is gated and never mutates."""
|
||||
|
||||
def setUp(self) -> None:
|
||||
self._tmp = tempfile.TemporaryDirectory()
|
||||
self.db = ControlPlaneDB(os.path.join(self._tmp.name, "cp.sqlite3"))
|
||||
self.db.upsert_dependency_edge(**_edge_kwargs(org="O", repo="R"))
|
||||
|
||||
def tearDown(self) -> None:
|
||||
self._tmp.cleanup()
|
||||
|
||||
def _call(self, *, read_block=None, **kwargs):
|
||||
with patch(
|
||||
"gitea_mcp_server._profile_operation_gate", return_value=read_block
|
||||
), patch(
|
||||
"gitea_mcp_server._resolve", return_value=("h", "O", "R")
|
||||
), patch(
|
||||
"gitea_mcp_server._permission_block_report", return_value={"blocked": True}
|
||||
), patch(
|
||||
"gitea_mcp_server._control_plane_db_or_error", return_value=(self.db, [])
|
||||
):
|
||||
return srv.gitea_list_dependency_edges(remote="prgs", **kwargs)
|
||||
|
||||
def test_returns_stored_edges(self) -> None:
|
||||
result = self._call()
|
||||
self.assertTrue(result["success"])
|
||||
self.assertTrue(result["read_only"])
|
||||
self.assertEqual(result["count"], 1)
|
||||
self.assertEqual(result["edges"][0]["target_number"], 628)
|
||||
self.assertEqual(len(result["edge_types"]), 7)
|
||||
|
||||
def test_reverse_lookup_filter(self) -> None:
|
||||
self.assertEqual(self._call(target_number=628)["count"], 1)
|
||||
self.assertEqual(self._call(target_number=999)["count"], 0)
|
||||
|
||||
def test_without_read_permission_it_fails_closed(self) -> None:
|
||||
result = self._call(read_block=["gitea.read not allowed"])
|
||||
self.assertFalse(result["success"])
|
||||
self.assertEqual(result["edges"], [])
|
||||
self.assertIn("permission_report", result)
|
||||
|
||||
def test_invalid_filter_fails_closed(self) -> None:
|
||||
result = self._call(edge_type="not_a_real_type")
|
||||
self.assertFalse(result["success"])
|
||||
self.assertEqual(result["edges"], [])
|
||||
self.assertTrue(any("fail closed" in r for r in result["reasons"]))
|
||||
|
||||
def test_tool_is_documented_in_the_inventory(self) -> None:
|
||||
"""The #781 drift guard requires a registered tool to be documented."""
|
||||
repo_root = os.path.dirname(os.path.dirname(os.path.abspath(__file__)))
|
||||
doc = os.path.join(repo_root, mcp_tool_inventory.INVENTORY_DOC_PATH)
|
||||
with open(doc, "r", encoding="utf-8") as handle:
|
||||
documented = mcp_tool_inventory.parse_documented_inventory(handle.read())
|
||||
self.assertIn("gitea_list_dependency_edges", documented)
|
||||
|
||||
|
||||
if __name__ == "__main__": # pragma: no cover
|
||||
unittest.main()
|
||||
Reference in New Issue
Block a user