← Files Adaptive Codex OrchestratorARCHIVED FILE

hooks/runtime.py

17.6 KB · Sep 30, 2026 · 23:14 UTC

↓ Download file

#!/usr/bin/env python3
"""Current Codex command-hook entry point for Adaptive Codex Orchestrator.

The runtime is intentionally offline, standard-library only, and silent except
for one JSON object on stdout.  It never reads transcript files.
"""

from __future__ import annotations

import json
import os
from pathlib import Path
import re
import sys
import time
from typing import Mapping, MutableMapping, Optional, Tuple


HOOK_DIR = Path(__file__).resolve().parent
if str(HOOK_DIR) not in sys.path:
    sys.path.insert(0, str(HOOK_DIR))

from lib.command_parser import ParseResult, parse_command
from lib.hook_output import context_output, encode_output, empty_output
from lib.policy_loader import (
    ACTIVE_POLICY_REVISION,
    PolicyError,
    compact_reminder,
    load_model_policy,
    worker_contract,
)
from lib.project_identity import project_id
from lib.state_store import (
    SESSION_TTL_SECONDS,
    StateStore,
    StateUnavailableError,
    effective_state,
)


MAX_INPUT_BYTES = 256 * 1024
SUPPORTED_EVENTS = frozenset(
    {"SessionStart", "UserPromptSubmit", "SubagentStart", "Stop", "SessionEnd"}
)
_SAFE_ID_RE = re.compile(r"^[A-Za-z0-9._:@-]{1,256}$")
_SAFE_MODEL_RE = re.compile(r"^[A-Za-z0-9._:-]{1,128}$")


def _safe_identifier(value: object) -> Optional[str]:
    if isinstance(value, str) and _SAFE_ID_RE.fullmatch(value):
        return value
    return None


def _safe_model(value: object) -> Optional[str]:
    if isinstance(value, str) and _SAFE_MODEL_RE.fullmatch(value):
        return value.casefold()
    return None


def _looks_like_sol(model: Optional[str]) -> bool:
    if not model:
        return False
    return model == "gpt-5.6" or model == "gpt-5.6-sol" or model.startswith(
        "gpt-5.6-sol-"
    )


def _compatibility(model: Optional[str]) -> Tuple[bool, str]:
    detected = model or "unavailable"
    return (not _looks_like_sol(model), detected)


def _compatibility_notice() -> str:
    return (
        "Ultra Orchestration is active in compatibility mode. The plugin does not "
        "change the parent model. For the intended behavior, select Sol with the "
        "Ultra setting manually. Ultra reasoning cannot be programmatically verified."
    )


def _storage_diagnostic(event: str) -> dict:
    return context_output(
        event,
        "Adaptive Codex Orchestrator could not access plugin-owned state. No persistent "
        "mode change was applied. Normal Codex behavior continues; invoke "
        "`$adaptive-orchestration <task>` for a current-task-only fallback.",
    )


def _status_text(
    effective: Mapping[str, object], model: Optional[str], configured_worker: str
) -> str:
    compatibility, detected = _compatibility(model)
    enabled = "ON" if effective.get("enabled") else "OFF"
    return "\n".join(
        (
            f"Ultra Orchestration: {enabled}",
            f"Effective scope: {effective.get('effective_scope', 'disabled')}",
            f"State source: {effective.get('state_source', 'disabled')}",
            f"Profile: {effective.get('profile', 'balanced')}",
            f"Configured worker: {configured_worker}",
            f"Parent model detected: {detected}",
            "Ultra reasoning verified: no",
            "Persistent hooks trusted: yes",
            f"Compatibility mode: {'yes' if compatibility else 'no'}",
        )
    )


def _event_context(event: str, *parts: str) -> dict:
    content = "\n\n".join(part.strip() for part in parts if part and part.strip())
    return context_output(event, content) if content else empty_output()


def _project(payload: Mapping[str, object]) -> Optional[str]:
    cwd = payload.get("cwd")
    if not isinstance(cwd, str) or not cwd or len(cwd) > 32768:
        return None
    try:
        return project_id(cwd)
    except (OSError, ValueError, UnicodeError):
        return None


def _session_notice_shown(state: Mapping[str, object], session_id: str) -> bool:
    sessions = state.get("sessions")
    if not isinstance(sessions, Mapping):
        return False
    session = sessions.get(session_id)
    return bool(
        isinstance(session, Mapping) and session.get("compatibility_notice_shown") is True
    )


def _active_policy_shown(state: Mapping[str, object], session_id: str) -> bool:
    sessions = state.get("sessions")
    if not isinstance(sessions, Mapping):
        return False
    session = sessions.get(session_id)
    return bool(
        isinstance(session, Mapping)
        and session.get("active_policy_revision") == ACTIVE_POLICY_REVISION
    )


def _mark_active_policy(store: StateStore, session_id: str) -> None:
    try:
        store.mark_active_policy(session_id, ACTIVE_POLICY_REVISION)
    except StateUnavailableError:
        # Context injection remains safe when the advisory de-duplication marker
        # cannot be written. A later turn may receive the compact policy again.
        pass


def _notice_if_needed(
    *,
    store: StateStore,
    state: Mapping[str, object],
    session_id: str,
    model: Optional[str],
) -> str:
    compatibility, _ = _compatibility(model)
    if not compatibility or _session_notice_shown(state, session_id):
        return ""
    try:
        store.mark_compatibility_notice(session_id)
        return _compatibility_notice()
    except StateUnavailableError:
        return (
            _compatibility_notice()
            + " The once-per-session marker could not be saved, so this notice may recur."
        )


def _active_context(
    *,
    plugin_root: Path,
    effective: Mapping[str, object],
    prefix: str = "",
    notice: str = "",
) -> str:
    policy = compact_reminder(plugin_root, str(effective["profile"]))
    return "\n\n".join(value for value in (prefix, notice, policy) if value)


def _handle_session_start(
    payload: Mapping[str, object],
    *,
    store: StateStore,
    state: Mapping[str, object],
    session_id: str,
    project_hash: str,
    plugin_root: Path,
) -> dict:
    effective = effective_state(state, session_id, project_hash)
    if not effective["enabled"]:
        return empty_output()
    model = _safe_model(payload.get("model"))
    notice = _notice_if_needed(
        store=store, state=state, session_id=session_id, model=model
    )
    source = payload.get("source")
    prefix = (
        "Codex context was compacted; restore the active orchestration policy below."
        if source == "compact"
        else ""
    )
    _mark_active_policy(store, session_id)
    return _event_context(
        "SessionStart",
        _active_context(
            plugin_root=plugin_root,
            effective=effective,
            prefix=prefix,
            notice=notice,
        ),
    )


def _apply_control(
    parsed: ParseResult,
    *,
    store: StateStore,
    session_id: str,
    project_hash: str,
    turn_id: Optional[str],
) -> dict:
    if parsed.intent in {"enable", "disable"}:
        return store.set_enabled(
            scope=parsed.scope or "session",
            enabled=parsed.intent == "enable",
            session_id=session_id,
            project_id=project_hash,
            turn_id=turn_id,
            profile=parsed.profile,
        )
    if parsed.intent == "set_profile" and parsed.profile:
        return store.set_profile(
            scope=parsed.scope or "session",
            profile=parsed.profile,
            session_id=session_id,
            project_id=project_hash,
            turn_id=turn_id,
        )
    raise ValueError("not a control intent")


def _handle_user_prompt(
    payload: Mapping[str, object],
    *,
    store: StateStore,
    state: Mapping[str, object],
    session_id: str,
    project_hash: str,
    plugin_root: Path,
) -> dict:
    prompt = payload.get("prompt")
    if not isinstance(prompt, str):
        return empty_output()
    turn_id = _safe_identifier(payload.get("turn_id"))
    model = _safe_model(payload.get("model"))
    configured_worker = str(load_model_policy(plugin_root)["primary_worker"])
    before = effective_state(state, session_id, project_hash)
    parsed = parse_command(prompt, mode_active=bool(before["enabled"]))

    # Status and ambiguous requests are explicitly read-only, including stale
    # one-shot cleanup and compatibility-notice bookkeeping.
    if parsed.intent == "status":
        return _event_context(
            "UserPromptSubmit", _status_text(before, model, configured_worker)
        )
    if parsed.ambiguous:
        return _event_context(
            "UserPromptSubmit",
            "Ultra Orchestration command was ambiguous; no state changed. "
            f"Reason: {parsed.reason}. State the desired scope/profile explicitly.",
        )

    # If Stop was skipped, consume a one-shot before handling a later turn.
    one_shot = state.get("one_shot", {})
    existing = one_shot.get(session_id) if isinstance(one_shot, Mapping) else None
    if isinstance(existing, Mapping) and (
        not turn_id
        or not existing.get("turn_id")
        or existing.get("turn_id") != turn_id
    ):
        state = store.clear_stale_one_shot(session_id, turn_id)
        before = effective_state(state, session_id, project_hash)
        parsed = parse_command(prompt, mode_active=bool(before["enabled"]))
        if parsed.ambiguous:
            # This ambiguity follows an automatic expiry, not a command change.
            return _event_context(
                "UserPromptSubmit",
                "A completed one-shot expired. The new control wording is ambiguous; no "
                f"new preference was applied. Reason: {parsed.reason}.",
            )

    if parsed.intent is None:
        if not before["enabled"]:
            return empty_output()
        notice = _notice_if_needed(
            store=store, state=state, session_id=session_id, model=model
        )
        if _active_policy_shown(state, session_id):
            return _event_context("UserPromptSubmit", notice)
        _mark_active_policy(store, session_id)
        return _event_context(
            "UserPromptSubmit", notice, compact_reminder(plugin_root, str(before["profile"]))
        )

    if parsed.scope == "project" and project_hash == "0" * 64:
        return _event_context(
            "UserPromptSubmit",
            "No project preference was changed because the host did not provide a usable "
            "working directory. Choose a session or one-shot scope instead.",
        )

    updated = _apply_control(
        parsed,
        store=store,
        session_id=session_id,
        project_hash=project_hash,
        turn_id=turn_id,
    )
    after = effective_state(updated, session_id, project_hash)
    action = {
        "enable": "enabled",
        "disable": "disabled",
        "set_profile": "profile updated",
    }[parsed.intent]
    applied = (
        f"Ultra Orchestration command applied: {action}; requested scope "
        f"{parsed.scope or 'session'}."
    )
    status = _status_text(after, model, configured_worker)

    if after["enabled"]:
        notice = _notice_if_needed(
            store=store, state=updated, session_id=session_id, model=model
        )
        # A new enable or active profile change gets the compact policy. A scoped
        # disable that leaves a higher-precedence ON state gets the same summary.
        if parsed.intent in {"enable", "set_profile"}:
            policy = _active_context(
                plugin_root=plugin_root, effective=after, notice=notice
            )
        else:
            policy = "\n\n".join(
                value
                for value in (
                    notice,
                    "The requested OFF preference was saved, but a higher-precedence ON "
                    "override remains effective for this turn. "
                    + compact_reminder(plugin_root, str(after["profile"])),
                )
                if value
            )
        _mark_active_policy(store, session_id)
        return _event_context("UserPromptSubmit", applied, status, policy)

    return _event_context(
        "UserPromptSubmit",
        applied,
        status,
        "Ultra Orchestration is OFF. Do not apply prior orchestration policy to this "
        "request unless the user explicitly invokes the current-task skill.",
    )


def _handle_subagent_start(
    payload: Mapping[str, object],
    *,
    state: Mapping[str, object],
    session_id: str,
    project_hash: str,
    plugin_root: Path,
) -> dict:
    effective = effective_state(state, session_id, project_hash)
    if not effective["enabled"]:
        return empty_output()
    policy = load_model_policy(plugin_root)
    model = _safe_model(payload.get("model"))
    if model and model != policy["primary_worker"].casefold():
        return empty_output()
    return _event_context(
        "SubagentStart",
        worker_contract(plugin_root, confirmed_spark=model == policy["primary_worker"].casefold()),
    )


def _has_stale_sessions(state: Mapping[str, object], session_id: str) -> bool:
    sessions = state.get("sessions", {})
    if not isinstance(sessions, Mapping):
        return False
    now = int(time.time())
    for key, entry in sessions.items():
        if key == session_id or not isinstance(entry, Mapping):
            continue
        last_seen = entry.get("last_seen")
        if isinstance(last_seen, int) and now - last_seen > SESSION_TTL_SECONDS:
            return True
    return False


def handle_event(
    payload: object,
    *,
    environ: Optional[Mapping[str, str]] = None,
    plugin_root: Optional[Path] = None,
) -> dict:
    """Handle one decoded hook event; exposed for offline fixture tests."""

    if not isinstance(payload, Mapping):
        return empty_output()
    if "schema_version" in payload and payload.get("schema_version") != 1:
        return empty_output()
    event = payload.get("hook_event_name")
    if event not in SUPPORTED_EVENTS:
        return empty_output()
    session_id = _safe_identifier(payload.get("session_id"))
    if not session_id:
        return empty_output()

    environ = environ if environ is not None else os.environ
    plugin_root = Path(plugin_root or environ.get("PLUGIN_ROOT") or HOOK_DIR.parent)
    data_value = environ.get("PLUGIN_DATA")
    if not data_value:
        if event == "UserPromptSubmit":
            parsed = parse_command(payload.get("prompt"), mode_active=False)
            if parsed.intent or parsed.ambiguous:
                return _storage_diagnostic("UserPromptSubmit")
        return empty_output()

    project_hash = _project(payload)
    if not project_hash:
        # Session/global operations remain safe. A fixed all-zero hash is never
        # persisted unless a project-scoped command reaches the explicit guard.
        project_hash = "0" * 64
    store = StateStore.from_data_dir(Path(data_value))
    try:
        # Cleanup hooks use one bounded lock acquisition so their three-second
        # host deadline cannot interrupt a second state-lock wait.
        if event == "Stop":
            store.clear_one_shot(session_id)
            return empty_output()
        if event == "SessionEnd":
            store.clear_session(session_id)
            return empty_output()
        state = store.read()
        if event == "SessionStart":
            sessions = state.get("sessions", {})
            if isinstance(sessions, Mapping) and session_id in sessions:
                state = store.touch_session(session_id)
            if _has_stale_sessions(state, session_id):
                state = store.prune_stale_sessions(keep_session=session_id)
            return _handle_session_start(
                payload,
                store=store,
                state=state,
                session_id=session_id,
                project_hash=project_hash,
                plugin_root=plugin_root,
            )
        if event == "UserPromptSubmit":
            return _handle_user_prompt(
                payload,
                store=store,
                state=state,
                session_id=session_id,
                project_hash=project_hash,
                plugin_root=plugin_root,
            )
        if event == "SubagentStart":
            return _handle_subagent_start(
                payload,
                state=state,
                session_id=session_id,
                project_hash=project_hash,
                plugin_root=plugin_root,
            )
    except (StateUnavailableError, PolicyError, OSError, ValueError, TypeError):
        if event == "UserPromptSubmit":
            parsed = parse_command(payload.get("prompt"), mode_active=False)
            if parsed.intent or parsed.ambiguous:
                return _storage_diagnostic("UserPromptSubmit")
        return empty_output()
    return empty_output()


def _read_stdin() -> object:
    raw = sys.stdin.buffer.read(MAX_INPUT_BYTES + 1)
    if len(raw) > MAX_INPUT_BYTES:
        return None
    try:
        return json.loads(raw.decode("utf-8"))
    except (UnicodeDecodeError, json.JSONDecodeError):
        return None


def main() -> int:
    try:
        result = handle_event(_read_stdin())
        sys.stdout.write(encode_output(result))
        sys.stdout.write("\n")
        return 0
    except Exception:
        # Hooks must never break ordinary Codex work. Avoid stderr/stdout details
        # that might expose host paths or event contents.
        sys.stdout.write("{}\n")
        return 0


if __name__ == "__main__":
    raise SystemExit(main())

SHA-256: 7d5f8ba11146730b47a5674a6d0b6b5d722e9ab3295aa48cf57129de0e58c89b