feat: route profiled ops runs through Glas

Assistant: codex
Assistant-Model: gpt-5.6-sol
Assistant-Session: 01a02b6f-7db1-7222-918b-e813a6bda38d
This commit is contained in:
tegwick 2026-08-22 23:55:29 +02:00
parent 68ab94581a
commit 0f01c3e421
17 changed files with 895 additions and 11 deletions

View file

@ -23,6 +23,11 @@ from rein_aharness.approaches import (
execute_approach,
select_approach,
)
from rein_aharness.glas_execution import (
GLAS_APPROACH,
GlasExecutionError,
execute_profiled_run,
)
from rein_aharness.ops_run_client import (
ActivityCoreOpsClient,
OpsRun,
@ -105,7 +110,7 @@ def process_one(
return ProcessResult(claimed=False, empty=True, reason="queue empty")
run = claimed[0]
approach = select_approach(run)
approach = GLAS_APPROACH if run.harness_profile_ref else select_approach(run)
logger.info(
"claimed run_id=%s approach=%s title=%r labels=%s",
run.id,
@ -140,6 +145,13 @@ def process_one(
detail={"dry_run": True},
)
if run.harness_profile_ref:
return _process_profiled_run(
client,
run,
report_to_hub=report_to_hub,
)
hb = _Heartbeat(client, run.id, cfg.lease_seconds)
hb.start()
try:
@ -207,6 +219,84 @@ def process_one(
)
def _process_profiled_run(
client: ActivityCoreOpsClient,
run: OpsRun,
*,
report_to_hub: bool,
) -> ProcessResult:
"""Execute a profiled row without consulting or falling back to legacy routing."""
hb = _Heartbeat(client, run.id, client.config.lease_seconds)
hb.start()
try:
try:
gateway_result = execute_profiled_run(
run,
client.config,
report_to_hub=report_to_hub,
)
except GlasExecutionError as exc:
reason = str(exc)
try:
out = client.fail(
run.id,
error=reason,
reopen=False,
result={"ok": False, "approach": GLAS_APPROACH, "reason": reason},
)
except OpsRunError as close_exc:
return ProcessResult(
claimed=True,
run_id=run.id,
approach=GLAS_APPROACH,
ok=False,
reason=f"close ops_run failed: {close_exc}; {reason}",
)
return ProcessResult(
claimed=True,
run_id=run.id,
approach=GLAS_APPROACH,
ok=False,
reason=reason,
ops_state=out.state,
)
finally:
hb.stop()
evidence = gateway_result["evidence"]
ok = gateway_result["ok"]
reason = str(evidence.get("error") or evidence.get("outcome") or "Glas execution failed")
try:
if ok:
out = client.complete(run.id, result=gateway_result)
else:
out = client.fail(
run.id,
error=reason,
reopen=False,
result=gateway_result,
)
except OpsRunError as exc:
return ProcessResult(
claimed=True,
run_id=run.id,
approach=GLAS_APPROACH,
ok=False,
reason=f"close ops_run failed: {exc}; gateway_ok={ok} {reason}",
detail={"execution_evidence": evidence},
)
return ProcessResult(
claimed=True,
run_id=run.id,
approach=GLAS_APPROACH,
ok=ok,
reason="" if ok else reason,
ops_state=out.state,
detail={"execution_evidence": evidence},
)
def poll_peek(client: ActivityCoreOpsClient | None = None) -> list[dict[str, Any]]:
"""List open ops_runs with selected approach (no claim)."""
client = client or ActivityCoreOpsClient()
@ -220,7 +310,10 @@ def poll_peek(client: ActivityCoreOpsClient | None = None) -> list[dict[str, Any
"state": run.state,
"labels": run.labels,
"target_repo": run.target_repo,
"approach": select_approach(run),
"approach": (
GLAS_APPROACH if run.harness_profile_ref else select_approach(run)
),
"harness_profile_ref": run.harness_profile_ref,
"created_at": run.raw.get("created_at"),
}
)

View file

@ -0,0 +1,91 @@
"""Profile-authoritative execution of activity-core ops runs through Glas.
Glas and sand-boxer are sibling runtime packages rather than hard dependencies
of rein-aharness's legacy paths. Imports therefore stay at the execution edge:
profile-absent coexistence continues to work, while a profiled row fails closed
with an actionable error if the governed runtime is not installed.
"""
from __future__ import annotations
from collections.abc import Callable
from typing import Any
from rein_aharness.ops_run_client import OpsRun, OpsRunConfig, resolve_ops_target
GLAS_APPROACH = "glas-profile"
_SCALAR_REFS = (
"correlation_id",
"assignment_ref",
"role_ref",
"duty_ref",
)
_LIST_REFS = ("goal_refs", "resource_envelope_refs")
class GlasExecutionError(RuntimeError):
"""The authoritative Glas invocation could not produce a GatewayResult."""
def _request_kwargs(run: OpsRun, config: OpsRunConfig, report_to_hub: bool) -> dict[str, Any]:
if not run.harness_profile_ref:
raise GlasExecutionError(f"ops_run {run.id} has no harness_profile_ref")
refs = run.execution_refs
kwargs: dict[str, Any] = {
"harness_profile_ref": run.harness_profile_ref,
"repo": str(resolve_ops_target(run, config)),
"title": run.title or "(untitled)",
"description": run.description or "",
"actor": config.worker_id,
"project": "rein-aharness",
"request_id": run.id,
"report_to_hub": report_to_hub,
}
for key in _SCALAR_REFS:
value = refs.get(key)
if isinstance(value, str) and value:
kwargs[key] = value
for key in _LIST_REFS:
value = refs.get(key)
if isinstance(value, list):
kwargs[key] = [item for item in value if isinstance(item, str)]
return kwargs
def execute_profiled_run(
run: OpsRun,
config: OpsRunConfig,
*,
report_to_hub: bool = True,
request_factory: Callable[..., Any] | None = None,
gateway: Callable[[Any], Any] | None = None,
) -> dict[str, Any]:
"""Invoke Glas and return its complete JSON-compatible GatewayResult."""
if request_factory is None or gateway is None:
try:
from glas_harness.contract import ExecutionRequest
from glas_harness.gateway import run_execution
except ImportError as exc:
raise GlasExecutionError(
"profiled ops_run requires glas-harness with sand-boxer support; "
"install the sibling checkouts into the claim-worker environment"
) from exc
request_factory = request_factory or ExecutionRequest
gateway = gateway or run_execution
try:
request = request_factory(**_request_kwargs(run, config, report_to_hub))
result = gateway(request)
raw = result.model_dump(mode="json") if hasattr(result, "model_dump") else result
except GlasExecutionError:
raise
except Exception as exc:
raise GlasExecutionError(f"Glas gateway invocation failed: {exc}") from exc
if not isinstance(raw, dict) or not isinstance(raw.get("ok"), bool):
raise GlasExecutionError("Glas gateway returned an invalid GatewayResult")
if not isinstance(raw.get("evidence"), dict):
raise GlasExecutionError("Glas GatewayResult is missing execution evidence")
return raw

View file

@ -28,6 +28,8 @@ _SAFE_RESPONSE_METADATA_KEYS = frozenset(
"created_at",
}
)
_SAFE_ERROR_KEYS = ("error", "provider_status", "provider", "model")
_MAX_ERROR_MESSAGE_CHARS = 400
class LLMConnectError(RuntimeError):
@ -60,9 +62,10 @@ class LLMConnectClient:
json=payload,
timeout=self.timeout_seconds,
)
resp.raise_for_status()
except httpx.HTTPError as exc:
raise LLMConnectError(f"llm-connect request failed: {exc}") from exc
if resp.status_code >= 400:
raise LLMConnectError(_llm_connect_error_text(resp))
try:
data = resp.json()
except ValueError as exc:
@ -74,6 +77,27 @@ class LLMConnectClient:
return content
def _llm_connect_error_text(resp: httpx.Response) -> str:
"""Keep llm-connect's safe cause without copying an upstream response blob."""
base = f"llm-connect returned HTTP {resp.status_code}"
try:
body = resp.json()
except ValueError:
return base
if not isinstance(body, dict):
return base
parts = [
f"{key}={body[key]}"
for key in _SAFE_ERROR_KEYS
if body.get(key) not in (None, "")
]
message = body.get("message")
if isinstance(message, str) and message.strip():
parts.append(f"message={message.strip()[:_MAX_ERROR_MESSAGE_CHARS]}")
return f"{base}: " + "; ".join(parts) if parts else base
def get_llm_connect_client() -> LLMConnectClient:
base_url = os.environ.get("LLM_CONNECT_URL", "").strip()
if not base_url:

View file

@ -56,6 +56,8 @@ class OpsRun:
source_id: str = ""
triggering_event_id: str = ""
approach_hint: str | None = None
harness_profile_ref: str | None = None
execution_refs: dict[str, Any] = field(default_factory=dict)
result: dict[str, Any] = field(default_factory=dict)
raw: dict[str, Any] = field(default_factory=dict)
@ -78,6 +80,8 @@ class OpsRun:
source_id=str(data.get("source_id") or ""),
triggering_event_id=str(data.get("triggering_event_id") or ""),
approach_hint=data.get("approach_hint"),
harness_profile_ref=data.get("harness_profile_ref"),
execution_refs=dict(data.get("execution_refs") or {}),
result=dict(data.get("result") or {}),
raw=data,
)