Saltar a contenido

ciel.runtime — Runtime de agentes y tools

Núcleo del runtime: DTOs de chat (ChatMessage, ChatChoice, ChatRequest, ChatResponse), contratos de provider/modelo (ChatProvider, ModelProvider), tool loop y despacho de herramientas (ToolProvider, DefaultToolDispatcher), resultados de ejecución (ToolLoopResult, AgentRuntimeResult, AgentContext) y el runtime concreto con trazas (DefaultAgentRuntime, AgentRuntime).

ChatMessage.content (multimodal)

ChatMessage.content es str | list[dict[str, Any]]:

  • str — texto plano (compatibilidad total con versiones previas).
  • list[dict]partes multimodales (text, image_url, input_audio). Los providers convierten estas partes a su formato nativo automáticamente.

ChatMessage.text() -> str concatena las partes de tipo "text" e ignora imágenes/audio (ver docs/api-reference/providers.md para los serializers por proveedor y ejemplos de partes).

ciel.runtime

ContentPart = Dict[str, Any] module-attribute

AgentContext dataclass

Source code in src/ciel/runtime/__init__.py
@dataclass(frozen=True)
class AgentContext:
    agent: str
    session_id: str
    tenant_id: Optional[str] = None
    metadata: Dict[str, Any] = field(default_factory=dict)

AgentRuntime

Async runtime contract for tool-loop execution and streaming.

Source code in src/ciel/runtime/__init__.py
class AgentRuntime:
    """Async runtime contract for tool-loop execution and streaming."""

    async def run_agent_loop(
        self,
        *,
        request: ChatRequest,
        tenant_id: Optional[str] = None,
        toolset: Optional[str] = None,
        limit: int = 32,
    ) -> AgentRuntimeResult:
        raise NotImplementedError

    async def stream_agent_loop(
        self,
        *,
        request: ChatRequest,
        tenant_id: Optional[str] = None,
        toolset: Optional[str] = None,
        limit: int = 32,
    ) -> AsyncIterator[ToolLoopResult]:
        raise NotImplementedError

AgentRuntimeResult dataclass

Source code in src/ciel/runtime/tools.py
@dataclass(frozen=True)
class AgentRuntimeResult:
    response: ChatResponse
    loop_results: Sequence[ToolLoopResult]
    tenant_id: Optional[str] = None
    metadata: Dict[str, Any] = field(default_factory=dict)

AuditEvent dataclass

Source code in src/ciel/observability/__init__.py
@dataclass
class AuditEvent:
    event: str
    session_id: Optional[str] = None
    agent: Optional[str] = None
    tool_call_id: Optional[str] = None
    data: Dict[str, Any] = None
    tenant_id: Optional[str] = None

    def __post_init__(self) -> None:
        if self.data is None:
            self.data = {}

ChatChoice dataclass

Source code in src/ciel/runtime/tools.py
@dataclass(frozen=True)
class ChatChoice:
    message: ChatMessage
    finish_reason: str
    usage: Optional[Dict[str, Any]] = None
    metadata: Dict[str, Any] = field(default_factory=dict)

ChatMessage dataclass

Source code in src/ciel/runtime/tools.py
@dataclass(frozen=True)
class ChatMessage:
    role: str
    content: ChatContent
    name: Optional[str] = None
    tool_call_id: Optional[str] = None
    tool_calls: Optional[list[dict[str, Any]]] = None
    metadata: Dict[str, Any] = field(default_factory=dict)

    def text(self) -> str:
        """Extract plain text from content, tolerant to multimodal parts.

        - ``str`` content is returned verbatim.
        - ``list`` content concatenates the ``text`` of every part whose type
          is ``"text"`` but drops images/audio/video, so consumers (CLI,
          compression, ``AgentResponse.text``) see only readable text.
        """
        content = self.content
        if isinstance(content, str):
            return content
        if not isinstance(content, list):
            return ""
        parts: list[str] = []
        for part in content:
            if isinstance(part, dict) and part.get("type") == "text":
                text = part.get("text", "")
                if isinstance(text, str):
                    parts.append(text)
        return "".join(parts)

text() -> str

Extract plain text from content, tolerant to multimodal parts.

  • str content is returned verbatim.
  • list content concatenates the text of every part whose type is "text" but drops images/audio/video, so consumers (CLI, compression, AgentResponse.text) see only readable text.
Source code in src/ciel/runtime/tools.py
def text(self) -> str:
    """Extract plain text from content, tolerant to multimodal parts.

    - ``str`` content is returned verbatim.
    - ``list`` content concatenates the ``text`` of every part whose type
      is ``"text"`` but drops images/audio/video, so consumers (CLI,
      compression, ``AgentResponse.text``) see only readable text.
    """
    content = self.content
    if isinstance(content, str):
        return content
    if not isinstance(content, list):
        return ""
    parts: list[str] = []
    for part in content:
        if isinstance(part, dict) and part.get("type") == "text":
            text = part.get("text", "")
            if isinstance(text, str):
                parts.append(text)
    return "".join(parts)

ChatProvider

Bases: ABC

Source code in src/ciel/providers/__init__.py
class ChatProvider(ABC):
    provider_name: str = "unknown"

    @abstractmethod
    async def complete(self, request: "ChatRequest") -> "ChatResponse":
        raise NotImplementedError

    @abstractmethod
    async def stream(self, request: "ChatRequest") -> Sequence["ChatResponse"]:
        raise NotImplementedError

    @abstractmethod
    async def models(self) -> Sequence[ModelInfo]:
        raise NotImplementedError

ChatRequest dataclass

Source code in src/ciel/runtime/tools.py
@dataclass(frozen=True)
class ChatRequest:
    messages: Sequence[ChatMessage]
    tools: Sequence[ToolSpec] = ()
    model: Optional[str] = None
    temperature: Optional[float] = None
    max_tokens: Optional[int] = None
    extra: Dict[str, Any] = field(default_factory=dict)

ChatResponse dataclass

Source code in src/ciel/runtime/tools.py
@dataclass(frozen=True)
class ChatResponse:
    choice: ChatChoice
    metadata: Dict[str, Any] = field(default_factory=dict)

CielError

Bases: Exception

Base Ciel error.

Source code in src/ciel/common/__init__.py
class CielError(Exception):
    """Base Ciel error."""

DefaultAgentRuntime

Concrete runtime wiring provider + tool execution with tracing.

Source code in src/ciel/runtime/__init__.py
class DefaultAgentRuntime:
    """Concrete runtime wiring provider + tool execution with tracing."""

    def __init__(
        self,
        *,
        provider: ChatProvider,
        dispatcher: DefaultToolDispatcher,
        registry: Optional[ProviderRegistry] = None,
        audit_sink: Optional[Any] = None,
        agent: str = "default",
        approval_policy: Optional[Any] = None,
    ) -> None:
        self.provider = provider
        self.dispatcher = dispatcher
        self.registry = registry
        self.audit_sink = audit_sink or NullAuditSink()
        self.agent = agent
        self.approval_policy = approval_policy

    async def _emit(self, event: AuditEvent, *, tenant_id: Optional[str] = None) -> AuditEvent:
        normalized = propagate(event, tenant_id=tenant_id)
        await self.audit_sink.write(normalized)
        return normalized

    async def run_agent_loop(
        self,
        *,
        request: ChatRequest,
        tenant_id: Optional[str] = None,
        toolset: Optional[str] = None,
        limit: int = 32,
    ) -> AgentRuntimeResult:
        """Run the agentic tool loop (multi-turn ReAct when tools are present).

        The loop issues up to ``limit`` model completions. After each completion
        that requests tool calls, the requested tools are dispatched and their
        results are appended as messages; the model is then called again. The
        loop stops when (a) the model returns no tool calls (``finish_reason``
        ``stop``), (b) ``limit`` turns are exhausted, or (c) there are no tools
        registered in the request.

        Backward compatibility: when ``limit <= 1`` or ``request.tools`` is
        empty, exactly one completion is produced and no tool messages are
        appended (the historical single-step behaviour). Every tool result from
        every turn is preserved in ``loop_results`` so the facade can surface
        the full ``tool_results`` list.
        """
        session_id = request.extra.get("session_id") or str(uuid.uuid4())
        await self._emit(
            AuditEvent(event="agent.loop.start", session_id=session_id, agent=self.agent, tenant_id=tenant_id),
            tenant_id=tenant_id,
        )

        messages: List[ChatMessage] = list(request.messages)
        loop_results: List[ToolLoopResult] = []
        finish_reason = "stop"
        has_tools = bool(request.tools)

        # Multi-turn only when tools exist and the caller allows more than one step.
        max_turns = max(1, limit) if (has_tools and limit > 1) else 1

        final_response = None
        for _ in range(max_turns):
            response = await self.provider.complete(
                ChatRequest(
                    messages=tuple(messages),
                    tools=request.tools,
                    model=request.model,
                    temperature=request.temperature,
                    max_tokens=request.max_tokens,
                    extra={**request.extra, "session_id": session_id, "tenant_id": tenant_id},
                )
            )
            final_response = response
            messages.append(response.choice.message)
            tool_calls = _extract_tool_calls(response)

            if not tool_calls:
                # Model is done; stop the loop.
                finish_reason = response.choice.finish_reason or "stop"
                break

            dispatch_results: List[ToolResult] = []
            for call in tool_calls:
                call_arguments = call.get("arguments") or {}
                decision = None
                if self.approval_policy is not None and call.get("name") not in {None, ""}:
                    try:
                        from ciel.security import ApprovalRequest, ApprovalPolicy
                        if isinstance(self.approval_policy, type):
                            policy = self.approval_policy()
                        else:
                            policy = self.approval_policy
                        if hasattr(policy, "evaluate"):
                            decision = policy.evaluate(
                                ApprovalRequest(
                                    request_id=call.get("id") or call.get("tool_call_id") or str(uuid.uuid4()),
                                    actor=tenant_id or "unknown",
                                    tool=call.get("name", ""),
                                    arguments=call_arguments,
                                    risk="medium",
                                    tenant=tenant_id,
                                )
                            )
                    except Exception:
                        decision = None
                if decision is not None and not decision.approved:
                    dispatch_results.append(
                        ToolResult(
                            id=call.get("id") or call.get("tool_call_id") or str(uuid.uuid4()),
                            name=call.get("name", ""),
                            error=f"ApprovalDenied: {decision.note or 'denied'}",
                            metadata={"tenant_id": tenant_id, "approval_decision": decision.note or "denied"},
                        )
                    )
                    continue
                dispatch_results.append(
                    await self.dispatcher.dispatch(
                        tenant_id=tenant_id,
                        toolset=toolset or self.dispatcher.default_toolset or "default",
                        name=call.get("name", ""),
                        arguments=call_arguments,
                        tool_call_id=call.get("id") or call.get("tool_call_id") or str(uuid.uuid4()),
                    )
                )

            tool_turn = ToolLoopResult(
                turn_id=str(uuid.uuid4()),
                messages=tuple(messages),
                tool_results=tuple(dispatch_results),
                finish_reason="tool_calls",
                tenant_id=tenant_id,
                metadata={"session_id": session_id, "toolset": toolset, "turn": len(loop_results) + 1},
            )
            loop_results.append(tool_turn)
            finish_reason = "tool_calls"
            await self._emit(
                AuditEvent(
                    event="agent.tool_calls.dispatched",
                    session_id=session_id,
                    agent=self.agent,
                    tool_call_id=dispatch_results[0].id if dispatch_results else None,
                    data={"tools": [r.name for r in dispatch_results], "turn": len(loop_results)},
                    tenant_id=tenant_id,
                ),
                tenant_id=tenant_id,
            )

            # Append tool results as tool messages so the next completion sees them.
            for res in dispatch_results:
                messages.append(
                    ChatMessage(
                        role="tool",
                        content=str(res.output) if res.error is None else (res.error or ""),
                        name=res.name,
                        tool_call_id=res.id,
                        metadata={"tenant_id": tenant_id, "tool_result": True, "error": res.error},
                    )
                )

        await self._emit(
            AuditEvent(event="agent.loop.end", session_id=session_id, agent=self.agent, tenant_id=tenant_id),
            tenant_id=tenant_id,
        )

        return AgentRuntimeResult(
            response=final_response,
            loop_results=tuple(loop_results),
            tenant_id=tenant_id,
            metadata={"session_id": session_id, "agent": self.agent},
        )

    async def stream_agent_loop(
        self,
        *,
        request: ChatRequest,
        tenant_id: Optional[str] = None,
        toolset: Optional[str] = None,
        limit: int = 32,
    ) -> AsyncIterator[ToolLoopResult]:
        result = await self.run_agent_loop(
            request=request,
            tenant_id=tenant_id,
            toolset=toolset,
            limit=limit,
        )
        for turn in result.loop_results:
            yield turn

    async def stream_tokens(
        self,
        *,
        request: ChatRequest,
        tenant_id: Optional[str] = None,
        toolset: Optional[str] = None,
    ) -> AsyncIterator[str]:
        """Stream incremental assistant tokens from the provider.

        Calls ``provider.stream`` (real SSE streaming) and re-emits the
        partial ``content`` of each incremental :class:`ChatResponse` as it
        arrives, so callers see the answer grow token by token.
        """
        chunks = await self.provider.stream(request=request)
        prior = ""
        for chunk in chunks:
            content = chunk.choice.message.text()
            if content != prior:
                yield content
                prior = content

run_agent_loop(*, request: ChatRequest, tenant_id: Optional[str] = None, toolset: Optional[str] = None, limit: int = 32) -> AgentRuntimeResult async

Run the agentic tool loop (multi-turn ReAct when tools are present).

The loop issues up to limit model completions. After each completion that requests tool calls, the requested tools are dispatched and their results are appended as messages; the model is then called again. The loop stops when (a) the model returns no tool calls (finish_reason stop), (b) limit turns are exhausted, or (c) there are no tools registered in the request.

Backward compatibility: when limit <= 1 or request.tools is empty, exactly one completion is produced and no tool messages are appended (the historical single-step behaviour). Every tool result from every turn is preserved in loop_results so the facade can surface the full tool_results list.

Source code in src/ciel/runtime/__init__.py
async def run_agent_loop(
    self,
    *,
    request: ChatRequest,
    tenant_id: Optional[str] = None,
    toolset: Optional[str] = None,
    limit: int = 32,
) -> AgentRuntimeResult:
    """Run the agentic tool loop (multi-turn ReAct when tools are present).

    The loop issues up to ``limit`` model completions. After each completion
    that requests tool calls, the requested tools are dispatched and their
    results are appended as messages; the model is then called again. The
    loop stops when (a) the model returns no tool calls (``finish_reason``
    ``stop``), (b) ``limit`` turns are exhausted, or (c) there are no tools
    registered in the request.

    Backward compatibility: when ``limit <= 1`` or ``request.tools`` is
    empty, exactly one completion is produced and no tool messages are
    appended (the historical single-step behaviour). Every tool result from
    every turn is preserved in ``loop_results`` so the facade can surface
    the full ``tool_results`` list.
    """
    session_id = request.extra.get("session_id") or str(uuid.uuid4())
    await self._emit(
        AuditEvent(event="agent.loop.start", session_id=session_id, agent=self.agent, tenant_id=tenant_id),
        tenant_id=tenant_id,
    )

    messages: List[ChatMessage] = list(request.messages)
    loop_results: List[ToolLoopResult] = []
    finish_reason = "stop"
    has_tools = bool(request.tools)

    # Multi-turn only when tools exist and the caller allows more than one step.
    max_turns = max(1, limit) if (has_tools and limit > 1) else 1

    final_response = None
    for _ in range(max_turns):
        response = await self.provider.complete(
            ChatRequest(
                messages=tuple(messages),
                tools=request.tools,
                model=request.model,
                temperature=request.temperature,
                max_tokens=request.max_tokens,
                extra={**request.extra, "session_id": session_id, "tenant_id": tenant_id},
            )
        )
        final_response = response
        messages.append(response.choice.message)
        tool_calls = _extract_tool_calls(response)

        if not tool_calls:
            # Model is done; stop the loop.
            finish_reason = response.choice.finish_reason or "stop"
            break

        dispatch_results: List[ToolResult] = []
        for call in tool_calls:
            call_arguments = call.get("arguments") or {}
            decision = None
            if self.approval_policy is not None and call.get("name") not in {None, ""}:
                try:
                    from ciel.security import ApprovalRequest, ApprovalPolicy
                    if isinstance(self.approval_policy, type):
                        policy = self.approval_policy()
                    else:
                        policy = self.approval_policy
                    if hasattr(policy, "evaluate"):
                        decision = policy.evaluate(
                            ApprovalRequest(
                                request_id=call.get("id") or call.get("tool_call_id") or str(uuid.uuid4()),
                                actor=tenant_id or "unknown",
                                tool=call.get("name", ""),
                                arguments=call_arguments,
                                risk="medium",
                                tenant=tenant_id,
                            )
                        )
                except Exception:
                    decision = None
            if decision is not None and not decision.approved:
                dispatch_results.append(
                    ToolResult(
                        id=call.get("id") or call.get("tool_call_id") or str(uuid.uuid4()),
                        name=call.get("name", ""),
                        error=f"ApprovalDenied: {decision.note or 'denied'}",
                        metadata={"tenant_id": tenant_id, "approval_decision": decision.note or "denied"},
                    )
                )
                continue
            dispatch_results.append(
                await self.dispatcher.dispatch(
                    tenant_id=tenant_id,
                    toolset=toolset or self.dispatcher.default_toolset or "default",
                    name=call.get("name", ""),
                    arguments=call_arguments,
                    tool_call_id=call.get("id") or call.get("tool_call_id") or str(uuid.uuid4()),
                )
            )

        tool_turn = ToolLoopResult(
            turn_id=str(uuid.uuid4()),
            messages=tuple(messages),
            tool_results=tuple(dispatch_results),
            finish_reason="tool_calls",
            tenant_id=tenant_id,
            metadata={"session_id": session_id, "toolset": toolset, "turn": len(loop_results) + 1},
        )
        loop_results.append(tool_turn)
        finish_reason = "tool_calls"
        await self._emit(
            AuditEvent(
                event="agent.tool_calls.dispatched",
                session_id=session_id,
                agent=self.agent,
                tool_call_id=dispatch_results[0].id if dispatch_results else None,
                data={"tools": [r.name for r in dispatch_results], "turn": len(loop_results)},
                tenant_id=tenant_id,
            ),
            tenant_id=tenant_id,
        )

        # Append tool results as tool messages so the next completion sees them.
        for res in dispatch_results:
            messages.append(
                ChatMessage(
                    role="tool",
                    content=str(res.output) if res.error is None else (res.error or ""),
                    name=res.name,
                    tool_call_id=res.id,
                    metadata={"tenant_id": tenant_id, "tool_result": True, "error": res.error},
                )
            )

    await self._emit(
        AuditEvent(event="agent.loop.end", session_id=session_id, agent=self.agent, tenant_id=tenant_id),
        tenant_id=tenant_id,
    )

    return AgentRuntimeResult(
        response=final_response,
        loop_results=tuple(loop_results),
        tenant_id=tenant_id,
        metadata={"session_id": session_id, "agent": self.agent},
    )

stream_tokens(*, request: ChatRequest, tenant_id: Optional[str] = None, toolset: Optional[str] = None) -> AsyncIterator[str] async

Stream incremental assistant tokens from the provider.

Calls provider.stream (real SSE streaming) and re-emits the partial content of each incremental :class:ChatResponse as it arrives, so callers see the answer grow token by token.

Source code in src/ciel/runtime/__init__.py
async def stream_tokens(
    self,
    *,
    request: ChatRequest,
    tenant_id: Optional[str] = None,
    toolset: Optional[str] = None,
) -> AsyncIterator[str]:
    """Stream incremental assistant tokens from the provider.

    Calls ``provider.stream`` (real SSE streaming) and re-emits the
    partial ``content`` of each incremental :class:`ChatResponse` as it
    arrives, so callers see the answer grow token by token.
    """
    chunks = await self.provider.stream(request=request)
    prior = ""
    for chunk in chunks:
        content = chunk.choice.message.text()
        if content != prior:
            yield content
            prior = content

DefaultToolDispatcher

Dispatch tool requests to a configured ToolProvider.

Source code in src/ciel/runtime/__init__.py
class DefaultToolDispatcher:
    """Dispatch tool requests to a configured ToolProvider."""

    provider: ToolProvider
    default_toolset: Optional[str]

    def __init__(self, provider: ToolProvider, default_toolset: Optional[str] = None) -> None:
        self.provider = provider
        self.default_toolset = default_toolset or getattr(provider.registry, "default_toolset", None)

    async def dispatch(
        self,
        *,
        tenant_id: Optional[str] = None,
        toolset: Optional[str] = None,
        name: str,
        arguments: Dict[str, Any],
        tool_call_id: str,
    ) -> ToolResult:
        result = await self.provider.execute(
            tenant_id=tenant_id,
            toolset=toolset or self.default_toolset or "default",
            name=name,
            arguments=arguments,
            tool_call_id=tool_call_id,
        )
        result.metadata.setdefault("tenant_id", tenant_id)
        return result

    async def dispatch_all(
        self,
        *,
        tenant_id: Optional[str] = None,
        toolset: Optional[str] = None,
        calls: Sequence[Dict[str, Any]],
    ):
        results: List[ToolResult] = []
        for call in calls:
            call_tenant_id = tenant_id or call.get("metadata", {}).get("tenant_id")
            results.append(
                await self.dispatch(
                    tenant_id=call_tenant_id,
                    toolset=toolset or self.default_toolset or "default",
                    name=call["name"],
                    arguments=call.get("arguments", {}),
                    tool_call_id=call.get("id") or call.get("tool_call_id") or str(uuid.uuid4()),
                )
            )
        return results

InMemoryAuditSink

Bases: AuditSink

Source code in src/ciel/observability/__init__.py
class InMemoryAuditSink(AuditSink):
    def __init__(self) -> None:
        self.events: List[AuditEvent] = []

    async def write(self, event: AuditEvent) -> None:
        self.events.append(event)

ModelProvider

Model/provider contract for completions (legacy alias).

Source code in src/ciel/runtime/__init__.py
class ModelProvider:
    """Model/provider contract for completions (legacy alias)."""

    async def complete(self, request: ChatRequest) -> ChatResponse:
        raise NotImplementedError

    async def stream(self, request: ChatRequest) -> Sequence[ChatResponse]:
        raise NotImplementedError

NullAuditSink

Bases: AuditSink

Source code in src/ciel/observability/__init__.py
class NullAuditSink(AuditSink):
    async def write(self, event: AuditEvent) -> None:
        return

ProviderRegistry

Source code in src/ciel/providers/__init__.py
class ProviderRegistry:
    def __init__(self) -> None:
        self._providers: Dict[str, ChatProvider] = {}
        self._configs: Dict[str, Dict[str, Any]] = {}

    def register(self, name: str, provider: ChatProvider, *, config: Optional[Dict[str, Any]] = None) -> None:
        self._providers[name] = provider
        self._configs[name] = config or {}

    def get(self, name: str) -> ChatProvider:
        if name not in self._providers:
            raise _domain_error(f"Provider not registered: {name}")
        return self._providers[name]

    def available(self) -> Sequence[str]:
        return list(self._providers.keys())

SkillError

Bases: Exception

Base error for skill library operations.

Source code in src/ciel/runtime/skills_lib.py
class SkillError(Exception):
    """Base error for skill library operations."""

SkillLibrary

Writable, in-memory skill store that wraps a :class:SkillRegistry.

The registry remains the source of truth for disk-loaded skills; the library layer adds creation, registration, update (with version bump) and removal of skills that live only in memory. Tenant isolation is supported via the optional tenant_id key on each stored :class:Skill.

Source code in src/ciel/runtime/skills_lib.py
class SkillLibrary:
    """Writable, in-memory skill store that wraps a :class:`SkillRegistry`.

    The registry remains the source of truth for disk-loaded skills; the
    library layer adds creation, registration, update (with version bump) and
    removal of skills that live only in memory. Tenant isolation is supported
    via the optional ``tenant_id`` key on each stored :class:`Skill`.
    """

    def __init__(self, registry: Optional[SkillRegistry] = None) -> None:
        self.registry = registry or SkillRegistry()
        # name -> list of Skill (newest last). We keep a list so update() can
        # preserve previous_versions for the evolution tree.
        self._skills: Dict[str, List[Skill]] = {}

    # -- backed by the passive registry (backward-compatible) -----------------

    def load_from_disk(self) -> List[Skill]:
        """Discover skills from the registry roots and index them in-memory."""
        found = self.registry.discover()
        for skill in found:
            self._skills.setdefault(skill.name, []).append(skill)
        return found

    # -- writable store --------------------------------------------------------

    def _sha256(self, content: str) -> str:
        import hashlib

        return hashlib.sha256(content.encode("utf-8")).hexdigest()

    def create_from_code(
        self,
        *,
        name: str,
        description: str,
        code: str,
        category: Optional[str] = None,
        tenant_id: Optional[str] = None,
        metadata: Optional[Dict[str, Any]] = None,
    ) -> Skill:
        """Compile ``code`` (syntax check) and store it as a new skill.

        Raises :class:`SkillError` if the code does not compile. The skill is
        NOT executed here — execution belongs to :class:`SkillVerifier`.
        """
        try:
            compile(code, f"<skill:{name}>", "exec")
        except SyntaxError as exc:  # noqa: BLE001
            raise SkillError(f"skill '{name}' has invalid syntax: {exc}") from exc

        content = code
        skill = Skill(
            name=name,
            description=description,
            content=content,
            category=category,
            metadata={
                **(metadata or {}),
                **({"tenant_id": tenant_id} if tenant_id else {}),
                "sha256": self._sha256(content),
            },
            sha256=self._sha256(content),
        )
        self._skills.setdefault(name, []).append(skill)
        return skill

    def register(self, skill: Skill) -> Skill:
        """Register an already-built :class:`Skill` (e.g. disk-loaded)."""
        if skill.sha256 is None:
            skill.sha256 = self._sha256(skill.content)
        self._skills.setdefault(skill.name, []).append(skill)
        return skill

    def get(self, name: str) -> Optional[Skill]:
        versions = self._skills.get(name)
        if not versions:
            # Fall back to the disk registry (backward-compat lookup).
            return self.registry.get(name)
        return versions[-1]

    def list_skills(self, *, category: Optional[str] = None, tenant_id: Optional[str] = None) -> List[Skill]:
        out: List[Skill] = []
        for versions in self._skills.values():
            latest = versions[-1]
            if category is not None and latest.category != category:
                continue
            if tenant_id is not None and latest.metadata.get("tenant_id") != tenant_id:
                continue
            out.append(latest)
        # Also include any registry-only skills not yet indexed in memory.
        for skill in self.registry.list_skills(category=category):
            if skill.name not in self._skills:
                if tenant_id is None or skill.metadata.get("tenant_id") == tenant_id:
                    out.append(skill)
        return out

    def history(self, name: str) -> List[Skill]:
        """Return every stored version of ``name`` (oldest first)."""
        return list(self._skills.get(name, []))

    def remove(self, name: str) -> bool:
        if name in self._skills:
            del self._skills[name]
            return True
        return False

    def update(
        self,
        *,
        name: str,
        description: Optional[str] = None,
        code: Optional[str] = None,
        category: Optional[str] = None,
        bump: str = "patch",
    ) -> Skill:
        """Create a new version of an existing skill, preserving history.

        ``bump`` is one of ``major``/``minor``/``patch`` (semantic) and is
        recorded in ``metadata.version``. The previous version is preserved in
        ``history(name)``.
        """
        versions = self._skills.get(name)
        if not versions:
            disk = self.registry.get(name)
            if disk is None:
                raise SkillError(f"cannot update unknown skill '{name}'")
            versions = [disk]
            self._skills[name] = versions

        previous = versions[-1]
        if code is None:
            code = previous.content
        else:
            try:
                compile(code, f"<skill:{name}>", "exec")
            except SyntaxError as exc:  # noqa: BLE001
                raise SkillError(f"skill '{name}' update has invalid syntax: {exc}") from exc

        new_version = self._next_version(previous.metadata.get("version"), bump)
        updated = Skill(
            name=name,
            description=description if description is not None else previous.description,
            content=code,
            category=category if category is not None else previous.category,
            metadata={
                **previous.metadata,
                "version": new_version,
                "previous_version": previous.metadata.get("version"),
                "sha256": self._sha256(code),
            },
            sha256=self._sha256(code),
        )
        versions.append(updated)
        return updated

    @staticmethod
    def _next_version(current: Optional[str], bump: str) -> str:
        major, minor, patch = (0, 0, 0)
        if current:
            parts = (current.split(".") + ["0", "0", "0"])[:3]
            try:
                major, minor, patch = (int(p) for p in parts)
            except ValueError:
                major, minor, patch = (0, 0, 0)
        if bump == "major":
            major, minor, patch = major + 1, 0, 0
        elif bump == "minor":
            minor, patch = minor + 1, 0
        else:
            patch += 1
        return f"{major}.{minor}.{patch}"

create_from_code(*, name: str, description: str, code: str, category: Optional[str] = None, tenant_id: Optional[str] = None, metadata: Optional[Dict[str, Any]] = None) -> Skill

Compile code (syntax check) and store it as a new skill.

Raises :class:SkillError if the code does not compile. The skill is NOT executed here — execution belongs to :class:SkillVerifier.

Source code in src/ciel/runtime/skills_lib.py
def create_from_code(
    self,
    *,
    name: str,
    description: str,
    code: str,
    category: Optional[str] = None,
    tenant_id: Optional[str] = None,
    metadata: Optional[Dict[str, Any]] = None,
) -> Skill:
    """Compile ``code`` (syntax check) and store it as a new skill.

    Raises :class:`SkillError` if the code does not compile. The skill is
    NOT executed here — execution belongs to :class:`SkillVerifier`.
    """
    try:
        compile(code, f"<skill:{name}>", "exec")
    except SyntaxError as exc:  # noqa: BLE001
        raise SkillError(f"skill '{name}' has invalid syntax: {exc}") from exc

    content = code
    skill = Skill(
        name=name,
        description=description,
        content=content,
        category=category,
        metadata={
            **(metadata or {}),
            **({"tenant_id": tenant_id} if tenant_id else {}),
            "sha256": self._sha256(content),
        },
        sha256=self._sha256(content),
    )
    self._skills.setdefault(name, []).append(skill)
    return skill

history(name: str) -> List[Skill]

Return every stored version of name (oldest first).

Source code in src/ciel/runtime/skills_lib.py
def history(self, name: str) -> List[Skill]:
    """Return every stored version of ``name`` (oldest first)."""
    return list(self._skills.get(name, []))

load_from_disk() -> List[Skill]

Discover skills from the registry roots and index them in-memory.

Source code in src/ciel/runtime/skills_lib.py
def load_from_disk(self) -> List[Skill]:
    """Discover skills from the registry roots and index them in-memory."""
    found = self.registry.discover()
    for skill in found:
        self._skills.setdefault(skill.name, []).append(skill)
    return found

register(skill: Skill) -> Skill

Register an already-built :class:Skill (e.g. disk-loaded).

Source code in src/ciel/runtime/skills_lib.py
def register(self, skill: Skill) -> Skill:
    """Register an already-built :class:`Skill` (e.g. disk-loaded)."""
    if skill.sha256 is None:
        skill.sha256 = self._sha256(skill.content)
    self._skills.setdefault(skill.name, []).append(skill)
    return skill

update(*, name: str, description: Optional[str] = None, code: Optional[str] = None, category: Optional[str] = None, bump: str = 'patch') -> Skill

Create a new version of an existing skill, preserving history.

bump is one of major/minor/patch (semantic) and is recorded in metadata.version. The previous version is preserved in history(name).

Source code in src/ciel/runtime/skills_lib.py
def update(
    self,
    *,
    name: str,
    description: Optional[str] = None,
    code: Optional[str] = None,
    category: Optional[str] = None,
    bump: str = "patch",
) -> Skill:
    """Create a new version of an existing skill, preserving history.

    ``bump`` is one of ``major``/``minor``/``patch`` (semantic) and is
    recorded in ``metadata.version``. The previous version is preserved in
    ``history(name)``.
    """
    versions = self._skills.get(name)
    if not versions:
        disk = self.registry.get(name)
        if disk is None:
            raise SkillError(f"cannot update unknown skill '{name}'")
        versions = [disk]
        self._skills[name] = versions

    previous = versions[-1]
    if code is None:
        code = previous.content
    else:
        try:
            compile(code, f"<skill:{name}>", "exec")
        except SyntaxError as exc:  # noqa: BLE001
            raise SkillError(f"skill '{name}' update has invalid syntax: {exc}") from exc

    new_version = self._next_version(previous.metadata.get("version"), bump)
    updated = Skill(
        name=name,
        description=description if description is not None else previous.description,
        content=code,
        category=category if category is not None else previous.category,
        metadata={
            **previous.metadata,
            "version": new_version,
            "previous_version": previous.metadata.get("version"),
            "sha256": self._sha256(code),
        },
        sha256=self._sha256(code),
    )
    versions.append(updated)
    return updated

SkillVerificationError

Bases: SkillError

Raised when a skill fails verification.

Source code in src/ciel/runtime/skills_lib.py
class SkillVerificationError(SkillError):
    """Raised when a skill fails verification."""

SkillVerificationResult dataclass

Outcome of :meth:SkillVerifier.verify.

Source code in src/ciel/runtime/skills_lib.py
@dataclass
class SkillVerificationResult:
    """Outcome of :meth:`SkillVerifier.verify`."""

    passed: bool
    skill: str
    attempts: int = 0
    error: Optional[str] = None
    traceback: Optional[str] = None
    expected: Optional[Any] = None
    got: Optional[Any] = None

SkillVerifier

Offline verifier: syntax check + executable test cases.

A test case is a dict {"call": {...}, "expect": <value>}. The verifier executes the skill code in an isolated namespace, looks up a callable named after the skill (or the first callable defined), invokes it with call arguments and compares the result to expect.

Source code in src/ciel/runtime/skills_lib.py
class SkillVerifier:
    """Offline verifier: syntax check + executable test cases.

    A test case is a dict ``{"call": {...}, "expect": <value>}``. The verifier
    executes the skill code in an isolated namespace, looks up a callable named
    after the skill (or the first callable defined), invokes it with ``call``
    arguments and compares the result to ``expect``.
    """

    def __init__(self, library: Optional[SkillLibrary] = None) -> None:
        self.library = library

    def _resolve_callable(self, skill: Skill, namespace: Dict[str, Any]) -> Any:
        if skill.name in namespace and callable(namespace[skill.name]):
            return namespace[skill.name]
        callables = [v for v in namespace.values() if callable(v) and not v.__module__ == "builtins"]
        if not callables:
            raise SkillVerificationError(f"skill '{skill.name}' defines no callable to invoke")
        return callables[0]

    def verify(self, skill: Skill, *, test_cases: Sequence[Dict[str, Any]]) -> SkillVerificationResult:
        # 1) Syntax validation first.
        try:
            compile(skill.content, f"<skill:{skill.name}>", "exec")
        except SyntaxError as exc:  # noqa: BLE001
            return SkillVerificationResult(
                passed=False,
                skill=skill.name,
                attempts=0,
                error=f"syntax error: {exc}",
                traceback=str(exc),
            )

        namespace: Dict[str, Any] = {}
        try:
            exec(skill.content, namespace)  # noqa: S102 — offline, trusted-by-construction
        except Exception as exc:  # noqa: BLE001
            return SkillVerificationResult(
                passed=False,
                skill=skill.name,
                attempts=0,
                error=f"load error: {type(exc).__name__}: {exc}",
                traceback=repr(exc),
            )

        fn = self._resolve_callable(skill, namespace)

        last_traceback: Optional[str] = None
        for attempt, case in enumerate(test_cases, start=1):
            call_args = case.get("call", {}) or {}
            expected = case.get("expect")
            try:
                got = fn(**call_args)
            except Exception as exc:  # noqa: BLE001
                last_traceback = repr(exc)
                return SkillVerificationResult(
                    passed=False,
                    skill=skill.name,
                    attempts=attempt,
                    error=f"case {attempt} raised {type(exc).__name__}: {exc}",
                    traceback=last_traceback,
                    expected=expected,
                )
            if got != expected:
                return SkillVerificationResult(
                    passed=False,
                    skill=skill.name,
                    attempts=attempt,
                    error=f"case {attempt}: expected {expected!r}, got {got!r}",
                    expected=expected,
                    got=got,
                )
        return SkillVerificationResult(
            passed=True,
            skill=skill.name,
            attempts=len(test_cases),
        )

    def verify_by_name(self, name: str, *, test_cases: Sequence[Dict[str, Any]]) -> SkillVerificationResult:
        if self.library is None:
            raise SkillVerificationError("SkillVerifier was built without a library; pass the skill directly")
        skill = self.library.get(name)
        if skill is None:
            raise SkillVerificationError(f"unknown skill '{name}'")
        return self.verify(skill, test_cases=test_cases)

StaticToolProvider

Bases: ToolProvider

Source code in src/ciel/runtime/__init__.py
class StaticToolProvider(ToolProvider):
    def __init__(self, tools: Mapping[str, Sequence[ToolSpec]], *, require_tenant: bool = False) -> None:
        registry = ToolRegistry(default_toolset="default")
        for toolset, values in tools.items():
            registry.register_toolset(
                ToolsetSchema(
                    name=toolset,
                    description="",
                    tools=tuple(values) if not isinstance(values, tuple) else values,
                    require_tenant=require_tenant,
                )
            )
        super().__init__(registry=registry, require_tenant_on_execution=require_tenant)

TenantRequired

Bases: CielError

Tenant context is required for the requested operation.

Source code in src/ciel/common/__init__.py
class TenantRequired(CielError):
    """Tenant context is required for the requested operation."""

Tool dataclass

Source code in src/ciel/runtime/tools.py
@dataclass
class Tool:
    spec: ToolSpec
    callable_: Any = None
    metadata: Dict[str, Any] = field(default_factory=dict)
    required_tenant: bool = False

ToolExecutionContext

Source code in src/ciel/runtime/tools.py
class ToolExecutionContext:
    def __init__(self, *, tenant_id, toolset, name, tool_call_id, arguments):
        if not tenant_id:
            raise ValueError(f"tenant_id is required to execute tool '{name}' in toolset '{toolset}'.")
        self.context = ToolCallContext(
            tenant_id=tenant_id,
            toolset=toolset,
            tool_name=name,
            tool_call_id=tool_call_id,
            arguments=arguments,
        )

    @property
    def tenant_id(self) -> Optional[str]:
        return self.context.tenant_id

    @property
    def toolset(self) -> str:
        return self.context.toolset

    @property
    def name(self) -> str:
        return self.context.tool_name

ToolLoopResult dataclass

Source code in src/ciel/runtime/tools.py
@dataclass(frozen=True)
class ToolLoopResult:
    turn_id: str
    messages: Sequence[ChatMessage]
    tool_results: Sequence[ToolResult]
    finish_reason: str
    tenant_id: Optional[str] = None
    metadata: Dict[str, Any] = field(default_factory=dict)

ToolProvider dataclass

Concrete tool provider used by the runtime's tool dispatcher.

Source code in src/ciel/runtime/__init__.py
@dataclass(frozen=True)
class ToolProvider:
    """Concrete tool provider used by the runtime's tool dispatcher."""

    registry: ToolRegistry
    require_tenant_on_execution: bool = True

    async def tool_specs(self, tenant_id: Optional[str], toolset: str) -> Sequence[ToolSpec]:
        return tuple(getattr(self.registry, "_toolsets", {}).get(toolset, ToolsetSchema(name=toolset, description="")).tools)

    async def execute(
        self,
        *,
        tenant_id: Optional[str],
        toolset: Optional[str],
        name: str,
        arguments: Dict[str, Any],
        tool_call_id: str,
    ) -> ToolResult:
        target_toolset = toolset or (self.registry.default_toolset or "default")
        if self.require_tenant_on_execution and not tenant_id:
            raise TenantRequired(f"tenant_id is required to execute tool '{name}' in toolset '{target_toolset}'.")
        tool = self.registry.get_tool(toolset=target_toolset, name=name)
        if tool is None:
            return ToolResult(id=tool_call_id, name=name, error=f"unknown tool: {target_toolset}.{name}", metadata={"tenant_id": tenant_id})
        if getattr(tool, "required_tenant", False) and not tenant_id:
            raise TenantRequired(f"tool '{name}' requires tenant_id, but none was provided.")
        if tool.callable_ is None:
            output = {"arguments": arguments, "description": tool.spec.description}
            return ToolResult(id=tool_call_id, name=name, output=output, metadata={"tenant_id": tenant_id})
        # Official tool callable contract:
        #   callable_(arguments: dict, *, tool_call_id: str, tenant_id: str | None) -> ToolResult | dict | Any
        try:
            result = tool.callable_(arguments, tool_call_id=tool_call_id, tenant_id=tenant_id)
            if inspect.isawaitable(result):
                result = await result
        except Exception as exc:  # noqa: BLE001 — surface tool errors as ToolResult
            return ToolResult(id=tool_call_id, name=name, error=f"{type(exc).__name__}: {exc}", metadata={"tenant_id": tenant_id})
        if isinstance(result, ToolResult):
            if not result.metadata.get("tenant_id"):
                result.metadata["tenant_id"] = tenant_id
            return result
        return ToolResult(id=tool_call_id, name=name, output=result, metadata={"tenant_id": tenant_id})

ToolRegistry

Source code in src/ciel/runtime/tools.py
class ToolRegistry:
    def __init__(self, *, default_toolset: Optional[str] = None) -> None:
        self._toolsets: Dict[str, ToolsetSchema] = {}
        self._tools: Dict[str, Dict[str, Tool]] = {}
        self._tenant_tools: Dict[str, Dict[str, Dict[str, Tool]]] = {}
        self.default_toolset = default_toolset

    def register_toolset(self, schema: ToolsetSchema) -> None:
        self._toolsets[schema.name] = schema
        self._tools.setdefault(schema.name, {})
        self._tenant_tools.setdefault(schema.name, {})
        for tool in schema.tools:
            self._tools[schema.name][tool.name] = Tool(spec=tool, required_tenant=schema.require_tenant)

    def register_tool(self, toolset, tool, *, tenant_id=None):
        if isinstance(toolset, ToolsetSchema):
            schema = toolset
            toolset = schema.name
            self._toolsets[schema.name] = schema
            self._tools.setdefault(schema.name, {})
            self._tenant_tools.setdefault(schema.name, {})
            for tool in schema.tools:
                self._tools[schema.name][tool.name] = Tool(spec=tool, required_tenant=schema.require_tenant)
            return
        schema = self._toolsets.get(toolset)
        require_tenant = schema.require_tenant if schema is not None else tool.required_tenant
        tool_obj = Tool(spec=tool.spec, callable_=tool.callable_, metadata=tool.metadata, required_tenant=tool.required_tenant or require_tenant)
        if schema is None:
            schema = ToolsetSchema(name=toolset, description="", require_tenant=require_tenant)
            self._toolsets[toolset] = schema
        self._tools.setdefault(toolset, {})
        self._tools[toolset][tool.spec.name] = tool_obj
        # Keep the schema's tool list in sync so get_toolset_schema/export_schema reflect registered tools.
        self._toolsets[toolset] = replace(schema, tools=tuple(t.spec for t in self._tools[toolset].values()))
        if tenant_id:
            self._tenant_tools.setdefault(toolset, {}).setdefault(tenant_id, {})[tool.spec.name] = tool_obj

    def get_toolset(self, name):
        return self._toolsets.get(name)

    def get_tool(self, toolset: str, name: str, tenant_id: Optional[str] = None):
        if tenant_id:
            tenant_tools = self._tenant_tools.get(toolset, {}).get(tenant_id)
            if tenant_tools is None:
                raise ValueError(f"tenant_id='{tenant_id}' has no mapped tools for toolset='{toolset}'")
            tool = tenant_tools.get(name)
            if tool is not None:
                return tool
        tools = self._tools.get(toolset)
        if not tools:
            return None
        return tools.get(name)

    def toolset_names(self) -> Sequence[str]:
        return tuple(self._toolsets.keys())

    def tool_names(self, toolset: str) -> Sequence[str]:
        return tuple(self._tools.get(toolset, {}).keys())

    def export_schema(self, toolset: str) -> Dict[str, Any]:
        schema = self._toolsets.get(toolset)
        if schema is None:
            raise KeyError(f"unknown toolset: {toolset}")
        return schema.to_json()

    async def lookup(self, *, tenant_id: Optional[str], toolset: str) -> Sequence[Tool]:
        schema = self._toolsets.get(toolset)
        if schema is None:
            return ()
        effective_tenant = schema.tenant_for(caller_tenant_id=tenant_id)
        if effective_tenant:
            tools = self._tenant_tools.get(toolset, {}).get(effective_tenant)
            if tools:
                return tuple(tools.values())
        return tuple(self._tools.get(toolset, {}).values())

ToolResult dataclass

Source code in src/ciel/runtime/tools.py
@dataclass
class ToolResult:
    id: str
    name: str
    output: Any = None
    error: Optional[str] = None
    usage: Optional[Dict[str, Any]] = None
    duration_ms: Optional[int] = None
    metadata: Dict[str, Any] = field(default_factory=dict)

ToolSpec dataclass

Source code in src/ciel/runtime/tools.py
@dataclass(frozen=True)
class ToolSpec:
    name: str
    description: str
    parameters: Mapping[str, Any]
    strict: bool = False
    metadata: Dict[str, Any] = field(default_factory=dict)

ToolsetSchema dataclass

Source code in src/ciel/runtime/tools.py
@dataclass
class ToolsetSchema:
    name: str
    description: str
    version: str = "1.0.0"
    tenants: Sequence[str] = ()
    default_tenant: Optional[str] = None
    require_tenant: bool = False
    tools: Sequence[ToolSpec] = ()

    def tenant_for(self, *, caller_tenant_id: Optional[str]) -> Optional[str]:
        if self.require_tenant and not caller_tenant_id:
            raise ValueError(f"Toolset '{self.name}' requires tenant_id; none was provided.")
        return caller_tenant_id or self.default_tenant

    def to_json(self) -> Dict[str, Any]:
        payload: Dict[str, Any] = {
            "name": self.name,
            "description": self.description,
            "version": self.version,
            "tools": [
                {
                    "name": tool.name,
                    "description": tool.description,
                    "parameters": tool.parameters,
                    "strict": tool.strict,
                    **( {"metadata": tool.metadata} if tool.metadata else {} ),
                }
                for tool in self.tools
            ],
        }
        if self.require_tenant:
            payload["require_tenant"] = True
        if self.tenants:
            payload["tenants"] = list(self.tenants)
        return payload

_extract_tool_calls(response: ChatResponse) -> List[Dict[str, Any]]

Source code in src/ciel/runtime/__init__.py
def _extract_tool_calls(response: ChatResponse) -> List[Dict[str, Any]]:
    metadata = response.metadata or {}
    raw = metadata.get("tool_calls")
    if isinstance(raw, list):
        return raw
    message_tool_calls = getattr(response.choice.message, "tool_calls", None)
    if isinstance(message_tool_calls, list):
        return message_tool_calls
    return []

assert_tenant_event(event: AuditEvent) -> None

Source code in src/ciel/observability/__init__.py
def assert_tenant_event(event: AuditEvent) -> None:
    if event.tenant_id is None:
        raise ValueError("AuditEvent requires tenant_id for multi-tenancy tracing")

propagate(event: AuditEvent, *, tenant_id: Optional[str] = None) -> AuditEvent

Source code in src/ciel/observability/__init__.py
def propagate(event: AuditEvent, *, tenant_id: Optional[str] = None) -> AuditEvent:
    normalized_tenant_id = tenant_id or event.tenant_id
    if normalized_tenant_id is None:
        raise ValueError(
            "propagate() requires tenant_id to be passed explicitly or present on event"
        )
    if tenant_id is not None:
        event.tenant_id = normalized_tenant_id
    return event

ChatContent = 'str | list[ContentPart]' module-attribute

ContentPart = Dict[str, Any] module-attribute

__all__ = ['ToolSpec', 'Tool', 'ToolResult', 'ToolProvider', 'StaticToolProvider', 'DefaultToolDispatcher', 'ToolCallContext', 'ToolsetSchema', 'ToolExecutionContext', 'ToolRegistry', 'TenantAwareToolProvider', 'ChatMessage', 'ChatChoice', 'ChatRequest', 'ChatResponse', 'ModelProvider', 'ToolLoopResult', 'AgentRuntimeResult', 'AgentContext', 'AgentRuntime'] module-attribute

AgentContext dataclass

Source code in src/ciel/runtime/tools.py
@dataclass(frozen=True)
class AgentContext:
    agent: str
    session_id: str
    tenant_id: Optional[str] = None
    metadata: Dict[str, Any] = field(default_factory=dict)

AgentRuntime

Async runtime contract for tool-loop execution and streaming.

Source code in src/ciel/runtime/tools.py
class AgentRuntime:
    """Async runtime contract for tool-loop execution and streaming."""

    async def run_agent_loop(
        self,
        *,
        request: ChatRequest,
        tenant_id: Optional[str] = None,
        toolset: Optional[str] = None,
        limit: int = 32,
    ) -> AgentRuntimeResult:
        raise NotImplementedError

    async def stream_agent_loop(
        self,
        *,
        request: ChatRequest,
        tenant_id: Optional[str] = None,
        toolset: Optional[str] = None,
        limit: int = 32,
    ):
        raise NotImplementedError

AgentRuntimeResult dataclass

Source code in src/ciel/runtime/tools.py
@dataclass(frozen=True)
class AgentRuntimeResult:
    response: ChatResponse
    loop_results: Sequence[ToolLoopResult]
    tenant_id: Optional[str] = None
    metadata: Dict[str, Any] = field(default_factory=dict)

ChatChoice dataclass

Source code in src/ciel/runtime/tools.py
@dataclass(frozen=True)
class ChatChoice:
    message: ChatMessage
    finish_reason: str
    usage: Optional[Dict[str, Any]] = None
    metadata: Dict[str, Any] = field(default_factory=dict)

ChatMessage dataclass

Source code in src/ciel/runtime/tools.py
@dataclass(frozen=True)
class ChatMessage:
    role: str
    content: ChatContent
    name: Optional[str] = None
    tool_call_id: Optional[str] = None
    tool_calls: Optional[list[dict[str, Any]]] = None
    metadata: Dict[str, Any] = field(default_factory=dict)

    def text(self) -> str:
        """Extract plain text from content, tolerant to multimodal parts.

        - ``str`` content is returned verbatim.
        - ``list`` content concatenates the ``text`` of every part whose type
          is ``"text"`` but drops images/audio/video, so consumers (CLI,
          compression, ``AgentResponse.text``) see only readable text.
        """
        content = self.content
        if isinstance(content, str):
            return content
        if not isinstance(content, list):
            return ""
        parts: list[str] = []
        for part in content:
            if isinstance(part, dict) and part.get("type") == "text":
                text = part.get("text", "")
                if isinstance(text, str):
                    parts.append(text)
        return "".join(parts)

text() -> str

Extract plain text from content, tolerant to multimodal parts.

  • str content is returned verbatim.
  • list content concatenates the text of every part whose type is "text" but drops images/audio/video, so consumers (CLI, compression, AgentResponse.text) see only readable text.
Source code in src/ciel/runtime/tools.py
def text(self) -> str:
    """Extract plain text from content, tolerant to multimodal parts.

    - ``str`` content is returned verbatim.
    - ``list`` content concatenates the ``text`` of every part whose type
      is ``"text"`` but drops images/audio/video, so consumers (CLI,
      compression, ``AgentResponse.text``) see only readable text.
    """
    content = self.content
    if isinstance(content, str):
        return content
    if not isinstance(content, list):
        return ""
    parts: list[str] = []
    for part in content:
        if isinstance(part, dict) and part.get("type") == "text":
            text = part.get("text", "")
            if isinstance(text, str):
                parts.append(text)
    return "".join(parts)

ChatRequest dataclass

Source code in src/ciel/runtime/tools.py
@dataclass(frozen=True)
class ChatRequest:
    messages: Sequence[ChatMessage]
    tools: Sequence[ToolSpec] = ()
    model: Optional[str] = None
    temperature: Optional[float] = None
    max_tokens: Optional[int] = None
    extra: Dict[str, Any] = field(default_factory=dict)

ChatResponse dataclass

Source code in src/ciel/runtime/tools.py
@dataclass(frozen=True)
class ChatResponse:
    choice: ChatChoice
    metadata: Dict[str, Any] = field(default_factory=dict)

DefaultToolDispatcher

Dispatch tool requests to a configured ToolProvider.

Source code in src/ciel/runtime/tools.py
class DefaultToolDispatcher:
    """Dispatch tool requests to a configured ToolProvider."""

    def __init__(self, provider: ToolProvider, default_toolset: Optional[str] = None) -> None:
        self.provider = provider
        self.default_toolset = default_toolset or getattr(provider.registry, "default_toolset", None) or "default"

    async def dispatch(
        self,
        *,
        toolset: Optional[str] = None,
        name: str,
        arguments: Dict[str, Any],
        tool_call_id: str,
        tenant_id: Optional[str] = None,
    ) -> ToolResult:
        result = await self.provider.execute(
            toolset=toolset or self.default_toolset,
            name=name,
            arguments=arguments,
            tool_call_id=tool_call_id,
            tenant_id=tenant_id,
        )
        result.metadata.setdefault("tenant_id", tenant_id)
        return result

    async def dispatch_all(
        self,
        *,
        toolset: Optional[str] = None,
        calls: Sequence[Dict[str, Any]],
        base_tenant_id: Optional[str] = None,
    ) -> Sequence[ToolResult]:
        results = []
        for call in calls:
            call_tenant_id = call.get("metadata", {}).get("tenant_id") or base_tenant_id
            results.append(
                await self.dispatch(
                    toolset=toolset or self.default_toolset,
                    name=call["name"],
                    arguments=call.get("arguments", {}),
                    tool_call_id=call.get("id") or call.get("tool_call_id"),
                    tenant_id=call_tenant_id,
                )
            )
        return results

ModelProvider

Model/provider contract for completions.

Source code in src/ciel/runtime/tools.py
class ModelProvider:
    """Model/provider contract for completions."""

    async def complete(self, request: ChatRequest) -> ChatResponse:
        raise NotImplementedError

    async def stream(self, request: ChatRequest) -> Sequence[ChatResponse]:
        raise NotImplementedError

StaticToolProvider

Bases: ToolProvider

Source code in src/ciel/runtime/tools.py
class StaticToolProvider(ToolProvider):
    def __init__(self, tools, *, require_tenant=False):
        registry = ToolRegistry(default_toolset="default")
        for toolset_name, values in tools.items():
            specs = tuple(values) if not isinstance(values, tuple) else values
            registry.register_toolset(
                ToolsetSchema(
                    name=toolset_name,
                    description="",
                    tools=specs,
                    require_tenant=require_tenant,
                )
            )
        self.registry = registry
        self.require_tenant = require_tenant

    async def tool_specs(self, tenant_id, toolset):
        return tuple(self.registry._toolsets.get(toolset, ToolsetSchema(name=toolset, description="")).tools)

    async def execute(self, *, toolset, name, arguments, tool_call_id, tenant_id=None):
        target_toolset = toolset or self.registry.default_toolset or "default"
        if self.require_tenant and not tenant_id:
            raise ValueError(f"tenant_id is required to execute tool '{name}'.")
        tool = self.registry.get_tool(toolset=target_toolset, name=name)
        if tool is None:
            return ToolResult(id=tool_call_id, name=name, error=f"tool not found: {target_toolset}.{name}", metadata={"tenant_id": tenant_id})
        try:
            result = tool.callable_(arguments, tool_call_id=tool_call_id, tenant_id=tenant_id)
            if asyncio.iscoroutine(result):
                result = await result
        except Exception as exc:  # noqa: BLE001 — surface tool errors as ToolResult
            return ToolResult(id=tool_call_id, name=name, error=f"{type(exc).__name__}: {exc}", metadata={"tenant_id": tenant_id})
        if isinstance(result, ToolResult):
            return result
        return ToolResult(id=tool_call_id, name=name, output=result, metadata={"tenant_id": tenant_id})

TenantAwareToolProvider

Tenant-aware provider contract.

Source code in src/ciel/runtime/tools.py
class TenantAwareToolProvider:
    """Tenant-aware provider contract."""

    async def tool_specs(self, tenant_id, toolset):
        raise NotImplementedError

    async def execute(self, *, context):
        raise NotImplementedError

Tool dataclass

Source code in src/ciel/runtime/tools.py
@dataclass
class Tool:
    spec: ToolSpec
    callable_: Any = None
    metadata: Dict[str, Any] = field(default_factory=dict)
    required_tenant: bool = False

ToolCallContext dataclass

Source code in src/ciel/runtime/tools.py
@dataclass
class ToolCallContext:
    tenant_id: Optional[str]
    toolset: str
    tool_name: str
    tool_call_id: str
    arguments: Dict[str, Any]

    def at(self, *, tenant_id: Optional[str]) -> ToolCallContext:
        return ToolCallContext(
            tenant_id=tenant_id,
            toolset=self.toolset,
            tool_name=self.tool_name,
            tool_call_id=self.tool_call_id,
            arguments=self.arguments,
        )

ToolExecutionContext

Source code in src/ciel/runtime/tools.py
class ToolExecutionContext:
    def __init__(self, *, tenant_id, toolset, name, tool_call_id, arguments):
        if not tenant_id:
            raise ValueError(f"tenant_id is required to execute tool '{name}' in toolset '{toolset}'.")
        self.context = ToolCallContext(
            tenant_id=tenant_id,
            toolset=toolset,
            tool_name=name,
            tool_call_id=tool_call_id,
            arguments=arguments,
        )

    @property
    def tenant_id(self) -> Optional[str]:
        return self.context.tenant_id

    @property
    def toolset(self) -> str:
        return self.context.toolset

    @property
    def name(self) -> str:
        return self.context.tool_name

ToolLoopResult dataclass

Source code in src/ciel/runtime/tools.py
@dataclass(frozen=True)
class ToolLoopResult:
    turn_id: str
    messages: Sequence[ChatMessage]
    tool_results: Sequence[ToolResult]
    finish_reason: str
    tenant_id: Optional[str] = None
    metadata: Dict[str, Any] = field(default_factory=dict)

ToolProvider

Contract: discover tools and execute tool calls.

Source code in src/ciel/runtime/tools.py
class ToolProvider:
    """Contract: discover tools and execute tool calls."""

    async def tool_specs(self, tenant_id, toolset):
        raise NotImplementedError

    async def execute(self, *, tenant_id, toolset, name, arguments, tool_call_id):
        raise NotImplementedError

ToolRegistry

Source code in src/ciel/runtime/tools.py
class ToolRegistry:
    def __init__(self, *, default_toolset: Optional[str] = None) -> None:
        self._toolsets: Dict[str, ToolsetSchema] = {}
        self._tools: Dict[str, Dict[str, Tool]] = {}
        self._tenant_tools: Dict[str, Dict[str, Dict[str, Tool]]] = {}
        self.default_toolset = default_toolset

    def register_toolset(self, schema: ToolsetSchema) -> None:
        self._toolsets[schema.name] = schema
        self._tools.setdefault(schema.name, {})
        self._tenant_tools.setdefault(schema.name, {})
        for tool in schema.tools:
            self._tools[schema.name][tool.name] = Tool(spec=tool, required_tenant=schema.require_tenant)

    def register_tool(self, toolset, tool, *, tenant_id=None):
        if isinstance(toolset, ToolsetSchema):
            schema = toolset
            toolset = schema.name
            self._toolsets[schema.name] = schema
            self._tools.setdefault(schema.name, {})
            self._tenant_tools.setdefault(schema.name, {})
            for tool in schema.tools:
                self._tools[schema.name][tool.name] = Tool(spec=tool, required_tenant=schema.require_tenant)
            return
        schema = self._toolsets.get(toolset)
        require_tenant = schema.require_tenant if schema is not None else tool.required_tenant
        tool_obj = Tool(spec=tool.spec, callable_=tool.callable_, metadata=tool.metadata, required_tenant=tool.required_tenant or require_tenant)
        if schema is None:
            schema = ToolsetSchema(name=toolset, description="", require_tenant=require_tenant)
            self._toolsets[toolset] = schema
        self._tools.setdefault(toolset, {})
        self._tools[toolset][tool.spec.name] = tool_obj
        # Keep the schema's tool list in sync so get_toolset_schema/export_schema reflect registered tools.
        self._toolsets[toolset] = replace(schema, tools=tuple(t.spec for t in self._tools[toolset].values()))
        if tenant_id:
            self._tenant_tools.setdefault(toolset, {}).setdefault(tenant_id, {})[tool.spec.name] = tool_obj

    def get_toolset(self, name):
        return self._toolsets.get(name)

    def get_tool(self, toolset: str, name: str, tenant_id: Optional[str] = None):
        if tenant_id:
            tenant_tools = self._tenant_tools.get(toolset, {}).get(tenant_id)
            if tenant_tools is None:
                raise ValueError(f"tenant_id='{tenant_id}' has no mapped tools for toolset='{toolset}'")
            tool = tenant_tools.get(name)
            if tool is not None:
                return tool
        tools = self._tools.get(toolset)
        if not tools:
            return None
        return tools.get(name)

    def toolset_names(self) -> Sequence[str]:
        return tuple(self._toolsets.keys())

    def tool_names(self, toolset: str) -> Sequence[str]:
        return tuple(self._tools.get(toolset, {}).keys())

    def export_schema(self, toolset: str) -> Dict[str, Any]:
        schema = self._toolsets.get(toolset)
        if schema is None:
            raise KeyError(f"unknown toolset: {toolset}")
        return schema.to_json()

    async def lookup(self, *, tenant_id: Optional[str], toolset: str) -> Sequence[Tool]:
        schema = self._toolsets.get(toolset)
        if schema is None:
            return ()
        effective_tenant = schema.tenant_for(caller_tenant_id=tenant_id)
        if effective_tenant:
            tools = self._tenant_tools.get(toolset, {}).get(effective_tenant)
            if tools:
                return tuple(tools.values())
        return tuple(self._tools.get(toolset, {}).values())

ToolResult dataclass

Source code in src/ciel/runtime/tools.py
@dataclass
class ToolResult:
    id: str
    name: str
    output: Any = None
    error: Optional[str] = None
    usage: Optional[Dict[str, Any]] = None
    duration_ms: Optional[int] = None
    metadata: Dict[str, Any] = field(default_factory=dict)

ToolSpec dataclass

Source code in src/ciel/runtime/tools.py
@dataclass(frozen=True)
class ToolSpec:
    name: str
    description: str
    parameters: Mapping[str, Any]
    strict: bool = False
    metadata: Dict[str, Any] = field(default_factory=dict)

ToolsetSchema dataclass

Source code in src/ciel/runtime/tools.py
@dataclass
class ToolsetSchema:
    name: str
    description: str
    version: str = "1.0.0"
    tenants: Sequence[str] = ()
    default_tenant: Optional[str] = None
    require_tenant: bool = False
    tools: Sequence[ToolSpec] = ()

    def tenant_for(self, *, caller_tenant_id: Optional[str]) -> Optional[str]:
        if self.require_tenant and not caller_tenant_id:
            raise ValueError(f"Toolset '{self.name}' requires tenant_id; none was provided.")
        return caller_tenant_id or self.default_tenant

    def to_json(self) -> Dict[str, Any]:
        payload: Dict[str, Any] = {
            "name": self.name,
            "description": self.description,
            "version": self.version,
            "tools": [
                {
                    "name": tool.name,
                    "description": tool.description,
                    "parameters": tool.parameters,
                    "strict": tool.strict,
                    **( {"metadata": tool.metadata} if tool.metadata else {} ),
                }
                for tool in self.tools
            ],
        }
        if self.require_tenant:
            payload["require_tenant"] = True
        if self.tenants:
            payload["tenants"] = list(self.tenants)
        return payload

Built-in tools shipped with Ciel.

All tools are defined as ToolSpec + callable and registered into the builtins toolset by register_builtin_tools. Network/sandbox tools are safe to import offline; their side effects only happen at execution time and can be sandboxed via ciel.sandbox.SandboxContext.

BUILTIN_TOOLS: tuple[Tool, ...] = (ECHO_TOOL, DATETIME_TOOL, HTTP_GET_TOOL, FILE_READ_TOOL, SHELL_TOOL) module-attribute

BUILTIN_TOOLSET = ToolsetSchema(name='builtins', description='Ciel built-in tools (echo, datetime, http_get, file_read, shell).', tools=(tuple((t.spec) for t in BUILTIN_TOOLS))) module-attribute

DATETIME_TOOL = Tool(spec=(ToolSpec(name='datetime', description='Return current UTC time.', parameters={'format': {'type': 'string'}})), callable_=_datetime) module-attribute

ECHO_TOOL = Tool(spec=(ToolSpec(name='echo', description='Echo back the provided text.', parameters={'text': {'type': 'string'}})), callable_=_echo) module-attribute

FILE_READ_TOOL = Tool(spec=(ToolSpec(name='file_read', description='Read a local file (sandboxed).', parameters={'path': {'type': 'string'}})), callable_=_file_read) module-attribute

HTTP_GET_TOOL = Tool(spec=(ToolSpec(name='http_get', description='GET a URL (requires network).', parameters={'url': {'type': 'string'}, 'timeout': {'type': 'number'}})), callable_=_http_get) module-attribute

SHELL_TOOL = Tool(spec=(ToolSpec(name='shell', description='Run a shell command (disabled by default policy).', parameters={'command': {'type': 'string'}})), callable_=_shell) module-attribute

SandboxContext dataclass

Source code in src/ciel/sandbox/__init__.py
@dataclass
class SandboxContext:
    policy: Optional[SandboxPolicy] = None
    executor: Optional[SandboxExecutor] = None

    def __post_init__(self) -> None:
        if self.policy is None:
            self.policy = SandboxPolicy()
        if self.executor is None:
            self.executor = SandboxExecutor()

    def evaluate(self, capability: str, command: Optional[str] = None) -> bool:
        if capability == "terminal":
            if not self.policy.allow_terminal:
                return False
            if command:
                if self.policy.denied_commands and command in self.policy.denied_commands:
                    return False
                if self.policy.allowed_commands and command not in self.policy.allowed_commands:
                    return False
            return True
        if capability == "file_write":
            return self.policy.allow_file_write
        if capability == "file_read":
            return self.policy.allow_file_read
        return False

    def execute(self, command: str, arguments: Optional[Dict[str, Any]] = None) -> str:
        arguments = arguments or {}
        if not self.evaluate("terminal", command=command):
            raise SandboxBlockedError("terminal", f"command '{command}' denied")
        # Ejecución REAL vía el executor (reemplaza el antiguo stub).
        full = command
        args = arguments.get("args")
        if args:
            full = command + " " + (args if isinstance(args, str) else " ".join(map(str, args)))
        result = self.executor.run(full)
        if result.exit_code != 0 and result.stderr:
            return result.stderr
        return result.stdout

    def write_file(self, path: str, content: str) -> str:
        if not self.evaluate("file_write"):
            raise SandboxBlockedError("file_write", f"write to '{path}' denied")
        from pathlib import Path

        p = Path(path)
        p.parent.mkdir(parents=True, exist_ok=True)
        p.write_text(content, encoding="utf-8")
        return f"wrote {len(content)} bytes to {path}"

    def read_file(self, path: str) -> str:
        if not self.evaluate("file_read"):
            raise SandboxBlockedError("file_read", f"read from '{path}' denied")
        from pathlib import Path

        p = Path(path)
        if not p.is_file():
            raise FileNotFoundError(path)
        return p.read_text(encoding="utf-8")

SandboxPolicy dataclass

Source code in src/ciel/sandbox/__init__.py
@dataclass
class SandboxPolicy:
    allow_file_read: bool = True
    allow_file_write: bool = False
    allow_terminal: bool = False
    allowed_commands: set[str] = field(default_factory=set)
    denied_commands: set[str] = field(default_factory=set)

Tool dataclass

Source code in src/ciel/runtime/tools.py
@dataclass
class Tool:
    spec: ToolSpec
    callable_: Any = None
    metadata: Dict[str, Any] = field(default_factory=dict)
    required_tenant: bool = False

ToolResult dataclass

Source code in src/ciel/runtime/tools.py
@dataclass
class ToolResult:
    id: str
    name: str
    output: Any = None
    error: Optional[str] = None
    usage: Optional[Dict[str, Any]] = None
    duration_ms: Optional[int] = None
    metadata: Dict[str, Any] = field(default_factory=dict)

ToolSpec dataclass

Source code in src/ciel/runtime/tools.py
@dataclass(frozen=True)
class ToolSpec:
    name: str
    description: str
    parameters: Mapping[str, Any]
    strict: bool = False
    metadata: Dict[str, Any] = field(default_factory=dict)

ToolsetSchema dataclass

Source code in src/ciel/runtime/tools.py
@dataclass
class ToolsetSchema:
    name: str
    description: str
    version: str = "1.0.0"
    tenants: Sequence[str] = ()
    default_tenant: Optional[str] = None
    require_tenant: bool = False
    tools: Sequence[ToolSpec] = ()

    def tenant_for(self, *, caller_tenant_id: Optional[str]) -> Optional[str]:
        if self.require_tenant and not caller_tenant_id:
            raise ValueError(f"Toolset '{self.name}' requires tenant_id; none was provided.")
        return caller_tenant_id or self.default_tenant

    def to_json(self) -> Dict[str, Any]:
        payload: Dict[str, Any] = {
            "name": self.name,
            "description": self.description,
            "version": self.version,
            "tools": [
                {
                    "name": tool.name,
                    "description": tool.description,
                    "parameters": tool.parameters,
                    "strict": tool.strict,
                    **( {"metadata": tool.metadata} if tool.metadata else {} ),
                }
                for tool in self.tools
            ],
        }
        if self.require_tenant:
            payload["require_tenant"] = True
        if self.tenants:
            payload["tenants"] = list(self.tenants)
        return payload

_datetime(arguments: Dict[str, Any], *, tool_call_id: str = '', tenant_id: Optional[str] = None) -> ToolResult

Source code in src/ciel/runtime/tools_builtins.py
def _datetime(arguments: Dict[str, Any], *, tool_call_id: str = "", tenant_id: Optional[str] = None) -> ToolResult:
    fmt = str(arguments.get("format", "%Y-%m-%dT%H:%M:%SZ"))
    now = datetime.datetime.now(datetime.timezone.utc).strftime(fmt)
    return ToolResult(id=tool_call_id, name="datetime", output={"now": now})

_echo(arguments: Dict[str, Any], *, tool_call_id: str = '', tenant_id: Optional[str] = None) -> ToolResult

Source code in src/ciel/runtime/tools_builtins.py
def _echo(arguments: Dict[str, Any], *, tool_call_id: str = "", tenant_id: Optional[str] = None) -> ToolResult:
    text = str(arguments.get("text", ""))
    return ToolResult(id=tool_call_id, name="echo", output={"echo": text})

_file_read(arguments: Dict[str, Any], *, tool_call_id: str = '', tenant_id: Optional[str] = None) -> ToolResult

Source code in src/ciel/runtime/tools_builtins.py
def _file_read(arguments: Dict[str, Any], *, tool_call_id: str = "", tenant_id: Optional[str] = None) -> ToolResult:
    policy = SandboxContext(policy=SandboxPolicy(allow_file_read=True))
    try:
        content = policy.read_file(str(arguments.get("path", "")))
        return ToolResult(id=tool_call_id, name="file_read", output={"content": content})
    except Exception as exc:
        return ToolResult(id=tool_call_id, name="file_read", error=str(exc))

_http_get(arguments: Dict[str, Any], *, tool_call_id: str = '', tenant_id: Optional[str] = None, client: Optional[httpx.AsyncClient] = None) -> ToolResult async

Source code in src/ciel/runtime/tools_builtins.py
async def _http_get(arguments: Dict[str, Any], *, tool_call_id: str = "", tenant_id: Optional[str] = None, client: Optional[httpx.AsyncClient] = None) -> ToolResult:
    url = str(arguments.get("url", ""))
    if not url:
        return ToolResult(id=tool_call_id, name="http_get", error="missing 'url'")
    ctx_client = client or httpx.AsyncClient(timeout=float(arguments.get("timeout", 30.0)))
    own = client is None
    try:
        resp = await ctx_client.get(url)
        body = resp.text
        return ToolResult(id=tool_call_id, name="http_get", output={"status": resp.status_code, "body": body[:4000]})
    except Exception as exc:  # pragma: no cover - network path
        return ToolResult(id=tool_call_id, name="http_get", error=f"request failed: {exc}")
    finally:
        if own:
            await ctx_client.aclose()

_shell(arguments: Dict[str, Any], *, tool_call_id: str = '', tenant_id: Optional[str] = None) -> ToolResult

Source code in src/ciel/runtime/tools_builtins.py
def _shell(arguments: Dict[str, Any], *, tool_call_id: str = "", tenant_id: Optional[str] = None) -> ToolResult:
    policy = SandboxContext(policy=SandboxPolicy(allow_terminal=False))
    try:
        out = policy.execute(str(arguments.get("command", "")))
        return ToolResult(id=tool_call_id, name="shell", output={"output": out})
    except Exception as exc:
        return ToolResult(id=tool_call_id, name="shell", error=str(exc))

register_builtin_tools(registry) -> None

Register all built-in tools into a ToolRegistry instance.

Source code in src/ciel/runtime/tools_builtins.py
def register_builtin_tools(registry) -> None:
    """Register all built-in tools into a ``ToolRegistry`` instance."""
    for tool in BUILTIN_TOOLS:
        registry.register_tool(BUILTIN_TOOLSET.name, tool)

__all__ = ['MemoryStore', 'MemoryEntry', 'init_schema', '_fts5_available'] module-attribute

MemoryEntry dataclass

Source code in src/ciel/runtime/memory.py
@dataclass
class MemoryEntry:
    tenant_id: Optional[str]
    session_id: str
    key: str
    value: Any = None
    value_json: str = "null"

MemoryStore

Bases: SqliteStateBackend

Alias retrocompatible a :class:SqliteStateBackend (SQLite en disco).

MemoryStore es SQLite en disco desde F5 — el nombre engaña. Para F15 se refactorizó para heredar de StateBackend de modo que cualquier store que hoy recibe un MemoryStore puede recibir también un PostgresStateBackend compartido sin cambios de API.

Se sigue construyendo con la misma firma: MemoryStore(db_path).

Source code in src/ciel/runtime/memory.py
class MemoryStore(SqliteStateBackend):
    """Alias retrocompatible a :class:`SqliteStateBackend` (SQLite en disco).

    ``MemoryStore`` es SQLite en disco desde F5 — el nombre engaña. Para
    F15 se refactorizó para heredar de ``StateBackend`` de modo que cualquier
    store que hoy recibe un ``MemoryStore`` puede recibir también un
    ``PostgresStateBackend`` compartido sin cambios de API.

    Se sigue construyendo con la misma firma: ``MemoryStore(db_path)``.
    """

    def __init__(self, db_path: str) -> None:
        super().__init__(db_path)

SqliteStateBackend

Bases: StateBackend

Backend SQLite en disco (default offline). Hereda el esquema de F15-.

Mantiene FTS5 cuando está disponible; si no, search degrada a lista vacía (comportamiento idéntico al MemoryStore original).

Source code in src/ciel/runtime/state_backend.py
class SqliteStateBackend(StateBackend):
    """Backend SQLite en disco (default offline). Hereda el esquema de F15-.

    Mantiene FTS5 cuando está disponible; si no, ``search`` degrada a lista
    vacía (comportamiento idéntico al ``MemoryStore`` original).
    """

    backend_type = "sqlite"

    def __init__(self, db_path: str) -> None:
        import sqlite3

        self.db_path = db_path
        self.conn = sqlite3.connect(db_path)
        self.conn.row_factory = sqlite3.Row
        init_schema(self.conn)

    def set(
        self,
        *,
        tenant_id: Optional[str],
        session_id: str,
        key: str,
        value: Any,
    ) -> None:
        tenant_value = self._sentinel(tenant_id)
        now = __import__("datetime").datetime.now(
            __import__("datetime").timezone.utc
        ).isoformat()
        self.conn.execute(
            """
            INSERT INTO memory (tenant_id, session_id, key, value_json, created_at, updated_at)
            VALUES (?, ?, ?, ?, ?, ?)
            ON CONFLICT(tenant_id, session_id, key) DO UPDATE SET value_json = excluded.value_json, updated_at = excluded.updated_at
            """,
            (tenant_value, session_id, key, self._dump(value), now, now),
        )
        self.conn.commit()

    def get(self, *, tenant_id: Optional[str], session_id: str, key: str) -> Optional[Any]:
        import sqlite3

        if tenant_id is None:
            row = self.conn.execute(
                "SELECT value_json FROM memory WHERE tenant_id = ? AND session_id = ? AND key = ?",
                ("__none__", session_id, key),
            ).fetchone()
        else:
            row = self.conn.execute(
                "SELECT value_json FROM memory WHERE tenant_id = ? AND session_id = ? AND key = ?",
                (tenant_id, session_id, key),
            ).fetchone()
        if row is None:
            return None
        return self._normalize_row(row["value_json"])

    def delete(self, *, tenant_id: Optional[str], session_id: str, key: str) -> None:
        if tenant_id is None:
            self.conn.execute(
                "DELETE FROM memory WHERE tenant_id = ? AND session_id = ? AND key = ?",
                ("__none__", session_id, key),
            )
        else:
            self.conn.execute(
                "DELETE FROM memory WHERE tenant_id = ? AND session_id = ? AND key = ?",
                (tenant_id, session_id, key),
            )
        self.conn.commit()

    def search(self, query: str, *, limit: int = 10) -> List[Dict[str, Any]]:
        import sqlite3

        try:
            rows = self.conn.execute(
                "SELECT key, value_json FROM memory_fts WHERE memory_fts MATCH ? LIMIT ?",
                (query, limit),
            ).fetchall()
        except sqlite3.OperationalError:
            rows = []
        results: List[Dict[str, Any]] = []
        for row in rows:
            try:
                results.append({"key": row["key"], "value": json.loads(row["value_json"])})
            except TypeError:
                results.append({"key": row["key"], "value": None})
        return results

    def record_tool_execution(
        self,
        *,
        tenant_id: Optional[str],
        session_id: str,
        toolset: str,
        tool_name: str,
        arguments: Any,
        started_at: str,
        finished_at: str,
        duration_ms: int,
        output: Any = None,
        error: Optional[str] = None,
    ) -> None:
        tenant_value = self._sentinel(tenant_id)
        self.conn.execute(
            """
            INSERT INTO tool_execution_log (tenant_id, session_id, toolset, tool_name, arguments_json, output_json, error, started_at, finished_at, duration_ms)
            VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
            """,
            (
                tenant_value,
                session_id,
                toolset,
                tool_name,
                self._dump(arguments),
                self._dump(output),
                error,
                started_at,
                finished_at,
                duration_ms,
            ),
        )
        self.conn.commit()

    # --- memoria episódica (Fase 17, tenant-filtered) ----------------------
    def memory_append(
        self,
        *,
        tenant_id: Optional[str],
        session_id: str,
        memory_id: str,
        value: Any,
    ) -> None:
        tenant_value = self._sentinel(tenant_id)
        now = datetime.now(timezone.utc).isoformat()
        self.conn.execute(
            """
            INSERT INTO memory_episodic (tenant_id, session_id, memory_id, value_json, created_at)
            VALUES (?, ?, ?, ?, ?)
            ON CONFLICT(tenant_id, session_id, memory_id) DO UPDATE SET value_json = excluded.value_json, created_at = excluded.created_at
            """,
            (tenant_value, session_id, memory_id, self._dump(value), now),
        )
        self.conn.commit()

    def memory_get(
        self,
        *,
        tenant_id: Optional[str],
        session_id: str,
        memory_id: str,
    ) -> Optional[dict]:
        tenant_value = self._sentinel(tenant_id)
        row = self.conn.execute(
            "SELECT value_json FROM memory_episodic WHERE tenant_id = ? AND session_id = ? AND memory_id = ?",
            (tenant_value, session_id, memory_id),
        ).fetchone()
        if row is None:
            return None
        return self._normalize_row(row["value_json"])

    def memory_get_recent(
        self,
        *,
        tenant_id: Optional[str],
        session_id: str,
        limit: int = 8,
    ) -> List[dict]:
        tenant_value = self._sentinel(tenant_id)
        rows = self.conn.execute(
            """
            SELECT value_json FROM memory_episodic
            WHERE tenant_id = ? AND session_id = ?
            ORDER BY id DESC LIMIT ?
            """,
            (tenant_value, session_id, limit),
        ).fetchall()
        out: List[dict] = []
        for row in rows:
            parsed = self._normalize_row(row["value_json"])
            if isinstance(parsed, dict):
                out.append(parsed)
        return out

    def memory_search_tenant(
        self,
        *,
        tenant_id: Optional[str],
        session_id: Optional[str],
        query: str,
        limit: int = 5,
    ) -> List[dict]:
        tenant_value = self._sentinel(tenant_id)
        # Búsqueda por substring (LIKE) filtrada ESTRICTAMENTE por tenant_id
        # (riesgo de fuga cross-tenant mitigado). No depende del tokenizer FTS5,
        # así que funciona siempre (offline-safe) y es determinista.
        # El value_json guarda el contenido con escape JSON (p.ej. la ñ como
        # \u00f1), por eso normalizamos la query al mismo formato que produce
        # json.dumps para que el LIKE coincida con caracteres no-ASCII.
        import json as _json

        # El value_json guarda el contenido con escape JSON (p.ej. la ñ como
        # \u00f1), por eso normalizamos la query al mismo formato que produce
        # json.dumps. Escapamos los metacaracteres de LIKE (% _ \) para que la
        # búsqueda sea literal y segura (sin usar ESCAPE, que interferiría con
        # la barra invertida del escape JSON).
        raw = _json.dumps(query).strip(chr(34))
        # Escapamos solo los metacaracteres de LIKE (% _) para que la búsqueda
        # sea literal. NO escapamos la barra invertida: el valor JSON almacena
        # \u00f1 con una barra literal y, sin cláusula ESCAPE, LIKE trata la
        # barra como carácter normal (coincide con el valor almacenado).
        escaped = raw.replace("%", "\\%").replace("_", "\\_")
        like = f"%{escaped}%"
        try:
            if session_id is not None:
                rows = self.conn.execute(
                    """
                    SELECT value_json FROM memory_episodic
                    WHERE tenant_id = ? AND session_id = ? AND value_json LIKE ?
                    ORDER BY id DESC LIMIT ?
                    """,
                    (tenant_value, session_id, like, limit),
                ).fetchall()
            else:
                rows = self.conn.execute(
                    """
                    SELECT value_json FROM memory_episodic
                    WHERE tenant_id = ? AND value_json LIKE ?
                    ORDER BY id DESC LIMIT ?
                    """,
                    (tenant_value, like, limit),
                ).fetchall()
        except Exception:  # query/SQL inesperado => degrada a [].
            return []
        out: List[dict] = []
        for row in rows:
            parsed = self._normalize_row(row["value_json"])
            if isinstance(parsed, dict):
                out.append(parsed)
        return out

    def memory_clear_session(
        self, *, tenant_id: Optional[str], session_id: str
    ) -> None:
        tenant_value = self._sentinel(tenant_id)
        self.conn.execute(
            "DELETE FROM memory_episodic WHERE tenant_id = ? AND session_id = ?",
            (tenant_value, session_id),
        )
        self.conn.commit()

    def close(self) -> None:
        self.conn.close()

_fts5_available(conn) -> bool

Source code in src/ciel/runtime/state_backend.py
def _fts5_available(conn) -> bool:  # type: ignore[no-untyped-def]
    row = conn.execute(
        "SELECT * FROM pragma_compile_options WHERE compile_options LIKE '%FTS5%'"
    ).fetchone()
    return bool(row)

init_schema(conn) -> None

Source code in src/ciel/runtime/state_backend.py
def init_schema(conn) -> None:  # type: ignore[no-untyped-def]
    conn.execute(
        """
        CREATE TABLE IF NOT EXISTS memory (
            id INTEGER PRIMARY KEY AUTOINCREMENT,
            tenant_id TEXT,
            session_id TEXT,
            key TEXT,
            value_json TEXT,
            created_at TEXT,
            updated_at TEXT,
            UNIQUE(tenant_id, session_id, key)
        );
        """
    )
    conn.execute(
        """
        CREATE TABLE IF NOT EXISTS tool_execution_log (
            id INTEGER PRIMARY KEY AUTOINCREMENT,
            tenant_id TEXT,
            session_id TEXT,
            toolset TEXT,
            tool_name TEXT,
            arguments_json TEXT,
            output_json TEXT,
            error TEXT,
            started_at TEXT,
            finished_at TEXT,
            duration_ms INTEGER
        );
        """
    )
    # --- memoria episódica (Fase 17) ----------------------------------------
    conn.execute(
        """
        CREATE TABLE IF NOT EXISTS memory_episodic (
            id INTEGER PRIMARY KEY AUTOINCREMENT,
            tenant_id TEXT,
            session_id TEXT,
            memory_id TEXT,
            value_json TEXT,
            created_at TEXT,
            UNIQUE(tenant_id, session_id, memory_id)
        );
        """
    )
    if _fts5_available(conn):
        try:
            conn.execute(
                """
                CREATE VIRTUAL TABLE IF NOT EXISTS memory_episodic_fts USING fts5(
                    memory_id,
                    value_json,
                    content='memory_episodic',
                    content_rowid='id',
                    tokenize='trigram'
                );
                """
            )
            conn.execute(
                """
                CREATE TRIGGER IF NOT EXISTS me_ai AFTER INSERT ON memory_episodic BEGIN
                    INSERT INTO memory_episodic_fts(rowid, memory_id, value_json) VALUES (new.id, new.memory_id, new.value_json);
                END;
                """
            )
            conn.execute(
                """
                CREATE TRIGGER IF NOT EXISTS me_ad AFTER DELETE ON memory_episodic BEGIN
                    INSERT INTO memory_episodic_fts(memory_episodic_fts, rowid, memory_id, value_json) VALUES ('delete', old.id, old.memory_id, old.value_json);
                END;
                """
            )
        except Exception:  # pragma: no cover - FTS5 init best-effort
            pass
    conn.commit()
    if _fts5_available(conn):
        try:
            conn.execute(
                """
                CREATE VIRTUAL TABLE IF NOT EXISTS memory_fts USING fts5(
                    key,
                    value_json,
                    content='memory',
                    content_rowid='id'
                );
                """
            )
            conn.execute(
                """
                CREATE TRIGGER IF NOT EXISTS memory_ai AFTER INSERT ON memory BEGIN
                    INSERT INTO memory_fts(rowid, key, value_json) VALUES (new.id, new.key, new.value_json);
                END;
                """
            )
            conn.execute(
                """
                CREATE TRIGGER IF NOT EXISTS memory_ad AFTER DELETE ON memory BEGIN
                    INSERT INTO memory_fts(memory_fts, rowid, key, value_json) VALUES ('delete', old.id, old.key, old.value_json);
                END;
                """
            )
            conn.execute(
                """
                CREATE TRIGGER IF NOT EXISTS memory_au AFTER UPDATE ON memory BEGIN
                    INSERT INTO memory_fts(memory_fts, rowid, key, value_json) VALUES ('delete', old.id, old.key, old.value_json);
                    INSERT INTO memory_fts(rowid, key, value_json) VALUES (new.id, new.key, new.value_json);
                END;
                """
            )
        except Exception:  # pragma: no cover - FTS5 init best-effort
            pass
    conn.commit()