import hashlib
import re
import time
import uuid
from datetime import UTC, datetime, timedelta
from typing import Annotated, Literal

from fastapi import APIRouter, Cookie, Depends, Header, Request, Response
from fastapi.encoders import jsonable_encoder
from pydantic import BaseModel, Field, field_validator, model_validator
from sqlalchemy import func, select

from .auth import AuthDep, SessionDep, _audit, _require_csrf
from .authorization import authorize, permissions_for
from .computer_client import (
    ComputerClient,
    ComputerWorkerActionResult,
    ComputerWorkerSessionState,
    get_computer_client,
)
from .computer_tools import COMPUTER_TOOLS
from .errors import AppError
from .health import ComputerHealthDep
from .models import (
    Agent,
    AgentPermission,
    AgentTool,
    ApprovalRequest,
    Computer,
    ComputerAction,
    ComputerArtifact,
    ComputerCheckpoint,
    ComputerControlLease,
    ComputerProfile,
    ComputerSession,
    DeviceSession,
    Execution,
    IdempotencyRecord,
    WorkspaceControlState,
)
from .policy import ActionRequest, AutonomyLevel, Decision, evaluate_action
from .security import canonical_payload_hash, payload_matches

router = APIRouter(prefix="/api/v1/computer", tags=["computer"])
ComputerClientDep = Annotated[ComputerClient, Depends(get_computer_client)]
TERMINAL_STATUSES = {"stopped", "failed"}
ACTIVE_STATUSES = {
    "starting", "ready", "ai_controlled", "human_controlled", "awaiting_human",
    "paused", "recovering", "stopping",
}
LEASE_DURATION = timedelta(minutes=5)


class ProfileInput(BaseModel):
    agent_id: uuid.UUID
    profile_key: str = Field(min_length=2, max_length=80)
    name: str = Field(min_length=1, max_length=160)
    retention_days: int = Field(default=30, ge=1, le=3650)

    @field_validator("profile_key")
    @classmethod
    def validate_key(cls, value: str) -> str:
        value = value.strip().lower()
        if not re.fullmatch(r"[a-z][a-z0-9_-]+", value):
            raise ValueError("Profile key must be a stable lowercase identifier")
        return value

    @field_validator("name")
    @classmethod
    def validate_name(cls, value: str) -> str:
        value = value.strip()
        if not value:
            raise ValueError("Profile name cannot be blank")
        return value


class ProfileUpdateInput(BaseModel):
    name: str | None = Field(default=None, min_length=1, max_length=160)
    retention_days: int | None = Field(default=None, ge=1, le=3650)
    status: Literal["active", "disabled"] | None = None

    @field_validator("name")
    @classmethod
    def validate_name(cls, value: str | None) -> str | None:
        if value is None:
            return None
        value = value.strip()
        if not value:
            raise ValueError("Profile name cannot be blank")
        return value


class SessionInput(BaseModel):
    task_summary: str = Field(min_length=1, max_length=20_000)
    profile_id: uuid.UUID
    agent_id: uuid.UUID | None = None
    execution_id: uuid.UUID | None = None

    @field_validator("task_summary")
    @classmethod
    def validate_task(cls, value: str) -> str:
        value = value.strip()
        if not value:
            raise ValueError("Task summary cannot be blank")
        return value


class BrowserActionRequest(BaseModel):
    tool: Literal[
        "observe", "navigate", "click", "type", "press", "wait_for", "extract",
        "tab_list", "tab_open", "tab_switch", "tab_close", "screenshot",
    ]
    arguments: dict = Field(default_factory=dict)

    @model_validator(mode="after")
    def validate_size(self):
        if len(self.arguments) > 20:
            raise ValueError("Too many action arguments")
        return self


class ActionReconciliationInput(BaseModel):
    outcome: Literal["succeeded", "failed"]
    evidence_summary: str = Field(min_length=3, max_length=1000)


class ActionDecisionInput(BaseModel):
    approve: bool
    arguments: dict = Field(default_factory=dict)


class HumanPointerRequest(BaseModel):
    action: Literal["click", "move", "scroll"]
    x: float = Field(ge=0, le=10_000)
    y: float = Field(ge=0, le=10_000)
    delta_x: float = Field(default=0, ge=-10_000, le=10_000)
    delta_y: float = Field(default=0, ge=-10_000, le=10_000)


class HumanKeyboardRequest(BaseModel):
    text: str | None = Field(default=None, max_length=20_000)
    key: str | None = Field(default=None, min_length=1, max_length=80)

    @model_validator(mode="after")
    def require_one_input(self):
        if (self.text is None) == (self.key is None):
            raise ValueError("Provide either text or key")
        return self


class EmergencyStopInput(BaseModel):
    reason: str = Field(min_length=1, max_length=500)

    @field_validator("reason")
    @classmethod
    def validate_reason(cls, value: str) -> str:
        value = value.strip()
        if not value:
            raise ValueError("Emergency stop reason cannot be blank")
        return value


def _profile_data(profile: ComputerProfile) -> dict:
    return {
        "id": str(profile.id),
        "agent_id": str(profile.agent_id) if profile.agent_id else None,
        "computer_id": str(profile.computer_id) if profile.computer_id else None,
        "profile_key": profile.profile_key,
        "name": profile.name,
        "status": profile.status,
        "retention_days": profile.retention_days,
        "version": profile.version,
        "revoked_at": profile.revoked_at,
        "created_at": profile.created_at,
        "updated_at": profile.updated_at,
    }


def _session_data(record: ComputerSession) -> dict:
    return {
        "id": str(record.id),
        "requested_by_user_id": str(record.requested_by_user_id),
        "agent_id": str(record.agent_id),
        "execution_id": str(record.execution_id) if record.execution_id else None,
        "profile_id": str(record.profile_id),
        "task_summary": record.task_summary,
        "status": record.status,
        "control_owner": record.control_owner,
        "worker_id": record.worker_id,
        "current_url": record.current_url,
        "active_tab_id": record.active_tab_id,
        "current_action_summary": record.current_action_summary,
        "upcoming_action_summary": record.upcoming_action_summary,
        "failure_code": record.failure_code,
        "failure_message": record.failure_message,
        "lease_version": record.lease_version,
        "version": record.version,
        "started_at": record.started_at,
        "last_active_at": record.last_active_at,
        "paused_at": record.paused_at,
        "stopped_at": record.stopped_at,
        "created_at": record.created_at,
        "updated_at": record.updated_at,
    }


async def _csrf(
    session: SessionDep, auth: AuthDep, cookie: str | None, header: str | None,
) -> None:
    _require_csrf(await session.get(DeviceSession, auth.session_id), cookie, header)


async def _control_state(
    session: SessionDep, workspace_id: uuid.UUID, *, lock: bool = False,
) -> WorkspaceControlState:
    statement = select(WorkspaceControlState).where(
        WorkspaceControlState.workspace_id == workspace_id
    )
    if lock:
        statement = statement.with_for_update()
    state = await session.scalar(statement)
    if not state:
        raise AppError(
            "COMPUTER_CONTROL_STATE_MISSING",
            "Computer control state is unavailable.",
            503,
        )
    return state


async def _load_profile(
    session: SessionDep, workspace_id: uuid.UUID, profile_id: uuid.UUID, *, lock: bool = False,
) -> ComputerProfile:
    statement = select(ComputerProfile).where(
        ComputerProfile.workspace_id == workspace_id,
        ComputerProfile.id == profile_id,
    )
    if lock:
        statement = statement.with_for_update()
    profile = await session.scalar(statement)
    if not profile:
        raise AppError("COMPUTER_PROFILE_NOT_FOUND", "The computer profile is unavailable.", 404)
    return profile


async def _load_computer_session(
    session: SessionDep, workspace_id: uuid.UUID, session_id: uuid.UUID, *, lock: bool = False,
) -> ComputerSession:
    statement = select(ComputerSession).where(
        ComputerSession.workspace_id == workspace_id,
        ComputerSession.id == session_id,
    )
    if lock:
        statement = statement.with_for_update()
    record = await session.scalar(statement)
    if not record:
        raise AppError("COMPUTER_SESSION_NOT_FOUND", "The computer session is unavailable.", 404)
    return record


async def _active_lease(
    session: SessionDep, workspace_id: uuid.UUID, session_id: uuid.UUID, *, lock: bool = False,
) -> ComputerControlLease | None:
    statement = (
        select(ComputerControlLease)
        .where(
            ComputerControlLease.workspace_id == workspace_id,
            ComputerControlLease.session_id == session_id,
            ComputerControlLease.released_at.is_(None),
        )
        .order_by(ComputerControlLease.fencing_token.desc())
        .limit(1)
    )
    if lock:
        statement = statement.with_for_update()
    return await session.scalar(statement)


def _release_lease(lease: ComputerControlLease | None, reason: str, now: datetime) -> None:
    if lease and lease.released_at is None:
        lease.released_at = now
        lease.release_reason = reason[:160]


def _heartbeat_lease(lease: ComputerControlLease, *, now: datetime | None = None) -> None:
    current = now or datetime.now(UTC)
    expires_at = lease.expires_at
    if expires_at.tzinfo is None:
        expires_at = expires_at.replace(tzinfo=UTC)
    if current >= expires_at:
        raise AppError("COMPUTER_LEASE_EXPIRED", "The control lease has expired.", 409)
    lease.heartbeat_at = current
    lease.expires_at = current + LEASE_DURATION


def _issue_lease(
    db_session: SessionDep,
    record: ComputerSession,
    *,
    owner: Literal["ai", "human", "system"],
    actor_user_id: uuid.UUID | None,
    now: datetime,
) -> ComputerControlLease:
    record.lease_version += 1
    lease = ComputerControlLease(
        id=uuid.uuid4(),
        workspace_id=record.workspace_id,
        session_id=record.id,
        owner=owner,
        actor_user_id=actor_user_id,
        fencing_token=record.lease_version,
        issued_at=now,
        expires_at=now + LEASE_DURATION,
        heartbeat_at=now,
    )
    db_session.add(lease)
    return lease


async def _checkpoint(
    db_session: SessionDep,
    record: ComputerSession,
    checkpoint_type: str,
    *,
    worker_state: ComputerWorkerSessionState | None = None,
    state_redacted: dict | None = None,
    action_id: uuid.UUID | None = None,
) -> None:
    sequence = await db_session.scalar(select(func.max(ComputerCheckpoint.sequence)).where(
        ComputerCheckpoint.workspace_id == record.workspace_id,
        ComputerCheckpoint.session_id == record.id,
    ))
    db_session.add(ComputerCheckpoint(
        id=uuid.uuid4(), workspace_id=record.workspace_id, session_id=record.id,
        action_id=action_id,
        sequence=(sequence or 0) + 1, checkpoint_type=checkpoint_type,
        current_url=(worker_state.current_url if worker_state else record.current_url),
        active_tab_id=(worker_state.active_tab_id if worker_state else record.active_tab_id),
        tabs=(worker_state.tabs if worker_state else []),
        state_redacted=state_redacted or {},
    ))


def _apply_worker_state(record: ComputerSession, worker: ComputerWorkerSessionState) -> None:
    record.worker_id = worker.worker_id
    record.current_url = worker.current_url
    record.active_tab_id = worker.active_tab_id
    record.last_active_at = datetime.now(UTC)


async def _persist_worker_failure(
    db_session: SessionDep,
    *,
    workspace_id: uuid.UUID,
    session_id: uuid.UUID,
    actor_id: uuid.UUID,
    request_id: str,
    error: AppError,
    status: Literal["paused", "failed"],
    event_type: str,
) -> None:
    await db_session.rollback()
    record = await _load_computer_session(db_session, workspace_id, session_id, lock=True)
    now = datetime.now(UTC)
    lease = await _active_lease(db_session, workspace_id, session_id, lock=True)
    _release_lease(lease, error.code, now)
    record.status = status
    record.control_owner = "none"
    record.failure_code = error.code
    record.failure_message = error.message
    record.paused_at = now if status == "paused" else record.paused_at
    record.stopped_at = now if status == "failed" else record.stopped_at
    record.version += 1
    await _checkpoint(
        db_session, record, "failure", state_redacted={"error_code": error.code}
    )
    await _audit(
        db_session, workspace_id=workspace_id, actor_id=actor_id,
        event_type=event_type, request_id=request_id,
        data={"computer_session_id": str(session_id), "error_code": error.code},
    )
    await db_session.commit()


@router.get("/status")
async def computer_status(
    auth: AuthDep, session: SessionDep, probe: ComputerHealthDep,
):
    await authorize(
        session, workspace_id=auth.workspace_id, user_id=auth.user_id,
        permission="computer.observe",
    )
    state = await _control_state(session, auth.workspace_id)
    return {
        "success": True,
        "service_state": await probe.status(),
        "emergency_stopped": state.emergency_stopped,
        "reason": state.reason,
        "stopped_at": state.stopped_at,
        "resumed_at": state.resumed_at,
        "version": state.version,
    }


@router.get("/profiles")
async def list_profiles(auth: AuthDep, session: SessionDep):
    await authorize(
        session, workspace_id=auth.workspace_id, user_id=auth.user_id,
        permission="computer.observe",
    )
    profiles = list((await session.scalars(select(ComputerProfile).where(
        ComputerProfile.workspace_id == auth.workspace_id
    ).order_by(ComputerProfile.name))).all())
    return {"success": True, "profiles": [_profile_data(profile) for profile in profiles]}


@router.post("/profiles", status_code=201)
async def create_profile(
    payload: ProfileInput, request: Request, auth: AuthDep, session: SessionDep,
    hayva_csrf: Annotated[str | None, Cookie()] = None,
    x_csrf_token: Annotated[str | None, Header()] = None,
):
    await authorize(
        session, workspace_id=auth.workspace_id, user_id=auth.user_id,
        permission="settings.modify",
    )
    await _csrf(session, auth, hayva_csrf, x_csrf_token)
    if await session.scalar(select(ComputerProfile.id).where(
        ComputerProfile.workspace_id == auth.workspace_id,
        ComputerProfile.profile_key == payload.profile_key,
    )):
        raise AppError("COMPUTER_PROFILE_CONFLICT", "This computer profile key exists.", 409)
    agent = await session.scalar(select(Agent).where(
        Agent.workspace_id == auth.workspace_id,
        Agent.id == payload.agent_id,
        Agent.computer_id.is_not(None),
    ))
    if not agent:
        raise AppError(
            "AGENT_DEDICATED_COMPUTER_REQUIRED",
            "The agent has no dedicated computer for this browser profile.",
            409,
        )
    if await session.scalar(select(ComputerProfile.id).where(
        ComputerProfile.workspace_id == auth.workspace_id,
        ComputerProfile.agent_id == agent.id,
    )):
        raise AppError(
            "COMPUTER_PROFILE_CONFLICT",
            "The agent already owns its private browser profile.",
            409,
        )
    values = payload.model_dump(exclude={"agent_id"})
    profile = ComputerProfile(
        id=uuid.uuid4(), workspace_id=auth.workspace_id,
        agent_id=agent.id, computer_id=agent.computer_id,
        created_by_user_id=auth.user_id, storage_key=uuid.uuid4().hex,
        status="active", version=1, **values,
    )
    session.add(profile)
    await _audit(
        session, workspace_id=auth.workspace_id, actor_id=auth.user_id,
        event_type="computer.profile_created", request_id=request.state.request_id,
        data={"computer_profile_id": str(profile.id), "profile_key": profile.profile_key},
    )
    await session.commit()
    await session.refresh(profile)
    return {"success": True, "profile": _profile_data(profile)}


@router.put("/profiles/{profile_id}")
async def update_profile(
    profile_id: uuid.UUID, payload: ProfileUpdateInput, request: Request,
    auth: AuthDep, session: SessionDep,
    hayva_csrf: Annotated[str | None, Cookie()] = None,
    x_csrf_token: Annotated[str | None, Header()] = None,
):
    await authorize(
        session, workspace_id=auth.workspace_id, user_id=auth.user_id,
        permission="settings.modify",
    )
    await _csrf(session, auth, hayva_csrf, x_csrf_token)
    profile = await _load_profile(session, auth.workspace_id, profile_id, lock=True)
    if profile.status == "revoked":
        raise AppError("COMPUTER_PROFILE_REVOKED", "A revoked profile cannot be changed.", 409)
    for key, value in payload.model_dump(exclude_unset=True).items():
        setattr(profile, key, value)
    profile.version += 1
    await _audit(
        session, workspace_id=auth.workspace_id, actor_id=auth.user_id,
        event_type="computer.profile_updated", request_id=request.state.request_id,
        data={"computer_profile_id": str(profile.id), "fields": sorted(payload.model_fields_set)},
    )
    await session.commit()
    await session.refresh(profile)
    return {"success": True, "profile": _profile_data(profile)}


@router.post("/profiles/{profile_id}/revoke")
async def revoke_profile(
    profile_id: uuid.UUID, request: Request, auth: AuthDep, session: SessionDep,
    hayva_csrf: Annotated[str | None, Cookie()] = None,
    x_csrf_token: Annotated[str | None, Header()] = None,
):
    await authorize(
        session, workspace_id=auth.workspace_id, user_id=auth.user_id,
        permission="settings.modify",
    )
    await _csrf(session, auth, hayva_csrf, x_csrf_token)
    profile = await _load_profile(session, auth.workspace_id, profile_id, lock=True)
    active_count = await session.scalar(select(func.count(ComputerSession.id)).where(
        ComputerSession.workspace_id == auth.workspace_id,
        ComputerSession.profile_id == profile.id,
        ComputerSession.status.in_(ACTIVE_STATUSES),
    ))
    if active_count:
        raise AppError(
            "COMPUTER_PROFILE_IN_USE",
            "Stop active sessions before revoking this computer profile.",
            409,
        )
    if profile.status != "revoked":
        profile.status = "revoked"
        profile.revoked_at = datetime.now(UTC)
        profile.version += 1
        await _audit(
            session, workspace_id=auth.workspace_id, actor_id=auth.user_id,
            event_type="computer.profile_revoked", request_id=request.state.request_id,
            data={"computer_profile_id": str(profile.id)},
        )
        await session.commit()
        await session.refresh(profile)
    return {"success": True, "profile": _profile_data(profile)}


@router.get("/sessions")
async def list_sessions(auth: AuthDep, session: SessionDep):
    await authorize(
        session, workspace_id=auth.workspace_id, user_id=auth.user_id,
        permission="computer.observe",
    )
    records = list((await session.scalars(select(ComputerSession).where(
        ComputerSession.workspace_id == auth.workspace_id
    ).order_by(ComputerSession.created_at.desc()).limit(100))).all())
    return {"success": True, "sessions": [_session_data(record) for record in records]}


@router.get("/sessions/{session_id}")
async def get_session(session_id: uuid.UUID, auth: AuthDep, session: SessionDep):
    await authorize(
        session, workspace_id=auth.workspace_id, user_id=auth.user_id,
        permission="computer.observe",
    )
    record = await _load_computer_session(session, auth.workspace_id, session_id)
    checkpoints = list((await session.scalars(select(ComputerCheckpoint).where(
        ComputerCheckpoint.workspace_id == auth.workspace_id,
        ComputerCheckpoint.session_id == session_id,
    ).order_by(ComputerCheckpoint.sequence.desc()).limit(50))).all())
    actions = list((await session.scalars(select(ComputerAction).where(
        ComputerAction.workspace_id == auth.workspace_id,
        ComputerAction.session_id == session_id,
    ).order_by(ComputerAction.created_at.desc()).limit(100))).all())
    artifacts = list((await session.scalars(select(ComputerArtifact).where(
        ComputerArtifact.workspace_id == auth.workspace_id,
        ComputerArtifact.session_id == session_id,
    ).order_by(ComputerArtifact.created_at.desc()).limit(100))).all())
    lease = await _active_lease(session, auth.workspace_id, session_id)
    return {
        "success": True,
        "session": _session_data(record),
        "control_lease": ({
            "owner": lease.owner, "expires_at": lease.expires_at,
            "heartbeat_at": lease.heartbeat_at,
        } if lease else None),
        "checkpoints": [{
            "id": str(item.id), "sequence": item.sequence,
            "checkpoint_type": item.checkpoint_type, "current_url": item.current_url,
            "active_tab_id": item.active_tab_id, "tabs": item.tabs,
            "state": item.state_redacted, "created_at": item.created_at,
        } for item in checkpoints],
        "actions": [{
            "id": str(item.id), "actor_type": item.actor_type,
            "tool": item.tool_name, "tool_name": item.tool_name, "risk": item.risk,
            "access_type": item.access_type, "status": item.status,
            "result": item.result_redacted, "verification": item.verification,
            "error_code": item.error_code, "created_at": item.created_at,
        } for item in actions],
        "artifacts": [{
            "id": str(item.id), "kind": item.kind, "status": item.status,
            "mime_type": item.mime_type, "size_bytes": item.size_bytes,
            "metadata": item.metadata_redacted, "created_at": item.created_at,
        } for item in artifacts],
    }


@router.post("/sessions", status_code=201)
async def start_session(
    payload: SessionInput, request: Request, auth: AuthDep, session: SessionDep,
    client: ComputerClientDep,
    hayva_csrf: Annotated[str | None, Cookie()] = None,
    x_csrf_token: Annotated[str | None, Header()] = None,
):
    await authorize(
        session, workspace_id=auth.workspace_id, user_id=auth.user_id,
        permission="computer.control",
    )
    await _csrf(session, auth, hayva_csrf, x_csrf_token)
    state = await _control_state(session, auth.workspace_id, lock=True)
    if state.emergency_stopped:
        raise AppError("COMPUTER_EMERGENCY_STOPPED", "Computer activity is emergency-stopped.", 409)
    profile = await _load_profile(session, auth.workspace_id, payload.profile_id)
    if profile.status != "active":
        raise AppError("COMPUTER_PROFILE_INACTIVE", "An active computer profile is required.", 409)
    agent_statement = select(Agent).where(
        Agent.workspace_id == auth.workspace_id, Agent.status == "active",
    )
    if payload.agent_id:
        agent_statement = agent_statement.where(Agent.id == payload.agent_id)
    else:
        agent_statement = agent_statement.where(Agent.name == "Personal Executive Assistant")
    agent = await session.scalar(agent_statement)
    if not agent:
        raise AppError("AGENT_UNAVAILABLE", "An active computer agent is required.", 409)
    if profile.agent_id != agent.id or profile.computer_id != agent.computer_id:
        raise AppError(
            "COMPUTER_PROFILE_AGENT_MISMATCH",
            "An agent may use only the browser profile inside its own computer.",
            403,
        )
    dedicated = await session.scalar(select(Computer).where(
        Computer.workspace_id == auth.workspace_id,
        Computer.id == agent.computer_id,
    ))
    if not dedicated or dedicated.status not in {"running", "idle"}:
        raise AppError(
            "COMPUTER_NOT_RUNNING",
            "Start the agent's dedicated computer before opening its browser runtime.",
            409,
        )
    if payload.execution_id and not await session.scalar(select(Execution.id).where(
        Execution.workspace_id == auth.workspace_id,
        Execution.id == payload.execution_id,
    )):
        raise AppError("EXECUTION_NOT_FOUND", "The execution is unavailable.", 404)
    now = datetime.now(UTC)
    record = ComputerSession(
        id=uuid.uuid4(), workspace_id=auth.workspace_id,
        requested_by_user_id=auth.user_id, agent_id=agent.id,
        execution_id=payload.execution_id, profile_id=profile.id,
        task_summary=payload.task_summary, status="starting", control_owner="system",
        lease_version=0, version=1,
    )
    session.add(record)
    lease = _issue_lease(session, record, owner="ai", actor_user_id=None, now=now)
    await _audit(
        session, workspace_id=auth.workspace_id, actor_id=auth.user_id,
        event_type="computer.session_start_requested", request_id=request.state.request_id,
        data={"computer_session_id": str(record.id), "profile_id": str(profile.id)},
    )
    await session.commit()
    try:
        worker = await client.start_session(
            workspace_id=auth.workspace_id, session_id=record.id, profile_id=profile.id,
            profile_storage_key=profile.storage_key, fencing_token=lease.fencing_token,
            request_id=request.state.request_id,
        )
    except AppError as error:
        await _persist_worker_failure(
            session, workspace_id=auth.workspace_id, session_id=record.id,
            actor_id=auth.user_id, request_id=request.state.request_id, error=error,
            status="failed", event_type="computer.session_start_failed",
        )
        raise
    if worker.status not in {"ready", "ai_controlled"}:
        error = AppError(
            "COMPUTER_START_STATE_INVALID",
            "The browser worker did not start into an active state.",
            502,
        )
        await _persist_worker_failure(
            session, workspace_id=auth.workspace_id, session_id=record.id,
            actor_id=auth.user_id, request_id=request.state.request_id, error=error,
            status="failed", event_type="computer.session_start_failed",
        )
        raise error
    record = await _load_computer_session(session, auth.workspace_id, record.id, lock=True)
    state = await _control_state(session, auth.workspace_id, lock=True)
    if state.emergency_stopped or record.status != "starting":
        await session.rollback()
        raise AppError("COMPUTER_START_INTERRUPTED", "Computer startup was interrupted.", 409)
    _apply_worker_state(record, worker)
    record.status = "ai_controlled"
    record.control_owner = "ai"
    record.started_at = now
    record.failure_code = None
    record.failure_message = None
    record.version += 1
    await _checkpoint(session, record, "session_started", worker_state=worker)
    await _audit(
        session, workspace_id=auth.workspace_id, actor_id=auth.user_id,
        event_type="computer.session_started", request_id=request.state.request_id,
        data={"computer_session_id": str(record.id), "worker_id": worker.worker_id},
    )
    await session.commit()
    await session.refresh(record)
    return {"success": True, "session": _session_data(record)}


def _validate_browser_arguments(tool: str, arguments: dict) -> None:
    rules: dict[str, tuple[set[str], set[str]]] = {
        "observe": (set(), set()),
        "navigate": ({"url"}, {"url"}),
        "click": ({"selector"}, {"selector"}),
        "type": ({"selector", "text"}, {"selector", "text"}),
        "press": ({"key", "selector"}, {"key"}),
        "wait_for": ({"selector", "state", "timeout_ms"}, {"selector"}),
        "extract": ({"selector", "attribute"}, {"selector"}),
        "tab_list": (set(), set()),
        "tab_open": ({"url"}, set()),
        "tab_switch": ({"tab_id"}, {"tab_id"}),
        "tab_close": ({"tab_id"}, {"tab_id"}),
        "screenshot": (set(), set()),
    }
    allowed, required = rules[tool]
    if set(arguments) - allowed or required - set(arguments):
        raise AppError("COMPUTER_ACTION_ARGUMENTS_INVALID", "Action arguments are invalid.", 422)
    string_limits = {
        "url": 4096, "selector": 1000, "text": 20_000,
        "key": 80, "state": 20, "attribute": 80, "tab_id": 160,
    }
    for key, limit in string_limits.items():
        if key in arguments and (
            not isinstance(arguments[key], str)
            or (key != "text" and not arguments[key].strip())
            or len(arguments[key]) > limit
        ):
            raise AppError(
                "COMPUTER_ACTION_ARGUMENTS_INVALID", "Action arguments are invalid.", 422
            )
    if "timeout_ms" in arguments and (
        not isinstance(arguments["timeout_ms"], int)
        or not 100 <= arguments["timeout_ms"] <= 30_000
    ):
        raise AppError("COMPUTER_ACTION_ARGUMENTS_INVALID", "Action arguments are invalid.", 422)


def _redact_browser_arguments(arguments: dict) -> dict:
    redacted: dict = {}
    for key, value in arguments.items():
        if key == "text":
            redacted[key] = {
                "redacted": True,
                "character_count": len(value) if isinstance(value, str) else None,
            }
        elif key == "url" and isinstance(value, str):
            redacted[key] = value.split("?", 1)[0].split("#", 1)[0][:4096]
        else:
            redacted[key] = value
    return redacted


def _redact_browser_result(tool: str, result: dict) -> dict:
    if tool == "observe":
        return {
            "url": result.get("url"),
            "title": result.get("title"),
            "interactive_element_count": len(result.get("interactive_elements", [])),
        }
    if tool == "extract":
        value = result.get("value")
        return {
            "value_redacted": True,
            "value_length": len(value) if isinstance(value, str) else None,
        }
    if tool == "screenshot":
        artifact = result.get("artifact") or {}
        return {"artifact_id": artifact.get("id")}
    return result


def _action_data(action: ComputerAction, approval: ApprovalRequest | None = None) -> dict:
    return {
        "id": str(action.id),
        "session_id": str(action.session_id),
        "tool": action.tool_name,
        "risk": action.risk,
        "access_type": action.access_type,
        "status": action.status,
        "parameters": action.parameters_redacted,
        "result": action.result_redacted,
        "verification": action.verification,
        "error_code": action.error_code,
        "duration_ms": action.duration_ms,
        "approval": ({
            "id": str(approval.id), "status": approval.status,
            "expires_at": approval.expires_at,
        } if approval else None),
        "created_at": action.created_at,
        "completed_at": action.completed_at,
    }


async def _authorize_browser_action(
    db_session: SessionDep,
    *,
    auth: AuthDep,
    record: ComputerSession,
    tool: str,
) -> Decision:
    definition = COMPUTER_TOOLS[tool]
    granted = await permissions_for(
        db_session, workspace_id=auth.workspace_id, user_id=auth.user_id
    )
    if definition.permission not in granted or "computer.control" not in granted:
        raise AppError("PERMISSION_DENIED", "You do not have permission for this action.", 403)
    agent = await db_session.scalar(select(Agent).where(
        Agent.workspace_id == auth.workspace_id,
        Agent.id == record.agent_id,
    ))
    if not agent or agent.status != "active" or agent.autonomy_level == 0:
        raise AppError("AGENT_DISABLED", "The computer agent is disabled.", 409)
    tool_granted = await db_session.scalar(select(AgentTool.tool_name).where(
        AgentTool.workspace_id == auth.workspace_id,
        AgentTool.agent_id == agent.id,
        AgentTool.tool_name == f"browser.{tool}",
    ))
    permission_granted = await db_session.scalar(select(AgentPermission.permission_key).where(
        AgentPermission.workspace_id == auth.workspace_id,
        AgentPermission.agent_id == agent.id,
        AgentPermission.permission_key == definition.permission,
    ))
    if not tool_granted or not permission_granted:
        raise AppError("AGENT_TOOL_REVOKED", "The browser tool grant is not active.", 403)
    decision = evaluate_action(
        ActionRequest(
            permission=definition.permission,
            risk=definition.risk,
            is_write=definition.is_write,
        ),
        AutonomyLevel(agent.autonomy_level),
        granted,
    )
    if decision == Decision.DENY:
        raise AppError("COMPUTER_ACTION_DENIED", "The browser action was denied.", 403)
    if decision == Decision.DRAFT_ONLY:
        return Decision.REQUIRE_APPROVAL
    return decision


async def _record_worker_artifact(
    db_session: SessionDep,
    *,
    record: ComputerSession,
    action: ComputerAction,
    result: dict,
) -> uuid.UUID | None:
    raw = result.get("artifact")
    if not isinstance(raw, dict):
        return None
    try:
        artifact_id = uuid.UUID(str(raw["id"]))
        storage_key = str(raw["storage_key"])
        sha256 = str(raw["sha256"])
        mime_type = str(raw["mime_type"])
        size_bytes = int(raw["size_bytes"])
        metadata = raw.get("metadata") if isinstance(raw.get("metadata"), dict) else {}
    except (KeyError, TypeError, ValueError) as error:
        raise AppError(
            "COMPUTER_ARTIFACT_INVALID", "The browser artifact metadata is invalid.", 502
        ) from error
    expected_prefix = f"{record.workspace_id}/{record.id}/"
    if (
        not storage_key.startswith(expected_prefix)
        or ".." in storage_key
        or len(storage_key) > 255
        or not re.fullmatch(r"[a-f0-9]{64}", sha256)
        or mime_type != "image/jpeg"
        or not 0 <= size_bytes <= 20_000_000
    ):
        raise AppError(
            "COMPUTER_ARTIFACT_INVALID", "The browser artifact metadata is invalid.", 502
        )
    retention_days = await db_session.scalar(select(ComputerProfile.retention_days).where(
        ComputerProfile.workspace_id == record.workspace_id,
        ComputerProfile.id == record.profile_id,
    ))
    if retention_days is None:
        raise AppError(
            "COMPUTER_PROFILE_NOT_FOUND", "The browser profile is unavailable.", 404
        )
    db_session.add(ComputerArtifact(
        id=artifact_id, workspace_id=record.workspace_id, session_id=record.id,
        action_id=action.id, kind="screenshot", status="available",
        storage_key=storage_key, sha256=sha256, mime_type=mime_type,
        size_bytes=size_bytes, metadata_redacted=metadata,
        expires_at=datetime.now(UTC) + timedelta(days=retention_days),
    ))
    return artifact_id


async def _execute_browser_action(
    db_session: SessionDep,
    *,
    auth: AuthDep,
    action_id: uuid.UUID,
    arguments: dict,
    client: ComputerClient,
) -> tuple[ComputerAction, ComputerWorkerActionResult]:
    action = await db_session.scalar(select(ComputerAction).where(
        ComputerAction.workspace_id == auth.workspace_id,
        ComputerAction.id == action_id,
    ).with_for_update())
    if not action:
        raise AppError("COMPUTER_ACTION_NOT_FOUND", "The browser action is unavailable.", 404)
    record = await _load_computer_session(
        db_session, auth.workspace_id, action.session_id, lock=True
    )
    state = await _control_state(db_session, auth.workspace_id, lock=True)
    lease = await _active_lease(db_session, auth.workspace_id, record.id, lock=True)
    if state.emergency_stopped:
        raise AppError("COMPUTER_EMERGENCY_STOPPED", "Computer activity is emergency-stopped.", 409)
    if record.status != "ai_controlled" or not lease or lease.owner != "ai":
        raise AppError("COMPUTER_AI_CONTROL_INACTIVE", "AI control is not active.", 409)
    _heartbeat_lease(lease)
    if action.status not in {"pending", "awaiting_approval"}:
        raise AppError("COMPUTER_ACTION_NOT_RUNNABLE", "The browser action cannot run.", 409)
    action.status = "running"
    action.started_at = datetime.now(UTC)
    record.current_action_summary = f"Browser {action.tool_name}"
    record.version += 1
    await _checkpoint(
        db_session, record, "before_action", action_id=action.id,
        state_redacted={"tool": action.tool_name},
    )
    await _audit(
        db_session, workspace_id=auth.workspace_id, actor_id=auth.user_id,
        event_type="computer.action_started", request_id=str(action.id),
        data={"computer_session_id": str(record.id), "action_id": str(action.id),
              "tool": action.tool_name},
    )
    await db_session.commit()
    started = time.perf_counter()
    try:
        worker = await client.action(
            session_id=record.id, fencing_token=lease.fencing_token,
            tool=action.tool_name, arguments=arguments, request_id=str(action.id),
        )
    except AppError as error:
        await db_session.rollback()
        action = await db_session.scalar(select(ComputerAction).where(
            ComputerAction.workspace_id == auth.workspace_id,
            ComputerAction.id == action_id,
        ).with_for_update())
        if action:
            outcome_unknown = (
                action.access_type == "write"
                and error.code in {
                    "COMPUTER_SERVICE_UNAVAILABLE",
                    "COMPUTER_SERVICE_FAILED",
                    "COMPUTER_ACTION_TIMEOUT",
                    "COMPUTER_SERVICE_RESPONSE_INVALID",
                }
            )
            action.status = "unknown" if outcome_unknown else "failed"
            action.error_code = (
                "COMPUTER_ACTION_OUTCOME_UNKNOWN" if outcome_unknown else error.code
            )
            action.duration_ms = int((time.perf_counter() - started) * 1000)
            action.completed_at = datetime.now(UTC)
            record = await _load_computer_session(
                db_session, auth.workspace_id, action.session_id, lock=True
            )
            record.current_action_summary = None
            record.failure_code = action.error_code
            record.failure_message = (
                "The browser action outcome must be reconciled by a human."
                if outcome_unknown else error.message
            )
            if outcome_unknown:
                active = await _active_lease(
                    db_session, auth.workspace_id, record.id, lock=True
                )
                _release_lease(active, "outcome_unknown", datetime.now(UTC))
                record.status = "paused"
                record.control_owner = "none"
                record.paused_at = datetime.now(UTC)
            record.version += 1
            await _audit(
                db_session, workspace_id=auth.workspace_id, actor_id=auth.user_id,
                event_type=(
                    "computer.action_outcome_unknown"
                    if outcome_unknown else "computer.action_failed"
                ),
                request_id=str(action.id),
                data={"computer_session_id": str(record.id), "action_id": str(action.id),
                      "error_code": error.code},
            )
            await db_session.commit()
        if action and action.status == "unknown":
            raise AppError(
                "COMPUTER_ACTION_OUTCOME_UNKNOWN",
                "The browser action outcome is unknown and requires human reconciliation.",
                409,
            ) from error
        raise
    action = await db_session.scalar(select(ComputerAction).where(
        ComputerAction.workspace_id == auth.workspace_id,
        ComputerAction.id == action_id,
    ).with_for_update())
    if not action:
        raise AppError("COMPUTER_ACTION_NOT_FOUND", "The browser action is unavailable.", 404)
    record = await _load_computer_session(
        db_session, auth.workspace_id, action.session_id, lock=True
    )
    current_lease = await _active_lease(db_session, auth.workspace_id, record.id, lock=True)
    if not current_lease or current_lease.fencing_token != lease.fencing_token:
        action.status = "unknown"
        action.error_code = "COMPUTER_LEASE_CHANGED"
        action.completed_at = datetime.now(UTC)
        record.current_action_summary = None
        await db_session.commit()
        raise AppError(
            "COMPUTER_ACTION_OUTCOME_UNKNOWN",
            "Control changed before the browser action could be verified.",
            409,
        )
    screenshot_artifact_id = await _record_worker_artifact(
        db_session, record=record, action=action, result=worker.result
    )
    action.status = "succeeded"
    action.result_redacted = _redact_browser_result(action.tool_name, worker.result)
    action.verification = worker.verification
    action.duration_ms = int((time.perf_counter() - started) * 1000)
    action.completed_at = datetime.now(UTC)
    action.error_code = None
    _apply_worker_state(record, worker.session)
    record.current_action_summary = None
    record.failure_code = None
    record.failure_message = None
    record.version += 1
    await _checkpoint(
        db_session, record, "after_action", action_id=action.id, worker_state=worker.session,
        state_redacted={"tool": action.tool_name, "verified": True},
    )
    if screenshot_artifact_id:
        checkpoint = await db_session.scalar(select(ComputerCheckpoint).where(
            ComputerCheckpoint.workspace_id == auth.workspace_id,
            ComputerCheckpoint.session_id == record.id,
        ).order_by(ComputerCheckpoint.sequence.desc()).limit(1))
        if checkpoint:
            checkpoint.screenshot_artifact_id = screenshot_artifact_id
    await _audit(
        db_session, workspace_id=auth.workspace_id, actor_id=auth.user_id,
        event_type="computer.action_succeeded", request_id=str(action.id),
        data={"computer_session_id": str(record.id), "action_id": str(action.id),
              "tool": action.tool_name, "verified": True},
    )
    await db_session.commit()
    await db_session.refresh(action)
    return action, worker


@router.post("/sessions/{session_id}/actions")
async def create_browser_action(
    session_id: uuid.UUID, payload: BrowserActionRequest, request: Request,
    response: Response, auth: AuthDep, session: SessionDep, client: ComputerClientDep,
    hayva_csrf: Annotated[str | None, Cookie()] = None,
    x_csrf_token: Annotated[str | None, Header()] = None,
    x_idempotency_key: Annotated[str | None, Header()] = None,
):
    await authorize(session, workspace_id=auth.workspace_id, user_id=auth.user_id,
                    permission="computer.control")
    await _csrf(session, auth, hayva_csrf, x_csrf_token)
    _validate_browser_arguments(payload.tool, payload.arguments)
    if not x_idempotency_key or not re.fullmatch(
        r"[A-Za-z0-9._:-]{8,120}", x_idempotency_key
    ):
        raise AppError(
            "IDEMPOTENCY_KEY_REQUIRED",
            "A valid idempotency key is required for browser actions.",
            422,
        )
    intent = {
        "session_id": str(session_id), "tool": payload.tool,
        "arguments": payload.arguments,
    }
    intent_hash = canonical_payload_hash(intent)
    idempotency_key = f"computer-action:{session_id}:{x_idempotency_key}"
    existing_idempotency = await session.scalar(select(IdempotencyRecord).where(
        IdempotencyRecord.workspace_id == auth.workspace_id,
        IdempotencyRecord.key == idempotency_key,
    ).with_for_update())
    if existing_idempotency:
        if existing_idempotency.request_hash != intent_hash:
            raise AppError(
                "IDEMPOTENCY_KEY_CONFLICT",
                "The idempotency key was already used for another request.",
                409,
            )
        stored = existing_idempotency.response_body or {}
        if existing_idempotency.status == "completed" and isinstance(stored.get("response"), dict):
            response.status_code = existing_idempotency.response_status or 200
            return stored["response"]
        raise AppError(
            "IDEMPOTENT_REQUEST_IN_PROGRESS",
            "The browser action request is already in progress.",
            409,
        )
    record = await _load_computer_session(session, auth.workspace_id, session_id, lock=True)
    if record.status != "ai_controlled":
        raise AppError("COMPUTER_AI_CONTROL_INACTIVE", "AI control is not active.", 409)
    lease = await _active_lease(session, auth.workspace_id, session_id, lock=True)
    if not lease or lease.owner != "ai":
        raise AppError("COMPUTER_LEASE_INVALID", "The AI control lease is unavailable.", 409)
    _heartbeat_lease(lease)
    decision = await _authorize_browser_action(
        session, auth=auth, record=record, tool=payload.tool
    )
    definition = COMPUTER_TOOLS[payload.tool]
    action = ComputerAction(
        id=uuid.uuid4(), workspace_id=auth.workspace_id, session_id=session_id,
        execution_id=record.execution_id, actor_type="agent", actor_id=record.agent_id,
        tool_name=payload.tool, parameters_redacted=_redact_browser_arguments(payload.arguments),
        risk=definition.risk.value, access_type=definition.access_type,
        action_intent_hash=intent_hash,
        status="awaiting_approval" if decision == Decision.REQUIRE_APPROVAL else "pending",
    )
    session.add(action)
    idempotency = IdempotencyRecord(
        id=uuid.uuid4(), workspace_id=auth.workspace_id, key=idempotency_key,
        request_hash=intent_hash, status="in_progress",
        response_body={"action_id": str(action.id)},
        expires_at=datetime.now(UTC) + timedelta(hours=24),
    )
    session.add(idempotency)
    if decision == Decision.REQUIRE_APPROVAL:
        approval = ApprovalRequest(
            id=uuid.uuid4(), workspace_id=auth.workspace_id, action_id=action.id,
            requested_by_type="agent", requested_by_id=record.agent_id,
            permission=definition.permission, risk=definition.risk.value,
            canonical_payload_hash=action.action_intent_hash,
            payload={
                "session_id": str(session_id), "tool": payload.tool,
                "arguments": action.parameters_redacted,
            },
            status="pending", expires_at=datetime.now(UTC) + timedelta(minutes=4),
        )
        session.add(approval)
        await _audit(
            session, workspace_id=auth.workspace_id, actor_id=auth.user_id,
            event_type="computer.action_approval_requested", request_id=request.state.request_id,
            data={"computer_session_id": str(session_id), "action_id": str(action.id),
                  "approval_id": str(approval.id), "tool": payload.tool},
        )
        await session.commit()
        await session.refresh(action)
        await session.refresh(approval)
        body = {
            "success": True, "approval_required": True,
            "action": _action_data(action, approval),
        }
        idempotency = await session.scalar(select(IdempotencyRecord).where(
            IdempotencyRecord.id == idempotency.id,
            IdempotencyRecord.workspace_id == auth.workspace_id,
        ).with_for_update())
        if idempotency:
            idempotency.status = "completed"
            idempotency.response_status = 202
            idempotency.response_body = {"response": jsonable_encoder(body)}
            await session.commit()
        response.status_code = 202
        return body
    await session.commit()
    try:
        action, worker = await _execute_browser_action(
            session, auth=auth, action_id=action.id, arguments=payload.arguments, client=client
        )
    except AppError as error:
        idempotency = await session.scalar(select(IdempotencyRecord).where(
            IdempotencyRecord.workspace_id == auth.workspace_id,
            IdempotencyRecord.key == idempotency_key,
        ).with_for_update())
        if idempotency:
            idempotency.status = "failed"
            idempotency.response_status = error.status_code
            idempotency.response_body = {
                "action_id": str(action.id), "error_code": error.code
            }
            await session.commit()
        raise
    body = {
        "success": True, "approval_required": False,
        "action": _action_data(action), "result": worker.result,
    }
    idempotency = await session.scalar(select(IdempotencyRecord).where(
        IdempotencyRecord.workspace_id == auth.workspace_id,
        IdempotencyRecord.key == idempotency_key,
    ).with_for_update())
    if idempotency:
        idempotency.status = "completed"
        idempotency.response_status = 200
        idempotency.response_body = {"response": jsonable_encoder(body)}
        await session.commit()
    return body


@router.get("/approvals")
async def list_computer_approvals(auth: AuthDep, session: SessionDep):
    await authorize(session, workspace_id=auth.workspace_id, user_id=auth.user_id,
                    permission="approvals.approve")
    approvals = list((await session.scalars(select(ApprovalRequest).where(
        ApprovalRequest.workspace_id == auth.workspace_id,
        ApprovalRequest.status == "pending",
    ).order_by(ApprovalRequest.created_at.desc()).limit(100))).all())
    return {"success": True, "approvals": [{
        "id": str(item.id), "action_id": str(item.action_id),
        "permission": item.permission, "risk": item.risk,
        "payload": item.payload, "status": item.status,
        "expires_at": item.expires_at, "created_at": item.created_at,
    } for item in approvals]}


@router.post("/actions/{action_id}/decision")
async def decide_computer_action(
    action_id: uuid.UUID, payload: ActionDecisionInput, request: Request,
    auth: AuthDep, session: SessionDep, client: ComputerClientDep,
    hayva_csrf: Annotated[str | None, Cookie()] = None,
    x_csrf_token: Annotated[str | None, Header()] = None,
):
    await authorize(session, workspace_id=auth.workspace_id, user_id=auth.user_id,
                    permission="approvals.approve")
    await _csrf(session, auth, hayva_csrf, x_csrf_token)
    approval = await session.scalar(select(ApprovalRequest).where(
        ApprovalRequest.workspace_id == auth.workspace_id,
        ApprovalRequest.action_id == action_id,
    ).with_for_update())
    action = await session.scalar(select(ComputerAction).where(
        ComputerAction.workspace_id == auth.workspace_id,
        ComputerAction.id == action_id,
    ).with_for_update())
    if not approval or not action:
        raise AppError("COMPUTER_APPROVAL_NOT_FOUND", "The approval is unavailable.", 404)
    if approval.status != "pending" or action.status != "awaiting_approval":
        raise AppError("COMPUTER_APPROVAL_DECIDED", "The approval has already been decided.", 409)
    now = datetime.now(UTC)
    expiry = approval.expires_at
    if expiry.tzinfo is None:
        expiry = expiry.replace(tzinfo=UTC)
    if now >= expiry:
        approval.status = "expired"
        approval.decided_at = now
        action.status = "blocked"
        action.error_code = "COMPUTER_APPROVAL_EXPIRED"
        action.completed_at = now
        await session.commit()
        raise AppError("COMPUTER_APPROVAL_EXPIRED", "The approval has expired.", 409)
    if not payload.approve:
        approval.status = "rejected"
        approval.decided_at = now
        approval.decided_by_user_id = auth.user_id
        action.status = "blocked"
        action.error_code = "COMPUTER_ACTION_REJECTED"
        action.completed_at = now
        await _audit(
            session, workspace_id=auth.workspace_id, actor_id=auth.user_id,
            event_type="computer.action_rejected", request_id=request.state.request_id,
            data={"action_id": str(action.id), "computer_session_id": str(action.session_id)},
        )
        await session.commit()
        await session.refresh(action)
        return {"success": True, "action": _action_data(action, approval)}
    _validate_browser_arguments(action.tool_name, payload.arguments)
    intent = {
        "session_id": str(action.session_id), "tool": action.tool_name,
        "arguments": payload.arguments,
    }
    if not payload_matches(intent, approval.canonical_payload_hash):
        raise AppError(
            "COMPUTER_APPROVAL_PAYLOAD_MISMATCH",
            "The approved action payload does not match the request.",
            409,
        )
    approval.status = "approved"
    approval.decided_at = now
    approval.decided_by_user_id = auth.user_id
    await _audit(
        session, workspace_id=auth.workspace_id, actor_id=auth.user_id,
        event_type="computer.action_approved", request_id=request.state.request_id,
        data={"action_id": str(action.id), "computer_session_id": str(action.session_id)},
    )
    await session.commit()
    action, worker = await _execute_browser_action(
        session, auth=auth, action_id=action.id, arguments=payload.arguments, client=client
    )
    return {"success": True, "action": _action_data(action, approval), "result": worker.result}


@router.post("/actions/{action_id}/reconcile")
async def reconcile_computer_action(
    action_id: uuid.UUID, payload: ActionReconciliationInput, request: Request,
    auth: AuthDep, session: SessionDep,
    hayva_csrf: Annotated[str | None, Cookie()] = None,
    x_csrf_token: Annotated[str | None, Header()] = None,
):
    await authorize(
        session, workspace_id=auth.workspace_id, user_id=auth.user_id,
        permission="computer.control",
    )
    await authorize(
        session, workspace_id=auth.workspace_id, user_id=auth.user_id,
        permission="approvals.approve",
    )
    await _csrf(session, auth, hayva_csrf, x_csrf_token)
    action = await session.scalar(select(ComputerAction).where(
        ComputerAction.workspace_id == auth.workspace_id,
        ComputerAction.id == action_id,
    ).with_for_update())
    if not action:
        raise AppError("COMPUTER_ACTION_NOT_FOUND", "The browser action is unavailable.", 404)
    if action.status != "unknown":
        raise AppError(
            "COMPUTER_ACTION_RECONCILIATION_INVALID",
            "Only an action with an unknown outcome can be reconciled.",
            409,
        )
    now = datetime.now(UTC)
    action.status = payload.outcome
    action.verification = {
        "status": "human_reconciled",
        "outcome": payload.outcome,
        "evidence_summary": payload.evidence_summary,
        "reconciled_by_user_id": str(auth.user_id),
        "reconciled_at": now.isoformat(),
    }
    action.error_code = None if payload.outcome == "succeeded" else "HUMAN_RECONCILED_FAILED"
    action.completed_at = now
    record = await _load_computer_session(
        session, auth.workspace_id, action.session_id, lock=True
    )
    record.failure_code = None if payload.outcome == "succeeded" else action.error_code
    record.failure_message = None if payload.outcome == "succeeded" else (
        "A human reconciled the browser action as failed."
    )
    record.version += 1
    await _checkpoint(
        session, record, "recovery", action_id=action.id,
        state_redacted={
            "outcome": payload.outcome,
            "evidence_summary": payload.evidence_summary,
        },
    )
    await _audit(
        session, workspace_id=auth.workspace_id, actor_id=auth.user_id,
        event_type="computer.action_reconciled", request_id=request.state.request_id,
        data={
            "computer_session_id": str(record.id), "action_id": str(action.id),
            "outcome": payload.outcome,
        },
    )
    await session.commit()
    await session.refresh(action)
    return {"success": True, "action": _action_data(action)}


async def _human_control_context(
    db_session: SessionDep, *, auth: AuthDep, session_id: uuid.UUID,
) -> tuple[ComputerSession, ComputerControlLease]:
    state = await _control_state(db_session, auth.workspace_id, lock=True)
    if state.emergency_stopped:
        raise AppError("COMPUTER_EMERGENCY_STOPPED", "Computer activity is emergency-stopped.", 409)
    record = await _load_computer_session(
        db_session, auth.workspace_id, session_id, lock=True
    )
    lease = await _active_lease(db_session, auth.workspace_id, session_id, lock=True)
    if record.status != "human_controlled" or not lease or lease.owner != "human":
        raise AppError("COMPUTER_HUMAN_CONTROL_INACTIVE", "Human control is not active.", 409)
    if lease.actor_user_id != auth.user_id:
        raise AppError("COMPUTER_LEASE_OWNED_BY_ANOTHER_USER", "Another user holds control.", 403)
    _heartbeat_lease(lease)
    return record, lease


async def _execute_human_input(
    db_session: SessionDep,
    *,
    auth: AuthDep,
    request: Request,
    session_id: uuid.UUID,
    tool_name: str,
    parameters_redacted: dict,
    worker_call,
) -> ComputerSession:
    record, lease = await _human_control_context(
        db_session, auth=auth, session_id=session_id
    )
    action = ComputerAction(
        id=uuid.uuid4(), workspace_id=auth.workspace_id, session_id=session_id,
        execution_id=record.execution_id, actor_type="user", actor_id=auth.user_id,
        tool_name=tool_name, parameters_redacted=parameters_redacted,
        risk="medium" if tool_name != "human.pointer.move" else "low",
        access_type="control", action_intent_hash=None,
        status="running", started_at=datetime.now(UTC),
    )
    db_session.add(action)
    await db_session.commit()
    started = time.perf_counter()
    try:
        worker = await worker_call(lease.fencing_token, str(action.id))
    except AppError as error:
        await db_session.rollback()
        action = await db_session.scalar(select(ComputerAction).where(
            ComputerAction.workspace_id == auth.workspace_id,
            ComputerAction.id == action.id,
        ).with_for_update())
        if action:
            action.status = "failed"
            action.error_code = error.code
            action.duration_ms = int((time.perf_counter() - started) * 1000)
            action.completed_at = datetime.now(UTC)
            await _audit(
                db_session, workspace_id=auth.workspace_id, actor_id=auth.user_id,
                event_type="computer.human_input_failed", request_id=request.state.request_id,
                data={"computer_session_id": str(session_id), "action_id": str(action.id),
                      "error_code": error.code},
            )
            await db_session.commit()
        raise
    action = await db_session.scalar(select(ComputerAction).where(
        ComputerAction.workspace_id == auth.workspace_id,
        ComputerAction.id == action.id,
    ).with_for_update())
    record = await _load_computer_session(
        db_session, auth.workspace_id, session_id, lock=True
    )
    current_lease = await _active_lease(
        db_session, auth.workspace_id, session_id, lock=True
    )
    if not action or not current_lease or current_lease.fencing_token != lease.fencing_token:
        raise AppError("COMPUTER_LEASE_CHANGED", "The human control lease changed.", 409)
    action.status = "succeeded"
    action.result_redacted = {"accepted": True}
    action.verification = {
        "status": "observed", "url": worker.current_url,
        "observed_at": worker.observed_at,
    }
    action.duration_ms = int((time.perf_counter() - started) * 1000)
    action.completed_at = datetime.now(UTC)
    _apply_worker_state(record, worker)
    record.version += 1
    await _audit(
        db_session, workspace_id=auth.workspace_id, actor_id=auth.user_id,
        event_type="computer.human_input", request_id=request.state.request_id,
        data={"computer_session_id": str(session_id), "action_id": str(action.id),
              "tool": tool_name},
    )
    await db_session.commit()
    await db_session.refresh(record)
    return record


@router.post("/sessions/{session_id}/human/pointer")
async def human_pointer(
    session_id: uuid.UUID, payload: HumanPointerRequest, request: Request,
    auth: AuthDep, session: SessionDep, client: ComputerClientDep,
    hayva_csrf: Annotated[str | None, Cookie()] = None,
    x_csrf_token: Annotated[str | None, Header()] = None,
):
    await authorize(session, workspace_id=auth.workspace_id, user_id=auth.user_id,
                    permission="computer.control")
    await _csrf(session, auth, hayva_csrf, x_csrf_token)
    raw_payload = payload.model_dump()

    async def call_worker(token: int, call_id: str):
        return await client.human_pointer(
            session_id=session_id, fencing_token=token,
            payload=raw_payload, request_id=call_id,
        )

    record = await _execute_human_input(
        session, auth=auth, request=request, session_id=session_id,
        tool_name=f"human.pointer.{payload.action}",
        parameters_redacted=raw_payload, worker_call=call_worker,
    )
    return {"success": True, "session": _session_data(record)}


@router.post("/sessions/{session_id}/human/keyboard")
async def human_keyboard(
    session_id: uuid.UUID, payload: HumanKeyboardRequest, request: Request,
    auth: AuthDep, session: SessionDep, client: ComputerClientDep,
    hayva_csrf: Annotated[str | None, Cookie()] = None,
    x_csrf_token: Annotated[str | None, Header()] = None,
):
    await authorize(session, workspace_id=auth.workspace_id, user_id=auth.user_id,
                    permission="computer.control")
    await _csrf(session, auth, hayva_csrf, x_csrf_token)
    raw_payload = payload.model_dump(exclude_none=True)
    redacted = (
        {"text": {"redacted": True, "character_count": len(payload.text)}}
        if payload.text is not None else {"key": payload.key}
    )

    async def call_worker(token: int, call_id: str):
        return await client.human_keyboard(
            session_id=session_id, fencing_token=token,
            payload=raw_payload, request_id=call_id,
        )

    record = await _execute_human_input(
        session, auth=auth, request=request, session_id=session_id,
        tool_name="human.keyboard", parameters_redacted=redacted,
        worker_call=call_worker,
    )
    return {"success": True, "session": _session_data(record)}


@router.get("/sessions/{session_id}/frame")
async def get_frame(
    session_id: uuid.UUID, request: Request, response: Response,
    auth: AuthDep, session: SessionDep, client: ComputerClientDep,
):
    await authorize(session, workspace_id=auth.workspace_id, user_id=auth.user_id,
                    permission="computer.observe")
    record = await _load_computer_session(session, auth.workspace_id, session_id)
    if record.status in TERMINAL_STATUSES or record.lease_version <= 0:
        raise AppError("COMPUTER_FRAME_UNAVAILABLE", "The live frame is unavailable.", 409)
    lease = await _active_lease(session, auth.workspace_id, session_id)
    fencing_token = lease.fencing_token if lease else record.lease_version
    frame = await client.frame(
        session_id=session_id, fencing_token=fencing_token,
        request_id=request.state.request_id,
    )
    response.headers["cache-control"] = "no-store"
    response.headers["x-content-type-options"] = "nosniff"
    return {"success": True, "frame": frame}


@router.get("/artifacts/{artifact_id}")
async def get_artifact(
    artifact_id: uuid.UUID, request: Request, auth: AuthDep,
    session: SessionDep, client: ComputerClientDep,
):
    await authorize(session, workspace_id=auth.workspace_id, user_id=auth.user_id,
                    permission="computer.observe")
    artifact = await session.scalar(select(ComputerArtifact).where(
        ComputerArtifact.workspace_id == auth.workspace_id,
        ComputerArtifact.id == artifact_id,
        ComputerArtifact.status == "available",
    ))
    if not artifact:
        raise AppError("COMPUTER_ARTIFACT_NOT_FOUND", "The artifact is unavailable.", 404)
    now = datetime.now(UTC)
    expires_at = artifact.expires_at
    if expires_at and expires_at.tzinfo is None:
        expires_at = expires_at.replace(tzinfo=UTC)
    if expires_at and now >= expires_at:
        raise AppError("COMPUTER_ARTIFACT_EXPIRED", "The artifact has expired.", 410)
    content = await client.artifact(
        workspace_id=auth.workspace_id, session_id=artifact.session_id,
        artifact_id=artifact.id, request_id=request.state.request_id,
    )
    if len(content) != artifact.size_bytes or not hashlib.sha256(content).hexdigest() == artifact.sha256:
        raise AppError(
            "COMPUTER_ARTIFACT_INTEGRITY_FAILED",
            "The artifact failed integrity verification.",
            502,
        )
    return Response(
        content=content, media_type=artifact.mime_type,
        headers={
            "cache-control": "no-store", "x-content-type-options": "nosniff",
            "content-security-policy": "default-src 'none'; sandbox",
            "content-disposition": f'inline; filename="computer-{artifact.id}.jpg"',
        },
    )


async def _prepare_worker_command(
    db_session: SessionDep,
    *, workspace_id: uuid.UUID, session_id: uuid.UUID,
    transition_status: Literal["paused", "recovering", "stopping"],
    lease_owner: Literal["ai", "system"] | None,
    actor_id: uuid.UUID,
    release_reason: str,
) -> tuple[ComputerSession, int]:
    record = await _load_computer_session(db_session, workspace_id, session_id, lock=True)
    current = await _active_lease(db_session, workspace_id, session_id, lock=True)
    now = datetime.now(UTC)
    token: int
    if lease_owner is not None:
        _release_lease(current, release_reason, now)
        issued = _issue_lease(
            db_session, record, owner=lease_owner,
            actor_user_id=None if lease_owner == "ai" else actor_id, now=now,
        )
        token = issued.fencing_token
    elif current:
        token = current.fencing_token
    else:
        issued = _issue_lease(
            db_session, record, owner="system", actor_user_id=actor_id, now=now,
        )
        token = issued.fencing_token
    record.status = transition_status
    record.control_owner = "system" if transition_status != "recovering" else "ai"
    record.version += 1
    await db_session.commit()
    return record, token


@router.post("/sessions/{session_id}/takeover")
async def take_over(
    session_id: uuid.UUID, request: Request, auth: AuthDep, session: SessionDep,
    client: ComputerClientDep,
    hayva_csrf: Annotated[str | None, Cookie()] = None,
    x_csrf_token: Annotated[str | None, Header()] = None,
):
    await authorize(session, workspace_id=auth.workspace_id, user_id=auth.user_id,
                    permission="computer.control")
    await _csrf(session, auth, hayva_csrf, x_csrf_token)
    record = await _load_computer_session(session, auth.workspace_id, session_id, lock=True)
    if record.status != "ai_controlled":
        raise AppError("COMPUTER_TAKEOVER_INVALID", "AI control is not active.", 409)
    lease = await _active_lease(session, auth.workspace_id, session_id, lock=True)
    if not lease or lease.owner != "ai":
        raise AppError("COMPUTER_LEASE_INVALID", "The AI control lease is unavailable.", 409)
    now = datetime.now(UTC)
    _heartbeat_lease(lease, now=now)
    old_token = lease.fencing_token
    _release_lease(lease, "human_takeover", now)
    record.status = "paused"
    record.control_owner = "system"
    record.paused_at = now
    record.version += 1
    await session.commit()
    try:
        worker = await client.command(
            session_id=session_id, command="pause", fencing_token=old_token,
            request_id=request.state.request_id,
        )
    except AppError as error:
        await _persist_worker_failure(
            session, workspace_id=auth.workspace_id, session_id=session_id,
            actor_id=auth.user_id, request_id=request.state.request_id, error=error,
            status="paused", event_type="computer.takeover_failed",
        )
        raise
    record = await _load_computer_session(session, auth.workspace_id, session_id, lock=True)
    now = datetime.now(UTC)
    human_lease = _issue_lease(
        session, record, owner="human", actor_user_id=auth.user_id, now=now
    )
    _apply_worker_state(record, worker)
    record.status = "human_controlled"
    record.control_owner = "human"
    record.failure_code = None
    record.failure_message = None
    record.version += 1
    await _checkpoint(
        session, record, "takeover", worker_state=worker,
        state_redacted={"lease_owner": "human"},
    )
    await _audit(
        session, workspace_id=auth.workspace_id, actor_id=auth.user_id,
        event_type="computer.human_takeover", request_id=request.state.request_id,
        data={"computer_session_id": str(session_id), "lease_version": human_lease.fencing_token},
    )
    await session.commit()
    await session.refresh(record)
    return {"success": True, "session": _session_data(record)}


@router.post("/sessions/{session_id}/recover")
async def recover_session(
    session_id: uuid.UUID, request: Request, auth: AuthDep, session: SessionDep,
    client: ComputerClientDep,
    hayva_csrf: Annotated[str | None, Cookie()] = None,
    x_csrf_token: Annotated[str | None, Header()] = None,
):
    await authorize(session, workspace_id=auth.workspace_id, user_id=auth.user_id,
                    permission="computer.control")
    await _csrf(session, auth, hayva_csrf, x_csrf_token)
    state = await _control_state(session, auth.workspace_id, lock=True)
    if state.emergency_stopped:
        raise AppError("COMPUTER_EMERGENCY_STOPPED", "Computer activity is emergency-stopped.", 409)
    record = await _load_computer_session(session, auth.workspace_id, session_id, lock=True)
    if record.status in {"stopped", "stopping", "human_controlled"}:
        raise AppError("COMPUTER_RECOVERY_INVALID", "The computer session cannot recover.", 409)
    profile = await _load_profile(session, auth.workspace_id, record.profile_id)
    if profile.status == "revoked":
        raise AppError("COMPUTER_PROFILE_REVOKED", "The computer profile was revoked.", 409)
    now = datetime.now(UTC)
    current = await _active_lease(session, auth.workspace_id, session_id, lock=True)
    _release_lease(current, "recovery", now)
    lease = _issue_lease(session, record, owner="ai", actor_user_id=None, now=now)
    record.status = "recovering"
    record.control_owner = "ai"
    record.failure_code = None
    record.failure_message = None
    record.version += 1
    await _audit(
        session, workspace_id=auth.workspace_id, actor_id=auth.user_id,
        event_type="computer.recovery_started", request_id=request.state.request_id,
        data={"computer_session_id": str(session_id), "lease_version": lease.fencing_token},
    )
    await session.commit()
    try:
        worker = await client.start_session(
            workspace_id=auth.workspace_id, session_id=session_id, profile_id=profile.id,
            profile_storage_key=profile.storage_key, fencing_token=lease.fencing_token,
            request_id=f"{request.state.request_id}.recover",
        )
    except AppError as error:
        await _persist_worker_failure(
            session, workspace_id=auth.workspace_id, session_id=session_id,
            actor_id=auth.user_id, request_id=request.state.request_id, error=error,
            status="paused", event_type="computer.recovery_failed",
        )
        raise
    if worker.status not in {"ready", "ai_controlled"}:
        error = AppError(
            "COMPUTER_RECOVERY_STATE_INVALID",
            "The browser worker did not recover into an active state.",
            502,
        )
        await _persist_worker_failure(
            session, workspace_id=auth.workspace_id, session_id=session_id,
            actor_id=auth.user_id, request_id=request.state.request_id, error=error,
            status="paused", event_type="computer.recovery_failed",
        )
        raise error
    record = await _load_computer_session(session, auth.workspace_id, session_id, lock=True)
    _apply_worker_state(record, worker)
    record.status = "ai_controlled"
    record.control_owner = "ai"
    record.paused_at = None
    record.failure_code = None
    record.failure_message = None
    record.version += 1
    await _checkpoint(
        session, record, "recovery", worker_state=worker,
        state_redacted={"persistent_profile_reopened": True},
    )
    await _audit(
        session, workspace_id=auth.workspace_id, actor_id=auth.user_id,
        event_type="computer.recovery_succeeded", request_id=request.state.request_id,
        data={"computer_session_id": str(session_id), "worker_id": worker.worker_id},
    )
    await session.commit()
    await session.refresh(record)
    return {"success": True, "session": _session_data(record), "recovered": True}


@router.post("/sessions/{session_id}/return-control")
async def return_control(
    session_id: uuid.UUID, request: Request, auth: AuthDep, session: SessionDep,
    client: ComputerClientDep,
    hayva_csrf: Annotated[str | None, Cookie()] = None,
    x_csrf_token: Annotated[str | None, Header()] = None,
):
    await authorize(session, workspace_id=auth.workspace_id, user_id=auth.user_id,
                    permission="computer.control")
    await _csrf(session, auth, hayva_csrf, x_csrf_token)
    state = await _control_state(session, auth.workspace_id)
    if state.emergency_stopped:
        raise AppError("COMPUTER_EMERGENCY_STOPPED", "Computer activity is emergency-stopped.", 409)
    record = await _load_computer_session(session, auth.workspace_id, session_id)
    lease = await _active_lease(session, auth.workspace_id, session_id)
    if record.status != "human_controlled" or not lease or lease.owner != "human":
        raise AppError("COMPUTER_RETURN_INVALID", "Human control is not active.", 409)
    if lease.actor_user_id != auth.user_id:
        raise AppError("COMPUTER_LEASE_OWNED_BY_ANOTHER_USER", "Another user holds control.", 403)
    _heartbeat_lease(lease)
    try:
        snapshot = await client.snapshot(
            session_id=session_id, fencing_token=lease.fencing_token,
            request_id=f"{request.state.request_id}.snapshot",
        )
    except AppError:
        await session.rollback()
        raise
    record = await _load_computer_session(session, auth.workspace_id, session_id, lock=True)
    current = await _active_lease(session, auth.workspace_id, session_id, lock=True)
    if not current or current.id != lease.id:
        raise AppError("COMPUTER_LEASE_CHANGED", "The control lease changed.", 409)
    now = datetime.now(UTC)
    _release_lease(current, "returned_to_ai", now)
    ai_lease = _issue_lease(session, record, owner="ai", actor_user_id=None, now=now)
    _apply_worker_state(record, snapshot)
    record.status = "recovering"
    record.control_owner = "ai"
    record.version += 1
    await session.commit()
    try:
        worker = await client.command(
            session_id=session_id, command="resume", fencing_token=ai_lease.fencing_token,
            request_id=f"{request.state.request_id}.resume",
        )
    except AppError as error:
        await _persist_worker_failure(
            session, workspace_id=auth.workspace_id, session_id=session_id,
            actor_id=auth.user_id, request_id=request.state.request_id, error=error,
            status="paused", event_type="computer.return_control_failed",
        )
        raise
    record = await _load_computer_session(session, auth.workspace_id, session_id, lock=True)
    _apply_worker_state(record, worker)
    record.status = "ai_controlled"
    record.control_owner = "ai"
    record.paused_at = None
    record.failure_code = None
    record.failure_message = None
    record.version += 1
    await _checkpoint(
        session, record, "return_control", worker_state=snapshot,
        state_redacted={"reobserved": True, "page_title": snapshot.page_title},
    )
    await _audit(
        session, workspace_id=auth.workspace_id, actor_id=auth.user_id,
        event_type="computer.control_returned_to_ai", request_id=request.state.request_id,
        data={"computer_session_id": str(session_id), "reobserved": True},
    )
    await session.commit()
    await session.refresh(record)
    return {"success": True, "session": _session_data(record), "reobserved": True}


async def _simple_command(
    *, command: Literal["pause", "resume", "stop"], session_id: uuid.UUID,
    request: Request, auth: AuthDep, db_session: SessionDep, client: ComputerClient,
) -> dict:
    record = await _load_computer_session(
        db_session, auth.workspace_id, session_id, lock=True
    )
    state = await _control_state(db_session, auth.workspace_id, lock=True)
    if command == "resume":
        if state.emergency_stopped:
            raise AppError("COMPUTER_EMERGENCY_STOPPED", "Computer activity is emergency-stopped.", 409)
        if record.status != "paused":
            raise AppError("COMPUTER_RESUME_INVALID", "The computer session is not paused.", 409)
        transition, owner = "recovering", "ai"
    elif command == "pause":
        if record.status not in {"ai_controlled", "human_controlled", "ready", "recovering"}:
            raise AppError("COMPUTER_PAUSE_INVALID", "The computer session cannot be paused.", 409)
        if record.status == "human_controlled":
            human_lease = await _active_lease(
                db_session, auth.workspace_id, session_id, lock=True
            )
            if (
                not human_lease
                or human_lease.owner != "human"
                or human_lease.actor_user_id != auth.user_id
            ):
                raise AppError(
                    "COMPUTER_HUMAN_CONTROL_OWNED",
                    "Only the active human controller can pause this session.",
                    403,
                )
            _heartbeat_lease(human_lease)
        transition, owner = "paused", None
    else:
        if record.status in TERMINAL_STATUSES:
            return {"success": True, "session": _session_data(record)}
        transition, owner = "stopping", None
    _, token = await _prepare_worker_command(
        db_session, workspace_id=auth.workspace_id, session_id=session_id,
        transition_status=transition, lease_owner=owner,
        actor_id=auth.user_id, release_reason=f"user_{command}",
    )
    try:
        worker = await client.command(
            session_id=session_id, command=command, fencing_token=token,
            request_id=request.state.request_id,
        )
    except AppError as error:
        await _persist_worker_failure(
            db_session, workspace_id=auth.workspace_id, session_id=session_id,
            actor_id=auth.user_id, request_id=request.state.request_id, error=error,
            status="paused" if command != "stop" else "failed",
            event_type=f"computer.{command}_failed",
        )
        raise
    record = await _load_computer_session(
        db_session, auth.workspace_id, session_id, lock=True
    )
    now = datetime.now(UTC)
    _apply_worker_state(record, worker)
    record.failure_code = None
    record.failure_message = None
    if command == "resume":
        record.status, record.control_owner, record.paused_at = "ai_controlled", "ai", None
        checkpoint_type = "recovery"
    elif command == "pause":
        active = await _active_lease(db_session, auth.workspace_id, session_id, lock=True)
        _release_lease(active, "paused", now)
        record.status, record.control_owner, record.paused_at = "paused", "none", now
        checkpoint_type = "pause"
    else:
        active = await _active_lease(db_session, auth.workspace_id, session_id, lock=True)
        _release_lease(active, "stopped", now)
        record.status, record.control_owner, record.stopped_at = "stopped", "none", now
        checkpoint_type = "completion"
    record.version += 1
    await _checkpoint(db_session, record, checkpoint_type, worker_state=worker)
    await _audit(
        db_session, workspace_id=auth.workspace_id, actor_id=auth.user_id,
        event_type={
            "pause": "computer.session_paused",
            "resume": "computer.session_resumed",
            "stop": "computer.session_stopped",
        }[command],
        request_id=request.state.request_id,
        data={"computer_session_id": str(session_id)},
    )
    await db_session.commit()
    await db_session.refresh(record)
    return {"success": True, "session": _session_data(record)}


@router.post("/sessions/{session_id}/pause")
async def pause_session(
    session_id: uuid.UUID, request: Request, auth: AuthDep, session: SessionDep,
    client: ComputerClientDep,
    hayva_csrf: Annotated[str | None, Cookie()] = None,
    x_csrf_token: Annotated[str | None, Header()] = None,
):
    await authorize(session, workspace_id=auth.workspace_id, user_id=auth.user_id,
                    permission="computer.control")
    await _csrf(session, auth, hayva_csrf, x_csrf_token)
    return await _simple_command(
        command="pause", session_id=session_id, request=request, auth=auth,
        db_session=session, client=client,
    )


@router.post("/sessions/{session_id}/resume")
async def resume_session(
    session_id: uuid.UUID, request: Request, auth: AuthDep, session: SessionDep,
    client: ComputerClientDep,
    hayva_csrf: Annotated[str | None, Cookie()] = None,
    x_csrf_token: Annotated[str | None, Header()] = None,
):
    await authorize(session, workspace_id=auth.workspace_id, user_id=auth.user_id,
                    permission="computer.control")
    await _csrf(session, auth, hayva_csrf, x_csrf_token)
    return await _simple_command(
        command="resume", session_id=session_id, request=request, auth=auth,
        db_session=session, client=client,
    )


@router.post("/sessions/{session_id}/stop")
async def stop_session(
    session_id: uuid.UUID, request: Request, auth: AuthDep, session: SessionDep,
    client: ComputerClientDep,
    hayva_csrf: Annotated[str | None, Cookie()] = None,
    x_csrf_token: Annotated[str | None, Header()] = None,
):
    await authorize(session, workspace_id=auth.workspace_id, user_id=auth.user_id,
                    permission="computer.control")
    await _csrf(session, auth, hayva_csrf, x_csrf_token)
    return await _simple_command(
        command="stop", session_id=session_id, request=request, auth=auth,
        db_session=session, client=client,
    )


@router.post("/emergency-stop")
async def emergency_stop(
    payload: EmergencyStopInput, request: Request, auth: AuthDep, session: SessionDep,
    client: ComputerClientDep,
    hayva_csrf: Annotated[str | None, Cookie()] = None,
    x_csrf_token: Annotated[str | None, Header()] = None,
):
    await authorize(session, workspace_id=auth.workspace_id, user_id=auth.user_id,
                    permission="computer.emergency_stop")
    await _csrf(session, auth, hayva_csrf, x_csrf_token)
    state = await _control_state(session, auth.workspace_id, lock=True)
    now = datetime.now(UTC)
    state.emergency_stopped = True
    state.reason = payload.reason
    state.stopped_by_user_id = auth.user_id
    state.stopped_at = now
    state.resumed_by_user_id = None
    state.resumed_at = None
    state.version += 1
    records = list((await session.scalars(select(ComputerSession).where(
        ComputerSession.workspace_id == auth.workspace_id,
        ComputerSession.status.in_(ACTIVE_STATUSES),
    ).with_for_update())).all())
    commands: list[tuple[uuid.UUID, int]] = []
    for record in records:
        lease = await _active_lease(session, auth.workspace_id, record.id, lock=True)
        _release_lease(lease, "emergency_stop", now)
        system_lease = _issue_lease(
            session, record, owner="system", actor_user_id=auth.user_id, now=now
        )
        record.status = "stopping"
        record.control_owner = "system"
        record.version += 1
        commands.append((record.id, system_lease.fencing_token))
    await _audit(
        session, workspace_id=auth.workspace_id, actor_id=auth.user_id,
        event_type="computer.emergency_stopped", request_id=request.state.request_id,
        data={"reason": payload.reason, "affected_sessions": len(commands)},
    )
    await session.commit()
    stopped, unverified = 0, 0
    for record_id, token in commands:
        try:
            worker = await client.command(
                session_id=record_id, command="stop", fencing_token=token,
                request_id=request.state.request_id,
            )
        except AppError as error:
            unverified += 1
            await _persist_worker_failure(
                session, workspace_id=auth.workspace_id, session_id=record_id,
                actor_id=auth.user_id, request_id=request.state.request_id, error=error,
                status="failed", event_type="computer.emergency_stop_unverified",
            )
            continue
        record = await _load_computer_session(session, auth.workspace_id, record_id, lock=True)
        lease = await _active_lease(session, auth.workspace_id, record_id, lock=True)
        _release_lease(lease, "emergency_stopped", datetime.now(UTC))
        _apply_worker_state(record, worker)
        record.status = "stopped"
        record.control_owner = "none"
        record.stopped_at = datetime.now(UTC)
        record.failure_code = None
        record.failure_message = None
        record.version += 1
        await _checkpoint(session, record, "completion", worker_state=worker,
                          state_redacted={"emergency_stop": True})
        await session.commit()
        stopped += 1
    return {
        "success": True,
        "emergency_stopped": True,
        "affected_sessions": len(commands),
        "verified_stopped_sessions": stopped,
        "unverified_sessions": unverified,
    }


@router.post("/emergency-resume")
async def emergency_resume(
    request: Request, auth: AuthDep, session: SessionDep,
    hayva_csrf: Annotated[str | None, Cookie()] = None,
    x_csrf_token: Annotated[str | None, Header()] = None,
):
    await authorize(session, workspace_id=auth.workspace_id, user_id=auth.user_id,
                    permission="computer.emergency_stop")
    await _csrf(session, auth, hayva_csrf, x_csrf_token)
    state = await _control_state(session, auth.workspace_id, lock=True)
    if state.emergency_stopped:
        state.emergency_stopped = False
        state.resumed_by_user_id = auth.user_id
        state.resumed_at = datetime.now(UTC)
        state.version += 1
        await _audit(
            session, workspace_id=auth.workspace_id, actor_id=auth.user_id,
            event_type="computer.emergency_stop_cleared", request_id=request.state.request_id,
            data={"sessions_auto_resumed": False},
        )
        await session.commit()
        await session.refresh(state)
    return {
        "success": True, "emergency_stopped": False,
        "sessions_auto_resumed": False, "version": state.version,
    }
