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

_DANGEROUS_TOOL_NAMES = frozenset({'shell', 'terminal', 'exec', 'file_write', 'write_file', 'code_exec', 'eval'}) 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."""

Curriculum dataclass

Plan versionado (curriculum) para un objetivo del agente.

Source code in src\ciel\runtime\curriculum.py
@dataclass
class Curriculum:
    """Plan versionado (curriculum) para un objetivo del agente."""

    goal: str = ""
    plan: List[str] = field(default_factory=list)
    tenant_id: Optional[str] = None
    version: str = INITIAL_VERSION
    created_at: Optional[str] = None
    sha256: str = ""
    previous_version: Optional[str] = None
    parent: Optional[str] = None
    metadata: Dict[str, Any] = field(default_factory=dict)

    def to_dict(self) -> Dict[str, Any]:
        return {
            "goal": self.goal,
            "plan": list(self.plan),
            "tenant_id": self.tenant_id,
            "version": self.version,
            "created_at": self.created_at,
            "sha256": self.sha256,
            "previous_version": self.previous_version,
            "parent": self.parent,
            "metadata": self.metadata,
        }

    @classmethod
    def from_dict(cls, data: Dict[str, Any]) -> "Curriculum":
        return cls(
            goal=data.get("goal", ""),
            plan=list(data.get("plan") or []),
            tenant_id=data.get("tenant_id"),
            version=data.get("version") or INITIAL_VERSION,
            created_at=data.get("created_at"),
            sha256=data.get("sha256", ""),
            previous_version=data.get("previous_version"),
            parent=data.get("parent"),
            metadata=data.get("metadata") or {},
        )

    def bump_version(self, kind: str = "patch") -> str:
        parts = (self.version.split(".") + ["0", "0", "0"])[:3]
        try:
            major, minor, patch = (int(p) for p in parts)
        except ValueError as exc:
            raise CurriculumError(f"invalid version string: {self.version!r}") from exc
        kind = (kind or "patch").lower()
        if kind == "major":
            return f"{major + 1}.0.0"
        if kind == "minor":
            return f"{major}.{minor + 1}.0"
        if kind == "patch":
            return f"{major}.{minor}.{patch + 1}"
        raise CurriculumError(f"unknown bump kind: {kind!r} (expected major/minor/patch)")

CurriculumRegistry

Registro multitenant de curricula versionados sobre un StateBackend.

Source code in src\ciel\runtime\curriculum.py
class CurriculumRegistry:
    """Registro multitenant de curricula versionados sobre un ``StateBackend``."""

    def __init__(self, backend: Any) -> None:
        self._backend = backend

    # --- escritura ----------------------------------------------------------
    def create(
        self,
        goal: str,
        plan: Sequence[str],
        *,
        tenant_id: Optional[str] = None,
        metadata: Optional[Dict[str, Any]] = None,
    ) -> Curriculum:
        """Crea la versión inicial ``0.0.0`` de un curriculum (o bumpea si existe)."""
        current = self.get(goal, tenant_id=tenant_id)
        if current is None:
            version = INITIAL_VERSION
            previous = None
        else:
            version = current.bump_version("patch")
            previous = current.version
        cur = Curriculum(
            goal=goal,
            plan=list(plan),
            tenant_id=tenant_id,
            version=version,
            created_at=_now_iso(),
            sha256=sha256_plan(goal, plan),
            previous_version=previous,
            parent=previous,
            metadata=metadata or {},
        )
        self._backend.curriculum_save(
            tenant_id=tenant_id,
            goal=goal,
            version=cur.version,
            plan_json=json.dumps(list(plan), ensure_ascii=False),
            value_json=json.dumps(cur.to_dict(), ensure_ascii=False),
            sha256=cur.sha256,
            previous_version=cur.previous_version,
            created_at=cur.created_at,
        )
        return cur

    # --- lectura ------------------------------------------------------------
    def get(
        self,
        goal: str,
        *,
        tenant_id: Optional[str] = None,
        version: Optional[str] = None,
    ) -> Optional[Curriculum]:
        row = self._backend.curriculum_get(tenant_id=tenant_id, goal=goal, version=version)
        if row is None:
            return None
        return Curriculum.from_dict(row)

    def history(self, goal: str, *, tenant_id: Optional[str] = None) -> List[Curriculum]:
        rows = self._backend.curriculum_get_history(tenant_id=tenant_id, goal=goal)
        return [Curriculum.from_dict(r) for r in rows]

    def evolution_tree(self, goal: str, *, tenant_id: Optional[str] = None) -> Dict[str, Any]:
        """Árbol de linaje, misma forma que ``prompt_versioning.evolution_tree``."""
        versions = self.history(goal, tenant_id=tenant_id)
        if not versions:
            raise CurriculumError(f"unknown curriculum: {goal!r} (tenant {tenant_id!r})")
        keys = [c.version for c in versions]
        nodes: Dict[str, Dict[str, Any]] = {}
        for idx, cur in enumerate(versions):
            raw_prev = cur.previous_version
            parent = raw_prev if raw_prev is not None else (keys[idx - 1] if idx > 0 else None)
            nodes[cur.version] = {
                "version": cur.version,
                "parent": parent,
                "previous_version": raw_prev,
                "children": [],
                "sha256": cur.sha256,
                "plan": list(cur.plan),
                "created_at": cur.created_at,
            }
        for key, node in nodes.items():
            if node["parent"] is not None and node["parent"] in nodes:
                nodes[node["parent"]]["children"].append(key)
        roots = [k for k, n in nodes.items() if n["parent"] is None]
        root = roots[0] if roots else (keys[0] if keys else None)
        return {"goal": goal, "root": root, "lineage": list(keys), "nodes": nodes}

create(goal: str, plan: Sequence[str], *, tenant_id: Optional[str] = None, metadata: Optional[Dict[str, Any]] = None) -> Curriculum

Crea la versión inicial 0.0.0 de un curriculum (o bumpea si existe).

Source code in src\ciel\runtime\curriculum.py
def create(
    self,
    goal: str,
    plan: Sequence[str],
    *,
    tenant_id: Optional[str] = None,
    metadata: Optional[Dict[str, Any]] = None,
) -> Curriculum:
    """Crea la versión inicial ``0.0.0`` de un curriculum (o bumpea si existe)."""
    current = self.get(goal, tenant_id=tenant_id)
    if current is None:
        version = INITIAL_VERSION
        previous = None
    else:
        version = current.bump_version("patch")
        previous = current.version
    cur = Curriculum(
        goal=goal,
        plan=list(plan),
        tenant_id=tenant_id,
        version=version,
        created_at=_now_iso(),
        sha256=sha256_plan(goal, plan),
        previous_version=previous,
        parent=previous,
        metadata=metadata or {},
    )
    self._backend.curriculum_save(
        tenant_id=tenant_id,
        goal=goal,
        version=cur.version,
        plan_json=json.dumps(list(plan), ensure_ascii=False),
        value_json=json.dumps(cur.to_dict(), ensure_ascii=False),
        sha256=cur.sha256,
        previous_version=cur.previous_version,
        created_at=cur.created_at,
    )
    return cur

evolution_tree(goal: str, *, tenant_id: Optional[str] = None) -> Dict[str, Any]

Árbol de linaje, misma forma que prompt_versioning.evolution_tree.

Source code in src\ciel\runtime\curriculum.py
def evolution_tree(self, goal: str, *, tenant_id: Optional[str] = None) -> Dict[str, Any]:
    """Árbol de linaje, misma forma que ``prompt_versioning.evolution_tree``."""
    versions = self.history(goal, tenant_id=tenant_id)
    if not versions:
        raise CurriculumError(f"unknown curriculum: {goal!r} (tenant {tenant_id!r})")
    keys = [c.version for c in versions]
    nodes: Dict[str, Dict[str, Any]] = {}
    for idx, cur in enumerate(versions):
        raw_prev = cur.previous_version
        parent = raw_prev if raw_prev is not None else (keys[idx - 1] if idx > 0 else None)
        nodes[cur.version] = {
            "version": cur.version,
            "parent": parent,
            "previous_version": raw_prev,
            "children": [],
            "sha256": cur.sha256,
            "plan": list(cur.plan),
            "created_at": cur.created_at,
        }
    for key, node in nodes.items():
        if node["parent"] is not None and node["parent"] in nodes:
            nodes[node["parent"]]["children"].append(key)
    roots = [k for k, n in nodes.items() if n["parent"] is None]
    root = roots[0] if roots else (keys[0] if keys else None)
    return {"goal": goal, "root": root, "lineage": list(keys), "nodes": nodes}

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.

Si se pasa sandbox (un :class:~ciel.sandbox.SandboxExecutor), las herramientas consideradas peligrosas (shell/exec/file_write/...) se enrutan por él y el backend real utilizado queda registrado en ToolResult.metadata["sandbox_backend"] + limits_applied (F-SB-11, Hueco A Fase 20). Sin sandbox, el dispatch es idéntico al anterior.

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

    Si se pasa ``sandbox`` (un :class:`~ciel.sandbox.SandboxExecutor`), las
    herramientas consideradas peligrosas (shell/exec/file_write/...) se enrutan
    por él y el backend real utilizado queda registrado en
    ``ToolResult.metadata["sandbox_backend"]`` + ``limits_applied`` (F-SB-11,
    Hueco A Fase 20). Sin ``sandbox``, el dispatch es idéntico al anterior.
    """

    provider: ToolProvider
    default_toolset: Optional[str]

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

    async def dispatch(
        self,
        *,
        tenant_id: Optional[str] = None,
        toolset: Optional[str] = None,
        name: str,
        arguments: Dict[str, Any],
        tool_call_id: str,
        sandbox: Any = None,
    ) -> ToolResult:
        active_sandbox = sandbox or self.sandbox
        if active_sandbox is not None and name in _DANGEROUS_TOOL_NAMES:
            return await self._dispatch_sandboxed(
                active_sandbox,
                name=name,
                arguments=arguments,
                tool_call_id=tool_call_id,
                tenant_id=tenant_id,
            )
        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_sandboxed(
        self,
        sandbox: Any,
        *,
        name: str,
        arguments: Dict[str, Any],
        tool_call_id: str,
        tenant_id: Optional[str],
    ) -> ToolResult:
        """Ejecuta una herramienta peligrosa por el sandbox y registra el backend.

        El sandbox decide si degradar a ``inprocess`` es seguro (Hueco A): si no,
        lanza :class:`~ciel.common.CielError`, que se captura y se devuelve como
        ``ToolResult.error`` (no crashea el agente) con el backend real en metadata.
        """
        from ciel.common import CielError

        command = arguments.get("command") or arguments.get("code") or name
        try:
            exec_result = sandbox.run(command, capability=name)
        except CielError as exc:
            return ToolResult(
                id=tool_call_id,
                name=name,
                error=str(exc),
                metadata={
                    "tenant_id": tenant_id,
                    "sandbox_backend": getattr(sandbox, "backend", "inprocess"),
                    "limits_applied": False,
                    "sandbox_rejected": True,
                },
            )
        backend = getattr(exec_result, "backend", "inprocess")
        limits_applied = getattr(exec_result, "limits_applied", False)
        if getattr(exec_result, "exit_code", 0) != 0 and getattr(exec_result, "stderr", ""):
            return ToolResult(
                id=tool_call_id,
                name=name,
                error=exec_result.stderr,
                metadata={
                    "tenant_id": tenant_id,
                    "sandbox_backend": backend,
                    "limits_applied": limits_applied,
                },
            )
        return ToolResult(
            id=tool_call_id,
            name=name,
            output=exec_result.stdout,
            metadata={
                "tenant_id": tenant_id,
                "sandbox_backend": backend,
                "limits_applied": limits_applied,
            },
        )

    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)

KGEdge dataclass

Arista dirigida entre dos nodos del knowledge graph.

Source code in src\ciel\runtime\knowledge_graph.py
@dataclass
class KGEdge:
    """Arista dirigida entre dos nodos del knowledge graph."""

    from_id: str
    to_id: str
    relation: str
    weight: float = 1.0

KGNode dataclass

Nodo del knowledge graph (aislado por tenant).

Source code in src\ciel\runtime\knowledge_graph.py
@dataclass
class KGNode:
    """Nodo del knowledge graph (aislado por tenant)."""

    id: str
    kind: str
    label: str
    tenant_id: Optional[str] = None
    embedding: Optional[List[float]] = None
    metadata: Dict[str, Any] = field(default_factory=dict)

KnowledgeGraph

Knowledge graph con persistencia SQLite + búsqueda semántica offline.

Parameters

backend: StateBackend SQLite existente (se reusa su conn), una ruta a fichero SQLite (str), o None para SQLite in-memory. embedding_provider: Provider de embeddings. Default: DeterministicEmbeddingProvider (hash-based, offline, sin red).

Source code in src\ciel\runtime\knowledge_graph.py
class KnowledgeGraph:
    """Knowledge graph con persistencia SQLite + búsqueda semántica offline.

    Parameters
    ----------
    backend:
        ``StateBackend`` SQLite existente (se reusa su ``conn``), una ruta a
        fichero SQLite (``str``), o ``None`` para SQLite in-memory.
    embedding_provider:
        Provider de embeddings. Default: ``DeterministicEmbeddingProvider``
        (hash-based, offline, sin red).
    """

    def __init__(
        self,
        backend: Optional[Any] = None,
        *,
        embedding_provider: Optional[EmbeddingProvider] = None,
    ) -> None:
        self._owns_conn = False
        if backend is None:
            self.conn = sqlite3.connect(":memory:")
            self._owns_conn = True
        elif isinstance(backend, str):
            self.conn = sqlite3.connect(backend)
            self._owns_conn = True
        elif isinstance(backend, StateBackend) and hasattr(backend, "conn"):
            self.conn = backend.conn  # SqliteStateBackend
        elif isinstance(backend, sqlite3.Connection):
            self.conn = backend
        else:
            raise TypeError(
                "KnowledgeGraph requiere un SqliteStateBackend, sqlite3.Connection, "
                f"ruta str o None; recibido: {type(backend)!r}"
            )
        self.conn.row_factory = sqlite3.Row
        init_kg_schema(self.conn)
        self._embeddings = embedding_provider or DeterministicEmbeddingProvider()
        self._index = InMemoryVectorStore()
        self._load_index()

    # --- helpers -----------------------------------------------------------
    @staticmethod
    def _sentinel(tenant_id: Optional[str]) -> str:
        return StateBackend._sentinel(tenant_id)

    @staticmethod
    def _now() -> str:
        return datetime.now(timezone.utc).isoformat()

    def _load_index(self) -> None:
        """Repuebla el índice vectorial en memoria desde SQLite (reapertura)."""
        rows = self.conn.execute(
            "SELECT tenant_id, node_id, label, embedding_json FROM kg_nodes"
        ).fetchall()
        for row in rows:
            emb = json.loads(row["embedding_json"]) if row["embedding_json"] else None
            if emb is None:
                emb = self._embeddings.embed_one(row["label"])
            tenant = None if row["tenant_id"] == "__none__" else row["tenant_id"]
            self._index.upsert(
                tenant_id=tenant,
                texts=[row["label"]],
                vectors=[emb],
                ids=[f"{row['tenant_id']}::{row['node_id']}"],
            )

    def _row_to_node(self, row: sqlite3.Row) -> KGNode:
        tenant = None if row["tenant_id"] == "__none__" else row["tenant_id"]
        return KGNode(
            id=row["node_id"],
            kind=row["kind"],
            label=row["label"],
            tenant_id=tenant,
            embedding=json.loads(row["embedding_json"]) if row["embedding_json"] else None,
            metadata=json.loads(row["metadata_json"]) if row["metadata_json"] else {},
        )

    # --- persistencia estilo F19 --------------------------------------------
    def kg_node_save(self, node: KGNode) -> KGNode:
        if node.embedding is None:
            node.embedding = self._embeddings.embed_one(node.label)
        tenant_value = self._sentinel(node.tenant_id)
        now = self._now()
        self.conn.execute(
            """
            INSERT INTO kg_nodes (tenant_id, node_id, kind, label, embedding_json, metadata_json, created_at, updated_at)
            VALUES (?, ?, ?, ?, ?, ?, ?, ?)
            ON CONFLICT(tenant_id, node_id) DO UPDATE SET
                kind = excluded.kind,
                label = excluded.label,
                embedding_json = excluded.embedding_json,
                metadata_json = excluded.metadata_json,
                updated_at = excluded.updated_at
            """,
            (
                tenant_value,
                node.id,
                node.kind,
                node.label,
                json.dumps(node.embedding),
                json.dumps(node.metadata),
                now,
                now,
            ),
        )
        self.conn.commit()
        self._index.upsert(
            tenant_id=node.tenant_id,
            texts=[node.label],
            vectors=[node.embedding],
            payloads=[{"node_id": node.id}],
            ids=[f"{tenant_value}::{node.id}"],
        )
        return node

    def kg_node_get(self, node_id: str, *, tenant_id: Optional[str] = None) -> Optional[KGNode]:
        row = self.conn.execute(
            "SELECT * FROM kg_nodes WHERE tenant_id = ? AND node_id = ?",
            (self._sentinel(tenant_id), node_id),
        ).fetchone()
        return self._row_to_node(row) if row else None

    def kg_edge_save(self, edge: KGEdge, *, tenant_id: Optional[str] = None) -> KGEdge:
        self.conn.execute(
            """
            INSERT INTO kg_edges (tenant_id, from_id, to_id, relation, weight, created_at)
            VALUES (?, ?, ?, ?, ?, ?)
            ON CONFLICT(tenant_id, from_id, to_id, relation) DO UPDATE SET
                weight = excluded.weight
            """,
            (
                self._sentinel(tenant_id),
                edge.from_id,
                edge.to_id,
                edge.relation,
                edge.weight,
                self._now(),
            ),
        )
        self.conn.commit()
        return edge

    def kg_edge_get(
        self,
        from_id: str,
        to_id: str,
        *,
        relation: Optional[str] = None,
        tenant_id: Optional[str] = None,
    ) -> List[KGEdge]:
        tenant_value = self._sentinel(tenant_id)
        if relation is not None:
            rows = self.conn.execute(
                "SELECT * FROM kg_edges WHERE tenant_id = ? AND from_id = ? AND to_id = ? AND relation = ?",
                (tenant_value, from_id, to_id, relation),
            ).fetchall()
        else:
            rows = self.conn.execute(
                "SELECT * FROM kg_edges WHERE tenant_id = ? AND from_id = ? AND to_id = ?",
                (tenant_value, from_id, to_id),
            ).fetchall()
        return [
            KGEdge(
                from_id=r["from_id"],
                to_id=r["to_id"],
                relation=r["relation"],
                weight=r["weight"],
            )
            for r in rows
        ]

    # --- API pública de grafo -----------------------------------------------
    def add_node(
        self,
        node_id: str,
        *,
        kind: str = "entity",
        label: Optional[str] = None,
        tenant_id: Optional[str] = None,
        metadata: Optional[Dict[str, Any]] = None,
    ) -> KGNode:
        node = KGNode(
            id=node_id,
            kind=kind,
            label=label or node_id,
            tenant_id=tenant_id,
            metadata=metadata or {},
        )
        return self.kg_node_save(node)

    def add_edge(
        self,
        from_id: str,
        to_id: str,
        *,
        relation: str = "related_to",
        weight: float = 1.0,
        tenant_id: Optional[str] = None,
    ) -> KGEdge:
        for nid in (from_id, to_id):
            if self.kg_node_get(nid, tenant_id=tenant_id) is None:
                raise KeyError(f"Nodo desconocido para tenant {tenant_id!r}: {nid!r}")
        edge = KGEdge(from_id=from_id, to_id=to_id, relation=relation, weight=weight)
        return self.kg_edge_save(edge, tenant_id=tenant_id)

    def neighbors(
        self,
        node_id: str,
        *,
        tenant_id: Optional[str] = None,
        direction: str = "both",
    ) -> List[Tuple[KGNode, KGEdge]]:
        """Vecinos por aristas explícitas. ``direction``: out | in | both."""
        tenant_value = self._sentinel(tenant_id)
        out: List[Tuple[KGNode, KGEdge]] = []
        if direction in ("out", "both"):
            rows = self.conn.execute(
                "SELECT * FROM kg_edges WHERE tenant_id = ? AND from_id = ?",
                (tenant_value, node_id),
            ).fetchall()
            for r in rows:
                node = self.kg_node_get(r["to_id"], tenant_id=tenant_id)
                if node is not None:
                    out.append(
                        (node, KGEdge(r["from_id"], r["to_id"], r["relation"], r["weight"]))
                    )
        if direction in ("in", "both"):
            rows = self.conn.execute(
                "SELECT * FROM kg_edges WHERE tenant_id = ? AND to_id = ?",
                (tenant_value, node_id),
            ).fetchall()
            for r in rows:
                node = self.kg_node_get(r["from_id"], tenant_id=tenant_id)
                if node is not None:
                    out.append(
                        (node, KGEdge(r["from_id"], r["to_id"], r["relation"], r["weight"]))
                    )
        return out

    def search(
        self,
        query: str,
        *,
        tenant_id: Optional[str] = None,
        top_k: int = 5,
    ) -> List[KGNode]:
        """Búsqueda semántica (coseno, offline) filtrada por tenant."""
        vector = self._embeddings.embed_one(query)
        records = self._index.query(tenant_id=tenant_id, vector=vector, top_k=top_k)
        out: List[KGNode] = []
        for rec in records:
            node_id = rec.payload.get("node_id") or rec.id.split("::", 1)[-1]
            node = self.kg_node_get(node_id, tenant_id=tenant_id)
            if node is not None:
                out.append(node)
        return out

    # --- interop opcional -----------------------------------------------------
    def to_networkx(self, *, tenant_id: Optional[str] = None):  # pragma: no cover
        """Exporta a ``networkx.DiGraph`` si networkx está instalado (opcional)."""
        try:
            import networkx as nx  # type: ignore
        except ImportError as exc:
            raise RuntimeError(
                "to_networkx requiere el paquete opcional 'networkx'."
            ) from exc
        g = nx.DiGraph()
        tenant_value = self._sentinel(tenant_id)
        for row in self.conn.execute(
            "SELECT * FROM kg_nodes WHERE tenant_id = ?", (tenant_value,)
        ).fetchall():
            node = self._row_to_node(row)
            g.add_node(node.id, kind=node.kind, label=node.label, **node.metadata)
        for row in self.conn.execute(
            "SELECT * FROM kg_edges WHERE tenant_id = ?", (tenant_value,)
        ).fetchall():
            g.add_edge(row["from_id"], row["to_id"], relation=row["relation"], weight=row["weight"])
        return g

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

neighbors(node_id: str, *, tenant_id: Optional[str] = None, direction: str = 'both') -> List[Tuple[KGNode, KGEdge]]

Vecinos por aristas explícitas. direction: out | in | both.

Source code in src\ciel\runtime\knowledge_graph.py
def neighbors(
    self,
    node_id: str,
    *,
    tenant_id: Optional[str] = None,
    direction: str = "both",
) -> List[Tuple[KGNode, KGEdge]]:
    """Vecinos por aristas explícitas. ``direction``: out | in | both."""
    tenant_value = self._sentinel(tenant_id)
    out: List[Tuple[KGNode, KGEdge]] = []
    if direction in ("out", "both"):
        rows = self.conn.execute(
            "SELECT * FROM kg_edges WHERE tenant_id = ? AND from_id = ?",
            (tenant_value, node_id),
        ).fetchall()
        for r in rows:
            node = self.kg_node_get(r["to_id"], tenant_id=tenant_id)
            if node is not None:
                out.append(
                    (node, KGEdge(r["from_id"], r["to_id"], r["relation"], r["weight"]))
                )
    if direction in ("in", "both"):
        rows = self.conn.execute(
            "SELECT * FROM kg_edges WHERE tenant_id = ? AND to_id = ?",
            (tenant_value, node_id),
        ).fetchall()
        for r in rows:
            node = self.kg_node_get(r["from_id"], tenant_id=tenant_id)
            if node is not None:
                out.append(
                    (node, KGEdge(r["from_id"], r["to_id"], r["relation"], r["weight"]))
                )
    return out

search(query: str, *, tenant_id: Optional[str] = None, top_k: int = 5) -> List[KGNode]

Búsqueda semántica (coseno, offline) filtrada por tenant.

Source code in src\ciel\runtime\knowledge_graph.py
def search(
    self,
    query: str,
    *,
    tenant_id: Optional[str] = None,
    top_k: int = 5,
) -> List[KGNode]:
    """Búsqueda semántica (coseno, offline) filtrada por tenant."""
    vector = self._embeddings.embed_one(query)
    records = self._index.query(tenant_id=tenant_id, vector=vector, top_k=top_k)
    out: List[KGNode] = []
    for rec in records:
        node_id = rec.payload.get("node_id") or rec.id.split("::", 1)[-1]
        node = self.kg_node_get(node_id, tenant_id=tenant_id)
        if node is not None:
            out.append(node)
    return out

to_networkx(*, tenant_id: Optional[str] = None)

Exporta a networkx.DiGraph si networkx está instalado (opcional).

Source code in src\ciel\runtime\knowledge_graph.py
def to_networkx(self, *, tenant_id: Optional[str] = None):  # pragma: no cover
    """Exporta a ``networkx.DiGraph`` si networkx está instalado (opcional)."""
    try:
        import networkx as nx  # type: ignore
    except ImportError as exc:
        raise RuntimeError(
            "to_networkx requiere el paquete opcional 'networkx'."
        ) from exc
    g = nx.DiGraph()
    tenant_value = self._sentinel(tenant_id)
    for row in self.conn.execute(
        "SELECT * FROM kg_nodes WHERE tenant_id = ?", (tenant_value,)
    ).fetchall():
        node = self._row_to_node(row)
        g.add_node(node.id, kind=node.kind, label=node.label, **node.metadata)
    for row in self.conn.execute(
        "SELECT * FROM kg_edges WHERE tenant_id = ?", (tenant_value,)
    ).fetchall():
        g.add_edge(row["from_id"], row["to_id"], relation=row["relation"], weight=row["weight"])
    return g

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,
        domain: Optional[str] = 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,
            domain=domain,
            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,
        domain: 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,
            domain=domain if domain is not None else getattr(previous, "domain", None),
            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:
        # Unificado (Fase 20): delega en SkillVersion.parse/bump — la única
        # implementación de semver vive en ciel.runtime.skill_versioning.
        from ciel.runtime.skill_versioning import SkillVersion

        try:
            base = SkillVersion.parse(current) if current else SkillVersion()
        except SkillError:
            base = SkillVersion()
        try:
            return base.bump(bump).version
        except SkillError:
            return base.bump("patch").version

create_from_code(*, name: str, description: str, code: str, category: Optional[str] = None, tenant_id: Optional[str] = None, metadata: Optional[Dict[str, Any]] = None, domain: Optional[str] = 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,
    domain: Optional[str] = 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,
        domain=domain,
        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, domain: 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,
    domain: 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,
        domain=domain if domain is not None else getattr(previous, "domain", None),
        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

SkillRecommender

Recomendador cross-domain de skills basado en embeddings (offline).

  • index() vectoriza Skill.content de la librería (por tenant).
  • recommend() devuelve skills relevantes al contexto. Si se conoce el dominio del contexto (context_domain), solo devuelve skills cuyo domain sea distinto (transfer learning entre dominios). Si no hay dominio, devuelve el top-k por similitud.
Source code in src\ciel\runtime\skill_recommender.py
class SkillRecommender:
    """Recomendador cross-domain de skills basado en embeddings (offline).

    - ``index()`` vectoriza ``Skill.content`` de la librería (por tenant).
    - ``recommend()`` devuelve skills relevantes al contexto. Si se conoce el
      dominio del contexto (``context_domain``), solo devuelve skills cuyo
      ``domain`` sea distinto (transfer learning entre dominios). Si no hay
      dominio, devuelve el top-k por similitud.
    """

    def __init__(
        self,
        library: SkillLibrary,
        *,
        embeddings: Optional[EmbeddingProvider] = None,
        store: Optional[VectorStore] = None,
    ) -> None:
        self.library = library
        self.embeddings = embeddings or DeterministicEmbeddingProvider()
        self.store = store or InMemoryVectorStore()

    # -- indexing --------------------------------------------------------

    def index(self, *, tenant_id: Optional[str] = None) -> int:
        """(Re)indexa las skills del tenant. Devuelve cuántas se indexaron."""
        skills = self.library.list_skills(tenant_id=tenant_id)
        if not skills:
            return 0
        texts = [f"{s.name}\n{s.description}\n{s.content}" for s in skills]
        vectors = self.embeddings.embed(texts)
        self.store.upsert(
            tenant_id=tenant_id,
            texts=texts,
            vectors=vectors,
            payloads=[
                {"skill_name": s.name, "domain": getattr(s, "domain", None)}
                for s in skills
            ],
            ids=[f"{tenant_id or '__none__'}:{s.name}" for s in skills],
        )
        return len(skills)

    # -- recommendation ---------------------------------------------------

    def recommend(
        self,
        context: str,
        *,
        tenant_id: Optional[str] = None,
        top_k: int = 3,
        context_domain: Optional[str] = None,
    ) -> List[Skill]:
        """Sugiere hasta ``top_k`` skills relevantes a ``context``.

        Cross-domain: si ``context_domain`` está definido, excluye skills de
        ese mismo dominio (solo transferencia desde OTROS dominios). Sin
        dominio, devuelve el top-k global por similitud coseno.
        """
        vector = self.embeddings.embed_one(context)
        # Pedimos de más para poder filtrar por dominio sin quedarnos cortos.
        records = self.store.query(tenant_id=tenant_id, vector=vector, top_k=max(top_k * 4, top_k))
        out: List[Skill] = []
        seen: set[str] = set()
        for record in records:
            name = record.payload.get("skill_name")
            if not name or name in seen:
                continue
            skill_domain = record.payload.get("domain")
            if context_domain is not None and skill_domain == context_domain:
                continue  # mismo dominio: no es transfer cross-domain
            skill = self.library.get(name)
            if skill is None:
                continue
            seen.add(name)
            out.append(skill)
            if len(out) >= top_k:
                break
        return out

index(*, tenant_id: Optional[str] = None) -> int

(Re)indexa las skills del tenant. Devuelve cuántas se indexaron.

Source code in src\ciel\runtime\skill_recommender.py
def index(self, *, tenant_id: Optional[str] = None) -> int:
    """(Re)indexa las skills del tenant. Devuelve cuántas se indexaron."""
    skills = self.library.list_skills(tenant_id=tenant_id)
    if not skills:
        return 0
    texts = [f"{s.name}\n{s.description}\n{s.content}" for s in skills]
    vectors = self.embeddings.embed(texts)
    self.store.upsert(
        tenant_id=tenant_id,
        texts=texts,
        vectors=vectors,
        payloads=[
            {"skill_name": s.name, "domain": getattr(s, "domain", None)}
            for s in skills
        ],
        ids=[f"{tenant_id or '__none__'}:{s.name}" for s in skills],
    )
    return len(skills)

recommend(context: str, *, tenant_id: Optional[str] = None, top_k: int = 3, context_domain: Optional[str] = None) -> List[Skill]

Sugiere hasta top_k skills relevantes a context.

Cross-domain: si context_domain está definido, excluye skills de ese mismo dominio (solo transferencia desde OTROS dominios). Sin dominio, devuelve el top-k global por similitud coseno.

Source code in src\ciel\runtime\skill_recommender.py
def recommend(
    self,
    context: str,
    *,
    tenant_id: Optional[str] = None,
    top_k: int = 3,
    context_domain: Optional[str] = None,
) -> List[Skill]:
    """Sugiere hasta ``top_k`` skills relevantes a ``context``.

    Cross-domain: si ``context_domain`` está definido, excluye skills de
    ese mismo dominio (solo transferencia desde OTROS dominios). Sin
    dominio, devuelve el top-k global por similitud coseno.
    """
    vector = self.embeddings.embed_one(context)
    # Pedimos de más para poder filtrar por dominio sin quedarnos cortos.
    records = self.store.query(tenant_id=tenant_id, vector=vector, top_k=max(top_k * 4, top_k))
    out: List[Skill] = []
    seen: set[str] = set()
    for record in records:
        name = record.payload.get("skill_name")
        if not name or name in seen:
            continue
        skill_domain = record.payload.get("domain")
        if context_domain is not None and skill_domain == context_domain:
            continue  # mismo dominio: no es transfer cross-domain
        skill = self.library.get(name)
        if skill is None:
            continue
        seen.add(name)
        out.append(skill)
        if len(out) >= top_k:
            break
    return out

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")

decompose_goal(goal: str, *, provider: Any = None) -> List[str]

Descompone goal en sub-tareas ordenadas.

Sin provider: heurística determinista offline (split por cláusulas). Con provider (ChatProvider, p. ej. MockProvider): se le pide el plan y se parsea; si la respuesta es vacía/inutilizable, fallback a la heurística offline.

Source code in src\ciel\runtime\curriculum.py
def decompose_goal(goal: str, *, provider: Any = None) -> List[str]:
    """Descompone ``goal`` en sub-tareas ordenadas.

    Sin ``provider``: heurística determinista offline (split por cláusulas).
    Con ``provider`` (ChatProvider, p. ej. ``MockProvider``): se le pide el
    plan y se parsea; si la respuesta es vacía/inutilizable, fallback a la
    heurística offline.
    """
    if not goal or not goal.strip():
        raise CurriculumError("goal vacío: nada que descomponer")
    if provider is None:
        return _heuristic_plan(goal)
    try:
        from ciel.runtime.tools import ChatMessage, ChatRequest

        request = ChatRequest(
            messages=(
                ChatMessage(
                    role="user",
                    content=(
                        "Descompón el siguiente objetivo en una lista JSON de "
                        f"sub-tareas ordenadas: {goal}"
                    ),
                ),
            )
        )
        try:
            loop = asyncio.get_running_loop()
        except RuntimeError:
            loop = None
        if loop is not None:  # pragma: no cover - llamada desde loop activo
            raise CurriculumError(
                "decompose_goal(provider=...) no puede llamarse desde un event "
                "loop activo; usa la heurística offline o llama fuera del loop"
            )
        resp = asyncio.run(provider.complete(request))
        text = getattr(getattr(resp, "choice", None), "message", None)
        text = getattr(text, "content", "") if text is not None else ""
        if isinstance(text, list):  # contenido multimodal
            text = "".join(p.get("text", "") for p in text if isinstance(p, dict))
        steps = _parse_provider_plan(text or "")
        if steps:
            return steps
    except CurriculumError:
        raise
    except Exception:
        pass  # provider roto => fallback offline
    return _heuristic_plan(goal)

install_curriculum_support(agent_cls: Any) -> Any

Engancha curriculum= / curriculum_config= sin reescribir api.py.

Idempotente (flag _curriculum_installed). Expone Agent.run_curriculum(goal, *, handler=None, plan=None, tenant_id=None) que descompone el goal (o usa plan), lo persiste vía :class:CurriculumRegistry y ejecuta las sub-tareas con AutonomousAgent.run_goal (EventLoop real de orquestación). Si no se pasa handler, cada sub-tarea se resuelve con Agent.arun.

Source code in src\ciel\runtime\curriculum.py
def install_curriculum_support(agent_cls: Any) -> Any:
    """Engancha ``curriculum=`` / ``curriculum_config=`` sin reescribir api.py.

    Idempotente (flag ``_curriculum_installed``). Expone
    ``Agent.run_curriculum(goal, *, handler=None, plan=None, tenant_id=None)``
    que descompone el goal (o usa ``plan``), lo persiste vía
    :class:`CurriculumRegistry` y ejecuta las sub-tareas con
    ``AutonomousAgent.run_goal`` (EventLoop real de orquestación). Si no se
    pasa ``handler``, cada sub-tarea se resuelve con ``Agent.arun``.
    """
    if getattr(agent_cls, "_curriculum_installed", False):
        return agent_cls
    original_init = agent_cls.__init__

    @functools.wraps(original_init)
    def _patched_init(self: Any, *args: Any, **kwargs: Any) -> None:
        cur_arg = kwargs.pop("curriculum", None)
        cur_config = kwargs.pop("curriculum_config", None)
        original_init(self, *args, **kwargs)

        config = cur_config or CurriculumConfig()
        if cur_arg is False or (cur_arg is None and cur_config is None):
            config = config.disabled()
        backend = config.backend
        if backend is None:
            backend = _default_backend()
        self._curriculum = AgentCurriculum(config=config, backend=backend)

    async def _run_curriculum(
        self: Any,
        goal: str,
        *,
        handler: Any = None,
        plan: Optional[Sequence[str]] = None,
        tenant_id: Optional[str] = None,
        provider: Any = None,
    ) -> Dict[str, Any]:
        """Descompone, persiste y ejecuta un curriculum. Devuelve dict resumen."""
        cur_state: Optional[AgentCurriculum] = getattr(self, "_curriculum", None)
        if cur_state is None:
            cur_state = AgentCurriculum(
                config=CurriculumConfig(), backend=_default_backend()
            )
            self._curriculum = cur_state
        eff_tenant = tenant_id or getattr(self, "tenant_id", None)
        steps = list(plan) if plan else decompose_goal(goal, provider=provider)
        curriculum = cur_state.registry.create(goal, steps, tenant_id=eff_tenant)

        if handler is None:
            async def handler(task: Any) -> Any:  # type: ignore[misc]
                resp = await self.arun(task.goal, tenant_id=eff_tenant)
                return getattr(resp, "text", resp)

        from ciel.orchestration.agent import AutonomousAgent

        runner = AutonomousAgent(
            name="curriculum",
            tenant_id=eff_tenant,
            max_attempts=cur_state.config.max_attempts,
        )
        tasks = await runner.run_goal(goal, handler, plan=steps)
        return {
            "goal": goal,
            "curriculum": curriculum.to_dict(),
            "tasks": [t.snapshot() for t in tasks],
            "succeeded": all(t.status == "succeeded" for t in tasks),
        }

    agent_cls.__init__ = _patched_init  # type: ignore[assignment]
    agent_cls.run_curriculum = _run_curriculum  # type: ignore[attr-defined]
    agent_cls._curriculum_installed = True  # type: ignore[attr-defined]
    return agent_cls

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

_DANGEROUS_TOOL_NAMES = frozenset({'shell', 'terminal', 'exec', 'file_write', 'write_file', 'code_exec', 'eval'}) 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.

Si se pasa sandbox (un :class:~ciel.sandbox.SandboxExecutor), las herramientas consideradas peligrosas (shell/exec/file_write/...) se enrutan por él y el backend real utilizado queda registrado en ToolResult.metadata["sandbox_backend"] + limits_applied (F-SB-11). Sin sandbox, el dispatch es idéntico al comportamiento anterior.

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

    Si se pasa ``sandbox`` (un :class:`~ciel.sandbox.SandboxExecutor`), las
    herramientas consideradas peligrosas (shell/exec/file_write/...) se enrutan
    por él y el backend real utilizado queda registrado en
    ``ToolResult.metadata["sandbox_backend"]`` + ``limits_applied`` (F-SB-11).
    Sin ``sandbox``, el dispatch es idéntico al comportamiento anterior.
    """

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

    async def dispatch(
        self,
        *,
        toolset: Optional[str] = None,
        name: str,
        arguments: Dict[str, Any],
        tool_call_id: str,
        tenant_id: Optional[str] = None,
        sandbox: Any = None,
    ) -> ToolResult:
        # Preferencia: kwarg de dispatch > kwarg del dispatcher.
        active_sandbox = sandbox or self.sandbox
        if active_sandbox is not None and name in _DANGEROUS_TOOL_NAMES:
            return await self._dispatch_sandboxed(
                active_sandbox,
                name=name,
                arguments=arguments,
                tool_call_id=tool_call_id,
                tenant_id=tenant_id,
            )
        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_sandboxed(
        self,
        sandbox: Any,
        *,
        name: str,
        arguments: Dict[str, Any],
        tool_call_id: str,
        tenant_id: Optional[str],
    ) -> ToolResult:
        """Ejecuta una herramienta peligrosa a través del sandbox y registra el backend.

        El sandbox decide si degradar a ``inprocess`` es seguro (Hueco A): si no,
        lanza :class:`~ciel.common.CielSecurityError`, que se captura y se
        devuelve como ``ToolResult.error`` (no crashea el agente) con el backend
        real en metadata.
        """
        from ciel.common import CielError

        command = arguments.get("command") or arguments.get("code") or name
        try:
            exec_result = sandbox.run(command, capability=name)
        except CielError as exc:
            return ToolResult(
                id=tool_call_id,
                name=name,
                error=str(exc),
                metadata={
                    "tenant_id": tenant_id,
                    "sandbox_backend": getattr(sandbox, "backend", "inprocess"),
                    "limits_applied": False,
                    "sandbox_rejected": True,
                },
            )
        backend = getattr(exec_result, "backend", "inprocess")
        limits_applied = getattr(exec_result, "limits_applied", False)
        if getattr(exec_result, "exit_code", 0) != 0 and getattr(exec_result, "stderr", ""):
            return ToolResult(
                id=tool_call_id,
                name=name,
                error=exec_result.stderr,
                metadata={
                    "tenant_id": tenant_id,
                    "sandbox_backend": backend,
                    "limits_applied": limits_applied,
                },
            )
        return ToolResult(
            id=tool_call_id,
            name=name,
            output=exec_result.stdout,
            metadata={
                "tenant_id": tenant_id,
                "sandbox_backend": backend,
                "limits_applied": limits_applied,
            },
        )

    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
    allow_insecure_inprocess: bool = False

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

    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, capability="terminal")
        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)
    # Si es False (default seguro) y el backend fuerte no está disponible para
    # una capacidad peligrosa, el executor RECHAZA la ejecución en inprocess en
    # lugar de degradar silenciosamente. Ponerlo en True es explícito y ruidoso.
    allow_insecure_inprocess: bool = False

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
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
645
646
647
648
649
650
651
652
653
654
655
656
657
658
659
660
661
662
663
664
665
666
667
668
669
670
671
672
673
674
675
676
677
678
679
680
681
682
683
684
685
686
687
688
689
690
691
692
693
694
695
696
697
698
699
700
701
702
703
704
705
706
707
708
709
710
711
712
713
714
715
716
717
718
719
720
721
722
723
724
725
726
727
728
729
730
731
732
733
734
735
736
737
738
739
740
741
742
743
744
745
746
747
748
749
750
751
752
753
754
755
756
757
758
759
760
761
762
763
764
765
766
767
768
769
770
771
772
773
774
775
776
777
778
779
780
781
782
783
784
785
786
787
788
789
790
791
792
793
794
795
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()

    # --- prompts versionados (Fase 19, tenant-filtered) ---------------------
    def prompt_save(
        self,
        *,
        tenant_id: Optional[str],
        name: str,
        version: str,
        prompt_text: str,
        value_json: str,
        sha256: str,
        previous_version: Optional[str],
        created_at: str,
    ) -> None:
        tenant_value = self._sentinel(tenant_id)
        self.conn.execute(
            """
            INSERT INTO prompt_versions (tenant_id, name, version, prompt_text, value_json, sha256, previous_version, created_at)
            VALUES (?, ?, ?, ?, ?, ?, ?, ?)
            ON CONFLICT(tenant_id, name, version) DO UPDATE SET
                prompt_text = excluded.prompt_text,
                value_json = excluded.value_json,
                sha256 = excluded.sha256,
                previous_version = excluded.previous_version,
                created_at = excluded.created_at
            """,
            (tenant_value, name, version, prompt_text, value_json, sha256, previous_version, created_at),
        )
        self.conn.commit()

    def prompt_get(
        self, *, tenant_id: Optional[str], name: str, version: Optional[str] = None
    ) -> Optional[dict]:
        tenant_value = self._sentinel(tenant_id)
        if version:
            row = self.conn.execute(
                "SELECT value_json FROM prompt_versions WHERE tenant_id = ? AND name = ? AND version = ?",
                (tenant_value, name, version),
            ).fetchone()
        else:
            row = self.conn.execute(
                "SELECT value_json FROM prompt_versions WHERE tenant_id = ? AND name = ? ORDER BY id DESC LIMIT 1",
                (tenant_value, name),
            ).fetchone()
        if row is None:
            return None
        return self._normalize_row(row["value_json"])

    def prompt_get_history(
        self, *, tenant_id: Optional[str], name: str
    ) -> List[dict]:
        tenant_value = self._sentinel(tenant_id)
        rows = self.conn.execute(
            "SELECT value_json FROM prompt_versions WHERE tenant_id = ? AND name = ? ORDER BY id ASC",
            (tenant_value, name),
        ).fetchall()
        out: List[dict] = []
        for row in rows:
            parsed = self._normalize_row(row["value_json"])
            if isinstance(parsed, dict):
                out.append(parsed)
        return out

    # --- log de estado cognitivo (Fase 19, tenant-filtered) ----------------
    def state_log_append(
        self,
        *,
        tenant_id: Optional[str],
        session_id: str,
        prompt_version: Optional[str],
        value_json: str,
        created_at: str,
    ) -> None:
        tenant_value = self._sentinel(tenant_id)
        self.conn.execute(
            """
            INSERT INTO cognitive_state_log (tenant_id, session_id, prompt_version, value_json, created_at)
            VALUES (?, ?, ?, ?, ?)
            """,
            (tenant_value, session_id, prompt_version, value_json, created_at),
        )
        self.conn.commit()

    def state_log_get_recent(
        self, *, tenant_id: Optional[str], session_id: str, limit: int = 16
    ) -> List[dict]:
        tenant_value = self._sentinel(tenant_id)
        rows = self.conn.execute(
            """
            SELECT value_json FROM cognitive_state_log
            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

    # --- curricula versionados (Fase 20 BLOQUE A, tenant-filtered) ----------
    def curriculum_save(
        self,
        *,
        tenant_id: Optional[str],
        goal: str,
        version: str,
        plan_json: str,
        value_json: str,
        sha256: str,
        previous_version: Optional[str],
        created_at: str,
    ) -> None:
        tenant_value = self._sentinel(tenant_id)
        self.conn.execute(
            """
            INSERT INTO curricula (tenant_id, goal, version, plan_json, value_json, sha256, previous_version, created_at)
            VALUES (?, ?, ?, ?, ?, ?, ?, ?)
            ON CONFLICT(tenant_id, goal, version) DO UPDATE SET
                plan_json = excluded.plan_json,
                value_json = excluded.value_json,
                sha256 = excluded.sha256,
                previous_version = excluded.previous_version,
                created_at = excluded.created_at
            """,
            (tenant_value, goal, version, plan_json, value_json, sha256, previous_version, created_at),
        )
        self.conn.commit()

    def curriculum_get(
        self, *, tenant_id: Optional[str], goal: str, version: Optional[str] = None
    ) -> Optional[dict]:
        tenant_value = self._sentinel(tenant_id)
        if version:
            row = self.conn.execute(
                "SELECT value_json FROM curricula WHERE tenant_id = ? AND goal = ? AND version = ?",
                (tenant_value, goal, version),
            ).fetchone()
        else:
            row = self.conn.execute(
                "SELECT value_json FROM curricula WHERE tenant_id = ? AND goal = ? ORDER BY id DESC LIMIT 1",
                (tenant_value, goal),
            ).fetchone()
        if row is None:
            return None
        return self._normalize_row(row["value_json"])

    def curriculum_get_history(
        self, *, tenant_id: Optional[str], goal: str
    ) -> List[dict]:
        tenant_value = self._sentinel(tenant_id)
        rows = self.conn.execute(
            "SELECT value_json FROM curricula WHERE tenant_id = ? AND goal = ? ORDER BY id ASC",
            (tenant_value, goal),
        ).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 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)
        );
        """
    )
    # --- prompts versionados (Fase 19, tenant-filtered) --------------------
    conn.execute(
        'CREATE TABLE IF NOT EXISTS prompt_versions (id INTEGER PRIMARY KEY AUTOINCREMENT, tenant_id TEXT, name TEXT, version TEXT, prompt_text TEXT, value_json TEXT, sha256 TEXT, previous_version TEXT, created_at TEXT, UNIQUE(tenant_id, name, version))'
    )
    # --- log de estado cognitivo (Fase 19, tenant-filtered) ----------------
    conn.execute(
        'CREATE TABLE IF NOT EXISTS cognitive_state_log (id INTEGER PRIMARY KEY AUTOINCREMENT, tenant_id TEXT, session_id TEXT, prompt_version TEXT, value_json TEXT, created_at TEXT)'
    )
    # --- curricula versionados (Fase 20 BLOQUE A, tenant-filtered) ----------
    conn.execute(
        'CREATE TABLE IF NOT EXISTS curricula (id INTEGER PRIMARY KEY AUTOINCREMENT, tenant_id TEXT, goal TEXT, version TEXT, plan_json TEXT, value_json TEXT, sha256 TEXT, previous_version TEXT, created_at TEXT, UNIQUE(tenant_id, goal, version))'
    )
    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()

Fase 19 — Prompt evolution versionado (Autonomía II, v0.13.0).

Este módulo aporta versionado semántico de prompts del agente (instructions sistemáticas) con persistencia offline en el StateBackend (SQLite en dev / Postgres en prod) y aislamiento estricto por tenant_id.

Molde: skill_versioning.py (Fase 12). A diferencia de los skills (que viven en memoria dentro de SkillLibrary), los prompts SÍ se persisten en SQLite a través del StateBackend para que sobrevivan reinicios y sean auditables.

Todo es network-free y API-key-free (offline-safe).

INITIAL_VERSION = '0.0.0' module-attribute

__all__ = ['PromptVersioningError', 'PromptVersion', 'PromptRegistry', 'sha256_text', 'INITIAL_VERSION'] module-attribute

PromptRegistry

Registro multitenant de prompts versionados sobre un StateBackend.

El backend es cualquier instancia de ciel.runtime.state_backend.StateBackend (SqliteStateBackend por defecto, PostgresStateBackend en prod). Todas las lecturas/escrituras filtran estrictamente por tenant_id.

Source code in src\ciel\runtime\prompt_versioning.py
class PromptRegistry:
    """Registro multitenant de prompts versionados sobre un ``StateBackend``.

    El ``backend`` es cualquier instancia de ``ciel.runtime.state_backend.StateBackend``
    (``SqliteStateBackend`` por defecto, ``PostgresStateBackend`` en prod). Todas
    las lecturas/escrituras filtran estrictamente por ``tenant_id``.
    """

    def __init__(self, backend: Any) -> None:
        # backend: ciel.runtime.state_backend.StateBackend
        self._backend = backend

    # --- escritura ----------------------------------------------------------

    def create(
        self,
        name: str,
        prompt_text: str,
        *,
        tenant_id: Optional[str] = None,
        changelog: str = "",
        metadata: Optional[Dict[str, Any]] = None,
    ) -> PromptVersion:
        """Crea la versión inicial ``0.0.0`` de un prompt.

        Lanza :class:`PromptVersioningError` si el nombre ya existe para el tenant.
        """
        existing = self._backend.prompt_get(tenant_id=tenant_id, name=name)
        if existing is not None:
            raise PromptVersioningError(
                f"el prompt {name!r} ya existe para tenant {tenant_id!r}; usa update()"
            )
        pv = PromptVersion(
            name=name,
            prompt_text=prompt_text,
            changelog=changelog,
            released_at=datetime.now(timezone.utc),
            sha256=sha256_text(prompt_text),
            previous_version=None,
            parent=None,
            metadata=metadata or {},
        )
        self._save(pv, tenant_id=tenant_id)
        return pv

    def update(
        self,
        name: str,
        prompt_text: str,
        *,
        tenant_id: Optional[str] = None,
        bump: str = "patch",
        changelog: str = "",
        metadata: Optional[Dict[str, Any]] = None,
    ) -> PromptVersion:
        """Bumpea un prompt existente y guarda la nueva versión."""
        current = self.get(name, tenant_id=tenant_id)
        if current is None:
            raise PromptVersioningError(
                f"no existe el prompt {name!r} para tenant {tenant_id!r}; usa create()"
            )
        new_ver = current.bump(bump)
        pv = PromptVersion(
            name=name,
            prompt_text=prompt_text,
            changelog=changelog,
            released_at=datetime.now(timezone.utc),
            sha256=sha256_text(prompt_text),
            previous_version=current.version,
            parent=current.version,
            metadata=metadata or current.metadata or {},
        )
        # conserva el major/minor/patch calculado por bump
        pv.major, pv.minor, pv.patch = new_ver.major, new_ver.minor, new_ver.patch
        self._save(pv, tenant_id=tenant_id)
        return pv

    def _save(self, pv: PromptVersion, *, tenant_id: Optional[str]) -> None:
        self._backend.prompt_save(
            tenant_id=tenant_id,
            name=pv.name,
            version=pv.version,
            prompt_text=pv.prompt_text,
            value_json=_dump(pv.to_dict()),
            sha256=pv.sha256,
            previous_version=pv.previous_version,
            created_at=pv.released_at.isoformat() if pv.released_at else _now_iso(),
        )

    # --- lectura ------------------------------------------------------------

    def get(
        self,
        name: str,
        *,
        tenant_id: Optional[str] = None,
        version: Optional[str] = None,
    ) -> Optional[PromptVersion]:
        row = self._backend.prompt_get(tenant_id=tenant_id, name=name, version=version)
        if row is None:
            return None
        return PromptVersion.from_dict(row)

    def history(self, name: str, *, tenant_id: Optional[str] = None) -> List[PromptVersion]:
        rows = self._backend.prompt_get_history(tenant_id=tenant_id, name=name)
        return [PromptVersion.from_dict(r) for r in rows]

    def evolution_tree(self, name: str, *, tenant_id: Optional[str] = None) -> Dict[str, Any]:
        """Devuelve el árbol de linaje de un prompt versionado.

        Misma forma que ``skill_versioning.evolution_tree``:

        ``{"name", "root", "lineage": [...], "nodes": {version: {...}}}``
        con ``parent`` / ``children`` / ``sha256`` / ``changelog`` / ``released_at``.
        """
        versions = self.history(name, tenant_id=tenant_id)
        if not versions:
            raise PromptVersioningError(f"unknown prompt: {name!r} (tenant {tenant_id!r})")

        keys: List[str] = [v.version for v in versions]
        nodes: Dict[str, Dict[str, Any]] = {}
        for idx, pv in enumerate(versions):
            key = pv.version
            raw_prev = pv.previous_version
            parent = raw_prev if raw_prev is not None else (keys[idx - 1] if idx > 0 else None)
            nodes[key] = {
                "version": key,
                "parent": parent,
                "previous_version": raw_prev,
                "children": [],
                "sha256": pv.sha256,
                "changelog": pv.changelog,
                "released_at": pv.released_at.isoformat() if pv.released_at else None,
                "prompt_text": pv.prompt_text,
            }

        for key, node in nodes.items():
            if node["parent"] is not None and node["parent"] in nodes:
                nodes[node["parent"]]["children"].append(key)

        roots = [k for k, n in nodes.items() if n["parent"] is None]
        root = roots[0] if roots else (keys[0] if keys else None)
        return {
            "name": name,
            "root": root,
            "lineage": list(keys),
            "nodes": nodes,
        }

create(name: str, prompt_text: str, *, tenant_id: Optional[str] = None, changelog: str = '', metadata: Optional[Dict[str, Any]] = None) -> PromptVersion

Crea la versión inicial 0.0.0 de un prompt.

Lanza :class:PromptVersioningError si el nombre ya existe para el tenant.

Source code in src\ciel\runtime\prompt_versioning.py
def create(
    self,
    name: str,
    prompt_text: str,
    *,
    tenant_id: Optional[str] = None,
    changelog: str = "",
    metadata: Optional[Dict[str, Any]] = None,
) -> PromptVersion:
    """Crea la versión inicial ``0.0.0`` de un prompt.

    Lanza :class:`PromptVersioningError` si el nombre ya existe para el tenant.
    """
    existing = self._backend.prompt_get(tenant_id=tenant_id, name=name)
    if existing is not None:
        raise PromptVersioningError(
            f"el prompt {name!r} ya existe para tenant {tenant_id!r}; usa update()"
        )
    pv = PromptVersion(
        name=name,
        prompt_text=prompt_text,
        changelog=changelog,
        released_at=datetime.now(timezone.utc),
        sha256=sha256_text(prompt_text),
        previous_version=None,
        parent=None,
        metadata=metadata or {},
    )
    self._save(pv, tenant_id=tenant_id)
    return pv

evolution_tree(name: str, *, tenant_id: Optional[str] = None) -> Dict[str, Any]

Devuelve el árbol de linaje de un prompt versionado.

Misma forma que skill_versioning.evolution_tree:

{"name", "root", "lineage": [...], "nodes": {version: {...}}} con parent / children / sha256 / changelog / released_at.

Source code in src\ciel\runtime\prompt_versioning.py
def evolution_tree(self, name: str, *, tenant_id: Optional[str] = None) -> Dict[str, Any]:
    """Devuelve el árbol de linaje de un prompt versionado.

    Misma forma que ``skill_versioning.evolution_tree``:

    ``{"name", "root", "lineage": [...], "nodes": {version: {...}}}``
    con ``parent`` / ``children`` / ``sha256`` / ``changelog`` / ``released_at``.
    """
    versions = self.history(name, tenant_id=tenant_id)
    if not versions:
        raise PromptVersioningError(f"unknown prompt: {name!r} (tenant {tenant_id!r})")

    keys: List[str] = [v.version for v in versions]
    nodes: Dict[str, Dict[str, Any]] = {}
    for idx, pv in enumerate(versions):
        key = pv.version
        raw_prev = pv.previous_version
        parent = raw_prev if raw_prev is not None else (keys[idx - 1] if idx > 0 else None)
        nodes[key] = {
            "version": key,
            "parent": parent,
            "previous_version": raw_prev,
            "children": [],
            "sha256": pv.sha256,
            "changelog": pv.changelog,
            "released_at": pv.released_at.isoformat() if pv.released_at else None,
            "prompt_text": pv.prompt_text,
        }

    for key, node in nodes.items():
        if node["parent"] is not None and node["parent"] in nodes:
            nodes[node["parent"]]["children"].append(key)

    roots = [k for k, n in nodes.items() if n["parent"] is None]
    root = roots[0] if roots else (keys[0] if keys else None)
    return {
        "name": name,
        "root": root,
        "lineage": list(keys),
        "nodes": nodes,
    }

update(name: str, prompt_text: str, *, tenant_id: Optional[str] = None, bump: str = 'patch', changelog: str = '', metadata: Optional[Dict[str, Any]] = None) -> PromptVersion

Bumpea un prompt existente y guarda la nueva versión.

Source code in src\ciel\runtime\prompt_versioning.py
def update(
    self,
    name: str,
    prompt_text: str,
    *,
    tenant_id: Optional[str] = None,
    bump: str = "patch",
    changelog: str = "",
    metadata: Optional[Dict[str, Any]] = None,
) -> PromptVersion:
    """Bumpea un prompt existente y guarda la nueva versión."""
    current = self.get(name, tenant_id=tenant_id)
    if current is None:
        raise PromptVersioningError(
            f"no existe el prompt {name!r} para tenant {tenant_id!r}; usa create()"
        )
    new_ver = current.bump(bump)
    pv = PromptVersion(
        name=name,
        prompt_text=prompt_text,
        changelog=changelog,
        released_at=datetime.now(timezone.utc),
        sha256=sha256_text(prompt_text),
        previous_version=current.version,
        parent=current.version,
        metadata=metadata or current.metadata or {},
    )
    # conserva el major/minor/patch calculado por bump
    pv.major, pv.minor, pv.patch = new_ver.major, new_ver.minor, new_ver.patch
    self._save(pv, tenant_id=tenant_id)
    return pv

PromptVersion dataclass

Versión semántica enriquecida de un prompt del agente.

Lleva major.minor.patch + changelog + released_at (timestamp) + el hash sha256 del texto, y la trazabilidad de linaje (previous_version / parent).

Source code in src\ciel\runtime\prompt_versioning.py
@dataclass
class PromptVersion:
    """Versión semántica enriquecida de un prompt del agente.

    Lleva ``major.minor.patch`` + ``changelog`` + ``released_at`` (timestamp) +
    el hash ``sha256`` del texto, y la trazabilidad de linaje
    (``previous_version`` / ``parent``).
    """

    major: int = 0
    minor: int = 0
    patch: int = 0
    name: str = ""
    prompt_text: str = ""
    changelog: str = ""
    released_at: Optional[datetime] = None
    sha256: str = ""
    previous_version: Optional[str] = None
    parent: Optional[str] = None
    metadata: Dict[str, Any] = field(default_factory=dict)

    # --- serialización ------------------------------------------------------

    @property
    def version(self) -> str:
        return f"{self.major}.{self.minor}.{self.patch}"

    @classmethod
    def parse(cls, version: str) -> "PromptVersion":
        parts = (version.split(".") + ["0", "0", "0"])[:3]
        try:
            major, minor, patch = (int(p) for p in parts)
        except ValueError as exc:
            raise PromptVersioningError(f"invalid version string: {version!r}") from exc
        return cls(major=major, minor=minor, patch=patch)

    def bump(self, kind: str = "patch") -> "PromptVersion":
        kind = (kind or "patch").lower()
        if kind == "major":
            return PromptVersion(self.major + 1, 0, 0)
        if kind == "minor":
            return PromptVersion(self.major, self.minor + 1, 0)
        if kind == "patch":
            return PromptVersion(self.major, self.minor, self.patch + 1)
        raise PromptVersioningError(f"unknown bump kind: {kind!r} (expected major/minor/patch)")

    def to_dict(self) -> Dict[str, Any]:
        return {
            "name": self.name,
            "version": self.version,
            "major": self.major,
            "minor": self.minor,
            "patch": self.patch,
            "prompt_text": self.prompt_text,
            "changelog": self.changelog,
            "released_at": self.released_at.isoformat() if self.released_at else None,
            "sha256": self.sha256,
            "previous_version": self.previous_version,
            "parent": self.parent,
            "metadata": self.metadata,
        }

    @classmethod
    def from_dict(cls, data: Dict[str, Any]) -> "PromptVersion":
        ver = cls.parse(data.get("version") or INITIAL_VERSION)
        ver.name = data.get("name", "")
        ver.prompt_text = data.get("prompt_text", "")
        ver.changelog = data.get("changelog", "")
        ver.sha256 = data.get("sha256", "")
        ver.previous_version = data.get("previous_version")
        ver.parent = data.get("parent")
        ver.metadata = data.get("metadata") or {}
        raw = data.get("released_at")
        if raw:
            try:
                ver.released_at = datetime.fromisoformat(raw)
            except (ValueError, TypeError):  # pragma: no cover - defensive
                ver.released_at = None
        return ver

PromptVersioningError

Bases: Exception

Error de versionado de prompts.

Source code in src\ciel\runtime\prompt_versioning.py
class PromptVersioningError(Exception):
    """Error de versionado de prompts."""

_dump(value: Any) -> str

Source code in src\ciel\runtime\prompt_versioning.py
def _dump(value: Any) -> str:
    import json

    try:
        return json.dumps(value, ensure_ascii=False)
    except TypeError:  # pragma: no cover - defensive
        return json.dumps({"repr": repr(value)})

_now_iso() -> str

Source code in src\ciel\runtime\prompt_versioning.py
def _now_iso() -> str:
    return datetime.now(timezone.utc).isoformat()

sha256_text(text: str) -> str

Hash determinista del texto del prompt (offline).

Source code in src\ciel\runtime\prompt_versioning.py
def sha256_text(text: str) -> str:
    """Hash determinista del texto del prompt (offline)."""
    return hashlib.sha256(text.encode("utf-8")).hexdigest()

Integración aditiva de self-reflection + learning-from-failure (Fase 19).

IGUAL que memory_agent_integration: NO reescribe api.py. Se invoca install_agent_reflection_support(Agent) al final de ciel/api.py y engancha:

  • Agent(reflection=...) / Agent(reflection_config=...) — habilita la reflexión post-run (opcional, offline-safe).
  • Tras cada arun/run, si el run tuvo fallos de tool, genera una lección determinista (sin red) y la persiste como memoria episódica role="lesson" (reutiliza el EpisodicStore de F17 → multitenant).
  • Exponer AgentResponse.reflection (property aditiva) con el resumen.

Degrada graceful: si no se instala, getattr(self, "_reflection", None) es None y el run no cambia.

__all__ = ['AgentReflection', 'ReflectionConfig', 'install_agent_reflection_support'] module-attribute

AgentReflection

State holder de reflexión para un Agent (aditivo, no invasivo).

Genera lecciones deterministas a partir de fallos de tool (sin red) y las persiste vía el EpisodicStore (memoria episódica F17, aislada por tenant). Si no hay store, opera en modo no-persistente (solo resumen).

Source code in src\ciel\runtime\reflection_agent_integration.py
class AgentReflection:
    """State holder de reflexión para un ``Agent`` (aditivo, no invasivo).

    Genera lecciones deterministas a partir de fallos de tool (sin red) y las
    persiste vía el ``EpisodicStore`` (memoria episódica F17, aislada por
    tenant). Si no hay store, opera en modo no-persistente (solo resumen).
    """

    def __init__(
        self,
        store: Any,
        config: Optional[ReflectionConfig] = None,
    ) -> None:
        # store: ciel.runtime.memory_episodic.EpisodicStore | None
        self.store = store
        self.config = config or ReflectionConfig()

    @property
    def enabled(self) -> bool:
        return self.config.enabled

    # ------------------------------------------------------------------ #
    def reflect(
        self,
        response: Any,
        *,
        tenant_id: Optional[str],
        session_id: str,
    ) -> Optional[dict]:
        """Reflexiona sobre un ``AgentResponse`` y (opcionalmente) persiste lección.

        Devuelve un dict resumen:
        ``{"had_failure", "failed_tools", "lessons_count", "lesson"}``
        o ``None`` si la reflexión está deshabilitada.
        """
        if not self.enabled:
            return None
        tool_results = response.tool_results
        failed = [r for r in tool_results if getattr(r, "error", None)]
        had_failure = bool(failed)
        lesson = None
        if had_failure:
            lesson = self._build_lesson(failed, response=response)
            if self.config.persist_lessons and self.store is not None:
                try:
                    self.store.append(
                        tenant_id=tenant_id,
                        session_id=session_id,
                        role="lesson",
                        content=lesson,
                        metadata={"kind": "reflection", "had_failure": True},
                    )
                except Exception:  # pragma: no cover - reflexión nunca rompe el run
                    pass
        return {
            "had_failure": had_failure,
            "failed_tools": [r.name for r in failed],
            "lessons_count": len(failed) if lesson else 0,
            "lesson": lesson,
        }

    def _build_lesson(self, failed: List[Any], *, response: Any) -> dict:
        """Lección determinista a partir de los fallos (sin red)."""
        details = []
        for r in failed:
            details.append(
                {
                    "tool": r.name,
                    "error": r.error,
                }
            )
        return {
            "type": "learning_from_failure",
            "created_at": datetime.now(timezone.utc).isoformat(),
            "summary": (
                f"El agente falló en {len(failed)} tool(s): "
                + ", ".join(d["tool"] for d in details)
                + ". En próximos runs, considera validar argumentos o usar un tool alternativo."
            ),
            "failures": details,
            "final_text": (response.text or "")[:500],
        }

    def lessons(self, *, tenant_id: Optional[str], session_id: str, limit: int = 5) -> List[dict]:
        """Recupera las lecciones persistidas (memoria episódica role='lesson')."""
        if self.store is None:
            return []
        recent = self.store.get_recent(
            tenant_id=tenant_id, session_id=session_id, limit=max(limit, self.config.max_lessons)
        )
        out: List[dict] = []
        for m in recent:
            if getattr(m, "role", None) == "lesson":
                content = m.content
                out.append(content if isinstance(content, dict) else {"content": content})
        return out

lessons(*, tenant_id: Optional[str], session_id: str, limit: int = 5) -> List[dict]

Recupera las lecciones persistidas (memoria episódica role='lesson').

Source code in src\ciel\runtime\reflection_agent_integration.py
def lessons(self, *, tenant_id: Optional[str], session_id: str, limit: int = 5) -> List[dict]:
    """Recupera las lecciones persistidas (memoria episódica role='lesson')."""
    if self.store is None:
        return []
    recent = self.store.get_recent(
        tenant_id=tenant_id, session_id=session_id, limit=max(limit, self.config.max_lessons)
    )
    out: List[dict] = []
    for m in recent:
        if getattr(m, "role", None) == "lesson":
            content = m.content
            out.append(content if isinstance(content, dict) else {"content": content})
    return out

reflect(response: Any, *, tenant_id: Optional[str], session_id: str) -> Optional[dict]

Reflexiona sobre un AgentResponse y (opcionalmente) persiste lección.

Devuelve un dict resumen: {"had_failure", "failed_tools", "lessons_count", "lesson"} o None si la reflexión está deshabilitada.

Source code in src\ciel\runtime\reflection_agent_integration.py
def reflect(
    self,
    response: Any,
    *,
    tenant_id: Optional[str],
    session_id: str,
) -> Optional[dict]:
    """Reflexiona sobre un ``AgentResponse`` y (opcionalmente) persiste lección.

    Devuelve un dict resumen:
    ``{"had_failure", "failed_tools", "lessons_count", "lesson"}``
    o ``None`` si la reflexión está deshabilitada.
    """
    if not self.enabled:
        return None
    tool_results = response.tool_results
    failed = [r for r in tool_results if getattr(r, "error", None)]
    had_failure = bool(failed)
    lesson = None
    if had_failure:
        lesson = self._build_lesson(failed, response=response)
        if self.config.persist_lessons and self.store is not None:
            try:
                self.store.append(
                    tenant_id=tenant_id,
                    session_id=session_id,
                    role="lesson",
                    content=lesson,
                    metadata={"kind": "reflection", "had_failure": True},
                )
            except Exception:  # pragma: no cover - reflexión nunca rompe el run
                pass
    return {
        "had_failure": had_failure,
        "failed_tools": [r.name for r in failed],
        "lessons_count": len(failed) if lesson else 0,
        "lesson": lesson,
    }

EpisodicStore

Store de memoria episódica sobre un StateBackend (aislado por tenant).

Operaciones: * append — persiste un turno. * get_recent — últimos N episodios de una sesión (orden cronológico). * get_by_id — recupera un episodio por su ID. * search — búsqueda por keywords filtrada POR TENANT (no cross-tenant). * clear_session — borra la memoria de una sesión. * as_context — serializa episodios recientes como texto para el system.

Source code in src\ciel\runtime\memory_episodic.py
class EpisodicStore:
    """Store de memoria episódica sobre un ``StateBackend`` (aislado por tenant).

    Operaciones:
    * ``append`` — persiste un turno.
    * ``get_recent`` — últimos N episodios de una sesión (orden cronológico).
    * ``get_by_id`` — recupera un episodio por su ID.
    * ``search`` — búsqueda por keywords filtrada POR TENANT (no cross-tenant).
    * ``clear_session`` — borra la memoria de una sesión.
    * ``as_context`` — serializa episodios recientes como texto para el system.
    """

    def __init__(self, backend: Any, *, namespace: str = "episodic") -> None:
        # backend: ciel.runtime.state_backend.StateBackend
        self._backend = backend
        self._ns = namespace

    # -- helpers internos ----------------------------------------------------
    def _key(self, session_id: str, memory_id: str) -> str:
        return f"{self._ns}:{session_id}:{memory_id}"

    # -- escritura -----------------------------------------------------------
    def append(
        self,
        *,
        tenant_id: Optional[str],
        session_id: str,
        role: str,
        content: Any,
        metadata: Optional[dict] = None,
    ) -> EpisodicMemory:
        memory_id = str(uuid.uuid4())
        mem = EpisodicMemory(
            id=memory_id,
            tenant_id=tenant_id,
            session_id=session_id,
            role=role,
            content=content,
            metadata=metadata or {},
        )
        self._backend.memory_append(
            tenant_id=tenant_id,
            session_id=session_id,
            memory_id=memory_id,
            value=mem.to_dict(),
        )
        return mem

    # -- lectura -------------------------------------------------------------
    def get_recent(
        self, *, tenant_id: Optional[str], session_id: str, limit: int = 8
    ) -> List[EpisodicMemory]:
        rows = self._backend.memory_get_recent(
            tenant_id=tenant_id, session_id=session_id, limit=limit
        )
        out: List[EpisodicMemory] = []
        for row in rows:
            try:
                out.append(EpisodicMemory.from_dict(row))
            except (KeyError, TypeError):  # pragma: no cover - defensive
                continue
        # Orden cronológico ascendente (más antiguo primero).
        out.sort(key=lambda m: m.created_at)
        return out

    def get_by_id(
        self, *, tenant_id: Optional[str], session_id: str, memory_id: str
    ) -> Optional[EpisodicMemory]:
        row = self._backend.memory_get(
            tenant_id=tenant_id, session_id=session_id, memory_id=memory_id
        )
        if row is None:
            return None
        try:
            return EpisodicMemory.from_dict(row)
        except (KeyError, TypeError):  # pragma: no cover - defensive
            return None

    def search(
        self,
        *,
        tenant_id: Optional[str],
        session_id: Optional[str] = None,
        query: str,
        limit: int = 5,
    ) -> List[EpisodicMemory]:
        rows = self._backend.memory_search_tenant(
            tenant_id=tenant_id,
            session_id=session_id,
            query=query,
            limit=limit,
        )
        out: List[EpisodicMemory] = []
        for row in rows:
            try:
                out.append(EpisodicMemory.from_dict(row))
            except (KeyError, TypeError):  # pragma: no cover - defensive
                continue
        return out

    def clear_session(self, *, tenant_id: Optional[str], session_id: str) -> None:
        self._backend.memory_clear_session(
            tenant_id=tenant_id, session_id=session_id
        )

    # -- serialización para el system prompt --------------------------------
    @staticmethod
    def format_as_context(memories: Sequence[EpisodicMemory]) -> str:
        if not memories:
            return ""
        lines: List[str] = []
        for m in memories:
            content = m.content
            if isinstance(content, (list, dict)):  # multimodal/structured
                import json

                try:
                    content = json.dumps(content, ensure_ascii=False)
                except (TypeError, ValueError):  # pragma: no cover
                    content = str(content)
            elif not isinstance(content, str):
                content = str(content)
            lines.append(f"[{m.role}] {content}")
        return "\n".join(lines)

    def as_context(
        self,
        *,
        tenant_id: Optional[str],
        session_id: str,
        recent: int = 8,
        query: Optional[str] = None,
        search_limit: int = 5,
    ) -> str:
        """Devuelve el texto de memoria para inyectar en el system prompt."""
        memories: List[EpisodicMemory] = list(
            self.get_recent(
                tenant_id=tenant_id, session_id=session_id, limit=recent
            )
        )
        if query:
            memories.extend(
                self.search(
                    tenant_id=tenant_id,
                    session_id=session_id,
                    query=query,
                    limit=search_limit,
                )
            )
        # Dedup por id preservando orden.
        seen = set()
        unique: List[EpisodicMemory] = []
        for m in memories:
            if m.id in seen:
                continue
            seen.add(m.id)
            unique.append(m)
        return self.format_as_context(unique)

as_context(*, tenant_id: Optional[str], session_id: str, recent: int = 8, query: Optional[str] = None, search_limit: int = 5) -> str

Devuelve el texto de memoria para inyectar en el system prompt.

Source code in src\ciel\runtime\memory_episodic.py
def as_context(
    self,
    *,
    tenant_id: Optional[str],
    session_id: str,
    recent: int = 8,
    query: Optional[str] = None,
    search_limit: int = 5,
) -> str:
    """Devuelve el texto de memoria para inyectar en el system prompt."""
    memories: List[EpisodicMemory] = list(
        self.get_recent(
            tenant_id=tenant_id, session_id=session_id, limit=recent
        )
    )
    if query:
        memories.extend(
            self.search(
                tenant_id=tenant_id,
                session_id=session_id,
                query=query,
                limit=search_limit,
            )
        )
    # Dedup por id preservando orden.
    seen = set()
    unique: List[EpisodicMemory] = []
    for m in memories:
        if m.id in seen:
            continue
        seen.add(m.id)
        unique.append(m)
    return self.format_as_context(unique)

ReflectionConfig dataclass

Configuración de reflexión del agente (degrada a 'sin reflexión').

Source code in src\ciel\runtime\reflection_agent_integration.py
@dataclass
class ReflectionConfig:
    """Configuración de reflexión del agente (degrada a 'sin reflexión')."""

    enabled: bool = True
    persist_lessons: bool = True
    max_lessons: int = 5

    @classmethod
    def disabled(cls) -> "ReflectionConfig":
        return cls(enabled=False)

_session_id_of(resp: AgentResponse) -> str

Source code in src\ciel\runtime\reflection_agent_integration.py
def _session_id_of(resp: AgentResponse) -> str:
    raw = getattr(resp, "raw", None)
    meta = getattr(raw, "metadata", None) or {}
    sid = meta.get("session_id")
    return sid if isinstance(sid, str) else "default"

install_agent_reflection_support(agent_cls: Any) -> Any

Engancha reflection= / reflection_config= en Agent sin reescribir api.py.

Idempotente: si ya se instaló, no re-envuelve (evita doble wrapper en tests).

Source code in src\ciel\runtime\reflection_agent_integration.py
def install_agent_reflection_support(agent_cls: Any) -> Any:
    """Engancha ``reflection=`` / ``reflection_config=`` en ``Agent`` sin reescribir api.py.

    Idempotente: si ya se instaló, no re-envuelve (evita doble wrapper en tests).
    """
    if getattr(agent_cls, "_reflection_installed", False):
        return agent_cls
    original_init = agent_cls.__init__
    original_arun = agent_cls.arun

    @functools.wraps(original_init)
    def _patched_init(self: Any, *args: Any, **kwargs: Any) -> None:
        reflection_arg = kwargs.pop("reflection", None)
        reflection_config = kwargs.pop("reflection_config", None)
        original_init(self, *args, **kwargs)

        config = reflection_config or ReflectionConfig()
        # Resolver el store de lecciones: reutiliza el EpisodicStore de memoria
        # si existe; si no, crea uno temporal offline (solo si se habilitó).
        store = None
        mem = getattr(self, "_memory", None)
        if mem is not None and getattr(mem, "store", None) is not None:
            store = mem.store
        elif reflection_arg is True or reflection_config is not None:
            from ciel.runtime.state_backend import SqliteStateBackend

            tmp = tempfile.mkdtemp(prefix="ciel-reflect-")
            be = SqliteStateBackend(str(Path(tmp) / "reflect.sqlite"))
            store = EpisodicStore(be)
        elif isinstance(reflection_arg, EpisodicStore):  # type: ignore[name-defined]
            store = reflection_arg
        if reflection_arg is False or (reflection_arg is None and reflection_config is None):
            config = config.disabled()
        self._reflection = AgentReflection(store=store, config=config)

    @functools.wraps(original_arun)
    async def _patched_arun(self: Any, prompt: Any, *, tenant_id=None, max_turns=10, limit=32) -> Any:
        resp = await original_arun(self, prompt, tenant_id=tenant_id, max_turns=max_turns, limit=limit)
        refl = getattr(self, "_reflection", None)
        if refl is not None and refl.enabled:
            sid = _session_id_of(resp)
            try:
                lesson = refl.reflect(resp, tenant_id=tenant_id, session_id=sid)
            except Exception:  # pragma: no cover - reflexión nunca rompe el run
                lesson = None
            try:
                resp.raw.metadata["reflection"] = lesson
            except Exception:  # pragma: no cover - defensivo
                pass
        return resp

    agent_cls.__init__ = _patched_init  # type: ignore[assignment]
    agent_cls.arun = _patched_arun  # type: ignore[assignment]
    agent_cls._reflection_installed = True  # type: ignore[attr-defined]
    return agent_cls

Introspección / estado cognitivo explicable (Fase 19, v0.13.0).

Aditivo y offline-safe. Engancha Agent(introspection=...) sin reescribir api.py y:

  • Inyecta un bloque [Estado cognitivo] en el system prompt (en _build_request) con el último snapshot conocido del agente, para que el modelo sea consciente de su propio estado (vuelve a inyectar el contexto de la introspección).
  • Registra un CognitiveSnapshot post-run en cognitive_state_log del StateBackend (aislado por tenant/session).
  • Expone Agent.introspect() para volcar los últimos snapshots.

Reusa EpisodicStore (conteo de turnos) y DeterministicEmbeddingProvider (embedding del estado, offline). El RAG (retrieved_context_ids) queda como None cuando el agente no usa RAG; el campo es opcional.

Degrada graceful: si no se instala, _cognitive no existe y el run no cambia.

_DEFAULT_COGNITIVE_BACKEND: Optional[StateBackend] = None module-attribute

__all__ = ['CognitiveSnapshot', 'IntrospectionReport', 'IntrospectionConfig', 'CognitiveState', 'install_cognitive_state_support'] module-attribute

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)

CognitiveSnapshot dataclass

Instantánea del estado cognitivo del agente en un run.

Source code in src\ciel\runtime\cognitive_state.py
@dataclass
class CognitiveSnapshot:
    """Instantánea del estado cognitivo del agente en un run."""

    tenant_id: Optional[str]
    session_id: str
    active_prompt_version: Optional[str] = None
    memory_turn_count: int = 0
    retrieved_context_ids: List[str] = field(default_factory=list)
    tool_calls: List[Dict[str, Any]] = field(default_factory=list)
    had_failure: bool = False
    confidence: float = 1.0
    rationale: str = ""
    created_at: str = field(default_factory=lambda: datetime.now(timezone.utc).isoformat())

    def to_dict(self) -> Dict[str, Any]:
        return {
            "tenant_id": self.tenant_id,
            "session_id": self.session_id,
            "active_prompt_version": self.active_prompt_version,
            "memory_turn_count": self.memory_turn_count,
            "retrieved_context_ids": list(self.retrieved_context_ids),
            "tool_calls": list(self.tool_calls),
            "had_failure": self.had_failure,
            "confidence": self.confidence,
            "rationale": self.rationale,
            "created_at": self.created_at,
        }

    @classmethod
    def from_dict(cls, data: Dict[str, Any]) -> "CognitiveSnapshot":
        return cls(
            tenant_id=data.get("tenant_id"),
            session_id=data.get("session_id", ""),
            active_prompt_version=data.get("active_prompt_version"),
            memory_turn_count=int(data.get("memory_turn_count", 0)),
            retrieved_context_ids=list(data.get("retrieved_context_ids") or []),
            tool_calls=list(data.get("tool_calls") or []),
            had_failure=bool(data.get("had_failure", False)),
            confidence=float(data.get("confidence", 1.0)),
            rationale=data.get("rationale", ""),
            created_at=data.get("created_at", ""),
        )

CognitiveState

State holder de introspección para un Agent (aditivo, offline).

Source code in src\ciel\runtime\cognitive_state.py
class CognitiveState:
    """State holder de introspección para un ``Agent`` (aditivo, offline)."""

    def __init__(
        self,
        store: Optional[EpisodicStore],
        backend: StateBackend,
        config: Optional[IntrospectionConfig] = None,
        session_id: Optional[str] = None,
    ) -> None:
        # store: EpisodicStore (para conteo de turnos); backend: StateBackend (log).
        # ``backend`` es obligatorio: si el agente no tiene memoria, el caller
        # debe pasar un ``StateBackend`` persistente (defaulteado a una instancia
        # por-proceso) para que los snapshots sobrevivan al run y sean consultables.
        # ``session_id`` es estable por agente (no transitorio por run) para que
        # la inyección y la introspección coincidan run a run.
        self.store = store
        self.backend = backend
        self.config = config or IntrospectionConfig()
        self._embedder = DeterministicEmbeddingProvider(dim=self.config.embedding_dim)
        self.session_id = session_id or str(uuid.uuid4())
        self.tenant_id: Optional[str] = None

    @property
    def enabled(self) -> bool:
        return self.config.enabled

    # ------------------------------------------------------------------ #
    def build_snapshot(
        self,
        *,
        tenant_id: Optional[str],
        session_id: str,
        response: Any,
        active_prompt_version: Optional[str] = None,
        retrieved_context_ids: Optional[List[str]] = None,
    ) -> CognitiveSnapshot:
        tool_results = response.tool_results if hasattr(response, "tool_results") else []
        had_failure = any(getattr(r, "error", None) for r in tool_results)
        tool_calls = [
            {"name": r.name, "error": getattr(r, "error", None)} for r in tool_results
        ]
        # conteo de turnos de memoria episódica (si hay store)
        memory_turn_count = 0
        if self.store is not None:
            recent = self.store.get_recent(
                tenant_id=tenant_id, session_id=session_id, limit=1000
            )
            memory_turn_count = len(recent)
        # confianza heurística determinista offline
        if had_failure:
            confidence = 0.3
            rationale = "fallo(s) de tool detectado(s) en el run"
        elif tool_calls:
            confidence = 0.8
            rationale = "run con uso de tools sin fallos"
        else:
            confidence = 1.0
            rationale = "run directo sin tools"
        return CognitiveSnapshot(
            tenant_id=tenant_id,
            session_id=session_id,
            active_prompt_version=active_prompt_version,
            memory_turn_count=memory_turn_count,
            retrieved_context_ids=list(retrieved_context_ids or []),
            tool_calls=tool_calls,
            had_failure=had_failure,
            confidence=confidence,
            rationale=rationale,
        )

    def record(self, snapshot: CognitiveSnapshot) -> None:
        """Persiste el snapshot en ``cognitive_state_log`` (tenant-filtered)."""
        try:
            self.backend.state_log_append(
                tenant_id=snapshot.tenant_id,
                session_id=snapshot.session_id,
                prompt_version=snapshot.active_prompt_version,
                value_json=_dump(snapshot.to_dict()),
                created_at=snapshot.created_at
                or datetime.now(timezone.utc).isoformat(),
            )
        except Exception:  # pragma: no cover - introspección nunca rompe el run
            pass

    def get_recent(
        self, *, tenant_id: Optional[str], session_id: str, limit: int = 16
    ) -> IntrospectionReport:
        snapshots: List[CognitiveSnapshot] = []
        try:
            rows = self.backend.state_log_get_recent(
                tenant_id=tenant_id, session_id=session_id, limit=limit
            )
            for row in rows:
                try:
                    snapshots.append(CognitiveSnapshot.from_dict(row))
                except Exception:  # pragma: no cover - defensivo
                    continue
        except Exception:  # pragma: no cover - defensivo
            pass
        return IntrospectionReport(
            tenant_id=tenant_id, session_id=session_id, snapshots=snapshots
        )

    def latest_block(self, *, tenant_id: Optional[str], session_id: str) -> Optional[str]:
        """Devuelve el bloque ``[Estado cognitivo]`` a inyectar, o None."""
        if not self.config.inject_into_prompt:
            return None
        report = self.get_recent(tenant_id=tenant_id, session_id=session_id, limit=1)
        snap = report.latest
        if snap is None:
            return None
        lines = [
            f"Versión de prompt activa: {snap.active_prompt_version or 'instructions (no versionado)'}",
            f"Turnos de memoria acumulados: {snap.memory_turn_count}",
            f"Tool calls en el último run: {len(snap.tool_calls)}",
            f"Falló en el último run: {'sí' if snap.had_failure else 'no'}",
            f"Confianza estimada: {snap.confidence:.2f}",
            f"Rationale: {snap.rationale}",
        ]
        return "[Estado cognitivo]\n" + "\n".join(lines)

    def embed_state(self, snapshot: CognitiveSnapshot) -> List[float]:
        """Embedding determinista del estado (offline)."""
        text = (
            f"failure={snapshot.had_failure} tools={len(snapshot.tool_calls)} "
            f"conf={snapshot.confidence:.2f} turns={snapshot.memory_turn_count}"
        )
        return self._embedder.embed_one(text)

embed_state(snapshot: CognitiveSnapshot) -> List[float]

Embedding determinista del estado (offline).

Source code in src\ciel\runtime\cognitive_state.py
def embed_state(self, snapshot: CognitiveSnapshot) -> List[float]:
    """Embedding determinista del estado (offline)."""
    text = (
        f"failure={snapshot.had_failure} tools={len(snapshot.tool_calls)} "
        f"conf={snapshot.confidence:.2f} turns={snapshot.memory_turn_count}"
    )
    return self._embedder.embed_one(text)

latest_block(*, tenant_id: Optional[str], session_id: str) -> Optional[str]

Devuelve el bloque [Estado cognitivo] a inyectar, o None.

Source code in src\ciel\runtime\cognitive_state.py
def latest_block(self, *, tenant_id: Optional[str], session_id: str) -> Optional[str]:
    """Devuelve el bloque ``[Estado cognitivo]`` a inyectar, o None."""
    if not self.config.inject_into_prompt:
        return None
    report = self.get_recent(tenant_id=tenant_id, session_id=session_id, limit=1)
    snap = report.latest
    if snap is None:
        return None
    lines = [
        f"Versión de prompt activa: {snap.active_prompt_version or 'instructions (no versionado)'}",
        f"Turnos de memoria acumulados: {snap.memory_turn_count}",
        f"Tool calls en el último run: {len(snap.tool_calls)}",
        f"Falló en el último run: {'sí' if snap.had_failure else 'no'}",
        f"Confianza estimada: {snap.confidence:.2f}",
        f"Rationale: {snap.rationale}",
    ]
    return "[Estado cognitivo]\n" + "\n".join(lines)

record(snapshot: CognitiveSnapshot) -> None

Persiste el snapshot en cognitive_state_log (tenant-filtered).

Source code in src\ciel\runtime\cognitive_state.py
def record(self, snapshot: CognitiveSnapshot) -> None:
    """Persiste el snapshot en ``cognitive_state_log`` (tenant-filtered)."""
    try:
        self.backend.state_log_append(
            tenant_id=snapshot.tenant_id,
            session_id=snapshot.session_id,
            prompt_version=snapshot.active_prompt_version,
            value_json=_dump(snapshot.to_dict()),
            created_at=snapshot.created_at
            or datetime.now(timezone.utc).isoformat(),
        )
    except Exception:  # pragma: no cover - introspección nunca rompe el run
        pass

DeterministicEmbeddingProvider

Bases: EmbeddingProvider

Embedding determinista offline (hash → vector). NO semántico real.

Útil para dev, tests y como fallback cuando no hay API key. Produce el mismo vector para el mismo texto (determinista), permitiendo búsqueda por similitud coseno coherente dentro de una corrida. NO usar en prod como única fuente de verdad semántica; combinar con BM25 (hybrid).

Source code in src\ciel\rag\embeddings.py
class DeterministicEmbeddingProvider(EmbeddingProvider):
    """Embedding determinista offline (hash → vector). NO semántico real.

    Útil para dev, tests y como fallback cuando no hay API key. Produce el
    mismo vector para el mismo texto (determinista), permitiendo búsqueda
    por similitud coseno coherente dentro de una corrida. NO usar en prod
    como única fuente de verdad semántica; combinar con BM25 (hybrid).
    """

    def __init__(self, dim: int = 256, seed: int = 0) -> None:
        self.dim = dim
        self._seed = seed

    def embed(self, texts: Sequence[str]) -> List[List[float]]:
        out: List[List[float]] = []
        for text in texts:
            vec = [0.0] * self.dim
            # Mezcla varias fragmentaciones del hash para dispersar.
            for window in range(0, max(1, len(text)), max(1, len(text) // self.dim or 1)):
                chunk = text[window : window + 8]
                h = hashlib.sha256(f"{self._seed}:{chunk}".encode("utf-8")).digest()
                for i in range(0, len(h), 4):
                    idx = (h[i] + h[i + 1] * 256) % self.dim
                    sign = 1.0 if (h[i + 2] & 1) == 0 else -1.0
                    vec[idx] += sign * ((h[i + 3] % 100) / 100.0)
            # Normaliza a norma 1 para coseno estable.
            norm = math.sqrt(sum(v * v for v in vec)) or 1.0
            out.append([v / norm for v in vec])
        return out

EpisodicStore

Store de memoria episódica sobre un StateBackend (aislado por tenant).

Operaciones: * append — persiste un turno. * get_recent — últimos N episodios de una sesión (orden cronológico). * get_by_id — recupera un episodio por su ID. * search — búsqueda por keywords filtrada POR TENANT (no cross-tenant). * clear_session — borra la memoria de una sesión. * as_context — serializa episodios recientes como texto para el system.

Source code in src\ciel\runtime\memory_episodic.py
class EpisodicStore:
    """Store de memoria episódica sobre un ``StateBackend`` (aislado por tenant).

    Operaciones:
    * ``append`` — persiste un turno.
    * ``get_recent`` — últimos N episodios de una sesión (orden cronológico).
    * ``get_by_id`` — recupera un episodio por su ID.
    * ``search`` — búsqueda por keywords filtrada POR TENANT (no cross-tenant).
    * ``clear_session`` — borra la memoria de una sesión.
    * ``as_context`` — serializa episodios recientes como texto para el system.
    """

    def __init__(self, backend: Any, *, namespace: str = "episodic") -> None:
        # backend: ciel.runtime.state_backend.StateBackend
        self._backend = backend
        self._ns = namespace

    # -- helpers internos ----------------------------------------------------
    def _key(self, session_id: str, memory_id: str) -> str:
        return f"{self._ns}:{session_id}:{memory_id}"

    # -- escritura -----------------------------------------------------------
    def append(
        self,
        *,
        tenant_id: Optional[str],
        session_id: str,
        role: str,
        content: Any,
        metadata: Optional[dict] = None,
    ) -> EpisodicMemory:
        memory_id = str(uuid.uuid4())
        mem = EpisodicMemory(
            id=memory_id,
            tenant_id=tenant_id,
            session_id=session_id,
            role=role,
            content=content,
            metadata=metadata or {},
        )
        self._backend.memory_append(
            tenant_id=tenant_id,
            session_id=session_id,
            memory_id=memory_id,
            value=mem.to_dict(),
        )
        return mem

    # -- lectura -------------------------------------------------------------
    def get_recent(
        self, *, tenant_id: Optional[str], session_id: str, limit: int = 8
    ) -> List[EpisodicMemory]:
        rows = self._backend.memory_get_recent(
            tenant_id=tenant_id, session_id=session_id, limit=limit
        )
        out: List[EpisodicMemory] = []
        for row in rows:
            try:
                out.append(EpisodicMemory.from_dict(row))
            except (KeyError, TypeError):  # pragma: no cover - defensive
                continue
        # Orden cronológico ascendente (más antiguo primero).
        out.sort(key=lambda m: m.created_at)
        return out

    def get_by_id(
        self, *, tenant_id: Optional[str], session_id: str, memory_id: str
    ) -> Optional[EpisodicMemory]:
        row = self._backend.memory_get(
            tenant_id=tenant_id, session_id=session_id, memory_id=memory_id
        )
        if row is None:
            return None
        try:
            return EpisodicMemory.from_dict(row)
        except (KeyError, TypeError):  # pragma: no cover - defensive
            return None

    def search(
        self,
        *,
        tenant_id: Optional[str],
        session_id: Optional[str] = None,
        query: str,
        limit: int = 5,
    ) -> List[EpisodicMemory]:
        rows = self._backend.memory_search_tenant(
            tenant_id=tenant_id,
            session_id=session_id,
            query=query,
            limit=limit,
        )
        out: List[EpisodicMemory] = []
        for row in rows:
            try:
                out.append(EpisodicMemory.from_dict(row))
            except (KeyError, TypeError):  # pragma: no cover - defensive
                continue
        return out

    def clear_session(self, *, tenant_id: Optional[str], session_id: str) -> None:
        self._backend.memory_clear_session(
            tenant_id=tenant_id, session_id=session_id
        )

    # -- serialización para el system prompt --------------------------------
    @staticmethod
    def format_as_context(memories: Sequence[EpisodicMemory]) -> str:
        if not memories:
            return ""
        lines: List[str] = []
        for m in memories:
            content = m.content
            if isinstance(content, (list, dict)):  # multimodal/structured
                import json

                try:
                    content = json.dumps(content, ensure_ascii=False)
                except (TypeError, ValueError):  # pragma: no cover
                    content = str(content)
            elif not isinstance(content, str):
                content = str(content)
            lines.append(f"[{m.role}] {content}")
        return "\n".join(lines)

    def as_context(
        self,
        *,
        tenant_id: Optional[str],
        session_id: str,
        recent: int = 8,
        query: Optional[str] = None,
        search_limit: int = 5,
    ) -> str:
        """Devuelve el texto de memoria para inyectar en el system prompt."""
        memories: List[EpisodicMemory] = list(
            self.get_recent(
                tenant_id=tenant_id, session_id=session_id, limit=recent
            )
        )
        if query:
            memories.extend(
                self.search(
                    tenant_id=tenant_id,
                    session_id=session_id,
                    query=query,
                    limit=search_limit,
                )
            )
        # Dedup por id preservando orden.
        seen = set()
        unique: List[EpisodicMemory] = []
        for m in memories:
            if m.id in seen:
                continue
            seen.add(m.id)
            unique.append(m)
        return self.format_as_context(unique)

as_context(*, tenant_id: Optional[str], session_id: str, recent: int = 8, query: Optional[str] = None, search_limit: int = 5) -> str

Devuelve el texto de memoria para inyectar en el system prompt.

Source code in src\ciel\runtime\memory_episodic.py
def as_context(
    self,
    *,
    tenant_id: Optional[str],
    session_id: str,
    recent: int = 8,
    query: Optional[str] = None,
    search_limit: int = 5,
) -> str:
    """Devuelve el texto de memoria para inyectar en el system prompt."""
    memories: List[EpisodicMemory] = list(
        self.get_recent(
            tenant_id=tenant_id, session_id=session_id, limit=recent
        )
    )
    if query:
        memories.extend(
            self.search(
                tenant_id=tenant_id,
                session_id=session_id,
                query=query,
                limit=search_limit,
            )
        )
    # Dedup por id preservando orden.
    seen = set()
    unique: List[EpisodicMemory] = []
    for m in memories:
        if m.id in seen:
            continue
        seen.add(m.id)
        unique.append(m)
    return self.format_as_context(unique)

IntrospectionConfig dataclass

Source code in src\ciel\runtime\cognitive_state.py
@dataclass
class IntrospectionConfig:
    enabled: bool = True
    inject_into_prompt: bool = True
    embedding_dim: int = 64

    @classmethod
    def disabled(cls) -> "IntrospectionConfig":
        return cls(enabled=False)

IntrospectionReport dataclass

Agrega snapshots de una sesión en un reporte introspectivo.

Source code in src\ciel\runtime\cognitive_state.py
@dataclass
class IntrospectionReport:
    """Agrega snapshots de una sesión en un reporte introspectivo."""

    tenant_id: Optional[str]
    session_id: str
    snapshots: List[CognitiveSnapshot] = field(default_factory=list)

    @property
    def latest(self) -> Optional[CognitiveSnapshot]:
        return self.snapshots[0] if self.snapshots else None

    def to_dict(self) -> Dict[str, Any]:
        return {
            "tenant_id": self.tenant_id,
            "session_id": self.session_id,
            "snapshot_count": len(self.snapshots),
            "snapshots": [s.to_dict() for s in self.snapshots],
        }

StateBackend

Bases: ABC

Interfaz mínima de persistencia compartida (multi-réplica).

Source code in src\ciel\runtime\state_backend.py
class StateBackend(abc.ABC):
    """Interfaz mínima de persistencia compartida (multi-réplica)."""

    backend_type: str = "abstract"

    @abc.abstractmethod
    def set(
        self,
        *,
        tenant_id: Optional[str],
        session_id: str,
        key: str,
        value: Any,
    ) -> None:
        ...

    @abc.abstractmethod
    def get(self, *, tenant_id: Optional[str], session_id: str, key: str) -> Optional[Any]:
        ...

    @abc.abstractmethod
    def delete(self, *, tenant_id: Optional[str], session_id: str, key: str) -> None:
        ...

    @abc.abstractmethod
    def search(self, query: str, *, limit: int = 10) -> List[Dict[str, Any]]:
        ...

    @abc.abstractmethod
    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:
        ...

    # --- memoria episódica (Fase 17, tenant-filtered) ----------------------
    # TODAS las lecturas filtran por tenant_id explícitamente para evitar
    # fuga cross-tenant (el ``search`` genérico NO filtra; no usarlo para memoria).
    @abc.abstractmethod
    def memory_append(
        self,
        *,
        tenant_id: Optional[str],
        session_id: str,
        memory_id: str,
        value: Any,
    ) -> None:
        ...

    @abc.abstractmethod
    def memory_get(
        self,
        *,
        tenant_id: Optional[str],
        session_id: str,
        memory_id: str,
    ) -> Optional[dict]:
        ...

    @abc.abstractmethod
    def memory_get_recent(
        self,
        *,
        tenant_id: Optional[str],
        session_id: str,
        limit: int = 8,
    ) -> List[dict]:
        ...

    @abc.abstractmethod
    def memory_search_tenant(
        self,
        *,
        tenant_id: Optional[str],
        session_id: Optional[str],
        query: str,
        limit: int = 5,
    ) -> List[dict]:
        ...

    @abc.abstractmethod
    def memory_clear_session(
        self, *, tenant_id: Optional[str], session_id: str
    ) -> None:
        ...

    # --- prompts versionados (Fase 19, tenant-filtered) -------------------
    @abc.abstractmethod
    def prompt_save(
        self,
        *,
        tenant_id: Optional[str],
        name: str,
        version: str,
        prompt_text: str,
        value_json: str,
        sha256: str,
        previous_version: Optional[str],
        created_at: str,
    ) -> None:
        ...

    @abc.abstractmethod
    def prompt_get(
        self, *, tenant_id: Optional[str], name: str, version: Optional[str] = None
    ) -> Optional[dict]:
        ...

    @abc.abstractmethod
    def prompt_get_history(
        self, *, tenant_id: Optional[str], name: str
    ) -> List[dict]:
        ...

    # --- log de estado cognitivo (Fase 19, tenant-filtered) ---------------
    @abc.abstractmethod
    def state_log_append(
        self,
        *,
        tenant_id: Optional[str],
        session_id: str,
        prompt_version: Optional[str],
        value_json: str,
        created_at: str,
    ) -> None:
        ...

    @abc.abstractmethod
    def state_log_get_recent(
        self, *, tenant_id: Optional[str], session_id: str, limit: int = 16
    ) -> List[dict]:
        ...

    # --- curricula versionados (Fase 20 BLOQUE A, tenant-filtered) ---------
    # NO abstractos para no romper backends de terceros: default NotImplementedError.
    def curriculum_save(
        self,
        *,
        tenant_id: Optional[str],
        goal: str,
        version: str,
        plan_json: str,
        value_json: str,
        sha256: str,
        previous_version: Optional[str],
        created_at: str,
    ) -> None:
        raise NotImplementedError(
            f"{type(self).__name__} no implementa curriculum_save (Fase 20)"
        )

    def curriculum_get(
        self, *, tenant_id: Optional[str], goal: str, version: Optional[str] = None
    ) -> Optional[dict]:
        raise NotImplementedError(
            f"{type(self).__name__} no implementa curriculum_get (Fase 20)"
        )

    def curriculum_get_history(
        self, *, tenant_id: Optional[str], goal: str
    ) -> List[dict]:
        raise NotImplementedError(
            f"{type(self).__name__} no implementa curriculum_get_history (Fase 20)"
        )

    @abc.abstractmethod
    def close(self) -> None:
        ...

    # --- readiness (usado por /readyz en F16) ------------------------------
    def is_ready(self) -> bool:
        """Devuelve True si el backend está conectado y migrado.

        El default asume listo; los backends remotos sobrescriben esto con una
        comprobación real de conectividad.
        """
        return True

    # --- utilidades compartidas --------------------------------------------
    @staticmethod
    def _sentinel(tenant_id: Optional[str]) -> Any:
        """SQLite no trata dos NULL como iguales para UNIQUE; normaliza None."""
        return tenant_id if tenant_id is not None else "__none__"

    @staticmethod
    def _dump(value: Any) -> str:
        try:
            return json.dumps(value)
        except TypeError:
            return json.dumps({"repr": repr(value)})

    @staticmethod
    def _normalize_row(value_json: Optional[str]) -> Optional[Any]:
        if value_json is None:
            return None
        try:
            return json.loads(value_json)
        except (TypeError, json.JSONDecodeError):
            return None

is_ready() -> bool

Devuelve True si el backend está conectado y migrado.

El default asume listo; los backends remotos sobrescriben esto con una comprobación real de conectividad.

Source code in src\ciel\runtime\state_backend.py
def is_ready(self) -> bool:
    """Devuelve True si el backend está conectado y migrado.

    El default asume listo; los backends remotos sobrescriben esto con una
    comprobación real de conectividad.
    """
    return True

SqliteStateBackend_memory() -> StateBackend

Source code in src\ciel\runtime\cognitive_state.py
def SqliteStateBackend_memory() -> StateBackend:
    from ciel.runtime.state_backend import SqliteStateBackend

    return SqliteStateBackend(":memory:")

_default_backend() -> StateBackend

Source code in src\ciel\runtime\cognitive_state.py
def _default_backend() -> StateBackend:
    global _DEFAULT_COGNITIVE_BACKEND
    if _DEFAULT_COGNITIVE_BACKEND is None:
        _DEFAULT_COGNITIVE_BACKEND = SqliteStateBackend_memory()
    return _DEFAULT_COGNITIVE_BACKEND

_dump(value: Any) -> str

Source code in src\ciel\runtime\cognitive_state.py
def _dump(value: Any) -> str:
    import json

    try:
        return json.dumps(value, ensure_ascii=False)
    except TypeError:  # pragma: no cover - defensivo
        return json.dumps({"repr": repr(value)})

_effective_tenant(agent: Any, tenant_id: Optional[str] = None) -> Optional[str]

Source code in src\ciel\runtime\cognitive_state.py
def _effective_tenant(agent: Any, tenant_id: Optional[str] = None) -> Optional[str]:
    if tenant_id:
        return tenant_id
    mem = getattr(agent, "_memory", None)
    if mem is not None and mem.enabled:
        return getattr(agent, "tenant_id", None) or "default"
    return getattr(agent, "tenant_id", None) or "default"

install_cognitive_state_support(agent_cls: Any) -> Any

Engancha introspection= / introspection_config= sin reescribir api.py.

Idempotente.

Source code in src\ciel\runtime\cognitive_state.py
def install_cognitive_state_support(agent_cls: Any) -> Any:
    """Engancha ``introspection=`` / ``introspection_config=`` sin reescribir api.py.

    Idempotente.
    """
    if getattr(agent_cls, "_cognitive_installed", False):
        return agent_cls
    original_init = agent_cls.__init__
    original_build = agent_cls._build_request
    original_arun = agent_cls.arun

    @functools.wraps(original_init)
    def _patched_init(self: Any, *args: Any, **kwargs: Any) -> None:
        intro_arg = kwargs.pop("introspection", None)
        intro_config = kwargs.pop("introspection_config", None)
        original_init(self, *args, **kwargs)

        config = intro_config or IntrospectionConfig()
        if intro_arg is False or (intro_arg is None and intro_config is None):
            config = config.disabled()

        # Resolver backend + store: reutiliza el EpisodicStore de memoria si existe
        # (así el log cognitivo vive en el mismo backend que la memoria del agente).
        store = None
        backend: StateBackend = _default_backend()
        session_id: Optional[str] = None
        mem = getattr(self, "_memory", None)
        if mem is not None and getattr(mem, "store", None) is not None:
            store = mem.store
            be = getattr(store, "_backend", None)
            if be is not None:
                backend = be
            session_id = mem.session_id() if mem.enabled else None
        if store is None:
            store = EpisodicStore(backend)
        if session_id is None:
            session_id = str(uuid.uuid4())
        self._cognitive = CognitiveState(
            store=store, backend=backend, config=config, session_id=session_id
        )

    @functools.wraps(original_build)
    def _patched_build(self: Any, prompt, *, tenant_id=None, session_id=None, **kw) -> Any:
        request = original_build(
            self, prompt, tenant_id=tenant_id, session_id=session_id, **kw
        )
        cog = getattr(self, "_cognitive", None)
        if cog is not None and cog.enabled:
            block = cog.latest_block(
                tenant_id=_effective_tenant(self, tenant_id), session_id=cog.session_id
            )
            if block:
                # ChatRequest es frozen -> reconstruimos, no reasignamos.
                new_messages = list(request.messages) + [
                    ChatMessage(role="system", content=block)
                ]
                request = ChatRequest(
                    messages=tuple(new_messages),
                    tools=request.tools,
                    model=request.model,
                    temperature=request.temperature,
                    max_tokens=request.max_tokens,
                    extra=request.extra,
                )
        return request

    @functools.wraps(original_arun)
    async def _patched_arun(
        self: Any, prompt, *, tenant_id=None, max_turns=10, limit=32
    ) -> Any:
        eff_tenant = _effective_tenant(self, tenant_id)
        resp = await original_arun(
            self, prompt, tenant_id=tenant_id, max_turns=max_turns, limit=limit
        )
        cog = getattr(self, "_cognitive", None)
        if cog is not None and cog.enabled:
            cog.tenant_id = eff_tenant
            try:
                snap = cog.build_snapshot(
                    tenant_id=eff_tenant,
                    session_id=cog.session_id,
                    response=resp,
                )
                cog.record(snap)
            except Exception:  # pragma: no cover - introspección nunca rompe el run
                pass
        return resp

    agent_cls.__init__ = _patched_init  # type: ignore[assignment]
    agent_cls._build_request = _patched_build  # type: ignore[assignment]
    agent_cls.arun = _patched_arun  # type: ignore[assignment]
    agent_cls._cognitive_installed = True  # type: ignore[attr-defined]

    # Helper de consulta expuesto en la instancia.
    def _introspect(self: Any, *, limit: int = 16) -> IntrospectionReport:
        cog = getattr(self, "_cognitive", None)
        if cog is None:
            return IntrospectionReport(tenant_id=None, session_id="", snapshots=[])
        return cog.get_recent(
            tenant_id=cog.tenant_id, session_id=cog.session_id, limit=limit
        )

    agent_cls.introspect = _introspect  # type: ignore[attr-defined]
    return agent_cls

Fase 20 BLOQUE A — Curriculum / autonomous goal setting (v0.13.x).

Un curriculum es la descomposición versionada de un goal en un plan de sub-tareas ordenadas, persistida en el StateBackend (tabla curricula) con aislamiento estricto por tenant_id y linaje (previous_version / parent), calcado del molde de prompt_versioning.py (Fase 19).

Piezas:

  • :class:Curriculum — dataclass con goal + plan + linaje + sha256.
  • :class:CurriculumRegistry — CRUD versionado sobre un StateBackend (curriculum_save / curriculum_get / curriculum_get_history).
  • :func:decompose_goal — descomposición determinista/offline por defecto (heurística de cláusulas); si se pasa un provider real (p. ej. MockProvider o un LLM), se le pide el plan y se parsea su respuesta.
  • :func:install_curriculum_support — integración ADITIVA con ciel.Agent (kwargs curriculum= / curriculum_config= + Agent.run_curriculum), idempotente, mismo patrón que install_cognitive_state_support.

Todo es network-free y API-key-free por defecto (offline-safe).

INITIAL_VERSION = '0.0.0' module-attribute

_SPLIT_PATTERN = re.compile('(?:;|\\.\\s+|\\.$|,\\s*luego\\s+|\\s+y\\s+luego\\s+|\\s+luego\\s+|\\s+y\\s+después\\s+|\\s+después\\s+|\\s+y\\s+|\\s+then\\s+|\\s+and\\s+then\\s+)', re.IGNORECASE) module-attribute

__all__ = ['CurriculumError', 'Curriculum', 'CurriculumConfig', 'CurriculumRegistry', 'AgentCurriculum', 'decompose_goal', 'install_curriculum_support', 'sha256_plan', 'INITIAL_VERSION'] module-attribute

AgentCurriculum

Estado de curriculum enganchado a una instancia de Agent.

Source code in src\ciel\runtime\curriculum.py
class AgentCurriculum:
    """Estado de curriculum enganchado a una instancia de Agent."""

    def __init__(self, *, config: CurriculumConfig, backend: Any) -> None:
        self.config = config
        self.backend = backend
        self.registry = CurriculumRegistry(backend)

    @property
    def enabled(self) -> bool:
        return bool(self.config.enabled)

Curriculum dataclass

Plan versionado (curriculum) para un objetivo del agente.

Source code in src\ciel\runtime\curriculum.py
@dataclass
class Curriculum:
    """Plan versionado (curriculum) para un objetivo del agente."""

    goal: str = ""
    plan: List[str] = field(default_factory=list)
    tenant_id: Optional[str] = None
    version: str = INITIAL_VERSION
    created_at: Optional[str] = None
    sha256: str = ""
    previous_version: Optional[str] = None
    parent: Optional[str] = None
    metadata: Dict[str, Any] = field(default_factory=dict)

    def to_dict(self) -> Dict[str, Any]:
        return {
            "goal": self.goal,
            "plan": list(self.plan),
            "tenant_id": self.tenant_id,
            "version": self.version,
            "created_at": self.created_at,
            "sha256": self.sha256,
            "previous_version": self.previous_version,
            "parent": self.parent,
            "metadata": self.metadata,
        }

    @classmethod
    def from_dict(cls, data: Dict[str, Any]) -> "Curriculum":
        return cls(
            goal=data.get("goal", ""),
            plan=list(data.get("plan") or []),
            tenant_id=data.get("tenant_id"),
            version=data.get("version") or INITIAL_VERSION,
            created_at=data.get("created_at"),
            sha256=data.get("sha256", ""),
            previous_version=data.get("previous_version"),
            parent=data.get("parent"),
            metadata=data.get("metadata") or {},
        )

    def bump_version(self, kind: str = "patch") -> str:
        parts = (self.version.split(".") + ["0", "0", "0"])[:3]
        try:
            major, minor, patch = (int(p) for p in parts)
        except ValueError as exc:
            raise CurriculumError(f"invalid version string: {self.version!r}") from exc
        kind = (kind or "patch").lower()
        if kind == "major":
            return f"{major + 1}.0.0"
        if kind == "minor":
            return f"{major}.{minor + 1}.0"
        if kind == "patch":
            return f"{major}.{minor}.{patch + 1}"
        raise CurriculumError(f"unknown bump kind: {kind!r} (expected major/minor/patch)")

CurriculumConfig dataclass

Configuración del soporte de curriculum en el Agent.

Source code in src\ciel\runtime\curriculum.py
@dataclass
class CurriculumConfig:
    """Configuración del soporte de curriculum en el Agent."""

    enabled: bool = True
    backend: Any = None  # StateBackend; default SQLite :memory:
    max_attempts: int = 5

    def disabled(self) -> "CurriculumConfig":
        return CurriculumConfig(enabled=False, backend=self.backend, max_attempts=self.max_attempts)

CurriculumError

Bases: Exception

Error de curriculum/goal setting.

Source code in src\ciel\runtime\curriculum.py
class CurriculumError(Exception):
    """Error de curriculum/goal setting."""

CurriculumRegistry

Registro multitenant de curricula versionados sobre un StateBackend.

Source code in src\ciel\runtime\curriculum.py
class CurriculumRegistry:
    """Registro multitenant de curricula versionados sobre un ``StateBackend``."""

    def __init__(self, backend: Any) -> None:
        self._backend = backend

    # --- escritura ----------------------------------------------------------
    def create(
        self,
        goal: str,
        plan: Sequence[str],
        *,
        tenant_id: Optional[str] = None,
        metadata: Optional[Dict[str, Any]] = None,
    ) -> Curriculum:
        """Crea la versión inicial ``0.0.0`` de un curriculum (o bumpea si existe)."""
        current = self.get(goal, tenant_id=tenant_id)
        if current is None:
            version = INITIAL_VERSION
            previous = None
        else:
            version = current.bump_version("patch")
            previous = current.version
        cur = Curriculum(
            goal=goal,
            plan=list(plan),
            tenant_id=tenant_id,
            version=version,
            created_at=_now_iso(),
            sha256=sha256_plan(goal, plan),
            previous_version=previous,
            parent=previous,
            metadata=metadata or {},
        )
        self._backend.curriculum_save(
            tenant_id=tenant_id,
            goal=goal,
            version=cur.version,
            plan_json=json.dumps(list(plan), ensure_ascii=False),
            value_json=json.dumps(cur.to_dict(), ensure_ascii=False),
            sha256=cur.sha256,
            previous_version=cur.previous_version,
            created_at=cur.created_at,
        )
        return cur

    # --- lectura ------------------------------------------------------------
    def get(
        self,
        goal: str,
        *,
        tenant_id: Optional[str] = None,
        version: Optional[str] = None,
    ) -> Optional[Curriculum]:
        row = self._backend.curriculum_get(tenant_id=tenant_id, goal=goal, version=version)
        if row is None:
            return None
        return Curriculum.from_dict(row)

    def history(self, goal: str, *, tenant_id: Optional[str] = None) -> List[Curriculum]:
        rows = self._backend.curriculum_get_history(tenant_id=tenant_id, goal=goal)
        return [Curriculum.from_dict(r) for r in rows]

    def evolution_tree(self, goal: str, *, tenant_id: Optional[str] = None) -> Dict[str, Any]:
        """Árbol de linaje, misma forma que ``prompt_versioning.evolution_tree``."""
        versions = self.history(goal, tenant_id=tenant_id)
        if not versions:
            raise CurriculumError(f"unknown curriculum: {goal!r} (tenant {tenant_id!r})")
        keys = [c.version for c in versions]
        nodes: Dict[str, Dict[str, Any]] = {}
        for idx, cur in enumerate(versions):
            raw_prev = cur.previous_version
            parent = raw_prev if raw_prev is not None else (keys[idx - 1] if idx > 0 else None)
            nodes[cur.version] = {
                "version": cur.version,
                "parent": parent,
                "previous_version": raw_prev,
                "children": [],
                "sha256": cur.sha256,
                "plan": list(cur.plan),
                "created_at": cur.created_at,
            }
        for key, node in nodes.items():
            if node["parent"] is not None and node["parent"] in nodes:
                nodes[node["parent"]]["children"].append(key)
        roots = [k for k, n in nodes.items() if n["parent"] is None]
        root = roots[0] if roots else (keys[0] if keys else None)
        return {"goal": goal, "root": root, "lineage": list(keys), "nodes": nodes}

create(goal: str, plan: Sequence[str], *, tenant_id: Optional[str] = None, metadata: Optional[Dict[str, Any]] = None) -> Curriculum

Crea la versión inicial 0.0.0 de un curriculum (o bumpea si existe).

Source code in src\ciel\runtime\curriculum.py
def create(
    self,
    goal: str,
    plan: Sequence[str],
    *,
    tenant_id: Optional[str] = None,
    metadata: Optional[Dict[str, Any]] = None,
) -> Curriculum:
    """Crea la versión inicial ``0.0.0`` de un curriculum (o bumpea si existe)."""
    current = self.get(goal, tenant_id=tenant_id)
    if current is None:
        version = INITIAL_VERSION
        previous = None
    else:
        version = current.bump_version("patch")
        previous = current.version
    cur = Curriculum(
        goal=goal,
        plan=list(plan),
        tenant_id=tenant_id,
        version=version,
        created_at=_now_iso(),
        sha256=sha256_plan(goal, plan),
        previous_version=previous,
        parent=previous,
        metadata=metadata or {},
    )
    self._backend.curriculum_save(
        tenant_id=tenant_id,
        goal=goal,
        version=cur.version,
        plan_json=json.dumps(list(plan), ensure_ascii=False),
        value_json=json.dumps(cur.to_dict(), ensure_ascii=False),
        sha256=cur.sha256,
        previous_version=cur.previous_version,
        created_at=cur.created_at,
    )
    return cur

evolution_tree(goal: str, *, tenant_id: Optional[str] = None) -> Dict[str, Any]

Árbol de linaje, misma forma que prompt_versioning.evolution_tree.

Source code in src\ciel\runtime\curriculum.py
def evolution_tree(self, goal: str, *, tenant_id: Optional[str] = None) -> Dict[str, Any]:
    """Árbol de linaje, misma forma que ``prompt_versioning.evolution_tree``."""
    versions = self.history(goal, tenant_id=tenant_id)
    if not versions:
        raise CurriculumError(f"unknown curriculum: {goal!r} (tenant {tenant_id!r})")
    keys = [c.version for c in versions]
    nodes: Dict[str, Dict[str, Any]] = {}
    for idx, cur in enumerate(versions):
        raw_prev = cur.previous_version
        parent = raw_prev if raw_prev is not None else (keys[idx - 1] if idx > 0 else None)
        nodes[cur.version] = {
            "version": cur.version,
            "parent": parent,
            "previous_version": raw_prev,
            "children": [],
            "sha256": cur.sha256,
            "plan": list(cur.plan),
            "created_at": cur.created_at,
        }
    for key, node in nodes.items():
        if node["parent"] is not None and node["parent"] in nodes:
            nodes[node["parent"]]["children"].append(key)
    roots = [k for k, n in nodes.items() if n["parent"] is None]
    root = roots[0] if roots else (keys[0] if keys else None)
    return {"goal": goal, "root": root, "lineage": list(keys), "nodes": nodes}

_default_backend() -> Any

Source code in src\ciel\runtime\curriculum.py
def _default_backend() -> Any:
    from ciel.runtime.state_backend import SqliteStateBackend

    return SqliteStateBackend(":memory:")

_heuristic_plan(goal: str) -> List[str]

Split determinista por frases/cláusulas ('; ', '.', ' y ', ' luego ', ...).

Source code in src\ciel\runtime\curriculum.py
def _heuristic_plan(goal: str) -> List[str]:
    """Split determinista por frases/cláusulas ('; ', '.', ' y ', ' luego ', ...)."""
    raw = _SPLIT_PATTERN.split(goal or "")
    steps = [s.strip(" .,;\t\n") for s in raw]
    steps = [s for s in steps if s]
    return steps or ([goal.strip()] if goal and goal.strip() else [])

_now_iso() -> str

Source code in src\ciel\runtime\curriculum.py
def _now_iso() -> str:
    return datetime.now(timezone.utc).isoformat()

_parse_provider_plan(text: str) -> List[str]

Parsea la respuesta de un provider a lista de pasos.

Acepta JSON (lista de strings) o texto con líneas/numeración/viñetas.

Source code in src\ciel\runtime\curriculum.py
def _parse_provider_plan(text: str) -> List[str]:
    """Parsea la respuesta de un provider a lista de pasos.

    Acepta JSON (lista de strings) o texto con líneas/numeración/viñetas.
    """
    text = (text or "").strip()
    if not text:
        return []
    try:
        data = json.loads(text)
        if isinstance(data, list):
            return [str(s).strip() for s in data if str(s).strip()]
    except (ValueError, TypeError):
        pass
    steps: List[str] = []
    for line in text.splitlines():
        cleaned = re.sub(r"^\s*(?:[-*•]|\d+[.)])\s*", "", line).strip()
        if cleaned:
            steps.append(cleaned)
    return steps

decompose_goal(goal: str, *, provider: Any = None) -> List[str]

Descompone goal en sub-tareas ordenadas.

Sin provider: heurística determinista offline (split por cláusulas). Con provider (ChatProvider, p. ej. MockProvider): se le pide el plan y se parsea; si la respuesta es vacía/inutilizable, fallback a la heurística offline.

Source code in src\ciel\runtime\curriculum.py
def decompose_goal(goal: str, *, provider: Any = None) -> List[str]:
    """Descompone ``goal`` en sub-tareas ordenadas.

    Sin ``provider``: heurística determinista offline (split por cláusulas).
    Con ``provider`` (ChatProvider, p. ej. ``MockProvider``): se le pide el
    plan y se parsea; si la respuesta es vacía/inutilizable, fallback a la
    heurística offline.
    """
    if not goal or not goal.strip():
        raise CurriculumError("goal vacío: nada que descomponer")
    if provider is None:
        return _heuristic_plan(goal)
    try:
        from ciel.runtime.tools import ChatMessage, ChatRequest

        request = ChatRequest(
            messages=(
                ChatMessage(
                    role="user",
                    content=(
                        "Descompón el siguiente objetivo en una lista JSON de "
                        f"sub-tareas ordenadas: {goal}"
                    ),
                ),
            )
        )
        try:
            loop = asyncio.get_running_loop()
        except RuntimeError:
            loop = None
        if loop is not None:  # pragma: no cover - llamada desde loop activo
            raise CurriculumError(
                "decompose_goal(provider=...) no puede llamarse desde un event "
                "loop activo; usa la heurística offline o llama fuera del loop"
            )
        resp = asyncio.run(provider.complete(request))
        text = getattr(getattr(resp, "choice", None), "message", None)
        text = getattr(text, "content", "") if text is not None else ""
        if isinstance(text, list):  # contenido multimodal
            text = "".join(p.get("text", "") for p in text if isinstance(p, dict))
        steps = _parse_provider_plan(text or "")
        if steps:
            return steps
    except CurriculumError:
        raise
    except Exception:
        pass  # provider roto => fallback offline
    return _heuristic_plan(goal)

install_curriculum_support(agent_cls: Any) -> Any

Engancha curriculum= / curriculum_config= sin reescribir api.py.

Idempotente (flag _curriculum_installed). Expone Agent.run_curriculum(goal, *, handler=None, plan=None, tenant_id=None) que descompone el goal (o usa plan), lo persiste vía :class:CurriculumRegistry y ejecuta las sub-tareas con AutonomousAgent.run_goal (EventLoop real de orquestación). Si no se pasa handler, cada sub-tarea se resuelve con Agent.arun.

Source code in src\ciel\runtime\curriculum.py
def install_curriculum_support(agent_cls: Any) -> Any:
    """Engancha ``curriculum=`` / ``curriculum_config=`` sin reescribir api.py.

    Idempotente (flag ``_curriculum_installed``). Expone
    ``Agent.run_curriculum(goal, *, handler=None, plan=None, tenant_id=None)``
    que descompone el goal (o usa ``plan``), lo persiste vía
    :class:`CurriculumRegistry` y ejecuta las sub-tareas con
    ``AutonomousAgent.run_goal`` (EventLoop real de orquestación). Si no se
    pasa ``handler``, cada sub-tarea se resuelve con ``Agent.arun``.
    """
    if getattr(agent_cls, "_curriculum_installed", False):
        return agent_cls
    original_init = agent_cls.__init__

    @functools.wraps(original_init)
    def _patched_init(self: Any, *args: Any, **kwargs: Any) -> None:
        cur_arg = kwargs.pop("curriculum", None)
        cur_config = kwargs.pop("curriculum_config", None)
        original_init(self, *args, **kwargs)

        config = cur_config or CurriculumConfig()
        if cur_arg is False or (cur_arg is None and cur_config is None):
            config = config.disabled()
        backend = config.backend
        if backend is None:
            backend = _default_backend()
        self._curriculum = AgentCurriculum(config=config, backend=backend)

    async def _run_curriculum(
        self: Any,
        goal: str,
        *,
        handler: Any = None,
        plan: Optional[Sequence[str]] = None,
        tenant_id: Optional[str] = None,
        provider: Any = None,
    ) -> Dict[str, Any]:
        """Descompone, persiste y ejecuta un curriculum. Devuelve dict resumen."""
        cur_state: Optional[AgentCurriculum] = getattr(self, "_curriculum", None)
        if cur_state is None:
            cur_state = AgentCurriculum(
                config=CurriculumConfig(), backend=_default_backend()
            )
            self._curriculum = cur_state
        eff_tenant = tenant_id or getattr(self, "tenant_id", None)
        steps = list(plan) if plan else decompose_goal(goal, provider=provider)
        curriculum = cur_state.registry.create(goal, steps, tenant_id=eff_tenant)

        if handler is None:
            async def handler(task: Any) -> Any:  # type: ignore[misc]
                resp = await self.arun(task.goal, tenant_id=eff_tenant)
                return getattr(resp, "text", resp)

        from ciel.orchestration.agent import AutonomousAgent

        runner = AutonomousAgent(
            name="curriculum",
            tenant_id=eff_tenant,
            max_attempts=cur_state.config.max_attempts,
        )
        tasks = await runner.run_goal(goal, handler, plan=steps)
        return {
            "goal": goal,
            "curriculum": curriculum.to_dict(),
            "tasks": [t.snapshot() for t in tasks],
            "succeeded": all(t.status == "succeeded" for t in tasks),
        }

    agent_cls.__init__ = _patched_init  # type: ignore[assignment]
    agent_cls.run_curriculum = _run_curriculum  # type: ignore[attr-defined]
    agent_cls._curriculum_installed = True  # type: ignore[attr-defined]
    return agent_cls

sha256_plan(goal: str, plan: Sequence[str]) -> str

Hash determinista de (goal, plan) — offline.

Source code in src\ciel\runtime\curriculum.py
def sha256_plan(goal: str, plan: Sequence[str]) -> str:
    """Hash determinista de (goal, plan) — offline."""
    payload = json.dumps({"goal": goal, "plan": list(plan)}, ensure_ascii=False, sort_keys=True)
    return hashlib.sha256(payload.encode("utf-8")).hexdigest()

Knowledge graph persistente y offline-safe (Fase 20 — Bloque B).

Grafo de conocimiento mínimo, ADITIVO (no toca la API existente):

  • :class:KGNode / :class:KGEdge — dataclasses propias (misma FORMA que GraphNode/GraphEdge de ciel.orchestration.graph pero sin importar de ahí para no colisionar con el grafo de orquestación).
  • :class:KnowledgeGraph — nodos + aristas con persistencia SQLite (tablas kg_nodes/kg_edges con índices, patrón F19 de state_backend) y búsqueda semántica offline vía DeterministicEmbeddingProvider + InMemoryVectorStore (Fase 17).

Decisiones de diseño: * Offline-safe primero: sin red, sin networkx obligatorio (import opcional en :meth:KnowledgeGraph.to_networkx). * Aislamiento estricto por tenant_id en TODAS las lecturas, reutilizando el sentinel de StateBackend._sentinel (SQLite no trata dos NULL como iguales en UNIQUE). * La persistencia acepta un StateBackend SQLite existente (reusa su conn), una ruta de fichero, o None (SQLite in-memory).

__all__ = ['KGNode', 'KGEdge', 'KnowledgeGraph', 'init_kg_schema'] module-attribute

DeterministicEmbeddingProvider

Bases: EmbeddingProvider

Embedding determinista offline (hash → vector). NO semántico real.

Útil para dev, tests y como fallback cuando no hay API key. Produce el mismo vector para el mismo texto (determinista), permitiendo búsqueda por similitud coseno coherente dentro de una corrida. NO usar en prod como única fuente de verdad semántica; combinar con BM25 (hybrid).

Source code in src\ciel\rag\embeddings.py
class DeterministicEmbeddingProvider(EmbeddingProvider):
    """Embedding determinista offline (hash → vector). NO semántico real.

    Útil para dev, tests y como fallback cuando no hay API key. Produce el
    mismo vector para el mismo texto (determinista), permitiendo búsqueda
    por similitud coseno coherente dentro de una corrida. NO usar en prod
    como única fuente de verdad semántica; combinar con BM25 (hybrid).
    """

    def __init__(self, dim: int = 256, seed: int = 0) -> None:
        self.dim = dim
        self._seed = seed

    def embed(self, texts: Sequence[str]) -> List[List[float]]:
        out: List[List[float]] = []
        for text in texts:
            vec = [0.0] * self.dim
            # Mezcla varias fragmentaciones del hash para dispersar.
            for window in range(0, max(1, len(text)), max(1, len(text) // self.dim or 1)):
                chunk = text[window : window + 8]
                h = hashlib.sha256(f"{self._seed}:{chunk}".encode("utf-8")).digest()
                for i in range(0, len(h), 4):
                    idx = (h[i] + h[i + 1] * 256) % self.dim
                    sign = 1.0 if (h[i + 2] & 1) == 0 else -1.0
                    vec[idx] += sign * ((h[i + 3] % 100) / 100.0)
            # Normaliza a norma 1 para coseno estable.
            norm = math.sqrt(sum(v * v for v in vec)) or 1.0
            out.append([v / norm for v in vec])
        return out

EmbeddingProvider

Bases: ABC

Contrato mínimo de embeddings.

Source code in src\ciel\rag\embeddings.py
class EmbeddingProvider(ABC):
    """Contrato mínimo de embeddings."""

    dim: int = 256

    @abstractmethod
    def embed(self, texts: Sequence[str]) -> List[List[float]]: ...

    def embed_one(self, text: str) -> List[float]:
        return self.embed([text])[0]

InMemoryVectorStore

Bases: VectorStore

Índice vectorial en memoria (offline, sin dependencias).

Aísla por tenant_id: query solo devuelve records del tenant pedido.

Source code in src\ciel\rag\vector_store.py
class InMemoryVectorStore(VectorStore):
    """Índice vectorial en memoria (offline, sin dependencias).

    Aísla por ``tenant_id``: ``query`` solo devuelve records del tenant pedido.
    """

    def __init__(self) -> None:
        self._records: List[VectorRecord] = []

    def upsert(
        self,
        *,
        tenant_id: Optional[str],
        texts: Sequence[str],
        vectors: Sequence[Sequence[float]],
        payloads: Optional[Sequence[Dict[str, Any]]] = None,
        ids: Optional[Sequence[str]] = None,
    ) -> List[str]:
        out: List[str] = []
        for i, text in enumerate(texts):
            rid = ids[i] if ids is not None else str(uuid.uuid4())
            payload = payloads[i] if payloads is not None else {}
            # Reemplazo si ya existe el id (idempotente).
            self._records = [r for r in self._records if r.id != rid]
            self._records.append(
                VectorRecord(
                    id=rid,
                    tenant_id=tenant_id,
                    vector=list(vectors[i]),
                    payload=payload,
                    text=text,
                )
            )
            out.append(rid)
        return out

    def query(
        self,
        *,
        tenant_id: Optional[str],
        vector: Sequence[float],
        top_k: int = 5,
    ) -> List[VectorRecord]:
        scored = [
            (cosine_similarity(vector, r.vector), r)
            for r in self._records
            if r.tenant_id == tenant_id
        ]
        scored.sort(key=lambda x: x[0], reverse=True)
        return [r for _, r in scored[:top_k]]

KGEdge dataclass

Arista dirigida entre dos nodos del knowledge graph.

Source code in src\ciel\runtime\knowledge_graph.py
@dataclass
class KGEdge:
    """Arista dirigida entre dos nodos del knowledge graph."""

    from_id: str
    to_id: str
    relation: str
    weight: float = 1.0

KGNode dataclass

Nodo del knowledge graph (aislado por tenant).

Source code in src\ciel\runtime\knowledge_graph.py
@dataclass
class KGNode:
    """Nodo del knowledge graph (aislado por tenant)."""

    id: str
    kind: str
    label: str
    tenant_id: Optional[str] = None
    embedding: Optional[List[float]] = None
    metadata: Dict[str, Any] = field(default_factory=dict)

KnowledgeGraph

Knowledge graph con persistencia SQLite + búsqueda semántica offline.

Parameters

backend: StateBackend SQLite existente (se reusa su conn), una ruta a fichero SQLite (str), o None para SQLite in-memory. embedding_provider: Provider de embeddings. Default: DeterministicEmbeddingProvider (hash-based, offline, sin red).

Source code in src\ciel\runtime\knowledge_graph.py
class KnowledgeGraph:
    """Knowledge graph con persistencia SQLite + búsqueda semántica offline.

    Parameters
    ----------
    backend:
        ``StateBackend`` SQLite existente (se reusa su ``conn``), una ruta a
        fichero SQLite (``str``), o ``None`` para SQLite in-memory.
    embedding_provider:
        Provider de embeddings. Default: ``DeterministicEmbeddingProvider``
        (hash-based, offline, sin red).
    """

    def __init__(
        self,
        backend: Optional[Any] = None,
        *,
        embedding_provider: Optional[EmbeddingProvider] = None,
    ) -> None:
        self._owns_conn = False
        if backend is None:
            self.conn = sqlite3.connect(":memory:")
            self._owns_conn = True
        elif isinstance(backend, str):
            self.conn = sqlite3.connect(backend)
            self._owns_conn = True
        elif isinstance(backend, StateBackend) and hasattr(backend, "conn"):
            self.conn = backend.conn  # SqliteStateBackend
        elif isinstance(backend, sqlite3.Connection):
            self.conn = backend
        else:
            raise TypeError(
                "KnowledgeGraph requiere un SqliteStateBackend, sqlite3.Connection, "
                f"ruta str o None; recibido: {type(backend)!r}"
            )
        self.conn.row_factory = sqlite3.Row
        init_kg_schema(self.conn)
        self._embeddings = embedding_provider or DeterministicEmbeddingProvider()
        self._index = InMemoryVectorStore()
        self._load_index()

    # --- helpers -----------------------------------------------------------
    @staticmethod
    def _sentinel(tenant_id: Optional[str]) -> str:
        return StateBackend._sentinel(tenant_id)

    @staticmethod
    def _now() -> str:
        return datetime.now(timezone.utc).isoformat()

    def _load_index(self) -> None:
        """Repuebla el índice vectorial en memoria desde SQLite (reapertura)."""
        rows = self.conn.execute(
            "SELECT tenant_id, node_id, label, embedding_json FROM kg_nodes"
        ).fetchall()
        for row in rows:
            emb = json.loads(row["embedding_json"]) if row["embedding_json"] else None
            if emb is None:
                emb = self._embeddings.embed_one(row["label"])
            tenant = None if row["tenant_id"] == "__none__" else row["tenant_id"]
            self._index.upsert(
                tenant_id=tenant,
                texts=[row["label"]],
                vectors=[emb],
                ids=[f"{row['tenant_id']}::{row['node_id']}"],
            )

    def _row_to_node(self, row: sqlite3.Row) -> KGNode:
        tenant = None if row["tenant_id"] == "__none__" else row["tenant_id"]
        return KGNode(
            id=row["node_id"],
            kind=row["kind"],
            label=row["label"],
            tenant_id=tenant,
            embedding=json.loads(row["embedding_json"]) if row["embedding_json"] else None,
            metadata=json.loads(row["metadata_json"]) if row["metadata_json"] else {},
        )

    # --- persistencia estilo F19 --------------------------------------------
    def kg_node_save(self, node: KGNode) -> KGNode:
        if node.embedding is None:
            node.embedding = self._embeddings.embed_one(node.label)
        tenant_value = self._sentinel(node.tenant_id)
        now = self._now()
        self.conn.execute(
            """
            INSERT INTO kg_nodes (tenant_id, node_id, kind, label, embedding_json, metadata_json, created_at, updated_at)
            VALUES (?, ?, ?, ?, ?, ?, ?, ?)
            ON CONFLICT(tenant_id, node_id) DO UPDATE SET
                kind = excluded.kind,
                label = excluded.label,
                embedding_json = excluded.embedding_json,
                metadata_json = excluded.metadata_json,
                updated_at = excluded.updated_at
            """,
            (
                tenant_value,
                node.id,
                node.kind,
                node.label,
                json.dumps(node.embedding),
                json.dumps(node.metadata),
                now,
                now,
            ),
        )
        self.conn.commit()
        self._index.upsert(
            tenant_id=node.tenant_id,
            texts=[node.label],
            vectors=[node.embedding],
            payloads=[{"node_id": node.id}],
            ids=[f"{tenant_value}::{node.id}"],
        )
        return node

    def kg_node_get(self, node_id: str, *, tenant_id: Optional[str] = None) -> Optional[KGNode]:
        row = self.conn.execute(
            "SELECT * FROM kg_nodes WHERE tenant_id = ? AND node_id = ?",
            (self._sentinel(tenant_id), node_id),
        ).fetchone()
        return self._row_to_node(row) if row else None

    def kg_edge_save(self, edge: KGEdge, *, tenant_id: Optional[str] = None) -> KGEdge:
        self.conn.execute(
            """
            INSERT INTO kg_edges (tenant_id, from_id, to_id, relation, weight, created_at)
            VALUES (?, ?, ?, ?, ?, ?)
            ON CONFLICT(tenant_id, from_id, to_id, relation) DO UPDATE SET
                weight = excluded.weight
            """,
            (
                self._sentinel(tenant_id),
                edge.from_id,
                edge.to_id,
                edge.relation,
                edge.weight,
                self._now(),
            ),
        )
        self.conn.commit()
        return edge

    def kg_edge_get(
        self,
        from_id: str,
        to_id: str,
        *,
        relation: Optional[str] = None,
        tenant_id: Optional[str] = None,
    ) -> List[KGEdge]:
        tenant_value = self._sentinel(tenant_id)
        if relation is not None:
            rows = self.conn.execute(
                "SELECT * FROM kg_edges WHERE tenant_id = ? AND from_id = ? AND to_id = ? AND relation = ?",
                (tenant_value, from_id, to_id, relation),
            ).fetchall()
        else:
            rows = self.conn.execute(
                "SELECT * FROM kg_edges WHERE tenant_id = ? AND from_id = ? AND to_id = ?",
                (tenant_value, from_id, to_id),
            ).fetchall()
        return [
            KGEdge(
                from_id=r["from_id"],
                to_id=r["to_id"],
                relation=r["relation"],
                weight=r["weight"],
            )
            for r in rows
        ]

    # --- API pública de grafo -----------------------------------------------
    def add_node(
        self,
        node_id: str,
        *,
        kind: str = "entity",
        label: Optional[str] = None,
        tenant_id: Optional[str] = None,
        metadata: Optional[Dict[str, Any]] = None,
    ) -> KGNode:
        node = KGNode(
            id=node_id,
            kind=kind,
            label=label or node_id,
            tenant_id=tenant_id,
            metadata=metadata or {},
        )
        return self.kg_node_save(node)

    def add_edge(
        self,
        from_id: str,
        to_id: str,
        *,
        relation: str = "related_to",
        weight: float = 1.0,
        tenant_id: Optional[str] = None,
    ) -> KGEdge:
        for nid in (from_id, to_id):
            if self.kg_node_get(nid, tenant_id=tenant_id) is None:
                raise KeyError(f"Nodo desconocido para tenant {tenant_id!r}: {nid!r}")
        edge = KGEdge(from_id=from_id, to_id=to_id, relation=relation, weight=weight)
        return self.kg_edge_save(edge, tenant_id=tenant_id)

    def neighbors(
        self,
        node_id: str,
        *,
        tenant_id: Optional[str] = None,
        direction: str = "both",
    ) -> List[Tuple[KGNode, KGEdge]]:
        """Vecinos por aristas explícitas. ``direction``: out | in | both."""
        tenant_value = self._sentinel(tenant_id)
        out: List[Tuple[KGNode, KGEdge]] = []
        if direction in ("out", "both"):
            rows = self.conn.execute(
                "SELECT * FROM kg_edges WHERE tenant_id = ? AND from_id = ?",
                (tenant_value, node_id),
            ).fetchall()
            for r in rows:
                node = self.kg_node_get(r["to_id"], tenant_id=tenant_id)
                if node is not None:
                    out.append(
                        (node, KGEdge(r["from_id"], r["to_id"], r["relation"], r["weight"]))
                    )
        if direction in ("in", "both"):
            rows = self.conn.execute(
                "SELECT * FROM kg_edges WHERE tenant_id = ? AND to_id = ?",
                (tenant_value, node_id),
            ).fetchall()
            for r in rows:
                node = self.kg_node_get(r["from_id"], tenant_id=tenant_id)
                if node is not None:
                    out.append(
                        (node, KGEdge(r["from_id"], r["to_id"], r["relation"], r["weight"]))
                    )
        return out

    def search(
        self,
        query: str,
        *,
        tenant_id: Optional[str] = None,
        top_k: int = 5,
    ) -> List[KGNode]:
        """Búsqueda semántica (coseno, offline) filtrada por tenant."""
        vector = self._embeddings.embed_one(query)
        records = self._index.query(tenant_id=tenant_id, vector=vector, top_k=top_k)
        out: List[KGNode] = []
        for rec in records:
            node_id = rec.payload.get("node_id") or rec.id.split("::", 1)[-1]
            node = self.kg_node_get(node_id, tenant_id=tenant_id)
            if node is not None:
                out.append(node)
        return out

    # --- interop opcional -----------------------------------------------------
    def to_networkx(self, *, tenant_id: Optional[str] = None):  # pragma: no cover
        """Exporta a ``networkx.DiGraph`` si networkx está instalado (opcional)."""
        try:
            import networkx as nx  # type: ignore
        except ImportError as exc:
            raise RuntimeError(
                "to_networkx requiere el paquete opcional 'networkx'."
            ) from exc
        g = nx.DiGraph()
        tenant_value = self._sentinel(tenant_id)
        for row in self.conn.execute(
            "SELECT * FROM kg_nodes WHERE tenant_id = ?", (tenant_value,)
        ).fetchall():
            node = self._row_to_node(row)
            g.add_node(node.id, kind=node.kind, label=node.label, **node.metadata)
        for row in self.conn.execute(
            "SELECT * FROM kg_edges WHERE tenant_id = ?", (tenant_value,)
        ).fetchall():
            g.add_edge(row["from_id"], row["to_id"], relation=row["relation"], weight=row["weight"])
        return g

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

neighbors(node_id: str, *, tenant_id: Optional[str] = None, direction: str = 'both') -> List[Tuple[KGNode, KGEdge]]

Vecinos por aristas explícitas. direction: out | in | both.

Source code in src\ciel\runtime\knowledge_graph.py
def neighbors(
    self,
    node_id: str,
    *,
    tenant_id: Optional[str] = None,
    direction: str = "both",
) -> List[Tuple[KGNode, KGEdge]]:
    """Vecinos por aristas explícitas. ``direction``: out | in | both."""
    tenant_value = self._sentinel(tenant_id)
    out: List[Tuple[KGNode, KGEdge]] = []
    if direction in ("out", "both"):
        rows = self.conn.execute(
            "SELECT * FROM kg_edges WHERE tenant_id = ? AND from_id = ?",
            (tenant_value, node_id),
        ).fetchall()
        for r in rows:
            node = self.kg_node_get(r["to_id"], tenant_id=tenant_id)
            if node is not None:
                out.append(
                    (node, KGEdge(r["from_id"], r["to_id"], r["relation"], r["weight"]))
                )
    if direction in ("in", "both"):
        rows = self.conn.execute(
            "SELECT * FROM kg_edges WHERE tenant_id = ? AND to_id = ?",
            (tenant_value, node_id),
        ).fetchall()
        for r in rows:
            node = self.kg_node_get(r["from_id"], tenant_id=tenant_id)
            if node is not None:
                out.append(
                    (node, KGEdge(r["from_id"], r["to_id"], r["relation"], r["weight"]))
                )
    return out

search(query: str, *, tenant_id: Optional[str] = None, top_k: int = 5) -> List[KGNode]

Búsqueda semántica (coseno, offline) filtrada por tenant.

Source code in src\ciel\runtime\knowledge_graph.py
def search(
    self,
    query: str,
    *,
    tenant_id: Optional[str] = None,
    top_k: int = 5,
) -> List[KGNode]:
    """Búsqueda semántica (coseno, offline) filtrada por tenant."""
    vector = self._embeddings.embed_one(query)
    records = self._index.query(tenant_id=tenant_id, vector=vector, top_k=top_k)
    out: List[KGNode] = []
    for rec in records:
        node_id = rec.payload.get("node_id") or rec.id.split("::", 1)[-1]
        node = self.kg_node_get(node_id, tenant_id=tenant_id)
        if node is not None:
            out.append(node)
    return out

to_networkx(*, tenant_id: Optional[str] = None)

Exporta a networkx.DiGraph si networkx está instalado (opcional).

Source code in src\ciel\runtime\knowledge_graph.py
def to_networkx(self, *, tenant_id: Optional[str] = None):  # pragma: no cover
    """Exporta a ``networkx.DiGraph`` si networkx está instalado (opcional)."""
    try:
        import networkx as nx  # type: ignore
    except ImportError as exc:
        raise RuntimeError(
            "to_networkx requiere el paquete opcional 'networkx'."
        ) from exc
    g = nx.DiGraph()
    tenant_value = self._sentinel(tenant_id)
    for row in self.conn.execute(
        "SELECT * FROM kg_nodes WHERE tenant_id = ?", (tenant_value,)
    ).fetchall():
        node = self._row_to_node(row)
        g.add_node(node.id, kind=node.kind, label=node.label, **node.metadata)
    for row in self.conn.execute(
        "SELECT * FROM kg_edges WHERE tenant_id = ?", (tenant_value,)
    ).fetchall():
        g.add_edge(row["from_id"], row["to_id"], relation=row["relation"], weight=row["weight"])
    return g

StateBackend

Bases: ABC

Interfaz mínima de persistencia compartida (multi-réplica).

Source code in src\ciel\runtime\state_backend.py
class StateBackend(abc.ABC):
    """Interfaz mínima de persistencia compartida (multi-réplica)."""

    backend_type: str = "abstract"

    @abc.abstractmethod
    def set(
        self,
        *,
        tenant_id: Optional[str],
        session_id: str,
        key: str,
        value: Any,
    ) -> None:
        ...

    @abc.abstractmethod
    def get(self, *, tenant_id: Optional[str], session_id: str, key: str) -> Optional[Any]:
        ...

    @abc.abstractmethod
    def delete(self, *, tenant_id: Optional[str], session_id: str, key: str) -> None:
        ...

    @abc.abstractmethod
    def search(self, query: str, *, limit: int = 10) -> List[Dict[str, Any]]:
        ...

    @abc.abstractmethod
    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:
        ...

    # --- memoria episódica (Fase 17, tenant-filtered) ----------------------
    # TODAS las lecturas filtran por tenant_id explícitamente para evitar
    # fuga cross-tenant (el ``search`` genérico NO filtra; no usarlo para memoria).
    @abc.abstractmethod
    def memory_append(
        self,
        *,
        tenant_id: Optional[str],
        session_id: str,
        memory_id: str,
        value: Any,
    ) -> None:
        ...

    @abc.abstractmethod
    def memory_get(
        self,
        *,
        tenant_id: Optional[str],
        session_id: str,
        memory_id: str,
    ) -> Optional[dict]:
        ...

    @abc.abstractmethod
    def memory_get_recent(
        self,
        *,
        tenant_id: Optional[str],
        session_id: str,
        limit: int = 8,
    ) -> List[dict]:
        ...

    @abc.abstractmethod
    def memory_search_tenant(
        self,
        *,
        tenant_id: Optional[str],
        session_id: Optional[str],
        query: str,
        limit: int = 5,
    ) -> List[dict]:
        ...

    @abc.abstractmethod
    def memory_clear_session(
        self, *, tenant_id: Optional[str], session_id: str
    ) -> None:
        ...

    # --- prompts versionados (Fase 19, tenant-filtered) -------------------
    @abc.abstractmethod
    def prompt_save(
        self,
        *,
        tenant_id: Optional[str],
        name: str,
        version: str,
        prompt_text: str,
        value_json: str,
        sha256: str,
        previous_version: Optional[str],
        created_at: str,
    ) -> None:
        ...

    @abc.abstractmethod
    def prompt_get(
        self, *, tenant_id: Optional[str], name: str, version: Optional[str] = None
    ) -> Optional[dict]:
        ...

    @abc.abstractmethod
    def prompt_get_history(
        self, *, tenant_id: Optional[str], name: str
    ) -> List[dict]:
        ...

    # --- log de estado cognitivo (Fase 19, tenant-filtered) ---------------
    @abc.abstractmethod
    def state_log_append(
        self,
        *,
        tenant_id: Optional[str],
        session_id: str,
        prompt_version: Optional[str],
        value_json: str,
        created_at: str,
    ) -> None:
        ...

    @abc.abstractmethod
    def state_log_get_recent(
        self, *, tenant_id: Optional[str], session_id: str, limit: int = 16
    ) -> List[dict]:
        ...

    # --- curricula versionados (Fase 20 BLOQUE A, tenant-filtered) ---------
    # NO abstractos para no romper backends de terceros: default NotImplementedError.
    def curriculum_save(
        self,
        *,
        tenant_id: Optional[str],
        goal: str,
        version: str,
        plan_json: str,
        value_json: str,
        sha256: str,
        previous_version: Optional[str],
        created_at: str,
    ) -> None:
        raise NotImplementedError(
            f"{type(self).__name__} no implementa curriculum_save (Fase 20)"
        )

    def curriculum_get(
        self, *, tenant_id: Optional[str], goal: str, version: Optional[str] = None
    ) -> Optional[dict]:
        raise NotImplementedError(
            f"{type(self).__name__} no implementa curriculum_get (Fase 20)"
        )

    def curriculum_get_history(
        self, *, tenant_id: Optional[str], goal: str
    ) -> List[dict]:
        raise NotImplementedError(
            f"{type(self).__name__} no implementa curriculum_get_history (Fase 20)"
        )

    @abc.abstractmethod
    def close(self) -> None:
        ...

    # --- readiness (usado por /readyz en F16) ------------------------------
    def is_ready(self) -> bool:
        """Devuelve True si el backend está conectado y migrado.

        El default asume listo; los backends remotos sobrescriben esto con una
        comprobación real de conectividad.
        """
        return True

    # --- utilidades compartidas --------------------------------------------
    @staticmethod
    def _sentinel(tenant_id: Optional[str]) -> Any:
        """SQLite no trata dos NULL como iguales para UNIQUE; normaliza None."""
        return tenant_id if tenant_id is not None else "__none__"

    @staticmethod
    def _dump(value: Any) -> str:
        try:
            return json.dumps(value)
        except TypeError:
            return json.dumps({"repr": repr(value)})

    @staticmethod
    def _normalize_row(value_json: Optional[str]) -> Optional[Any]:
        if value_json is None:
            return None
        try:
            return json.loads(value_json)
        except (TypeError, json.JSONDecodeError):
            return None

is_ready() -> bool

Devuelve True si el backend está conectado y migrado.

El default asume listo; los backends remotos sobrescriben esto con una comprobación real de conectividad.

Source code in src\ciel\runtime\state_backend.py
def is_ready(self) -> bool:
    """Devuelve True si el backend está conectado y migrado.

    El default asume listo; los backends remotos sobrescriben esto con una
    comprobación real de conectividad.
    """
    return True

init_kg_schema(conn: sqlite3.Connection) -> None

Crea las tablas kg_nodes/kg_edges si no existen (idempotente).

Source code in src\ciel\runtime\knowledge_graph.py
def init_kg_schema(conn: sqlite3.Connection) -> None:
    """Crea las tablas ``kg_nodes``/``kg_edges`` si no existen (idempotente)."""
    conn.execute(
        """
        CREATE TABLE IF NOT EXISTS kg_nodes (
            id INTEGER PRIMARY KEY AUTOINCREMENT,
            tenant_id TEXT,
            node_id TEXT,
            kind TEXT,
            label TEXT,
            embedding_json TEXT,
            metadata_json TEXT,
            created_at TEXT,
            updated_at TEXT,
            UNIQUE(tenant_id, node_id)
        );
        """
    )
    conn.execute(
        """
        CREATE TABLE IF NOT EXISTS kg_edges (
            id INTEGER PRIMARY KEY AUTOINCREMENT,
            tenant_id TEXT,
            from_id TEXT,
            to_id TEXT,
            relation TEXT,
            weight REAL DEFAULT 1.0,
            created_at TEXT,
            UNIQUE(tenant_id, from_id, to_id, relation)
        );
        """
    )
    conn.execute(
        "CREATE INDEX IF NOT EXISTS idx_kg_nodes_tenant ON kg_nodes(tenant_id, node_id)"
    )
    conn.execute(
        "CREATE INDEX IF NOT EXISTS idx_kg_edges_from ON kg_edges(tenant_id, from_id)"
    )
    conn.execute(
        "CREATE INDEX IF NOT EXISTS idx_kg_edges_to ON kg_edges(tenant_id, to_id)"
    )
    conn.commit()

__all__ = ['SkillRecommender'] module-attribute

DeterministicEmbeddingProvider

Bases: EmbeddingProvider

Embedding determinista offline (hash → vector). NO semántico real.

Útil para dev, tests y como fallback cuando no hay API key. Produce el mismo vector para el mismo texto (determinista), permitiendo búsqueda por similitud coseno coherente dentro de una corrida. NO usar en prod como única fuente de verdad semántica; combinar con BM25 (hybrid).

Source code in src\ciel\rag\embeddings.py
class DeterministicEmbeddingProvider(EmbeddingProvider):
    """Embedding determinista offline (hash → vector). NO semántico real.

    Útil para dev, tests y como fallback cuando no hay API key. Produce el
    mismo vector para el mismo texto (determinista), permitiendo búsqueda
    por similitud coseno coherente dentro de una corrida. NO usar en prod
    como única fuente de verdad semántica; combinar con BM25 (hybrid).
    """

    def __init__(self, dim: int = 256, seed: int = 0) -> None:
        self.dim = dim
        self._seed = seed

    def embed(self, texts: Sequence[str]) -> List[List[float]]:
        out: List[List[float]] = []
        for text in texts:
            vec = [0.0] * self.dim
            # Mezcla varias fragmentaciones del hash para dispersar.
            for window in range(0, max(1, len(text)), max(1, len(text) // self.dim or 1)):
                chunk = text[window : window + 8]
                h = hashlib.sha256(f"{self._seed}:{chunk}".encode("utf-8")).digest()
                for i in range(0, len(h), 4):
                    idx = (h[i] + h[i + 1] * 256) % self.dim
                    sign = 1.0 if (h[i + 2] & 1) == 0 else -1.0
                    vec[idx] += sign * ((h[i + 3] % 100) / 100.0)
            # Normaliza a norma 1 para coseno estable.
            norm = math.sqrt(sum(v * v for v in vec)) or 1.0
            out.append([v / norm for v in vec])
        return out

EmbeddingProvider

Bases: ABC

Contrato mínimo de embeddings.

Source code in src\ciel\rag\embeddings.py
class EmbeddingProvider(ABC):
    """Contrato mínimo de embeddings."""

    dim: int = 256

    @abstractmethod
    def embed(self, texts: Sequence[str]) -> List[List[float]]: ...

    def embed_one(self, text: str) -> List[float]:
        return self.embed([text])[0]

InMemoryVectorStore

Bases: VectorStore

Índice vectorial en memoria (offline, sin dependencias).

Aísla por tenant_id: query solo devuelve records del tenant pedido.

Source code in src\ciel\rag\vector_store.py
class InMemoryVectorStore(VectorStore):
    """Índice vectorial en memoria (offline, sin dependencias).

    Aísla por ``tenant_id``: ``query`` solo devuelve records del tenant pedido.
    """

    def __init__(self) -> None:
        self._records: List[VectorRecord] = []

    def upsert(
        self,
        *,
        tenant_id: Optional[str],
        texts: Sequence[str],
        vectors: Sequence[Sequence[float]],
        payloads: Optional[Sequence[Dict[str, Any]]] = None,
        ids: Optional[Sequence[str]] = None,
    ) -> List[str]:
        out: List[str] = []
        for i, text in enumerate(texts):
            rid = ids[i] if ids is not None else str(uuid.uuid4())
            payload = payloads[i] if payloads is not None else {}
            # Reemplazo si ya existe el id (idempotente).
            self._records = [r for r in self._records if r.id != rid]
            self._records.append(
                VectorRecord(
                    id=rid,
                    tenant_id=tenant_id,
                    vector=list(vectors[i]),
                    payload=payload,
                    text=text,
                )
            )
            out.append(rid)
        return out

    def query(
        self,
        *,
        tenant_id: Optional[str],
        vector: Sequence[float],
        top_k: int = 5,
    ) -> List[VectorRecord]:
        scored = [
            (cosine_similarity(vector, r.vector), r)
            for r in self._records
            if r.tenant_id == tenant_id
        ]
        scored.sort(key=lambda x: x[0], reverse=True)
        return [r for _, r in scored[:top_k]]

Skill dataclass

Source code in src\ciel\runtime\skills.py
@dataclass
class Skill:
    name: str
    description: str
    content: str
    category: Optional[str] = None
    metadata: Dict[str, Any] = field(default_factory=dict)
    sha256: Optional[str] = None
    domain: Optional[str] = None

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,
        domain: Optional[str] = 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,
            domain=domain,
            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,
        domain: 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,
            domain=domain if domain is not None else getattr(previous, "domain", None),
            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:
        # Unificado (Fase 20): delega en SkillVersion.parse/bump — la única
        # implementación de semver vive en ciel.runtime.skill_versioning.
        from ciel.runtime.skill_versioning import SkillVersion

        try:
            base = SkillVersion.parse(current) if current else SkillVersion()
        except SkillError:
            base = SkillVersion()
        try:
            return base.bump(bump).version
        except SkillError:
            return base.bump("patch").version

create_from_code(*, name: str, description: str, code: str, category: Optional[str] = None, tenant_id: Optional[str] = None, metadata: Optional[Dict[str, Any]] = None, domain: Optional[str] = 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,
    domain: Optional[str] = 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,
        domain=domain,
        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, domain: 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,
    domain: 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,
        domain=domain if domain is not None else getattr(previous, "domain", None),
        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

SkillRecommender

Recomendador cross-domain de skills basado en embeddings (offline).

  • index() vectoriza Skill.content de la librería (por tenant).
  • recommend() devuelve skills relevantes al contexto. Si se conoce el dominio del contexto (context_domain), solo devuelve skills cuyo domain sea distinto (transfer learning entre dominios). Si no hay dominio, devuelve el top-k por similitud.
Source code in src\ciel\runtime\skill_recommender.py
class SkillRecommender:
    """Recomendador cross-domain de skills basado en embeddings (offline).

    - ``index()`` vectoriza ``Skill.content`` de la librería (por tenant).
    - ``recommend()`` devuelve skills relevantes al contexto. Si se conoce el
      dominio del contexto (``context_domain``), solo devuelve skills cuyo
      ``domain`` sea distinto (transfer learning entre dominios). Si no hay
      dominio, devuelve el top-k por similitud.
    """

    def __init__(
        self,
        library: SkillLibrary,
        *,
        embeddings: Optional[EmbeddingProvider] = None,
        store: Optional[VectorStore] = None,
    ) -> None:
        self.library = library
        self.embeddings = embeddings or DeterministicEmbeddingProvider()
        self.store = store or InMemoryVectorStore()

    # -- indexing --------------------------------------------------------

    def index(self, *, tenant_id: Optional[str] = None) -> int:
        """(Re)indexa las skills del tenant. Devuelve cuántas se indexaron."""
        skills = self.library.list_skills(tenant_id=tenant_id)
        if not skills:
            return 0
        texts = [f"{s.name}\n{s.description}\n{s.content}" for s in skills]
        vectors = self.embeddings.embed(texts)
        self.store.upsert(
            tenant_id=tenant_id,
            texts=texts,
            vectors=vectors,
            payloads=[
                {"skill_name": s.name, "domain": getattr(s, "domain", None)}
                for s in skills
            ],
            ids=[f"{tenant_id or '__none__'}:{s.name}" for s in skills],
        )
        return len(skills)

    # -- recommendation ---------------------------------------------------

    def recommend(
        self,
        context: str,
        *,
        tenant_id: Optional[str] = None,
        top_k: int = 3,
        context_domain: Optional[str] = None,
    ) -> List[Skill]:
        """Sugiere hasta ``top_k`` skills relevantes a ``context``.

        Cross-domain: si ``context_domain`` está definido, excluye skills de
        ese mismo dominio (solo transferencia desde OTROS dominios). Sin
        dominio, devuelve el top-k global por similitud coseno.
        """
        vector = self.embeddings.embed_one(context)
        # Pedimos de más para poder filtrar por dominio sin quedarnos cortos.
        records = self.store.query(tenant_id=tenant_id, vector=vector, top_k=max(top_k * 4, top_k))
        out: List[Skill] = []
        seen: set[str] = set()
        for record in records:
            name = record.payload.get("skill_name")
            if not name or name in seen:
                continue
            skill_domain = record.payload.get("domain")
            if context_domain is not None and skill_domain == context_domain:
                continue  # mismo dominio: no es transfer cross-domain
            skill = self.library.get(name)
            if skill is None:
                continue
            seen.add(name)
            out.append(skill)
            if len(out) >= top_k:
                break
        return out

index(*, tenant_id: Optional[str] = None) -> int

(Re)indexa las skills del tenant. Devuelve cuántas se indexaron.

Source code in src\ciel\runtime\skill_recommender.py
def index(self, *, tenant_id: Optional[str] = None) -> int:
    """(Re)indexa las skills del tenant. Devuelve cuántas se indexaron."""
    skills = self.library.list_skills(tenant_id=tenant_id)
    if not skills:
        return 0
    texts = [f"{s.name}\n{s.description}\n{s.content}" for s in skills]
    vectors = self.embeddings.embed(texts)
    self.store.upsert(
        tenant_id=tenant_id,
        texts=texts,
        vectors=vectors,
        payloads=[
            {"skill_name": s.name, "domain": getattr(s, "domain", None)}
            for s in skills
        ],
        ids=[f"{tenant_id or '__none__'}:{s.name}" for s in skills],
    )
    return len(skills)

recommend(context: str, *, tenant_id: Optional[str] = None, top_k: int = 3, context_domain: Optional[str] = None) -> List[Skill]

Sugiere hasta top_k skills relevantes a context.

Cross-domain: si context_domain está definido, excluye skills de ese mismo dominio (solo transferencia desde OTROS dominios). Sin dominio, devuelve el top-k global por similitud coseno.

Source code in src\ciel\runtime\skill_recommender.py
def recommend(
    self,
    context: str,
    *,
    tenant_id: Optional[str] = None,
    top_k: int = 3,
    context_domain: Optional[str] = None,
) -> List[Skill]:
    """Sugiere hasta ``top_k`` skills relevantes a ``context``.

    Cross-domain: si ``context_domain`` está definido, excluye skills de
    ese mismo dominio (solo transferencia desde OTROS dominios). Sin
    dominio, devuelve el top-k global por similitud coseno.
    """
    vector = self.embeddings.embed_one(context)
    # Pedimos de más para poder filtrar por dominio sin quedarnos cortos.
    records = self.store.query(tenant_id=tenant_id, vector=vector, top_k=max(top_k * 4, top_k))
    out: List[Skill] = []
    seen: set[str] = set()
    for record in records:
        name = record.payload.get("skill_name")
        if not name or name in seen:
            continue
        skill_domain = record.payload.get("domain")
        if context_domain is not None and skill_domain == context_domain:
            continue  # mismo dominio: no es transfer cross-domain
        skill = self.library.get(name)
        if skill is None:
            continue
        seen.add(name)
        out.append(skill)
        if len(out) >= top_k:
            break
    return out

VectorStore

Bases: ABC

Contrato de índice vectorial (aislado por tenant_id).

Source code in src\ciel\rag\vector_store.py
class VectorStore(ABC):
    """Contrato de índice vectorial (aislado por tenant_id)."""

    @abstractmethod
    def upsert(
        self,
        *,
        tenant_id: Optional[str],
        texts: Sequence[str],
        vectors: Sequence[Sequence[float]],
        payloads: Optional[Sequence[Dict[str, Any]]] = None,
        ids: Optional[Sequence[str]] = None,
    ) -> List[str]: ...

    @abstractmethod
    def query(
        self,
        *,
        tenant_id: Optional[str],
        vector: Sequence[float],
        top_k: int = 5,
    ) -> List[VectorRecord]: ...