Files
openhuman/src/core/all.rs
T

938 lines
48 KiB
Rust
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
//! Registry and dispatch logic for all OpenHuman controllers.
//!
//! This module serves as the central hub for registering domain-specific
//! controllers (e.g., memory, skills, config) and providing a unified
//! interface for both the CLI and RPC layers to invoke them.
use std::future::Future;
use std::pin::Pin;
use std::sync::OnceLock;
use serde_json::{Map, Value};
use crate::core::ControllerSchema;
/// A pinned, boxed future returned by a controller handler.
pub type ControllerFuture = Pin<Box<dyn Future<Output = Result<Value, String>> + Send + 'static>>;
/// A function pointer type for controller handlers.
///
/// Handlers take a map of parameters and return a [`ControllerFuture`].
pub type ControllerHandler = fn(Map<String, Value>) -> ControllerFuture;
/// A function pointer type for domain-specific CLI handlers.
pub type CliHandler = fn(&[String]) -> anyhow::Result<()>;
/// A registered standalone CLI adapter for a domain.
#[derive(Clone)]
pub struct RegisteredCliAdapter {
pub namespace: &'static str,
pub handler: CliHandler,
}
/// A registered controller combining its schema and handler function.
#[derive(Clone)]
pub struct RegisteredController {
/// The schema defining the controller's identity and parameters.
pub schema: ControllerSchema,
/// The actual function that executes the controller's logic.
pub handler: ControllerHandler,
}
impl RegisteredController {
/// Returns the canonical RPC method name for this controller (e.g., `openhuman.memory_doc_put`).
pub fn rpc_method_name(&self) -> String {
rpc_method_name(&self.schema)
}
}
/// The global static registry of all controllers, initialized once on first access.
static REGISTRY: OnceLock<Vec<RegisteredController>> = OnceLock::new();
/// Internal-only controllers: registered for RPC dispatch but NOT in the agent-facing
/// schema catalog. These handlers are callable by trusted callers (e.g. the Tauri scanner)
/// but should not be advertised to agents via tool listings or schema discovery.
static INTERNAL_REGISTRY: OnceLock<Vec<RegisteredController>> = OnceLock::new();
/// The global static registry of standalone CLI adapters.
static CLI_ADAPTERS: OnceLock<Vec<RegisteredCliAdapter>> = OnceLock::new();
/// Returns a reference to the global controller registry.
///
/// This function initializes the registry if it hasn't been already,
/// performing validation to ensure no duplicates or missing handlers exist.
fn registry() -> &'static [RegisteredController] {
REGISTRY
.get_or_init(|| {
let registered = build_registered_controllers();
let declared = build_declared_controller_schemas();
validate_registry(&registered, &declared).unwrap_or_else(|err| {
panic!("invalid controller registry: {err}");
});
registered
})
.as_slice()
}
/// Returns a reference to the internal-only controller registry.
///
/// These controllers are callable over RPC but are NOT included in agent tool listings
/// or schema discovery endpoints.
fn internal_registry() -> &'static [RegisteredController] {
INTERNAL_REGISTRY
.get_or_init(build_internal_only_controllers)
.as_slice()
}
/// Returns a reference to the global CLI adapter registry.
fn cli_adapters() -> &'static [RegisteredCliAdapter] {
CLI_ADAPTERS.get_or_init(|| {
vec![RegisteredCliAdapter {
namespace: "voice",
handler: crate::openhuman::voice::cli::run_standalone_subcommand,
}]
})
}
/// Aggregates all controller implementations from across the codebase.
///
/// This function is responsible for collecting every domain-specific controller
/// registered in the system. It is used during the initialization of the
/// global [`REGISTRY`].
///
/// When adding a new domain/namespace, its `all_*_registered_controllers()`
/// function must be called here to make it available via RPC and CLI.
fn build_registered_controllers() -> Vec<RegisteredController> {
let mut controllers = Vec::new();
// Application information and capabilities
controllers.extend(crate::openhuman::about_app::all_about_app_registered_controllers());
// AgentBox marketplace adapter status
controllers.extend(crate::openhuman::agentbox::all_agentbox_registered_controllers());
// Core application shell state
controllers.extend(crate::openhuman::app_state::all_app_state_registered_controllers());
// Audio generation + podcast-style email delivery
controllers.extend(crate::openhuman::audio_toolkit::all_audio_toolkit_registered_controllers());
// Composio integration controllers
controllers.extend(crate::openhuman::composio::all_composio_registered_controllers());
// Scheduled job management
controllers.extend(crate::openhuman::cron::all_cron_registered_controllers());
// Proactive task ingestion from external tools (github/notion/linear/clickup)
controllers.extend(crate::openhuman::task_sources::all_task_sources_registered_controllers());
controllers.extend(crate::openhuman::dashboard::all_dashboard_registered_controllers());
// MCP client subsystem: Smithery registry browser, local server install/connect, tool dispatch
controllers.extend(crate::openhuman::mcp_registry::all_mcp_registry_registered_controllers());
// Webview APIs bridge — proxies connector calls (Gmail, …) through
// a WebSocket to the Tauri shell so curl reaches the live webview.
controllers.extend(crate::openhuman::webview_apis::all_webview_apis_registered_controllers());
// Agent definition and prompt inspection
controllers.extend(crate::openhuman::agent::all_agent_registered_controllers());
// Persistent agent profiles (flavours): name, soul, memory sources, skills, MCP, connectors.
controllers.extend(crate::openhuman::profiles::all_profiles_registered_controllers());
// User-facing agent registry: defaults, enablement, custom agents, tool policy.
controllers
.extend(crate::openhuman::agent_registry::all_agent_registry_registered_controllers());
// Local procedural operating experience for agent self-learning
controllers
.extend(crate::openhuman::agent_experience::all_agent_experience_registered_controllers());
// System and process health monitoring
controllers.extend(crate::openhuman::health::all_health_registered_controllers());
// Diagnostic tools
controllers.extend(crate::openhuman::doctor::all_doctor_registered_controllers());
// Secret storage and encryption
controllers.extend(crate::openhuman::encryption::all_encryption_registered_controllers());
// Keyring consent — user approval before local secret storage fallback
controllers
.extend(crate::openhuman::keyring_consent::all_keyring_consent_registered_controllers());
// Security policy metadata
controllers.extend(crate::openhuman::security::all_security_registered_controllers());
// Interactive approval workflow (#1339 — gate external-effect tool calls)
controllers.extend(crate::openhuman::approval::all_approval_registered_controllers());
// Agent-generated artifact storage, retrieval, and lifecycle management
controllers.extend(crate::openhuman::artifacts::all_artifacts_registered_controllers());
// Background heartbeat loop controls
controllers.extend(crate::openhuman::heartbeat::all_heartbeat_registered_controllers());
// Ad-hoc static directory HTTP hosting for local file sharing / previews
controllers.extend(crate::openhuman::http_host::all_http_host_registered_controllers());
// Token usage and billing cost tracking
controllers.extend(crate::openhuman::cost::all_cost_registered_controllers());
// x402 machine-payable API payment protocol
controllers.extend(crate::openhuman::x402::all_x402_registered_controllers());
// Inline autocomplete settings
controllers.extend(crate::openhuman::autocomplete::all_autocomplete_registered_controllers());
// External messaging channels (Web, Telegram, etc.)
controllers.extend(
crate::openhuman::channels::providers::web::all_web_channel_registered_controllers(),
);
controllers
.extend(crate::openhuman::channels::controllers::all_channels_registered_controllers());
// Persistent configuration management
controllers.extend(crate::openhuman::config::all_config_registered_controllers());
// Local sidecar reachability + backend Socket.IO state diagnostics (#1527)
controllers.extend(crate::openhuman::connectivity::all_connectivity_registered_controllers());
// User credentials and session management
controllers.extend(crate::openhuman::credentials::all_credentials_registered_controllers());
// Desktop service management
controllers.extend(crate::openhuman::service::all_service_registered_controllers());
// Data migration utilities
controllers.extend(crate::openhuman::migration::all_migration_registered_controllers());
// Saved council definitions for the desktop Model Council surface.
controllers
.extend(crate::openhuman::council_registry::all_council_registry_registered_controllers());
// Model Council: multi-model deliberation (parallel members + chair synthesis)
controllers.extend(crate::openhuman::model_council::all_model_council_registered_controllers());
// Background command monitors for agent-scoped event sources
controllers.extend(crate::openhuman::monitor::all_monitor_registered_controllers());
// Unified inference domain: text / vision / local runtime / cloud providers.
// (Formerly split across inference, local AI, and providers modules.)
controllers.extend(crate::openhuman::inference::all_inference_registered_controllers());
controllers.extend(crate::openhuman::inference::all_local_inference_registered_controllers());
// Embedding provider configuration and embed RPC.
controllers.extend(crate::openhuman::embeddings::all_embeddings_registered_controllers());
// People resolution and interaction scoring
controllers.extend(crate::openhuman::people::all_people_registered_controllers());
// Screen capture and UI analysis
controllers.extend(
crate::openhuman::screen_intelligence::all_screen_intelligence_registered_controllers(),
);
// Sandbox execution backends (Docker, local jail, policy, cleanup)
controllers.extend(crate::openhuman::sandbox::all_sandbox_registered_controllers());
// Backend Socket.IO bridge + related runtime plumbing
controllers.extend(crate::openhuman::socket::all_socket_registered_controllers());
// Managed Node.js runtime bridge (tool listing + dispatch)
controllers.extend(crate::openhuman::javascript::all_javascript_registered_controllers());
// Discovered SKILL.md skills and their bundled resources
controllers.extend(crate::openhuman::workflows::all_workflows_registered_controllers());
// Skill runtime: run/cancel/log skill executions and resolve Node/Python toolchains
controllers.extend(crate::openhuman::skill_runtime::all_skill_runtime_registered_controllers());
// Skill registry: browse, search, install from remote registries
controllers
.extend(crate::openhuman::skill_registry::all_skill_registry_registered_controllers());
// User workspace and file management
controllers.extend(crate::openhuman::workspace::all_workspace_registered_controllers());
// Workflow tool registry
controllers.extend(crate::openhuman::tools::all_tools_registered_controllers());
// Unified read-only registry across MCP stdio tools and controller-backed tools
controllers.extend(crate::openhuman::tool_registry::all_tool_registry_registered_controllers());
// Document and knowledge graph storage
controllers.extend(crate::openhuman::memory::all_memory_registered_controllers());
// Memory tree ingestion layer (#707 — canonicalised chunks with provenance)
controllers.extend(crate::openhuman::memory_tree::all_memory_tree_registered_controllers());
// Memory tree retrieval layer (#710 — LLM-callable read tools over the tree)
controllers.extend(crate::openhuman::memory_tree::all_retrieval_registered_controllers());
// Slack → memory-tree ingestion engine (per-message ingest, no bucketing)
controllers.extend(
crate::openhuman::composio::providers::slack::all_slack_memory_registered_controllers(),
);
// Per-connection memory sync status, controls, and progress (#1136)
controllers.extend(
crate::openhuman::memory_sync::sync_status::all_memory_sync_status_registered_controllers(),
);
// Memory sources — user-configured data connectors registry
controllers
.extend(crate::openhuman::memory_sources::all_memory_sources_registered_controllers());
// Memory diff — snapshot-based change tracking for memory sources
controllers.extend(crate::openhuman::memory_diff::all_memory_diff_registered_controllers());
// Link shortener for long tracking URLs — saves LLM tokens
controllers
.extend(crate::openhuman::redirect_links::all_redirect_links_registered_controllers());
// Referral and growth tracking
controllers.extend(crate::openhuman::referral::all_referral_registered_controllers());
// Billing and subscription management
controllers.extend(crate::openhuman::billing::all_billing_registered_controllers());
// Team and role management
controllers.extend(crate::openhuman::team::all_team_registered_controllers());
// E2E test support — `openhuman.test_reset` wipes sidecar state in-place.
// Gated behind the `e2e-test-support` cargo feature so shipped binaries
// never even register the destructive wipe RPC. Flipped on by the E2E
// build script (app/scripts/e2e-build.sh).
#[cfg(feature = "e2e-test-support")]
controllers.extend(crate::openhuman::test_support::all_test_support_registered_controllers());
// Local wallet metadata and onboarding status
controllers.extend(crate::openhuman::wallet::all_wallet_registered_controllers());
// High-level web3 surface (swaps / bridges / dapp calls) over the wallet
controllers.extend(crate::openhuman::web3::all_web3_registered_controllers());
// Local assistive surfaces over third-party provider apps
controllers.extend(
crate::openhuman::provider_surfaces::all_provider_surfaces_registered_controllers(),
);
// OS-level text input interactions
controllers.extend(crate::openhuman::text_input::all_text_input_registered_controllers());
// Voice transcription and synthesis
controllers.extend(crate::openhuman::voice::all_voice_registered_controllers());
// Background awareness and autonomous tasks
controllers.extend(crate::openhuman::subconscious::all_subconscious_registered_controllers());
// Webhook tunnel management
controllers.extend(crate::openhuman::webhooks::all_webhooks_registered_controllers());
// Core binary update management
controllers.extend(crate::openhuman::update::all_update_registered_controllers());
// Hierarchical knowledge summarization
controllers.extend(crate::openhuman::memory_tree::all_tree_summarizer_registered_controllers());
// Self-learning and user context enrichment
controllers.extend(crate::openhuman::learning::all_learning_registered_controllers());
// Conversation thread and message management
controllers.extend(crate::openhuman::threads::all_threads_registered_controllers());
// Per-thread todo list (agent task board CRUD over RPC)
controllers.extend(crate::openhuman::todos::all_todos_registered_controllers());
// Embedded webview native notifications
controllers.extend(
crate::openhuman::webview_notifications::all_webview_notifications_registered_controllers(),
);
// Integration notification ingest, triage, and per-provider settings
controllers.extend(crate::openhuman::notifications::all_notifications_registered_controllers());
// Google Meet call-join request validation (shell handles the webview)
controllers.extend(crate::openhuman::meet::all_meet_registered_controllers());
// Agent meetings — backend-delegated Meet bot via Socket.IO
controllers
.extend(crate::openhuman::agent_meetings::all_agent_meetings_registered_controllers());
// Live meet-agent loop: STT/LLM/TTS over the open call's audio.
controllers.extend(crate::openhuman::meet_agent::all_meet_agent_registered_controllers());
// Desktop companion — Clicky-style interaction loop.
controllers.extend(
crate::openhuman::desktop_companion::all_desktop_companion_registered_controllers(),
);
// Structured WhatsApp Web data — agent-facing read-only controllers (list/search).
// The write-path ingest controller is registered separately in build_internal_only_controllers.
controllers.extend(crate::openhuman::whatsapp_data::all_whatsapp_data_registered_controllers());
// Mobile device pairing and management
controllers.extend(crate::openhuman::devices::all_devices_registered_controllers());
// Durable agent session database — queryable index over transcripts, lineage, tool calls
controllers.extend(crate::openhuman::session_db::all_session_db_registered_controllers());
// Background agent command center — read-only grouped view over the run ledger
controllers
.extend(crate::openhuman::agent_orchestration::all_command_center_registered_controllers());
// Durable dynamic workflow runs — definitions + read surface over the run ledger
controllers
.extend(crate::openhuman::agent_orchestration::all_workflow_run_registered_controllers());
// Durable agent-team coordination — teams, members, dependency-aware task claiming, messaging
controllers
.extend(crate::openhuman::agent_orchestration::all_agent_team_registered_controllers());
// Git-worktree isolation manager — list / status / diff / remove worker worktrees (#3376)
controllers
.extend(crate::openhuman::agent_orchestration::all_worktree_registered_controllers());
controllers
}
/// Aggregates controllers that are registered for RPC routing but NOT exposed to agents.
///
/// These are write-path or internal-only handlers callable by trusted callers
/// (e.g. the Tauri scanner ingest path) that should not appear in agent tool listings.
fn build_internal_only_controllers() -> Vec<RegisteredController> {
let mut controllers = Vec::new();
// whatsapp_data ingest: scanner-side write path. Callable over RPC by the
// Tauri scanner but excluded from agent-facing schema discovery.
controllers.extend(crate::openhuman::whatsapp_data::all_whatsapp_data_internal_controllers());
// MCP write audit list: internal-only so the desktop UI/CLI can inspect
// local write history without exposing cross-client history as an MCP tool.
controllers.extend(crate::openhuman::mcp_audit::all_mcp_audit_internal_controllers());
controllers
}
/// Aggregates all controller schemas from across the codebase.
///
/// Similar to [`build_registered_controllers`], but only collects the metadata
/// (schema) for each controller. This is used for discovery and validation.
fn build_declared_controller_schemas() -> Vec<ControllerSchema> {
let mut schemas = Vec::new();
schemas.extend(crate::openhuman::about_app::all_about_app_controller_schemas());
schemas.extend(crate::openhuman::agentbox::all_agentbox_controller_schemas());
schemas.extend(crate::openhuman::app_state::all_app_state_controller_schemas());
schemas.extend(crate::openhuman::audio_toolkit::all_audio_toolkit_controller_schemas());
schemas.extend(crate::openhuman::composio::all_composio_controller_schemas());
schemas.extend(crate::openhuman::cron::all_cron_controller_schemas());
schemas.extend(crate::openhuman::task_sources::all_task_sources_controller_schemas());
schemas.extend(crate::openhuman::dashboard::all_dashboard_controller_schemas());
schemas.extend(crate::openhuman::mcp_registry::all_mcp_registry_controller_schemas());
schemas.extend(crate::openhuman::webview_apis::all_webview_apis_controller_schemas());
schemas.extend(crate::openhuman::agent::all_agent_controller_schemas());
schemas.extend(crate::openhuman::profiles::all_profiles_controller_schemas());
schemas.extend(crate::openhuman::agent_registry::all_agent_registry_controller_schemas());
schemas.extend(crate::openhuman::agent_experience::all_agent_experience_controller_schemas());
schemas.extend(crate::openhuman::health::all_health_controller_schemas());
schemas.extend(crate::openhuman::doctor::all_doctor_controller_schemas());
schemas.extend(crate::openhuman::encryption::all_encryption_controller_schemas());
schemas.extend(crate::openhuman::keyring_consent::all_keyring_consent_controller_schemas());
schemas.extend(crate::openhuman::security::all_security_controller_schemas());
schemas.extend(crate::openhuman::approval::all_approval_controller_schemas());
schemas.extend(crate::openhuman::artifacts::all_artifacts_controller_schemas());
schemas.extend(crate::openhuman::heartbeat::all_heartbeat_controller_schemas());
schemas.extend(crate::openhuman::http_host::all_http_host_controller_schemas());
schemas.extend(crate::openhuman::cost::all_cost_controller_schemas());
schemas.extend(crate::openhuman::x402::all_x402_controller_schemas());
schemas.extend(crate::openhuman::autocomplete::all_autocomplete_controller_schemas());
schemas
.extend(crate::openhuman::channels::providers::web::all_web_channel_controller_schemas());
schemas.extend(crate::openhuman::channels::controllers::all_channels_controller_schemas());
schemas.extend(crate::openhuman::config::all_config_controller_schemas());
schemas.extend(crate::openhuman::connectivity::all_connectivity_controller_schemas());
schemas.extend(crate::openhuman::credentials::all_credentials_controller_schemas());
schemas.extend(crate::openhuman::service::all_service_controller_schemas());
schemas.extend(crate::openhuman::migration::all_migration_controller_schemas());
schemas.extend(crate::openhuman::council_registry::all_council_registry_controller_schemas());
schemas.extend(crate::openhuman::model_council::all_model_council_controller_schemas());
schemas.extend(crate::openhuman::monitor::all_monitor_controller_schemas());
schemas.extend(crate::openhuman::inference::all_inference_controller_schemas());
schemas.extend(crate::openhuman::inference::all_local_inference_controller_schemas());
schemas.extend(crate::openhuman::embeddings::all_embeddings_controller_schemas());
schemas.extend(crate::openhuman::people::all_people_controller_schemas());
schemas.extend(
crate::openhuman::screen_intelligence::all_screen_intelligence_controller_schemas(),
);
schemas.extend(crate::openhuman::sandbox::all_sandbox_controller_schemas());
schemas.extend(crate::openhuman::socket::all_socket_controller_schemas());
schemas.extend(crate::openhuman::javascript::all_javascript_controller_schemas());
schemas.extend(crate::openhuman::workflows::all_workflows_controller_schemas());
schemas.extend(crate::openhuman::skill_runtime::all_skill_runtime_controller_schemas());
schemas.extend(crate::openhuman::skill_registry::all_skill_registry_controller_schemas());
schemas.extend(crate::openhuman::workspace::all_workspace_controller_schemas());
schemas.extend(crate::openhuman::tools::all_tools_controller_schemas());
schemas.extend(crate::openhuman::tool_registry::all_tool_registry_controller_schemas());
schemas.extend(crate::openhuman::memory::all_memory_controller_schemas());
schemas.extend(crate::openhuman::memory_tree::all_memory_tree_controller_schemas());
schemas.extend(crate::openhuman::memory_tree::all_retrieval_controller_schemas());
schemas.extend(
crate::openhuman::composio::providers::slack::all_slack_memory_controller_schemas(),
);
schemas.extend(
crate::openhuman::memory_sync::sync_status::all_memory_sync_status_controller_schemas(),
);
schemas.extend(crate::openhuman::memory_sources::all_memory_sources_controller_schemas());
schemas.extend(crate::openhuman::memory_diff::all_memory_diff_controller_schemas());
schemas.extend(crate::openhuman::redirect_links::all_redirect_links_controller_schemas());
schemas.extend(crate::openhuman::referral::all_referral_controller_schemas());
schemas.extend(crate::openhuman::billing::all_billing_controller_schemas());
schemas.extend(crate::openhuman::team::all_team_controller_schemas());
#[cfg(feature = "e2e-test-support")]
schemas.extend(crate::openhuman::test_support::all_test_support_controller_schemas());
schemas.extend(crate::openhuman::wallet::all_wallet_controller_schemas());
schemas.extend(crate::openhuman::web3::all_web3_controller_schemas());
schemas.extend(crate::openhuman::provider_surfaces::all_provider_surfaces_controller_schemas());
schemas.extend(crate::openhuman::text_input::all_text_input_controller_schemas());
schemas.extend(crate::openhuman::voice::all_voice_controller_schemas());
schemas.extend(crate::openhuman::subconscious::all_subconscious_controller_schemas());
schemas.extend(crate::openhuman::webhooks::all_webhooks_controller_schemas());
schemas.extend(crate::openhuman::update::all_update_controller_schemas());
schemas.extend(crate::openhuman::memory_tree::all_tree_summarizer_controller_schemas());
schemas.extend(crate::openhuman::learning::all_learning_controller_schemas());
// Conversation thread and message management
schemas.extend(crate::openhuman::threads::all_threads_controller_schemas());
// Per-thread todo list (agent task board CRUD over RPC)
schemas.extend(crate::openhuman::todos::all_todos_controller_schemas());
// Embedded webview native notifications
schemas.extend(
crate::openhuman::webview_notifications::all_webview_notifications_controller_schemas(),
);
// Integration notification ingest, triage, and per-provider settings
schemas.extend(crate::openhuman::notifications::all_notifications_controller_schemas());
// Google Meet call-join request validation
schemas.extend(crate::openhuman::meet::all_meet_controller_schemas());
// Agent meetings — backend-delegated Meet bot via Socket.IO
schemas.extend(crate::openhuman::agent_meetings::all_agent_meetings_controller_schemas());
// Live meet-agent listening + speaking loop
schemas.extend(crate::openhuman::meet_agent::all_meet_agent_controller_schemas());
// Desktop companion — Clicky-style interaction loop.
schemas.extend(crate::openhuman::desktop_companion::all_desktop_companion_controller_schemas());
// Structured WhatsApp Web data — local SQLite store, agent-queryable
schemas.extend(crate::openhuman::whatsapp_data::all_whatsapp_data_controller_schemas());
// Mobile device pairing and management
schemas.extend(crate::openhuman::devices::all_devices_controller_schemas());
// Durable agent session database
schemas.extend(crate::openhuman::session_db::all_session_db_controller_schemas());
// Background agent command center
schemas.extend(crate::openhuman::agent_orchestration::all_command_center_controller_schemas());
// Durable dynamic workflow runs
schemas.extend(crate::openhuman::agent_orchestration::all_workflow_run_controller_schemas());
// Durable agent-team coordination
schemas.extend(crate::openhuman::agent_orchestration::all_agent_team_controller_schemas());
// Git-worktree isolation manager (#3376)
schemas.extend(crate::openhuman::agent_orchestration::all_worktree_controller_schemas());
schemas
}
/// Returns a vector of all currently registered controllers.
pub fn all_registered_controllers() -> Vec<RegisteredController> {
registry().to_vec()
}
/// Returns a vector of all currently declared controller schemas.
pub fn all_controller_schemas() -> Vec<ControllerSchema> {
let _ = registry();
build_declared_controller_schemas()
}
/// Generates a standardized RPC method name from a controller schema.
pub fn rpc_method_name(schema: &ControllerSchema) -> String {
format!("openhuman.{}_{}", schema.namespace, schema.function)
}
/// Returns a human-readable description for a given namespace.
///
/// This is used for CLI help output.
pub fn namespace_description(namespace: &str) -> Option<&'static str> {
match namespace {
"about_app" => Some("Catalog the app's user-facing capabilities and where to find them."),
"agentbox" => Some("AgentBox marketplace adapter status — mode flag and GMI MaaS provider wiring."),
"ai" => Some("Agent-generated artifact storage, retrieval, and lifecycle management."),
"app_state" => Some("Expose core-owned app shell state for frontend polling."),
"auth" => Some("Manage app session and provider credentials."),
"agent_experience" => Some("Local procedural experience capture and retrieval for agents."),
"autocomplete" => Some("Inline autocomplete engine controls and style settings."),
"channels" => Some("Channel definitions, connections, and lifecycle management."),
"composio" => Some(
"Composio OAuth integrations proxied via the backend — toolkits, connections, tools, and actions."
),
"config" => Some("Read and update persisted runtime configuration."),
"connectivity" => Some(
"Connectivity diagnostics for the local sidecar, listening port, and backend Socket.IO state.",
),
"cron" => Some("Manage scheduled jobs and run history."),
"dashboard" => Some(
"Operator-facing dashboard aggregations: per-model health comparison rows.",
),
"mcp_clients" => Some(
"Browse the Smithery.ai MCP registry, install MCP servers locally, manage their stdio connections, and expose their tools to the agent.",
),
"mcp_setup" => Some(
"MCP setup agent surface: search registries, request secrets out-of-band (opaque refs, no raw values in agent context), test, and install + connect.",
),
"decrypt" => Some("Decrypt secure values managed by secret storage."),
"doctor" => Some("Run diagnostics for workspace and runtime health."),
"encrypt" => Some("Encrypt secure values managed by secret storage."),
"health" => Some("Process and component health snapshots."),
"inference" => Some("Connect to configured text, vision, and embedding inference runtimes."),
"migrate" => Some("Data migration utilities."),
"javascript" => Some("First-class JavaScript runtime bridge for listing and dispatching tools."),
"monitor" => Some("Start, inspect, read, and stop bounded background command monitors."),
"screen_intelligence" => Some("Screen capture, permissions, and accessibility automation."),
"security" => Some("Security policy and autonomy guardrail metadata."),
"service" => Some("Desktop service lifecycle management."),
"skill_registry" => Some("Browse, search, install, and uninstall skills from remote registries (OpenHuman, Hermes, OpenClaw)."),
"skill_runtime" => Some("Run installed skills, inspect run logs, and resolve Node/Python skill runtimes."),
"workflows" => Some("Discovered workflows (WORKFLOW.md/SKILL.md bundles) and their resources."),
"socket" => Some("Backend Socket.IO bridge controls."),
"memory" => Some("Document storage, vector search, key-value store, and knowledge graph."),
"memory_tree" => Some(
"Canonical chunk ingestion, provenance capture, and chunk retrieval for source-grounded memory.",
),
"memory_sync" => Some(
"Per-connection memory sync status, user enable toggle, and live progress for the desktop UI.",
),
"memory_sources" => Some(
"User-configured data connectors (Composio, folders, GitHub repos, RSS, web pages) that feed memory.",
),
"memory_diff" => Some(
"Snapshot-based change tracking for memory sources — capture state, compute diffs, and surface changes to agents.",
),
"redirect_links" => Some(
"Shorten long tracking URLs to `openhuman://link/<id>` placeholders (SQLite-backed) to save tokens in prompts, with round-trip rewrite helpers.",
),
"referral" => Some("Referral codes, stats, and apply flows via the hosted backend API."),
"run_ledger" => Some(
"Durable agent and workflow run state, child lineage, events, telemetry, and checkpoint references.",
),
"agent_work" => Some(
"Background agent command center — recent agent runs grouped by status (needs-input, working, completed, failed, stopped).",
),
"workflow_run" => Some(
"Durable dynamic workflow runs — declarative multi-agent definitions and the read surface over persisted runs.",
),
"agent_team" => Some(
"Durable agent-team coordination: teams, members, dependency-aware task claiming, and teammate messaging.",
),
"billing" => Some("Subscription plan, payment links, and credit top-up via the backend."),
"team" => Some("Team member management, invites, and role changes via the backend."),
"tool_registry" => Some(
"Read-only discovery for MCP stdio tools and controller-backed tools, including routes, schemas, version, allowed agents, and health.",
),
"test" => Some(
"E2E test support — wipe sidecar state in-place between specs.",
),
"wallet" => Some("Local wallet onboarding status and derived multi-chain account metadata."),
"web3_swap" => Some("Single-chain crypto swaps via deBridge, built on the local wallet."),
"web3_bridge" => Some("Cross-chain crypto bridges via deBridge DLN, built on the local wallet."),
"web3_dapp" => Some("Generic EVM dapp contract calls signed by the local wallet."),
"provider_surfaces" => Some(
"Local-first assistive surfaces for provider events, respond queues, and drafts.",
),
"voice" => Some("Speech-to-text and text-to-speech using local models."),
"subconscious" => Some("Periodic local-model background awareness loop."),
"text_input" => Some("Read, insert, and preview text in the OS-focused input field."),
"webhooks" => {
Some("Webhook tunnel registrations and captured request/response debug logs.")
}
"webview_apis" => Some(
"Typed connector APIs (Gmail, …) proxied over a loopback WebSocket to the Tauri shell so core-side JSON-RPC reaches live-webview CDP operations.",
),
"update" => {
Some("Self-update: check GitHub Releases for newer core binary and stage updates.")
}
"tree_summarizer" => {
Some("Hierarchical time-based summarization tree for background knowledge compression.")
}
"learning" => Some(
"User context enrichment — LinkedIn profile scraping and onboarding intelligence.",
),
"people" => {
Some("Contact resolution and recency × frequency × reciprocity × depth scoring.")
},
"notification" => Some(
"Integration notification ingest, triage scoring, listing, read-state, \
and per-provider routing settings.",
),
"meet" => Some(
"Validate Google Meet call-join requests and mint a request_id; the desktop \
shell opens the embedded CEF webview that joins the call as an anonymous guest.",
),
"meet_agent" => Some(
"Live agent loop for an open Google Meet call: shell streams inbound PCM, \
core runs VAD-segmented STT → LLM → TTS, shell pulls synthesized PCM back.",
),
"agent_meetings" => Some(
"Backend-delegated meeting bot (Google Meet, Zoom, Teams, Webex) via Socket.IO — join, leave, and harness response.",
),
"devices" => Some(
"Paired mobile device management — pairing channel creation, listing, and revocation.",
),
"whatsapp_data" => Some(
"Structured WhatsApp conversation and message store — list chats, read messages, and search across WhatsApp Web data.",
),
"companion" => Some(
"Desktop companion — Clicky-style hotkey-driven interaction loop with STT, LLM, TTS, and visual pointing.",
),
_ => None,
}
}
/// Returns the CLI handler for a given namespace, if one is registered.
pub fn cli_handler_for_namespace(namespace: &str) -> Option<CliHandler> {
cli_adapters()
.iter()
.find(|a| a.namespace == namespace)
.map(|a| a.handler)
}
/// Looks up an RPC method name based on namespace and function.
pub fn rpc_method_from_parts(namespace: &str, function: &str) -> Option<String> {
registry()
.iter()
.find(|r| r.schema.namespace == namespace && r.schema.function == function)
.map(|r| r.rpc_method_name())
}
/// Retrieves the schema for a specific RPC method.
///
/// Checks both the agent-facing registry and the internal registry so that
/// parameter validation still applies to internal-only methods (e.g. ingest).
pub fn schema_for_rpc_method(method: &str) -> Option<ControllerSchema> {
registry()
.iter()
.chain(internal_registry().iter())
.find(|r| r.rpc_method_name() == method)
.map(|r| r.schema.clone())
}
/// Validates that the provided parameters match the requirements of the controller schema.
///
/// # Errors
///
/// Returns an error message if required parameters are missing or if unknown parameters are provided.
pub fn validate_params(
schema: &ControllerSchema,
params: &Map<String, Value>,
) -> Result<(), String> {
for input in &schema.inputs {
if input.required && !params.contains_key(input.name) {
return Err(format!(
"missing required param '{}': {}",
input.name, input.comment
));
}
}
for key in params.keys() {
if !schema.inputs.iter().any(|f| f.name == key) {
return Err(format!(
"unknown param '{}' for {}.{}",
key, schema.namespace, schema.function
));
}
}
// Type-check each present param against its declared `TypeSchema`, so every
// controller gets uniform validation before dispatch rather than relying on
// each handler's `serde_json::from_value`. Absent (optional) params are
// already handled by the required-presence check above.
for input in &schema.inputs {
if let Some(value) = params.get(input.name) {
check_type(value, &input.ty).map_err(|expected| {
format!(
"invalid type for param '{}' in {}.{}: expected {}, got {}",
input.name,
schema.namespace,
schema.function,
expected,
json_type_name(value),
)
})?;
}
}
Ok(())
}
/// A short, human-readable name for the JSON kind of `value`, used in
/// `validate_params` type-mismatch errors.
fn json_type_name(value: &Value) -> &'static str {
match value {
Value::Null => "null",
Value::Bool(_) => "bool",
Value::Number(_) => "number",
Value::String(_) => "string",
Value::Array(_) => "array",
Value::Object(_) => "object",
}
}
/// Validate a JSON `value` against a declared [`TypeSchema`].
///
/// Returns `Ok(())` on a match, or `Err(expected)` where `expected` is a short
/// description of the type that was required. Unknown/opaque shapes
/// (`Json`, `Bytes`, `Ref`) accept any value — they are validated by the
/// handler's typed deserialization.
fn check_type(value: &Value, ty: &crate::core::TypeSchema) -> Result<(), &'static str> {
use crate::core::TypeSchema;
// JSON-RPC semantics (preserved from the prior presence-only check):
// an explicit `null` satisfies validation for any declared type. Required
// fields are checked for *presence*, not value; stronger contracts are
// enforced by the handler's typed deserialization.
if value.is_null() {
return Ok(());
}
match ty {
// Opaque / handler-validated shapes accept any JSON value.
//
// Structured types (`Object`/`Map`/`Ref`) are deliberately lenient: a
// struct field may have a custom `Deserialize` impl that accepts more
// than one JSON shape (e.g. `agent_registry.update`'s `subagents`
// accepts both `{ "allowlist": [...] }` and a legacy bare array), and
// the declared schema can only describe one of them. Strictly checking
// the JSON kind here would reject inputs the handler's `serde_json`
// deserialization accepts. We therefore keep pre-dispatch validation to
// scalar leaf types (where a type confusion would otherwise reach an
// `as_str()/as_u64()`-style accessor) and defer object/map shape
// validation to the handler.
TypeSchema::Json
| TypeSchema::Bytes
| TypeSchema::Ref(_)
| TypeSchema::Object { .. }
| TypeSchema::Map(_) => Ok(()),
TypeSchema::Bool => value.is_boolean().then_some(()).ok_or("bool"),
TypeSchema::String => value.is_string().then_some(()).ok_or("string"),
TypeSchema::I64 => value.is_i64().then_some(()).ok_or("integer"),
TypeSchema::U64 => value.is_u64().then_some(()).ok_or("unsigned integer"),
TypeSchema::F64 => {
// Accept any JSON number (ints are valid floats).
value.is_number().then_some(()).ok_or("number")
}
// `Option<T>` accepts null or a value matching the inner type.
TypeSchema::Option(inner) => {
if value.is_null() {
Ok(())
} else {
check_type(value, inner)
}
}
TypeSchema::Array(inner) => match value.as_array() {
Some(items) => {
for item in items {
check_type(item, inner)?;
}
Ok(())
}
None => Err("array"),
},
TypeSchema::Enum { variants } => match value.as_str() {
Some(s) if variants.contains(&s) => Ok(()),
Some(_) => Err("one of the allowed enum variants"),
None => Err("string"),
},
}
}
/// Attempts to invoke a registered RPC method by name.
///
/// Checks both the agent-facing controller registry and the internal-only registry,
/// so scanner-side write paths (e.g. `openhuman.whatsapp_data_ingest`) are routable
/// even though they are not included in agent tool listings.
///
/// Returns `None` if the method is not found in either registry.
pub async fn try_invoke_registered_rpc(
method: &str,
params: Map<String, Value>,
) -> Option<Result<Value, String>> {
for controller in registry() {
if controller.rpc_method_name() == method {
return Some((controller.handler)(params).await);
}
}
for controller in internal_registry() {
if controller.rpc_method_name() == method {
return Some((controller.handler)(params).await);
}
}
None
}
/// Validates the consistency of the controller registry.
///
/// Ensures that:
/// - There are no duplicate controllers or RPC methods.
/// - Every declared schema has a registered handler.
/// - Every registered handler has a declared schema.
/// - Namespaces and functions are not empty.
/// - Required input names are unique within a controller.
fn validate_registry(
registered: &[RegisteredController],
declared: &[ControllerSchema],
) -> Result<(), String> {
use std::collections::{BTreeMap, BTreeSet};
let mut errors: Vec<String> = Vec::new();
let mut declared_keys = BTreeSet::new();
let mut declared_rpc_methods = BTreeSet::new();
let mut registered_keys = BTreeSet::new();
let mut registered_rpc_methods = BTreeSet::new();
for schema in declared {
let key = format!("{}.{}", schema.namespace, schema.function);
if !declared_keys.insert(key.clone()) {
errors.push(format!("duplicate declared controller `{key}`"));
}
let rpc_method = rpc_method_name(schema);
if !declared_rpc_methods.insert(rpc_method.clone()) {
errors.push(format!("duplicate declared rpc method `{rpc_method}`"));
}
if schema.namespace.trim().is_empty() {
errors.push(format!(
"invalid declared controller `{key}`: namespace must not be empty"
));
}
if schema.function.trim().is_empty() {
errors.push(format!(
"invalid declared controller `{key}`: function must not be empty"
));
}
let mut required_inputs = BTreeSet::new();
let mut required_dupes: BTreeMap<String, usize> = BTreeMap::new();
for input in schema.inputs.iter().filter(|input| input.required) {
if !required_inputs.insert(input.name.to_string()) {
*required_dupes.entry(input.name.to_string()).or_default() += 1;
}
}
for (name, _) in required_dupes {
errors.push(format!(
"duplicate required input `{name}` in `{}`",
schema.method_name()
));
}
}
for controller in registered {
let key = format!(
"{}.{}",
controller.schema.namespace, controller.schema.function
);
if !registered_keys.insert(key.clone()) {
errors.push(format!("duplicate registered controller `{key}`"));
}
let rpc_method = controller.rpc_method_name();
if !registered_rpc_methods.insert(rpc_method.clone()) {
errors.push(format!("duplicate registered rpc method `{rpc_method}`"));
}
}
for key in declared_keys.difference(&registered_keys) {
errors.push(format!(
"declared controller `{key}` has no registered handler"
));
}
for key in registered_keys.difference(&declared_keys) {
errors.push(format!(
"registered controller `{key}` has no declared schema"
));
}
if errors.is_empty() {
Ok(())
} else {
Err(errors.join("; "))
}
}
#[derive(Debug, Clone)]
pub struct HttpMethodSchemaDefinition {
pub method: String,
pub namespace: &'static str,
pub function: &'static str,
pub description: &'static str,
pub inputs: Vec<crate::core::FieldSchema>,
pub outputs: Vec<crate::core::FieldSchema>,
}
pub fn all_http_method_schemas() -> Vec<HttpMethodSchemaDefinition> {
let mut methods = vec![
HttpMethodSchemaDefinition {
method: "core.ping".to_string(),
namespace: "core",
function: "ping",
description: "Liveness probe for the core JSON-RPC server.",
inputs: vec![],
outputs: vec![crate::core::FieldSchema {
name: "ok",
ty: crate::core::TypeSchema::Bool,
comment: "Always true when the server is reachable.",
required: true,
}],
},
HttpMethodSchemaDefinition {
method: "core.version".to_string(),
namespace: "core",
function: "version",
description: "Returns the core binary version.",
inputs: vec![],
outputs: vec![crate::core::FieldSchema {
name: "version",
ty: crate::core::TypeSchema::String,
comment: "Semantic version string for the running core binary.",
required: true,
}],
},
];
methods.extend(
all_controller_schemas()
.into_iter()
.map(|schema| HttpMethodSchemaDefinition {
method: rpc_method_name(&schema),
namespace: schema.namespace,
function: schema.function,
description: schema.description,
inputs: schema.inputs,
outputs: schema.outputs,
}),
);
methods
}
#[cfg(test)]
#[path = "all_tests.rs"]
mod tests;