Skip to content

Core API Reference

effero.core

Cognition core: orchestrator/planner, memory, and model router.

This is the "brain" layer described in the project README's Architecture section. Submodules:

  • agent -- the plan -> act -> observe -> replan loop
  • memory -- working / episodic / semantic memory
  • planner -- task decomposition and replanning
  • router -- routes LLM calls to a local or cloud backend

EventBus

High-performance asynchronous pub/sub event bus with compiled pattern routing.

Source code in src/effero/core/event_bus.py
class EventBus:
    """High-performance asynchronous pub/sub event bus with compiled pattern routing."""

    def __init__(self, max_history: int = 1000):
        self._max_history = max_history
        self._history: deque[Event] = deque(maxlen=max_history)
        self._subscribers: dict[str, list[EventHandler]] = {}
        self._exact_subscribers: dict[str, list[EventHandler]] = {}
        self._pattern_subscribers: list[tuple[str, re.Pattern[str], list[EventHandler]]] = []

    def subscribe(self, topic_pattern: str, handler: EventHandler) -> None:
        """Subscribe handler to exact topic or wildcard pattern."""
        if topic_pattern not in self._subscribers:
            self._subscribers[topic_pattern] = []
            if any(ch in topic_pattern for ch in ("*", "?", "[")):
                compiled = re.compile(fnmatch.translate(topic_pattern))
                self._pattern_subscribers.append((topic_pattern, compiled, self._subscribers[topic_pattern]))
            else:
                self._exact_subscribers[topic_pattern] = self._subscribers[topic_pattern]

        self._subscribers[topic_pattern].append(handler)

    def publish(self, event: Event) -> None:
        """Publish event and dispatch asynchronously to matched subscribers."""
        self._history.append(event)

        # Collect unique handlers across exact matches and wildcard patterns
        target_handlers: list[EventHandler] = []

        # 1. O(1) exact match lookup
        if event.topic in self._exact_subscribers:
            target_handlers.extend(self._exact_subscribers[event.topic])

        # 2. Pre-compiled regex evaluation for wildcard subscriptions
        for _pattern_str, regex, handlers in self._pattern_subscribers:
            if regex.match(event.topic):
                target_handlers.extend(handlers)

        if not target_handlers:
            return

        try:
            loop = asyncio.get_running_loop()
        except RuntimeError:
            return

        for handler in target_handlers:
            loop.create_task(self._dispatch(handler, event))

    def get_history(self, topic_pattern: str = "*", limit: int = 100) -> list[Event]:
        """Retrieve recent chronological history filtered by topic."""
        if limit <= 0:
            return []

        if topic_pattern == "*":
            # Direct slice without pattern checks
            hist = list(self._history)
            return hist[-limit:]

        if not any(ch in topic_pattern for ch in ("*", "?", "[")):
            matches = [e for e in self._history if e.topic == topic_pattern]
            return matches[-limit:]

        regex = re.compile(fnmatch.translate(topic_pattern))
        matches = [e for e in self._history if regex.match(e.topic)]
        return matches[-limit:]

    async def _dispatch(self, handler: EventHandler, event: Event) -> None:
        try:
            await handler(event)
        except Exception as e:
            logger.error(f"Error in event handler for topic {event.topic}: {e}")

    @staticmethod
    def _match_topic(pattern: str, topic: str) -> bool:
        return fnmatch.fnmatch(topic, pattern)

subscribe(topic_pattern, handler)

Subscribe handler to exact topic or wildcard pattern.

Source code in src/effero/core/event_bus.py
def subscribe(self, topic_pattern: str, handler: EventHandler) -> None:
    """Subscribe handler to exact topic or wildcard pattern."""
    if topic_pattern not in self._subscribers:
        self._subscribers[topic_pattern] = []
        if any(ch in topic_pattern for ch in ("*", "?", "[")):
            compiled = re.compile(fnmatch.translate(topic_pattern))
            self._pattern_subscribers.append((topic_pattern, compiled, self._subscribers[topic_pattern]))
        else:
            self._exact_subscribers[topic_pattern] = self._subscribers[topic_pattern]

    self._subscribers[topic_pattern].append(handler)

publish(event)

Publish event and dispatch asynchronously to matched subscribers.

Source code in src/effero/core/event_bus.py
def publish(self, event: Event) -> None:
    """Publish event and dispatch asynchronously to matched subscribers."""
    self._history.append(event)

    # Collect unique handlers across exact matches and wildcard patterns
    target_handlers: list[EventHandler] = []

    # 1. O(1) exact match lookup
    if event.topic in self._exact_subscribers:
        target_handlers.extend(self._exact_subscribers[event.topic])

    # 2. Pre-compiled regex evaluation for wildcard subscriptions
    for _pattern_str, regex, handlers in self._pattern_subscribers:
        if regex.match(event.topic):
            target_handlers.extend(handlers)

    if not target_handlers:
        return

    try:
        loop = asyncio.get_running_loop()
    except RuntimeError:
        return

    for handler in target_handlers:
        loop.create_task(self._dispatch(handler, event))

get_history(topic_pattern='*', limit=100)

Retrieve recent chronological history filtered by topic.

Source code in src/effero/core/event_bus.py
def get_history(self, topic_pattern: str = "*", limit: int = 100) -> list[Event]:
    """Retrieve recent chronological history filtered by topic."""
    if limit <= 0:
        return []

    if topic_pattern == "*":
        # Direct slice without pattern checks
        hist = list(self._history)
        return hist[-limit:]

    if not any(ch in topic_pattern for ch in ("*", "?", "[")):
        matches = [e for e in self._history if e.topic == topic_pattern]
        return matches[-limit:]

    regex = re.compile(fnmatch.translate(topic_pattern))
    matches = [e for e in self._history if regex.match(e.topic)]
    return matches[-limit:]

EpisodicMemory

Source code in src/effero/core/memory/episodic.py
class EpisodicMemory:
    def __init__(self, session_dir: Path | str | None = None):
        if session_dir is None:
            session_dir = Path.cwd() / "sessions"
        self.session_dir = Path(session_dir)
        self.session_dir.mkdir(parents=True, exist_ok=True)
        session_id = str(int(time.time()))
        self.session_path = self.session_dir / f"session_{session_id}.jsonl"

    def get_session_path(self) -> Path:
        return self.session_path

    def record(self, event_type: str, data: dict[str, Any]) -> None:
        entry = {
            "timestamp": time.time(),
            "type": event_type,
            "data": data,
        }
        try:
            with open(self.session_path, "a") as f:
                f.write(json.dumps(entry) + "\n")
        except Exception as e:
            logger.error(f"Failed to write episodic memory: {e}")

    def load_history(self, limit: int = 100) -> list[dict[str, Any]]:
        """Load the most recent episodic events from disk without reading the full file into RAM."""
        if not self.session_path.exists() or limit <= 0:
            return []

        entries: list[dict[str, Any]] = []
        chunk_size = 8192

        try:
            with open(self.session_path, "rb") as f:
                f.seek(0, 2)
                file_size = f.tell()
                pos = file_size
                buffer = b""

                while pos > 0 and len(entries) < limit:
                    read_size = min(chunk_size, pos)
                    pos -= read_size
                    f.seek(pos)
                    chunk = f.read(read_size)
                    buffer = chunk + buffer

                    lines = buffer.split(b"\n")
                    buffer = lines[0]

                    for raw_line in reversed(lines[1:]):
                        line_str = raw_line.decode("utf-8", errors="replace").strip()
                        if line_str and len(entries) < limit:
                            try:
                                entries.append(json.loads(line_str))
                            except json.JSONDecodeError:
                                continue

                if pos == 0 and buffer.strip() and len(entries) < limit:
                    line_str = buffer.decode("utf-8", errors="replace").strip()
                    if line_str:
                        try:
                            entries.append(json.loads(line_str))
                        except json.JSONDecodeError:
                            pass

            entries.reverse()
        except Exception as e:
            logger.error(f"Failed to read episodic memory: {e}")

        return entries

load_history(limit=100)

Load the most recent episodic events from disk without reading the full file into RAM.

Source code in src/effero/core/memory/episodic.py
def load_history(self, limit: int = 100) -> list[dict[str, Any]]:
    """Load the most recent episodic events from disk without reading the full file into RAM."""
    if not self.session_path.exists() or limit <= 0:
        return []

    entries: list[dict[str, Any]] = []
    chunk_size = 8192

    try:
        with open(self.session_path, "rb") as f:
            f.seek(0, 2)
            file_size = f.tell()
            pos = file_size
            buffer = b""

            while pos > 0 and len(entries) < limit:
                read_size = min(chunk_size, pos)
                pos -= read_size
                f.seek(pos)
                chunk = f.read(read_size)
                buffer = chunk + buffer

                lines = buffer.split(b"\n")
                buffer = lines[0]

                for raw_line in reversed(lines[1:]):
                    line_str = raw_line.decode("utf-8", errors="replace").strip()
                    if line_str and len(entries) < limit:
                        try:
                            entries.append(json.loads(line_str))
                        except json.JSONDecodeError:
                            continue

            if pos == 0 and buffer.strip() and len(entries) < limit:
                line_str = buffer.decode("utf-8", errors="replace").strip()
                if line_str:
                    try:
                        entries.append(json.loads(line_str))
                    except json.JSONDecodeError:
                        pass

        entries.reverse()
    except Exception as e:
        logger.error(f"Failed to read episodic memory: {e}")

    return entries

SemanticMemory

Vector-store backed semantic recall memory with spatial dimension indexing.

Stores knowledge chunks, actions, and observations, enabling nearest-neighbor semantic search to retrieve contextually relevant memories for planning.

Source code in src/effero/core/memory/semantic.py
class SemanticMemory:
    """Vector-store backed semantic recall memory with spatial dimension indexing.

    Stores knowledge chunks, actions, and observations, enabling nearest-neighbor
    semantic search to retrieve contextually relevant memories for planning.
    """

    def __init__(
        self,
        embedder: EmbeddingProvider | None = None,
        persistence_path: str | Path | None = None,
    ):
        self.embedder: EmbeddingProvider = embedder or LightweightTFIDFEmbedding()
        self.persistence_path = Path(persistence_path) if persistence_path else None
        self._records: dict[str, MemoryRecord] = {}
        self._dim_index: dict[int, set[str]] = {}

        if self.persistence_path and self.persistence_path.exists():
            self.load()

    def _index_record(self, record_id: str, vec: list[float]) -> None:
        """Register active non-zero dimensions of an embedding in the inverted index."""
        for i, val in enumerate(vec):
            if val > 1e-5:
                self._dim_index.setdefault(i, set()).add(record_id)

    def store(
        self,
        text: str,
        metadata: dict[str, Any] | None = None,
        record_id: str | None = None,
    ) -> str:
        """Store a text snippet and its metadata in semantic vector memory."""
        rec_id = record_id or str(uuid.uuid4())
        vec = self.embedder.embed(text)
        record = MemoryRecord(
            id=rec_id,
            text=text,
            embedding=vec,
            metadata=metadata or {},
        )
        self._records[rec_id] = record
        self._index_record(rec_id, vec)

        if self.persistence_path:
            self.save()

        return rec_id

    def search(
        self,
        query: str,
        top_k: int = 5,
        min_score: float = 0.0,
        filter_metadata: dict[str, Any] | None = None,
    ) -> list[MemoryRecord]:
        """Search memory for records most similar to the query."""
        if not self._records or not query.strip() or top_k <= 0:
            return []

        query_vec = self.embedder.embed(query)

        # Fast candidate pruning via inverted dimension index if knowledge base is large
        if len(self._records) > 50:
            active_dims = [i for i, v in enumerate(query_vec) if v > 1e-5]
            if active_dims:
                candidate_ids: set[str] = set()
                for d in active_dims:
                    matched = self._dim_index.get(d)
                    if matched:
                        candidate_ids.update(matched)
                candidates = [self._records[cid] for cid in candidate_ids if cid in self._records]
            else:
                candidates = list(self._records.values())
        else:
            candidates = list(self._records.values())

        scored: list[MemoryRecord] = []
        for record in candidates:
            if filter_metadata:
                match = all(record.metadata.get(k) == v for k, v in filter_metadata.items())
                if not match:
                    continue

            sim = cosine_similarity(query_vec, record.embedding)
            if sim >= min_score:
                rec_copy = MemoryRecord(
                    id=record.id,
                    text=record.text,
                    embedding=record.embedding,
                    metadata=record.metadata,
                    score=sim,
                    timestamp=record.timestamp,
                )
                scored.append(rec_copy)

        return heapq.nlargest(top_k, scored, key=lambda r: r.score)

    def delete(self, record_id: str) -> bool:
        """Delete a record by ID."""
        if record_id in self._records:
            del self._records[record_id]
            for dim_set in self._dim_index.values():
                dim_set.discard(record_id)
            if self.persistence_path:
                self.save()
            return True
        return False

    def clear(self) -> None:
        """Clear all stored semantic memories."""
        self._records.clear()
        self._dim_index.clear()
        if self.persistence_path and self.persistence_path.exists():
            try:
                self.persistence_path.unlink()
            except OSError as e:
                logger.error(f"Error removing memory persistence file: {e}")

    def count(self) -> int:
        """Total number of stored memory records."""
        return len(self._records)

    def save(self, path: str | Path | None = None) -> None:
        """Persist memory records to JSON file."""
        target = Path(path) if path else self.persistence_path
        if not target:
            return

        target.parent.mkdir(parents=True, exist_ok=True)
        data = [r.to_dict() for r in self._records.values()]
        target.write_text(json.dumps(data, indent=2), encoding="utf-8")

    def load(self, path: str | Path | None = None) -> None:
        """Load memory records from JSON file."""
        target = Path(path) if path else self.persistence_path
        if not target or not target.exists():
            return

        try:
            content = target.read_text(encoding="utf-8")
            data = json.loads(content)
            self._records = {item["id"]: MemoryRecord.from_dict(item) for item in data}
            self._dim_index.clear()
            for rec in self._records.values():
                self._index_record(rec.id, rec.embedding)
        except Exception as e:
            logger.error(f"Failed to load semantic memory from {target}: {e}")

store(text, metadata=None, record_id=None)

Store a text snippet and its metadata in semantic vector memory.

Source code in src/effero/core/memory/semantic.py
def store(
    self,
    text: str,
    metadata: dict[str, Any] | None = None,
    record_id: str | None = None,
) -> str:
    """Store a text snippet and its metadata in semantic vector memory."""
    rec_id = record_id or str(uuid.uuid4())
    vec = self.embedder.embed(text)
    record = MemoryRecord(
        id=rec_id,
        text=text,
        embedding=vec,
        metadata=metadata or {},
    )
    self._records[rec_id] = record
    self._index_record(rec_id, vec)

    if self.persistence_path:
        self.save()

    return rec_id

search(query, top_k=5, min_score=0.0, filter_metadata=None)

Search memory for records most similar to the query.

Source code in src/effero/core/memory/semantic.py
def search(
    self,
    query: str,
    top_k: int = 5,
    min_score: float = 0.0,
    filter_metadata: dict[str, Any] | None = None,
) -> list[MemoryRecord]:
    """Search memory for records most similar to the query."""
    if not self._records or not query.strip() or top_k <= 0:
        return []

    query_vec = self.embedder.embed(query)

    # Fast candidate pruning via inverted dimension index if knowledge base is large
    if len(self._records) > 50:
        active_dims = [i for i, v in enumerate(query_vec) if v > 1e-5]
        if active_dims:
            candidate_ids: set[str] = set()
            for d in active_dims:
                matched = self._dim_index.get(d)
                if matched:
                    candidate_ids.update(matched)
            candidates = [self._records[cid] for cid in candidate_ids if cid in self._records]
        else:
            candidates = list(self._records.values())
    else:
        candidates = list(self._records.values())

    scored: list[MemoryRecord] = []
    for record in candidates:
        if filter_metadata:
            match = all(record.metadata.get(k) == v for k, v in filter_metadata.items())
            if not match:
                continue

        sim = cosine_similarity(query_vec, record.embedding)
        if sim >= min_score:
            rec_copy = MemoryRecord(
                id=record.id,
                text=record.text,
                embedding=record.embedding,
                metadata=record.metadata,
                score=sim,
                timestamp=record.timestamp,
            )
            scored.append(rec_copy)

    return heapq.nlargest(top_k, scored, key=lambda r: r.score)

delete(record_id)

Delete a record by ID.

Source code in src/effero/core/memory/semantic.py
def delete(self, record_id: str) -> bool:
    """Delete a record by ID."""
    if record_id in self._records:
        del self._records[record_id]
        for dim_set in self._dim_index.values():
            dim_set.discard(record_id)
        if self.persistence_path:
            self.save()
        return True
    return False

clear()

Clear all stored semantic memories.

Source code in src/effero/core/memory/semantic.py
def clear(self) -> None:
    """Clear all stored semantic memories."""
    self._records.clear()
    self._dim_index.clear()
    if self.persistence_path and self.persistence_path.exists():
        try:
            self.persistence_path.unlink()
        except OSError as e:
            logger.error(f"Error removing memory persistence file: {e}")

count()

Total number of stored memory records.

Source code in src/effero/core/memory/semantic.py
def count(self) -> int:
    """Total number of stored memory records."""
    return len(self._records)

save(path=None)

Persist memory records to JSON file.

Source code in src/effero/core/memory/semantic.py
def save(self, path: str | Path | None = None) -> None:
    """Persist memory records to JSON file."""
    target = Path(path) if path else self.persistence_path
    if not target:
        return

    target.parent.mkdir(parents=True, exist_ok=True)
    data = [r.to_dict() for r in self._records.values()]
    target.write_text(json.dumps(data, indent=2), encoding="utf-8")

load(path=None)

Load memory records from JSON file.

Source code in src/effero/core/memory/semantic.py
def load(self, path: str | Path | None = None) -> None:
    """Load memory records from JSON file."""
    target = Path(path) if path else self.persistence_path
    if not target or not target.exists():
        return

    try:
        content = target.read_text(encoding="utf-8")
        data = json.loads(content)
        self._records = {item["id"]: MemoryRecord.from_dict(item) for item in data}
        self._dim_index.clear()
        for rec in self._records.values():
            self._index_record(rec.id, rec.embedding)
    except Exception as e:
        logger.error(f"Failed to load semantic memory from {target}: {e}")

WorkingMemory dataclass

Source code in src/effero/core/memory/working.py
@dataclass
class WorkingMemory:
    max_messages: int = 100
    messages: list[dict[str, Any]] = field(default_factory=list)
    telemetry_streams: dict[str, TimeSeriesRingBuffer[Any]] = field(default_factory=dict)

    def add_message(self, role: str, content: str, **kwargs: Any) -> None:
        msg = {"role": role, "content": content}
        msg.update(kwargs)
        self.messages.append(msg)
        self._trim()

    def add_tool_result(self, tool_call_id: str, name: str, result: str) -> None:
        msg = {
            "role": "tool",
            "tool_call_id": tool_call_id,
            "name": name,
            "content": result,
        }
        self.messages.append(msg)
        self._trim()

    def record_metric(
        self,
        metric_name: str,
        value: Any,
        timestamp: float | None = None,
        max_capacity: int = 500,
    ) -> None:
        """Append a timestamped observation to named telemetry stream."""
        if metric_name not in self.telemetry_streams:
            self.telemetry_streams[metric_name] = TimeSeriesRingBuffer(max_capacity=max_capacity)
        self.telemetry_streams[metric_name].append(value, timestamp=timestamp)

    def get_latest_metric(self, metric_name: str) -> Any | None:
        """Retrieve most recent observation for a named telemetry metric."""
        stream = self.telemetry_streams.get(metric_name)
        if stream is None:
            return None
        pt = stream.get_latest()
        return pt.value if pt else None

    def get_metric_series(self, metric_name: str, last_seconds: float | None = None) -> list[tuple[float, Any]]:
        """Retrieve chronological (timestamp, value) observations for a named metric."""
        stream = self.telemetry_streams.get(metric_name)
        if stream is None:
            return []
        if last_seconds is not None:
            points = stream.get_last_n_seconds(last_seconds)
        else:
            points = list(stream._buffer)
        return [(pt.timestamp, pt.value) for pt in points]

    def interpolate_metric(self, metric_name: str, timestamp: float) -> float | None:
        """Estimate numeric scalar value at specific historical timestamp."""
        stream = self.telemetry_streams.get(metric_name)
        if stream is None:
            return None
        return stream.interpolate_numeric(timestamp)

    def rolling_metric(self, metric_name: str, seconds: float) -> RollingWindowView[Any] | None:
        """Return a rolling window statistical view for a named metric."""
        stream = self.telemetry_streams.get(metric_name)
        if stream is None:
            return None
        return stream.rolling(seconds)

    def get_context(self) -> list[dict[str, Any]]:
        return list(self.messages)

    def clear(self) -> None:
        self.messages.clear()
        for stream in self.telemetry_streams.values():
            stream.clear()

    def to_dict(self) -> dict[str, Any]:
        return {
            "messages": self.messages,
            "telemetry_metrics": {k: len(v) for k, v in self.telemetry_streams.items()},
        }

    def _trim(self) -> None:
        if len(self.messages) <= self.max_messages:
            return

        system_msgs = [m for m in self.messages if m.get("role") == "system"]
        num_to_keep = self.max_messages - len(system_msgs)

        if num_to_keep <= 0:
            self.messages[:] = system_msgs
            return

        non_system_to_keep: list[dict[str, Any]] = []
        for m in reversed(self.messages):
            if m.get("role") != "system":
                non_system_to_keep.append(m)
                if len(non_system_to_keep) >= num_to_keep:
                    break

        non_system_to_keep.reverse()
        self.messages[:] = system_msgs + non_system_to_keep

record_metric(metric_name, value, timestamp=None, max_capacity=500)

Append a timestamped observation to named telemetry stream.

Source code in src/effero/core/memory/working.py
def record_metric(
    self,
    metric_name: str,
    value: Any,
    timestamp: float | None = None,
    max_capacity: int = 500,
) -> None:
    """Append a timestamped observation to named telemetry stream."""
    if metric_name not in self.telemetry_streams:
        self.telemetry_streams[metric_name] = TimeSeriesRingBuffer(max_capacity=max_capacity)
    self.telemetry_streams[metric_name].append(value, timestamp=timestamp)

get_latest_metric(metric_name)

Retrieve most recent observation for a named telemetry metric.

Source code in src/effero/core/memory/working.py
def get_latest_metric(self, metric_name: str) -> Any | None:
    """Retrieve most recent observation for a named telemetry metric."""
    stream = self.telemetry_streams.get(metric_name)
    if stream is None:
        return None
    pt = stream.get_latest()
    return pt.value if pt else None

get_metric_series(metric_name, last_seconds=None)

Retrieve chronological (timestamp, value) observations for a named metric.

Source code in src/effero/core/memory/working.py
def get_metric_series(self, metric_name: str, last_seconds: float | None = None) -> list[tuple[float, Any]]:
    """Retrieve chronological (timestamp, value) observations for a named metric."""
    stream = self.telemetry_streams.get(metric_name)
    if stream is None:
        return []
    if last_seconds is not None:
        points = stream.get_last_n_seconds(last_seconds)
    else:
        points = list(stream._buffer)
    return [(pt.timestamp, pt.value) for pt in points]

interpolate_metric(metric_name, timestamp)

Estimate numeric scalar value at specific historical timestamp.

Source code in src/effero/core/memory/working.py
def interpolate_metric(self, metric_name: str, timestamp: float) -> float | None:
    """Estimate numeric scalar value at specific historical timestamp."""
    stream = self.telemetry_streams.get(metric_name)
    if stream is None:
        return None
    return stream.interpolate_numeric(timestamp)

rolling_metric(metric_name, seconds)

Return a rolling window statistical view for a named metric.

Source code in src/effero/core/memory/working.py
def rolling_metric(self, metric_name: str, seconds: float) -> RollingWindowView[Any] | None:
    """Return a rolling window statistical view for a named metric."""
    stream = self.telemetry_streams.get(metric_name)
    if stream is None:
        return None
    return stream.rolling(seconds)

effero.core.router.base

Base interface for LLM backends.

LLMRequest dataclass

A request to an LLM backend.

Source code in src/effero/core/router/base.py
@dataclass
class LLMRequest:
    """A request to an LLM backend."""

    messages: list[dict[str, Any]]
    tools: list[dict[str, Any]] | None = None
    temperature: float = 0.7
    max_tokens: int = 4096
    system: str | None = None

ToolCall dataclass

A tool call returned by the LLM.

Source code in src/effero/core/router/base.py
@dataclass
class ToolCall:
    """A tool call returned by the LLM."""

    id: str
    name: str
    arguments: dict[str, Any]

LLMResponse dataclass

A response from an LLM backend.

Source code in src/effero/core/router/base.py
@dataclass
class LLMResponse:
    """A response from an LLM backend."""

    content: str | None = None
    tool_calls: list[ToolCall] = field(default_factory=list)
    usage: dict[str, int] = field(default_factory=dict)
    model: str = ""
    backend: str = ""
    confidence: float | None = None

LLMBackend

Bases: ABC

Abstract base class for LLM backends.

Source code in src/effero/core/router/base.py
class LLMBackend(ABC):
    """Abstract base class for LLM backends."""

    @abstractmethod
    async def complete(self, request: LLMRequest) -> LLMResponse:
        """Send a completion request to the LLM."""
        ...

    @abstractmethod
    async def is_available(self) -> bool:
        """Check if this backend is reachable and configured."""
        ...

complete(request) abstractmethod async

Send a completion request to the LLM.

Source code in src/effero/core/router/base.py
@abstractmethod
async def complete(self, request: LLMRequest) -> LLMResponse:
    """Send a completion request to the LLM."""
    ...

is_available() abstractmethod async

Check if this backend is reachable and configured.

Source code in src/effero/core/router/base.py
@abstractmethod
async def is_available(self) -> bool:
    """Check if this backend is reachable and configured."""
    ...