Skip to content

Safety API Reference

effero.safety

Runtime guardrail engine: policy-as-code checks that run independently of the LLM, between the skill layer and the Device Abstraction Layer.

Evaluates declarative policies against live telemetry and device facts via the high-performance out-of-process Rust safety kernel daemon (TCP 127.0.0.1:9400).

ApprovalHandler

Bases: ABC

Abstract interface for receiving and resolving HITL approval requests.

Source code in src/effero/safety/approval.py
class ApprovalHandler(ABC):
    """Abstract interface for receiving and resolving HITL approval requests."""

    @abstractmethod
    async def request_approval(self, request: ApprovalRequest) -> ApprovalResponse:
        """Present the approval request to a human operator and await decision."""
        ...

request_approval(request) abstractmethod async

Present the approval request to a human operator and await decision.

Source code in src/effero/safety/approval.py
@abstractmethod
async def request_approval(self, request: ApprovalRequest) -> ApprovalResponse:
    """Present the approval request to a human operator and await decision."""
    ...

ApprovalRequest dataclass

A request for human approval before executing a sensitive action.

Source code in src/effero/safety/approval.py
@dataclass
class ApprovalRequest:
    """A request for human approval before executing a sensitive action."""

    id: str = field(default_factory=lambda: str(uuid.uuid4()))
    skill_name: str = ""
    arguments: dict[str, Any] = field(default_factory=dict)
    reason: str = "Skill requires human authorization"
    safety_class: str = "act_with_approval"
    created_at: str = field(default_factory=lambda: datetime.now(UTC).isoformat())

ApprovalResponse dataclass

The human decision for an approval request.

Source code in src/effero/safety/approval.py
@dataclass
class ApprovalResponse:
    """The human decision for an approval request."""

    request_id: str
    approved: bool
    responder: str = "human"
    comment: str | None = None
    timestamp: str = field(default_factory=lambda: datetime.now(UTC).isoformat())

AutoApprovalHandler

Bases: ApprovalHandler

Approval handler that resolves requests automatically.

Supports: - All-or-nothing mode (default: approve_all=True or approve_all=False) - Graduated per-SafetyClass mode (graduated=True): - READ_ONLY / ACT_AUTONOMOUS: auto-approved (True) - ACT_RESTRICTED: auto-denied (False, simulation only) - ACT_WITH_APPROVAL: evaluates according to approve_approval_gated (default False)

Source code in src/effero/safety/approval.py
class AutoApprovalHandler(ApprovalHandler):
    """Approval handler that resolves requests automatically.

    Supports:
    - All-or-nothing mode (default: `approve_all=True` or `approve_all=False`)
    - Graduated per-SafetyClass mode (`graduated=True`):
      - READ_ONLY / ACT_AUTONOMOUS: auto-approved (True)
      - ACT_RESTRICTED: auto-denied (False, simulation only)
      - ACT_WITH_APPROVAL: evaluates according to `approve_approval_gated` (default False)
    """

    def __init__(
        self,
        approve_all: bool = True,
        graduated: bool = False,
        approve_approval_gated: bool = False,
        responder_name: str = "auto",
    ) -> None:
        self.approve_all = approve_all
        self.graduated = graduated
        self.approve_approval_gated = approve_approval_gated
        self.responder_name = responder_name
        self.history: list[tuple[ApprovalRequest, ApprovalResponse]] = []

    async def request_approval(self, request: ApprovalRequest) -> ApprovalResponse:
        norm_class = (request.safety_class or "act_with_approval").lower()
        if self.graduated:
            if norm_class in ("read_only", "act_autonomous"):
                approved = True
                comment = f"Auto-approved graduated {norm_class}"
            elif norm_class == "act_restricted":
                approved = False
                comment = "Denied act_restricted in graduated auto mode"
            else:
                approved = self.approve_approval_gated
                comment = f"Graduated auto-decision for {norm_class}: {approved}"
        else:
            approved = self.approve_all
            comment = "Automatic decision by policy/simulation"

        decision = ApprovalResponse(
            request_id=request.id,
            approved=approved,
            responder=self.responder_name,
            comment=comment,
        )
        self.history.append((request, decision))
        return decision

CallbackApprovalHandler

Bases: ApprovalHandler

Approval handler that dispatches to an async callback (e.g. Web UI, WebSocket, Telegram).

Source code in src/effero/safety/approval.py
class CallbackApprovalHandler(ApprovalHandler):
    """Approval handler that dispatches to an async callback (e.g. Web UI, WebSocket, Telegram)."""

    def __init__(
        self,
        callback: Callable[[ApprovalRequest], Any],
        timeout_seconds: float = 120.0,
    ):
        self.callback = callback
        self.timeout_seconds = timeout_seconds
        self._pending: dict[str, asyncio.Future[ApprovalResponse]] = {}

    async def request_approval(self, request: ApprovalRequest) -> ApprovalResponse:
        future = asyncio.get_running_loop().create_future()
        self._pending[request.id] = future

        try:
            result = self.callback(request)
            if hasattr(result, "__await__"):
                asyncio.create_task(result)

            return await asyncio.wait_for(future, timeout=self.timeout_seconds)
        except TimeoutError:
            logger.warning(f"Approval request {request.id} timed out; defaulting to deny.")
            return ApprovalResponse(
                request_id=request.id,
                approved=False,
                responder="system_timeout",
                comment="Approval timed out",
            )
        finally:
            self._pending.pop(request.id, None)

    def resolve(self, request_id: str, approved: bool, comment: str | None = None) -> bool:
        """Resolve a pending approval request externally."""
        if request_id in self._pending and not self._pending[request_id].done():
            self._pending[request_id].set_result(
                ApprovalResponse(
                    request_id=request_id,
                    approved=approved,
                    responder="remote_operator",
                    comment=comment,
                )
            )
            return True
        return False

resolve(request_id, approved, comment=None)

Resolve a pending approval request externally.

Source code in src/effero/safety/approval.py
def resolve(self, request_id: str, approved: bool, comment: str | None = None) -> bool:
    """Resolve a pending approval request externally."""
    if request_id in self._pending and not self._pending[request_id].done():
        self._pending[request_id].set_result(
            ApprovalResponse(
                request_id=request_id,
                approved=approved,
                responder="remote_operator",
                comment=comment,
            )
        )
        return True
    return False

ConsoleApprovalHandler

Bases: ApprovalHandler

Interactive CLI terminal approval prompt.

Source code in src/effero/safety/approval.py
class ConsoleApprovalHandler(ApprovalHandler):
    """Interactive CLI terminal approval prompt."""

    def __init__(self, timeout_seconds: float = 60.0):
        self.timeout_seconds = timeout_seconds

    async def request_approval(self, request: ApprovalRequest) -> ApprovalResponse:
        print("\n" + "=" * 60)
        print("🚨 EFFERO SAFETY KERNEL: APPROVAL REQUIRED")
        print("=" * 60)
        print(f"  Request ID:   {request.id}")
        print(f"  Skill:        {request.skill_name}")
        print(f"  Safety Class: {request.safety_class}")
        print(f"  Reason:       {request.reason}")
        print(f"  Arguments:    {request.arguments}")
        print("=" * 60)

        loop = asyncio.get_running_loop()
        try:
            # Prompt user without blocking async event loop
            user_input = await asyncio.wait_for(
                loop.run_in_executor(None, input, "Authorize this action? [y/N]: "),
                timeout=self.timeout_seconds,
            )
            approved = user_input.strip().lower() in ("y", "yes")
        except TimeoutError:
            print("\n[Timeout waiting for authorization: Denied by default]")
            approved = False
        except (KeyboardInterrupt, EOFError):
            approved = False

        return ApprovalResponse(
            request_id=request.id,
            approved=approved,
            responder="console_operator",
            comment="Approved via console" if approved else "Denied via console",
        )

GraduatedApprovalHandler

Bases: ApprovalHandler

Graduated approval handler: auto-approves safe classes, gates ACT_WITH_APPROVAL via delegate.

  • READ_ONLY: auto-approved
  • ACT_AUTONOMOUS: auto-approved
  • ACT_WITH_APPROVAL: delegates to secondary handler (e.g. Console or Callback)
  • ACT_RESTRICTED: denied by default (simulation/sandbox only)
Source code in src/effero/safety/approval.py
class GraduatedApprovalHandler(ApprovalHandler):
    """Graduated approval handler: auto-approves safe classes, gates ACT_WITH_APPROVAL via delegate.

    - READ_ONLY: auto-approved
    - ACT_AUTONOMOUS: auto-approved
    - ACT_WITH_APPROVAL: delegates to secondary handler (e.g. Console or Callback)
    - ACT_RESTRICTED: denied by default (simulation/sandbox only)
    """

    def __init__(self, delegate: ApprovalHandler | None = None):
        self.delegate = delegate or ConsoleApprovalHandler()

    async def request_approval(self, request: ApprovalRequest) -> ApprovalResponse:
        norm_class = request.safety_class.lower()

        if norm_class in ("read_only", "act_autonomous"):
            return ApprovalResponse(
                request_id=request.id,
                approved=True,
                responder="graduated_auto",
                comment=f"Auto-approved {norm_class} action",
            )
        elif norm_class == "act_restricted":
            return ApprovalResponse(
                request_id=request.id,
                approved=False,
                responder="graduated_auto",
                comment="Denied act_restricted action (simulation only)",
            )
        else:
            return await self.delegate.request_approval(request)

SafetyClient

Async TCP client for the effero-safety-kernel.

The safety kernel is an independent Rust process that evaluates proposed actions against a declarative policy before they reach hardware. Communication uses newline-delimited JSON over TCP.

Source code in src/effero/safety/client.py
class SafetyClient:
    """Async TCP client for the effero-safety-kernel.

    The safety kernel is an independent Rust process that evaluates
    proposed actions against a declarative policy before they reach hardware.
    Communication uses newline-delimited JSON over TCP.
    """

    def __init__(self, host: str = "127.0.0.1", port: int = 9400) -> None:
        self.host = host
        self.port = port
        self._reader: asyncio.StreamReader | None = None
        self._writer: asyncio.StreamWriter | None = None
        self._connected = False
        self._lock = asyncio.Lock()

    async def connect(self) -> None:
        """Connect to the safety kernel."""
        self._reader, self._writer = await asyncio.open_connection(self.host, self.port)
        self._connected = True
        logger.info(f"Connected to safety kernel at {self.host}:{self.port}")

    async def _send_request(self, request: dict[str, Any]) -> dict[str, Any]:
        """Send a JSON request and read the JSON response under concurrency lock."""
        if not self._connected:
            await self.connect()

        async with self._lock:
            assert self._writer is not None
            assert self._reader is not None
            self._writer.write((json.dumps(request) + "\n").encode())
            await self._writer.drain()

            line = await self._reader.readline()
            if not line:
                self._connected = False
                raise ConnectionError("Safety kernel closed the connection")

            return json.loads(line.decode())

    async def check_action(self, skill_name: str, facts: dict[str, Any] | None = None) -> dict[str, Any]:
        """Check whether an action is allowed by the safety policy.

        Args:
            skill_name: The skill being invoked (e.g., 'iot.lights.toggle')
            facts: Optional dict of numeric/boolean facts for condition evaluation

        Returns:
            dict with 'decision' key: 'allow', 'require_approval', 'deny', or 'limit'
        """
        request = {"type": "evaluate", "skill": skill_name, "facts": facts or {}}
        result = await self._send_request(request)
        logger.debug(f"Safety check for '{skill_name}': {result}")
        return result

    async def send_heartbeat(self) -> dict[str, Any]:
        """Send a heartbeat to prevent the safety kernel watchdog from tripping."""
        request = {"type": "heartbeat"}
        return await self._send_request(request)

    async def get_status(self) -> dict[str, Any]:
        """Query safety kernel and watchdog status."""
        request = {"type": "status"}
        return await self._send_request(request)

    async def reset_watchdog(self) -> dict[str, Any]:
        """Reset the watchdog after a timeout trip."""
        request = {"type": "reset_watchdog"}
        return await self._send_request(request)

    async def close(self) -> None:
        """Close the connection."""
        if self._writer:
            self._writer.close()
            await self._writer.wait_closed()
        self._connected = False

    @property
    def connected(self) -> bool:
        return self._connected

connect() async

Connect to the safety kernel.

Source code in src/effero/safety/client.py
async def connect(self) -> None:
    """Connect to the safety kernel."""
    self._reader, self._writer = await asyncio.open_connection(self.host, self.port)
    self._connected = True
    logger.info(f"Connected to safety kernel at {self.host}:{self.port}")

check_action(skill_name, facts=None) async

Check whether an action is allowed by the safety policy.

Parameters:

Name Type Description Default
skill_name str

The skill being invoked (e.g., 'iot.lights.toggle')

required
facts dict[str, Any] | None

Optional dict of numeric/boolean facts for condition evaluation

None

Returns:

Type Description
dict[str, Any]

dict with 'decision' key: 'allow', 'require_approval', 'deny', or 'limit'

Source code in src/effero/safety/client.py
async def check_action(self, skill_name: str, facts: dict[str, Any] | None = None) -> dict[str, Any]:
    """Check whether an action is allowed by the safety policy.

    Args:
        skill_name: The skill being invoked (e.g., 'iot.lights.toggle')
        facts: Optional dict of numeric/boolean facts for condition evaluation

    Returns:
        dict with 'decision' key: 'allow', 'require_approval', 'deny', or 'limit'
    """
    request = {"type": "evaluate", "skill": skill_name, "facts": facts or {}}
    result = await self._send_request(request)
    logger.debug(f"Safety check for '{skill_name}': {result}")
    return result

send_heartbeat() async

Send a heartbeat to prevent the safety kernel watchdog from tripping.

Source code in src/effero/safety/client.py
async def send_heartbeat(self) -> dict[str, Any]:
    """Send a heartbeat to prevent the safety kernel watchdog from tripping."""
    request = {"type": "heartbeat"}
    return await self._send_request(request)

get_status() async

Query safety kernel and watchdog status.

Source code in src/effero/safety/client.py
async def get_status(self) -> dict[str, Any]:
    """Query safety kernel and watchdog status."""
    request = {"type": "status"}
    return await self._send_request(request)

reset_watchdog() async

Reset the watchdog after a timeout trip.

Source code in src/effero/safety/client.py
async def reset_watchdog(self) -> dict[str, Any]:
    """Reset the watchdog after a timeout trip."""
    request = {"type": "reset_watchdog"}
    return await self._send_request(request)

close() async

Close the connection.

Source code in src/effero/safety/client.py
async def close(self) -> None:
    """Close the connection."""
    if self._writer:
        self._writer.close()
        await self._writer.wait_closed()
    self._connected = False

SafetyDaemonManager

Discovers, starts, monitors, and stops the Rust safety kernel daemon.

Source code in src/effero/safety/daemon.py
class SafetyDaemonManager:
    """Discovers, starts, monitors, and stops the Rust safety kernel daemon."""

    def __init__(
        self,
        host: str = "127.0.0.1",
        port: int = 9400,
        policy_path: str | Path | None = None,
        custom_binary_path: str | Path | None = None,
    ) -> None:
        self.host = host
        self.port = port
        self.policy_path = Path(policy_path) if policy_path else None
        self.custom_binary_path = Path(custom_binary_path) if custom_binary_path else None
        self._process: asyncio.subprocess.Process | None = None
        self._owns_process = False

    def find_binary(self) -> Path | None:
        """Find the effero-safety-kerneld binary on PATH or local workspace targets."""
        if self.custom_binary_path and self.custom_binary_path.exists():
            return self.custom_binary_path

        # Check system PATH
        which_path = shutil.which("effero-safety-kerneld")
        if which_path:
            return Path(which_path)

        # Check workspace target directories (target/release, target/debug)
        # Assuming repo root relative to this module
        current = Path(__file__).resolve().parents[3]  # repo root
        candidates = [
            current / "target" / "release" / "effero-safety-kerneld.exe",
            current / "target" / "release" / "effero-safety-kerneld",
            current / "target" / "debug" / "effero-safety-kerneld.exe",
            current / "target" / "debug" / "effero-safety-kerneld",
        ]
        for candidate in candidates:
            if candidate.exists() and candidate.is_file():
                return candidate

        return None

    async def is_running(self) -> bool:
        """Check if a safety kernel daemon is already accepting connections on host:port."""
        try:
            reader, writer = await asyncio.wait_for(
                asyncio.open_connection(self.host, self.port),
                timeout=0.5,
            )
            writer.close()
            await writer.wait_closed()
            return True
        except (TimeoutError, OSError):
            return False

    async def ensure_running(self) -> bool:
        """Ensure the safety daemon is running; spawn it if absent and binary is found."""
        if await self.is_running():
            logger.info(f"Safety kernel already running on {self.host}:{self.port}")
            return True

        binary = self.find_binary()
        if not binary:
            logger.warning("effero-safety-kerneld binary not found. Skipping auto-spawn.")
            return False

        args = [str(binary), "--port", str(self.port), "--host", self.host]
        if self.policy_path and self.policy_path.exists():
            args.extend(["--policy", str(self.policy_path)])

        logger.info(f"Auto-spawning safety kernel daemon: {' '.join(args)}")
        try:
            self._process = await asyncio.create_subprocess_exec(
                *args,
                stdout=asyncio.subprocess.PIPE,
                stderr=asyncio.subprocess.PIPE,
            )
            self._owns_process = True

            # Wait for daemon to bind and respond to TCP connections
            for _ in range(20):
                await asyncio.sleep(0.1)
                if await self.is_running():
                    logger.info(f"Safety kernel daemon ready on {self.host}:{self.port}")
                    return True

            logger.error("Safety kernel daemon started but did not respond on port in time")
            await self.stop()
            return False
        except Exception as e:
            logger.error(f"Failed to spawn safety kernel daemon: {e}")
            return False

    async def stop(self) -> None:
        """Stop the child safety daemon process if we started it."""
        if self._owns_process and self._process:
            logger.info("Stopping safety kernel daemon...")
            try:
                self._process.terminate()
                await asyncio.wait_for(self._process.wait(), timeout=2.0)
            except (ProcessLookupError, OSError):
                pass
            except Exception:
                try:
                    self._process.kill()
                except Exception:
                    pass
            self._process = None
            self._owns_process = False
            logger.info("Safety kernel daemon stopped")

find_binary()

Find the effero-safety-kerneld binary on PATH or local workspace targets.

Source code in src/effero/safety/daemon.py
def find_binary(self) -> Path | None:
    """Find the effero-safety-kerneld binary on PATH or local workspace targets."""
    if self.custom_binary_path and self.custom_binary_path.exists():
        return self.custom_binary_path

    # Check system PATH
    which_path = shutil.which("effero-safety-kerneld")
    if which_path:
        return Path(which_path)

    # Check workspace target directories (target/release, target/debug)
    # Assuming repo root relative to this module
    current = Path(__file__).resolve().parents[3]  # repo root
    candidates = [
        current / "target" / "release" / "effero-safety-kerneld.exe",
        current / "target" / "release" / "effero-safety-kerneld",
        current / "target" / "debug" / "effero-safety-kerneld.exe",
        current / "target" / "debug" / "effero-safety-kerneld",
    ]
    for candidate in candidates:
        if candidate.exists() and candidate.is_file():
            return candidate

    return None

is_running() async

Check if a safety kernel daemon is already accepting connections on host:port.

Source code in src/effero/safety/daemon.py
async def is_running(self) -> bool:
    """Check if a safety kernel daemon is already accepting connections on host:port."""
    try:
        reader, writer = await asyncio.wait_for(
            asyncio.open_connection(self.host, self.port),
            timeout=0.5,
        )
        writer.close()
        await writer.wait_closed()
        return True
    except (TimeoutError, OSError):
        return False

ensure_running() async

Ensure the safety daemon is running; spawn it if absent and binary is found.

Source code in src/effero/safety/daemon.py
async def ensure_running(self) -> bool:
    """Ensure the safety daemon is running; spawn it if absent and binary is found."""
    if await self.is_running():
        logger.info(f"Safety kernel already running on {self.host}:{self.port}")
        return True

    binary = self.find_binary()
    if not binary:
        logger.warning("effero-safety-kerneld binary not found. Skipping auto-spawn.")
        return False

    args = [str(binary), "--port", str(self.port), "--host", self.host]
    if self.policy_path and self.policy_path.exists():
        args.extend(["--policy", str(self.policy_path)])

    logger.info(f"Auto-spawning safety kernel daemon: {' '.join(args)}")
    try:
        self._process = await asyncio.create_subprocess_exec(
            *args,
            stdout=asyncio.subprocess.PIPE,
            stderr=asyncio.subprocess.PIPE,
        )
        self._owns_process = True

        # Wait for daemon to bind and respond to TCP connections
        for _ in range(20):
            await asyncio.sleep(0.1)
            if await self.is_running():
                logger.info(f"Safety kernel daemon ready on {self.host}:{self.port}")
                return True

        logger.error("Safety kernel daemon started but did not respond on port in time")
        await self.stop()
        return False
    except Exception as e:
        logger.error(f"Failed to spawn safety kernel daemon: {e}")
        return False

stop() async

Stop the child safety daemon process if we started it.

Source code in src/effero/safety/daemon.py
async def stop(self) -> None:
    """Stop the child safety daemon process if we started it."""
    if self._owns_process and self._process:
        logger.info("Stopping safety kernel daemon...")
        try:
            self._process.terminate()
            await asyncio.wait_for(self._process.wait(), timeout=2.0)
        except (ProcessLookupError, OSError):
            pass
        except Exception:
            try:
                self._process.kill()
            except Exception:
                pass
        self._process = None
        self._owns_process = False
        logger.info("Safety kernel daemon stopped")