Skip to content

Index

system

Top-level system composition: JarvisSystem, SystemBuilder, and helpers.

Classes

SystemBuilder

SystemBuilder(config: Optional[JarvisConfig] = None, *, config_path: Optional[Any] = None)

Config-driven fluent builder for JarvisSystem.

Source code in src/openjarvis/system/builder.py
def __init__(
    self,
    config: Optional[JarvisConfig] = None,
    *,
    config_path: Optional[Any] = None,
) -> None:
    if config is not None:
        self._config = config
    elif config_path is not None:
        from pathlib import Path

        self._config = load_config(Path(config_path))
    else:
        self._config = load_config()

    self._engine_key: Optional[str] = None
    self._engine_instance: Optional[InferenceEngine] = None
    self._engine_instance_key: Optional[str] = None
    self._model: Optional[str] = None
    self._agent_name: Optional[str] = None
    self._tool_names: Optional[List[str]] = None
    self._telemetry: Optional[bool] = None
    self._traces: Optional[bool] = None
    self._bus: Optional[EventBus] = None
    self._sandbox: Optional[bool] = None
    self._scheduler: Optional[bool] = None
    self._workflow: Optional[bool] = None
    self._sessions: Optional[bool] = None
    self._speech: Optional[bool] = None
    self._mcp_clients: List = []
    self._mcp_tools: List[BaseTool] = []
Functions
engine_instance
engine_instance(engine: InferenceEngine, key: str = 'openai-compat') -> SystemBuilder

Inject a pre-built engine instance, bypassing engine discovery.

Used by callers that must target one exact endpoint (e.g. jarvis eval --base-url). build() health-checks the instance and raises a loud error if it is unreachable — it never silently substitutes a different discovered engine.

Source code in src/openjarvis/system/builder.py
def engine_instance(
    self, engine: InferenceEngine, key: str = "openai-compat"
) -> SystemBuilder:
    """Inject a pre-built engine instance, bypassing engine discovery.

    Used by callers that must target one exact endpoint (e.g.
    ``jarvis eval --base-url``). ``build()`` health-checks the instance
    and raises a loud error if it is unreachable — it never silently
    substitutes a different discovered engine.
    """
    self._engine_instance = engine
    self._engine_instance_key = key
    return self
build
build() -> JarvisSystem

Construct a fully wired JarvisSystem.

Source code in src/openjarvis/system/builder.py
def build(self) -> JarvisSystem:
    """Construct a fully wired JarvisSystem."""
    # Discovery state belongs to one build only.  Once a system is
    # returned, that system owns the clients and adapters captured below;
    # retaining them here would make a reused builder hand closed clients
    # from an earlier system to the next one.
    self._clear_mcp_discovery_state(close_clients=True)
    try:
        system = self._build()
    except BaseException:
        # No system took ownership, so release any clients opened before
        # the build failed.
        self._clear_mcp_discovery_state(close_clients=True)
        raise
    self._clear_mcp_discovery_state(close_clients=False)
    return system

AgentRuntime dataclass

AgentRuntime(agent: Optional[BaseAgent] = None, agent_name: str = '', manager: Optional[AgentManager] = None, scheduler: Optional[AgentScheduler] = None, executor: Optional[AgentExecutor] = None)

Active agent and agent lifecycle managers.

Observability dataclass

Observability(telemetry_store: Optional[TelemetryStore] = None, trace_store: Optional[TraceStore] = None, trace_collector: Optional[TraceCollector] = None, gpu_monitor: Optional[GpuMonitor] = None)

Telemetry, traces, and hardware monitoring.

Scheduling dataclass

Scheduling(store: Optional[SchedulerStore] = None, runner: Optional[TaskScheduler] = None)

Task scheduler and its persistent store.

SecurityContext dataclass

SecurityContext(capability_policy: Optional[CapabilityPolicy] = None, audit_logger: Optional[AuditLogger] = None, boundary_guard: Optional[BoundaryGuard] = None)

Security policy, audit, and boundary enforcement.

JarvisSystem dataclass

JarvisSystem(config: JarvisConfig, bus: EventBus, engine: InferenceEngine, engine_key: str, model: str, agent: Optional[BaseAgent] = None, agent_name: str = '', tools: List[BaseTool] = list(), tool_executor: Optional[ToolExecutor] = None, memory_backend: Optional[MemoryBackend] = None, channel_backend: Optional[BaseChannel] = None, router: Optional[RouterPolicy] = None, mcp_server: Optional[MCPServer] = None, telemetry_store: Optional[TelemetryStore] = None, trace_store: Optional[TraceStore] = None, trace_collector: Optional[TraceCollector] = None, gpu_monitor: Optional[GpuMonitor] = None, scheduler_store: Optional[SchedulerStore] = None, scheduler: Optional[TaskScheduler] = None, container_runner: Optional[ContainerRunner] = None, workflow_engine: Optional[WorkflowEngine] = None, session_store: Optional[SessionStore] = None, capability_policy: Optional[CapabilityPolicy] = None, audit_logger: Optional[AuditLogger] = None, boundary_guard: Optional[BoundaryGuard] = None, operator_manager: Optional[OperatorManager] = None, agent_manager: Optional[AgentManager] = None, agent_scheduler: Optional[AgentScheduler] = None, agent_executor: Optional[AgentExecutor] = None, speech_backend: Optional[SpeechBackend] = None, skill_manager: Optional[SkillManager] = None, _learning_orchestrator: Optional[LearningOrchestrator] = None, _mcp_clients: List[MCPClient] = list(), mcp_tools: List[BaseTool] = list())

Fully wired system -- the single source of truth for primitive composition.

Functions
wire_channel
wire_channel(channel_bridge: Any) -> None

Register a message handler on channel_bridge that routes every incoming message through this system (agent or engine) and replies.

Sessions are isolated per "<channel>:<conversation_id>" key so each chat retains its own history.

PARAMETER DESCRIPTION
channel_bridge

A connected :class:~openjarvis.channels._stubs.BaseChannel instance whose on_message method accepts a callable.

TYPE: Any

Source code in src/openjarvis/system/core.py
def wire_channel(self, channel_bridge: Any) -> None:
    """Register a message handler on *channel_bridge* that routes every
    incoming message through this system (agent or engine) and replies.

    Sessions are isolated per ``"<channel>:<conversation_id>"`` key so
    each chat retains its own history.

    Parameters
    ----------
    channel_bridge:
        A connected :class:`~openjarvis.channels._stubs.BaseChannel`
        instance whose ``on_message`` method accepts a callable.
    """
    from openjarvis.core.types import Message
    from openjarvis.sessions.session import SessionStore

    if self.session_store is None:
        from pathlib import Path

        self.session_store = SessionStore(
            db_path=Path(self.config.sessions.db_path).expanduser(),
            max_age_hours=self.config.sessions.max_age_hours,
            consolidation_threshold=self.config.sessions.consolidation_threshold,
        )

    _system = self  # capture for closure

    def _on_channel_message(cm) -> None:
        session_key = f"{cm.channel}:{cm.conversation_id}"
        session = _system.session_store.get_or_create(
            session_key,
            channel=cm.channel,
            channel_user_id=cm.sender,
        )

        prior_msgs: List[Message] = []
        for sm in session.messages:
            try:
                role = Role(sm.role)
            except ValueError:
                role = Role.USER
            prior_msgs.append(Message(role=role, content=sm.content))

        reply = ""
        try:
            if _system.agent_name and _system.agent_name != "none":
                result = _system.ask(
                    cm.content,
                    context=False,
                    agent=_system.agent_name,
                    prior_messages=prior_msgs,
                )
                reply = result.get("content", "")
            else:
                result = _system.ask(
                    cm.content,
                    context=False,
                    prior_messages=prior_msgs,
                )
                reply = result.get("content", "")
        except Exception:
            logger.exception("Channel message handler error")
            reply = "Sorry, I encountered an error processing your message."

        try:
            _system.session_store.save_message(
                session.session_id,
                "user",
                cm.content,
                channel=cm.channel,
            )
            _system.session_store.save_message(
                session.session_id,
                "assistant",
                reply,
                channel=cm.channel,
            )
        except Exception:
            logger.debug("Session save error", exc_info=True)

        if reply:
            try:
                # Canonical channel send contract (see BaseChannel.send):
                # the first positional arg is the DESTINATION id, and the
                # `conversation_id=` kwarg is the inbound message id used as
                # a reply/thread reference.  ``cm.conversation_id`` holds the
                # real per-adapter destination (Discord/Slack channel id,
                # Telegram chat id, ...) while ``cm.channel`` is only the
                # channel TYPE label ("discord", "telegram", ...).  Passing
                # the type label as the destination produced HTTP 400s
                # (#515) and using the channel id as a reply reference
                # produced MESSAGE_REFERENCE_UNKNOWN_MESSAGE (#516).
                channel_bridge.send(
                    cm.conversation_id,
                    reply,
                    conversation_id=getattr(cm, "message_id", ""),
                )
            except Exception:
                logger.exception("Channel send error")

    channel_bridge.on_message(_on_channel_message)
close
close() -> None

Release resources.

Source code in src/openjarvis/system/core.py
def close(self) -> None:
    """Release resources."""
    if self.scheduler and hasattr(self.scheduler, "stop"):
        self.scheduler.stop()
    for resource in (
        self.scheduler_store,
        self.engine,
        self.gpu_monitor,
        self.telemetry_store,
        self.trace_store,
        self.memory_backend,
        self.session_store,
        self.channel_backend,
        self.workflow_engine,
        self.container_runner,
    ):
        if resource and hasattr(resource, "close"):
            resource.close()
    if self.agent_manager is not None:
        self.agent_manager.close()
    if self.agent_scheduler is not None:
        self.agent_scheduler.stop()
    self._close_mcp_clients()

QueryOrchestrator

QueryOrchestrator(system: OrchestratorDeps)
Source code in src/openjarvis/system/orchestrator.py
def __init__(self, system: OrchestratorDeps) -> None:
    self._system = system
Functions
ask
ask(query: str, *, context: bool = True, temperature: Optional[float] = None, max_tokens: Optional[int] = None, agent: Optional[str] = None, tools: Optional[List[str]] = None, system_prompt: Optional[str] = None, operator_id: Optional[str] = None, prior_messages: Optional[List[Message]] = None) -> Dict[str, Any]

Execute a query through the system and return a result dict.

Source code in src/openjarvis/system/orchestrator.py
def ask(
    self,
    query: str,
    *,
    context: bool = True,
    temperature: Optional[float] = None,
    max_tokens: Optional[int] = None,
    agent: Optional[str] = None,
    tools: Optional[List[str]] = None,
    system_prompt: Optional[str] = None,
    operator_id: Optional[str] = None,
    prior_messages: Optional[List[Message]] = None,
) -> Dict[str, Any]:
    """Execute a query through the system and return a result dict."""
    s = self._system
    if temperature is None:
        temperature = s.config.intelligence.temperature
    if max_tokens is None:
        max_tokens = s.config.intelligence.max_tokens

    messages = [Message(role=Role.USER, content=query)]

    if context and s.config.agent.context_from_memory:
        try:
            from openjarvis.memory import load_configured_facts
            from openjarvis.tools.storage.context import (
                ContextConfig,
                inject_context,
            )

            ctx_cfg = ContextConfig(
                top_k=s.config.memory.context_top_k,
                min_score=s.config.memory.context_min_score,
                max_context_tokens=s.config.memory.context_max_tokens,
            )
            facts = load_configured_facts(s.config)
            messages = inject_context(
                query,
                messages,
                s.memory_backend,
                config=ctx_cfg,
                facts=facts,
            )
        except Exception as exc:
            logger.warning("Failed to inject memory context: %s", exc)

    use_agent = agent or s.agent_name
    if not agent and use_agent != "none":
        detected = self._detect_agent_intent(query)
        if detected:
            use_agent = detected
    if use_agent and use_agent != "none":
        return self._run_agent(
            query,
            messages,
            use_agent,
            tools,
            temperature,
            max_tokens,
            system_prompt=system_prompt,
            operator_id=operator_id,
            prior_messages=prior_messages,
        )

    result = s.engine.generate(
        messages,
        model=s.model,
        temperature=temperature,
        max_tokens=max_tokens,
    )
    return {
        "content": result.get("content", ""),
        "usage": result.get("usage", {}),
        "model": s.model,
        "engine": s.engine_key,
    }

OrchestratorDeps

Bases: Protocol

Minimum surface of JarvisSystem that QueryOrchestrator depends on.

Tests can satisfy this with a lightweight class — no need to construct the full JarvisSystem dataclass or materialize every subsystem.