pub struct ServerState {Show 63 fields
pub journal_dir: PathBuf,
pub peer_audit_journal: PathBuf,
pub sessions: Mutex<HashMap<String, Arc<ClientSession>>>,
pub inference: OnceLock<Arc<InferenceEngine>>,
pub host: Arc<HostState>,
pub shared_memgine: Option<Arc<Mutex<MemgineEngine>>>,
pub identity_store: IdentityStore,
pub trajectory_store: Arc<TrajectoryStore>,
pub selfheal: Arc<SelfhealService>,
pub heal: Arc<HealService>,
pub voice_sessions: Arc<VoiceSessionRegistry>,
pub meetings: Arc<MeetingRegistry>,
pub a2ui: A2uiSurfaceStore,
pub ui_agent: Arc<UIImprovementAgent>,
pub ui_agent_oscillation: Arc<OscillationDetector>,
pub ui_agent_budget: Arc<IterationBudget>,
pub admission: Arc<InferenceAdmission>,
pub a2ui_route_auth: Mutex<HashMap<String, A2aRouteAuth>>,
pub supervisor: OnceLock<Arc<Supervisor>>,
pub declagents: OnceLock<Arc<DeclRegistry>>,
pub routing: OnceLock<Arc<RoutingStore>>,
pub sync: Mutex<Option<Arc<Mutex<SyncSubsystem>>>>,
pub org_sync: Mutex<Option<OrgSyncHolder>>,
pub observer_manifest_path: OnceLock<PathBuf>,
pub a2a_dispatcher: OnceLock<Arc<A2aDispatcher>>,
pub a2ui_subscribers: Mutex<HashMap<String, Arc<WsChannel>>>,
pub mcp_executor: Arc<McpToolExecutor>,
pub connectors: OnceLock<Arc<ConnectorManager>>,
pub channel_supervisor: OnceLock<Arc<ChannelSupervisor>>,
pub durable_tasks: Mutex<JoinSet<String>>,
pub auth_completion_owner_id: String,
pub auth_token: OnceLock<String>,
pub host_token: OnceLock<String>,
pub mobile_runtime_url: OnceLock<String>,
pub mobile_registration_url: OnceLock<String>,
pub parslee_session: OnceLock<ParsleeSession>,
pub attached_agents: Mutex<HashMap<String, String>>,
pub agent_memgines: Mutex<HashMap<String, Arc<Mutex<MemgineEngine>>>>,
pub namespace_memgines: Mutex<HashMap<String, Arc<Mutex<MemgineEngine>>>>,
pub coder_sessions: Mutex<CoderSessionMap>,
pub coder_disk_gc_base: Instant,
pub coder_disk_gc_at: AtomicU64,
pub coder_subscribers: Mutex<HashMap<(String, String), Arc<WsChannel>>>,
pub coder_watchers: Mutex<HashMap<String, (u64, Arc<WsChannel>)>>,
pub coder_discussions: Mutex<DiscussionMap>,
pub chat_sessions: Mutex<HashMap<String, ChatSession>>,
pub peer_guards: Mutex<HashMap<String, DeliveryGuard>>,
pub peer_standing: Mutex<HashMap<String, PeerStanding>>,
pub held_peer_messages: Mutex<VecDeque<HeldPeerMessage>>,
pub lan_discovery: Mutex<Option<LanDirectory>>,
pub peer_identity: Mutex<Option<Arc<PeerIdentity>>>,
pub peer_trust: PeerTrust,
pub chat_collectors: Mutex<HashMap<String, ChatCollector>>,
pub chat_goals: Mutex<HashMap<String, ChatGoalState>>,
pub runs: Mutex<HashMap<String, RunMeta>>,
pub run_subscribers: Mutex<HashMap<(String, String), RunTraceSubscriber>>,
pub browser_views: Arc<BrowserViewRegistry>,
pub run_store: RunStore,
pub mcp_url: OnceLock<String>,
pub approval_gate: ApprovalGate,
pub supervision: Arc<SupervisionRegistry>,
pub approval_ledger: Arc<RwLock<ApprovalLedger>>,
pub harness_measurer: RwLock<Option<Arc<dyn HarnessMeasurer>>>,
/* private fields */
}Expand description
Global server state shared across all connections.
Fields§
§journal_dir: PathBuf§peer_audit_journal: PathBufResolved once when the state is constructed. Peer-message writes must
use this path and never consult process-global CAR_HOME themselves.
sessions: Mutex<HashMap<String, Arc<ClientSession>>>§inference: OnceLock<Arc<InferenceEngine>>§host: Arc<HostState>When Some, create_session clones this handle into every new
ClientSession.memgine — embedders that want a single shared
memgine across all WS sessions set this. Standalone car-server
leaves it None, which gives each session its own engine
(preserving today’s behavior).
identity_store: IdentityStoreSource of truth for the assistant identity and user profile. Kept on the state so RPC tests/embedders can isolate it from the process environment.
trajectory_store: Arc<TrajectoryStore>Daemon-wide trajectory store, cloned into every session’s
Runtime so execution traces persist to one place.
Shared rather than per-session on purpose. A trajectory is evidence about how a tool behaves, and that evidence does not belong to the connection that happened to produce it — per-session stores would scatter the history across connections and reset it on every reconnect, leaving the derived success rates built from a handful of samples. This is the same reasoning that makes the approval ledger daemon-wide.
Until now nothing called Runtime::with_trajectory_store, so
persist_trajectory returned early on every execution and the store
was dead code: ToolFeedback::from_trajectories existed with no data
to read. Wiring it here is what makes the feedback loop real.
selfheal: Arc<SelfhealService>Daemon-wide self-healing detector/fixer + CAR_HOME ledger. Shared
across sessions so selfheal.run and the cadence cannot overlap and
every caller reads the same dismissals/detections.
heal: Arc<HealService>The self-healing repair loop: reads a configured tracker, runs a coder session, gates the result on a multi-model panel, opens a pull request.
Distinct from Self::selfheal, which detects and never writes. Shared
across sessions for the same reason: heal.run and the cadence must not
overlap, and both must read the same claim ledger.
voice_sessions: Arc<VoiceSessionRegistry>Process-wide voice session registry. Each
voice.transcribe_stream.start call registers its own per-client
WsVoiceEventSink so events route back to the originating WS
connection only.
meetings: Arc<MeetingRegistry>Process-wide meeting registry. Meeting ids are global; each
meeting binds to the originating client’s WS for upstream
events but persists transcripts to the resolved
.car/meetings/<id>/ regardless of which client started it.
a2ui: A2uiSurfaceStoreProcess-wide A2UI surface store. Agent-produced surfaces are visible to every host UI subscriber, independent of the WebSocket session that applied the update.
ui_agent: Arc<UIImprovementAgent>In-process UI-improvement agent. Invoked from
handle_a2ui_render_report with each inbound report; returned
Decision::Patch envelopes are applied via the standard
apply_a2ui_envelope path so all subscribers see the patch.
Arc so the agent’s interior DashMap state survives across
handler calls even when ServerState is cheap-cloned.
ui_agent_oscillation: Arc<OscillationDetector>Per-surface oscillation detector for the UI-improvement
loop. Sits between the agent’s Decision::Patch and the
apply path so A→B→A patch cycles get cooled down without
the agent itself having to track history. neo’s review:
“controllers use workqueue backoff; reconcilers stay
stateless.”
ui_agent_budget: Arc<IterationBudget>Per-surface iteration budget. Backstop against runaway
loops the oscillation detector misses — caps total agent-
driven patches per surface at DEFAULT_MAX_ITERATIONS.
admission: Arc<InferenceAdmission>Process-wide concurrency gate for inference RPC handlers. Sized
from host RAM at startup, overridable via
crate::admission::ENV_MAX_CONCURRENT. Without this, N
concurrent users multiply KV-cache and activation memory and
take the host out (#114-adjacent: filed alongside the daemon
always-on rework). The semaphore lives on ServerState so it
is shared across every WebSocket session in the same process.
a2ui_route_auth: Mutex<HashMap<String, A2aRouteAuth>>Server-side A2A continuation auth keyed by A2UI surface id.
Kept out of A2uiSurface.owner so host renderers never see
bearer/API-key material.
supervisor: OnceLock<Arc<Supervisor>>Lifecycle-managed agents — declarative manifest at
~/.car/agents.json driving spawn/restart/stop. Closes
Parslee-ai/car-releases#27. Lazy-initialized so embedders that
don’t want process supervision don’t pay the disk-touch cost
at server start.
declagents: OnceLock<Arc<DeclRegistry>>Declarative (in-daemon) agent registry — a parallel store to the
supervisor for agents the coder→agent loop builds. Lazy-initialized on
first use (~/.car/declagents.json).
routing: OnceLock<Arc<RoutingStore>>Routing learning state — per-agent success stats + agent→agent edge
weights for capability-similarity routing (~/.car/routing.json).
Lazy-initialized on first use, sibling to Self::declagents.
sync: Mutex<Option<Arc<Mutex<SyncSubsystem>>>>Multi-device sync + execution-lease subsystem (B6). Lazily opened on
first sync.*/lease.* contact, rooted at <journal_dir>/sync/ — a
daemon-held car_sync::SyncSession over an FsRelay (the single-user
two-device loopback works out of the box) plus an in-process
linearizable car_sync::InMemoryLeaseCoordinator. One subsystem per
daemon = one device in the sync fleet. The std::sync::Mutex<Option>
serializes the fallible first-open (the oplog journal holds an exclusive
advisory lock, so a double-open would fail) without holding the lock
across the subsystem’s async work — the subsystem itself lives behind a
tokio::sync::Mutex inside the Arc. A distributed relay + a
cross-daemon linearizable LeaseCoordinator backend is the documented
B6 follow-up; the traits are its contract.
org_sync: Mutex<Option<OrgSyncHolder>>The SHARED org-scope delivery subsystem — a SECOND device oplog whose relay
scope is org:{orgId}, so Scope::Shared{org} ops converge across an org’s
members. None unless org-scope is opted in (the same gated
PARSLEE_SYNC_ORG_SCOPE branch that resolves the org key builds it);
populated by open_sync_subsystem alongside the personal subsystem. Only
the Scope::Shared write sites + pump route to it — the personal path
(sync_subsystem() + all readers) is untouched. Slice 8a.
observer_manifest_path: OnceLock<PathBuf>Manifest path this daemon is observing but does NOT own.
Set by car-server when boot-time supervisor construction
fails with car_registry::supervisor::SupervisorError::AlreadyRunning
— another car-server process on the host holds the exclusive
lock on this manifest. In that state, supervisor() returns a
clear “observe-only” error so mutation handlers refuse
(preventing the duplicate-spawn bug from
Parslee-ai/car-releases#44), while read-only handlers
(agents.list, agents.health) fall back to
car_registry::supervisor::Supervisor::list_from_manifest /
car_registry::supervisor::Supervisor::health_from_manifest
so operators can still inspect what the primary daemon is
supervising.
a2a_dispatcher: OnceLock<Arc<A2aDispatcher>>In-core A2A dispatcher — embedders that consume car-server-core
get A2A reachability “for free” without standing up a separate
HTTP listener. Closes Parslee-ai/car-releases#28. Lazy-init so
the embedder can override the runtime / task store / agent card
via ServerStateConfig::with_a2a_runtime etc. before the
first dispatch.
a2ui_subscribers: Mutex<HashMap<String, Arc<WsChannel>>>WS clients subscribed to A2UI envelope events. After every
successful a2ui.apply / a2ui.ingest, the resulting
A2uiApplyResult is broadcast to every subscriber as an
a2ui.event JSON-RPC notification. Closes
Parslee-ai/car-releases#29. Subscribers register via the
a2ui/subscribe method and are auto-cleaned on WS disconnect.
mcp_executor: Arc<McpToolExecutor>Per-launch auth token. When Some, the WS dispatcher rejects
non-auth methods on unauthenticated sessions until the client
calls session.auth with the matching value. When None,
auth is disabled and every connection works as before. Set
at startup by car-server unless --no-auth is passed
(default flipped 2026-05); embedders that want to enable
auth call ServerState::install_auth_token. Closes
Parslee-ai/car-releases#32.
Shared MCP tool executor backing remote connectors. One per
process: every WS session’s runtime executor is a
car_engine::McpToolExecutor::share_with_fallback view over
this, so connector tools route here while non-connector tools
fall back to the per-session WsToolExecutor. Connectors
(their live sessions + routes) live on this shared instance, so
a connector enabled in one WS session is reachable from all.
connectors: OnceLock<Arc<ConnectorManager>>Remote MCP connector manager (~/.car/connectors.json).
Lazy-initialized over mcp_executor so
embedders that never touch connectors pay no disk cost. See
ServerState::connectors.
channel_supervisor: OnceLock<Arc<ChannelSupervisor>>Runtime approval-transport channel supervisor (Units 1/2/3). Lazy-
initialized at daemon boot by spawn_channel_pollers, which records the
channels it spawned at startup. The host-gated messaging.config.set
handler reaches it via ServerState to spawn a channel’s watcher the
instant the user enables it — no restart (U1). Also holds the per-channel
liveness messaging.status reads (U2) and the send path writes (U3).
None for embedders that never boot the channel pollers.
durable_tasks: Mutex<JoinSet<String>>Daemon-owned operations which must outlive the WebSocket request that
started them. The tasks hold no ClientSession or WsChannel; only
their oneshot response waiters are connection-scoped.
Two kinds of work live here. Credential-backed auth operations (attempt
completion and ordinary auth mutations), so a socket close can release
its sink immediately without cancelling a coordinator/keychain operation
halfway through its durable result. And coder.start, which registers a
session and provisions a worktree before a multi-minute contract
derivation: once that side effect exists on disk, the run’s lifetime
must not be owned by the board that asked for it.
auth_completion_owner_id: StringPer-daemon incarnation token persisted into redeeming auth leases. Awaiting-callback reservations intentionally have no daemon owner and survive a restart; a redeeming lease owned by a different incarnation is an orphan and is terminally failed by completion-status reconciliation.
auth_token: OnceLock<String>§host_token: OnceLock<String>Per-launch host token — a credential distinct from
auth_token, granting the host-management role (host-class
reads like cross-agent run traces). Critically it is never
served over GET /auth-token: a session becomes host-role only
by presenting this via session.auth { host_token }, and the
only way to obtain it is reading the 0600 host-token file,
which a different local user cannot. This is what stops any
authenticated local client from self-elevating to host and
reading every agent’s run traces (Parslee-ai/car#254). When
None, host-role can’t be granted (no host reads).
mobile_runtime_url: OnceLock<String>Mobile Parslee Core runtime URL this daemon can hand to consumer host apps.
Set by car-server at startup from the daemon bind address or an
explicit public URL override. The token remains the per-launch auth
token; this field only names the WebSocket endpoint.
mobile_registration_url: OnceLock<String>Explicit public mobile runtime URL to register after an authorized credential action succeeds. Absent for the default loopback URL, which is discoverable locally but must not be advertised to Parslee cloud.
parslee_session: OnceLock<ParsleeSession>Parslee cloud identity activated by an explicit credential-backed
request after car auth login has been completed.
attached_agents: Mutex<HashMap<String, String>>agent_id -> client_id map of currently-attached lifecycle
agents (#169). Populated by the session.auth handler when a
supervised child presents its agent_id + per-agent token;
drained on disconnect by remove_session. Single-claim:
a second connection presenting the same agent_id is
rejected so the daemon-side per-agent state stays unambiguous.
agent_memgines: Mutex<HashMap<String, Arc<Mutex<MemgineEngine>>>>agent_id -> persistent memgine map (#170). Lazy-loaded on
first connection per id from ~/.car/memory/agents/<id>.jsonl,
retained across daemon restart, surviving any single
disconnect/reconnect of the supervised child. Connections
that auth without an agent_id (browser, host, ad-hoc CLI)
keep the per-WS ephemeral memgine on ClientSession.memgine
— no behaviour change.
namespace_memgines: Mutex<HashMap<String, Arc<Mutex<MemgineEngine>>>>Daemon-owned memgines keyed by a host-declared memory namespace
(session.auth { memory_namespace }), lazily loaded from
~/.car/memory/memory-namespaces/<encoded-ns>.json — the filename is a
percent-encoding of the namespace, injective so two namespaces can never
share a snapshot file (#891).
A separate axis from agent_memgines, deliberately: one agent may work
across several namespaces, and two hosts may share a namespace without
sharing an identity (Parslee-ai/car-releases#79). A host that binds none
keeps today’s behaviour — the daemon’s shared graph — so the MCP-to-WS
shared knowledge base is untouched.
coder_sessions: Mutex<CoderSessionMap>Live coder sessions keyed by session_id (built-in coding agent).
Process-wide so a session outlives the WS connection that started
it; terminal sessions stay listed for coder.get/coder.list
history until daemon restart (snapshots persist under
~/.car/coder/).
coder_disk_gc_base: InstantMonotonic base for Self::coder_disk_gc_at, taken at construction.
coder_disk_gc_at: AtomicU64Seconds since Self::coder_disk_gc_base at the last coder state-dir
sweep. Rate-limits the disk GC amortized onto coder.start (car#1339)
so a burst of starts re-reads the directory once, not once each.
Monotonic rather than Unix seconds because this is an INTERVAL: a wall clock that steps backwards would stamp a future value and suppress every later sweep for the daemon’s lifetime, which is the bug this closes, reintroduced through the clock.
On the state rather than a process-global static so each ServerState —
every test builds its own — carries its own stamp, and one test’s sweep
cannot suppress another’s.
coder_subscribers: Mutex<HashMap<(String, String), Arc<WsChannel>>>Live coder.event subscribers keyed by (session_id, client_id)
— explicit fanout, same shape as run_subscribers. Lock order:
the session’s event buffer → this map; never the reverse.
coder_watchers: Mutex<HashMap<String, (u64, Arc<WsChannel>)>>coder.watch board subscribers keyed by client_id — one
subscription per connection covering EVERY session, so a board learns
about runs started by any other client (car code, CarHost, milo)
without polling. Fanned to as coder.session_changed. Lock order:
a session’s event buffer → coder_subscribers → this map; never the
reverse.
The value carries a registration generation alongside the channel.
A board re-registers periodically (its coder.watch is idempotent and
returns a fresh snapshot), and the fanout sheds a wedged watcher after
releasing this lock — so a shed that removed by client_id alone would
delete a registration made in that window and silently unwatch a healthy
board. The shed compares the generation and removes only the entry it
actually timed out on.
coder_discussions: Mutex<DiscussionMap>Live coder.discuss.* runtimes keyed by discussion_id. Model history
is checkpointed separately; recovery rebuilds this connection-owned
runtime after validating the saved repository and principal binding.
chat_sessions: Mutex<HashMap<String, ChatSession>>In-flight agents.chat sessions keyed by session_id. See
ChatSession for shape. Populated by agents.chat,
cleared on terminal agent.chat.event or
agents.chat.cancel. Disconnect cleanup happens in
remove_session — any in-flight session bound to either the
disconnecting host or agent client is dropped so subsequent
stray notifications from a respawned agent fall on the floor
rather than racing into a stale stream.
peer_guards: Mutex<HashMap<String, DeliveryGuard>>Per-recipient channel guards for peer messaging (agents.message).
Keyed by recipient agent id. Bounds the channel rather than the sender’s authority: size, identical repeats, per-sender rate, and queue depth. Two agents answering each other form a loop that no policy rule catches, because neither is misbehaving — this is what terminates it.
Entries are dropped in remove_session when the agent detaches, so a
restarted agent starts with a clean budget rather than inheriting the
rate history of the process that used to own its name.
peer_standing: Mutex<HashMap<String, PeerStanding>>Earned standing per sending principal — a peer key inbound, a local agent id outbound.
Deliberately not reaped on detach, unlike Self::peer_guards one
line above. That hook exists so a respawned agent does not inherit a
dead process’s rate and dedupe budget, which is channel state.
Standing is the opposite kind: clearing it on disconnect would let a
degraded agent wipe its record by reattaching, which is a cheaper
erasure than the renaming this design already refuses to allow.
Bounded by TTL and map cap instead, pruned on write.
held_peer_messages: Mutex<VecDeque<HeldPeerMessage>>Peer messages held for operator approval, oldest first.
A message is held when the sender’s read_only posture is
require_approval. Bounded at car_peers::HOLD_CAP; past that the
oldest is dropped, so a stream of held messages cannot grow without
limit. Kept flat rather than per-recipient because approval is an
operator action over the whole queue, not a per-agent one.
Deliberately in memory only: a held message is a live decision awaiting a human, and one that outlived a daemon restart would be delivered into a world that had moved on.
lan_discovery: Mutex<Option<LanDirectory>>Live view of CAR daemons discovered on the local network.
None when mDNS could not start (no multicast, a locked-down sandbox) —
LAN discovery is then simply absent rather than the daemon failing to
boot. A network with no peers and a network CAR cannot browse look the
same to a caller, which is why the status surface reports which it is.
peer_identity: Mutex<Option<Arc<PeerIdentity>>>This daemon’s ed25519 peer identity, and the keys it accepts.
The identity signs outbound cross-host requests; the trust set decides
which inbound signers are CAR. Both are None until the A2A surface
starts, because neither is meaningful without a network surface to
authenticate.
peer_trust: PeerTrust§chat_collectors: Mutex<HashMap<String, ChatCollector>>In-process collectors for agent.chat streams that have no host UI
to forward to — currently the A2A conversational bridge, which reverse-
calls agent.chat on a host session and aggregates the streamed deltas
into a single reply. Keyed by session_id. When a collector exists for a
session, try_forward_agent_chat_event feeds chunks here instead of
forwarding to a host channel. The collecting task owns the entry’s
lifetime (inserts before the reverse-call, removes when done/timed out).
chat_goals: Mutex<HashMap<String, ChatGoalState>>Standing deterministic chat goals keyed by session_id. This is a
host-facing status/control registry, separate from ephemeral
chat_sessions routing so a goal can be inspected after a turn finishes.
runs: Mutex<HashMap<String, RunMeta>>Agent runs keyed by run_id (agent run tracing, U1). Process-
wide (not per-session) so a run’s record outlives the WS
connection that produced it — the durable, connection-
independent boundary client_id cannot be (R1). Populated by
runs.start, made terminal by runs.complete, and swept to
Incomplete on a mid-run disconnect past the grace window
(R5). U2/U3 build the per-turn recorder and disk store on top
of this registry.
run_subscribers: Mutex<HashMap<(String, String), RunTraceSubscriber>>Live runs.trace.event subscribers keyed by (run_id, host_client_id) (agent run tracing, U4). Each value is a
crate::host::RunTraceSubscriber — the producer side of a
bounded channel whose dedicated drain task writes frames to that
connection’s WS. Two CarHost windows on one run register two
distinct entries (explicit fanout — the built-in notification
registry is single-subscriber-per-method).
Lock contract (invariant #1): runs.subscribe snapshots the
run’s turns AND inserts the subscriber while holding the
runs lock; the recorder (record_run_turns),
start_run, complete_run, and mark_run_incomplete append to
runs and push to run_subscribers while holding the SAME
runs lock — so snapshot/register and append/notify are
serialized. No turn appended in the snapshot/register window is
dropped (gap) or double-delivered (dup). The lock order is always
runs → run_subscribers; never the reverse.
browser_views: Arc<BrowserViewRegistry>Every browser the drawer can reach (browser.view.*), keyed by
conversation/agent-session — plus the one standing user session every
conversation without an agent-attached browser shares. Owns the
per-view snapshot/cursor/subscriber fanout; see
crate::browser_view.
run_store: RunStoreDisk-backed run-trace store (agent run tracing, U3). Source of
truth for REPLAY (U5) — persists each run’s RunStarted, turns,
and terminal record as JSONL under ~/.car/runs/{agent_id}/ so a
run survives daemon restarts (R4), bounded by retention/GC (R6)
and protected at rest with 0600/0700 perms + backup exclusion
(R14). The in-memory runs buffer stays the source
for the LIVE stream (U4); this store mirrors what was recorded.
Derived from journal_dir at construction (sibling runs/).
mcp_url: OnceLock<String>Bound MCP HTTP-streamable URL (e.g.
"http://127.0.0.1:9102/mcp") — car-server installs this
after binding the listener. Used by the
agents.invoke_external handler to default
InvokeOptions.mcp_endpoint so external agents
(Claude Code today) load the daemon’s CAR namespace via
--mcp-config automatically. None when MCP isn’t bound
(e.g. --mcp-bind disabled).
approval_gate: ApprovalGateApproval gate for high-risk WS methods (audit 2026-05). The
gate intercepts automation.run_applescript,
automation.shortcuts.run, messages.send, mail.send, the three
mutating calendar.* methods, and vision.ocr before they dispatch, raises a
host.create_approval for the user to act on, and waits
(with a timeout) for host.resolve_approval. Approve →
dispatch continues; deny / timeout → JSON-RPC error code
-32003. The set of gated methods and the wait timeout are
embedder-overridable via
ServerStateConfig::with_approval_gate.
supervision: Arc<SupervisionRegistry>Out-of-process supervision (proposal item 5). Holds the subscribed
supervisors and the intents parked on their verdicts; the matching
crate::supervision::SupervisionGate is registered as an
AdmissionGate on every session runtime alongside
StaticVerificationGate.
Daemon-wide rather than per-session, deliberately: a supervisor connects on its OWN WebSocket and supervises proposals submitted on OTHER sessions. A per-session registry would only ever see the supervisor’s own (empty) traffic.
Costs nothing until someone subscribes — the gate returns Allow
without building an intent while the subscriber set is empty.
approval_ledger: Arc<RwLock<ApprovalLedger>>The daemon-wide shared HITL approval ledger (kernel review C1) —
journal-backed at ~/.car/approvals.jsonl (override via
ServerStateConfig::with_approval_journal). ONE ledger for all
sessions: permission.approve/reject on any connection records
here, and every ledger consumer (permission.evaluate/pending,
evolution.run, cascade.run, skill.enforce_deployment/
ingest_governed) reads here — so an approval granted by a host
connection is visible to the agent connection that surfaced it and
survives daemon restart. Per-session gates keep only tier/classifier
state (see ClientSession::permission_gate). Falls back to an
in-memory ledger (loudly warned, approvals NOT restart-durable) only
when the journal is unopenable.
harness_measurer: RwLock<Option<Arc<dyn HarnessMeasurer>>>The in-process harness evaluator, when the binary installed one — the
thing that lets evolution.run grade a harness candidate ITSELF
instead of waiting for an operator to run car-bench-harness twice and
hand both files back in.
None on a build that did not install one (every embedder, and any
binary other than car-server): evolution.run with harness_measure
then ERRS rather than quietly falling back to HITL, because an opt-in
that silently does nothing would report an unattended cycle that never
measured anything.
A std::sync::RwLock on purpose — no caller holds it across an await;
they clone the Arc out and drop the guard.
Implementations§
Source§impl ServerState
impl ServerState
Sourcepub fn standalone(journal_dir: PathBuf) -> Self
pub fn standalone(journal_dir: PathBuf) -> Self
Constructor for the standalone car-server binary. Each WS
connection gets its own per-session memgine — matches the
pre-extraction default and is correct for a single-process
daemon serving one user at a time.
Embedders must not call this. It silently leaves
shared_memgine = None, which re-introduces the dual-memgine
bug U7 was created to prevent (one engine in the embedder, a
fresh one inside every WS session). Embedders use
ServerState::embedded instead, which makes the shared
engine handle a required argument so it cannot be forgotten.
Sourcepub fn sync_subsystem(&self) -> Result<Arc<Mutex<SyncSubsystem>>, String>
pub fn sync_subsystem(&self) -> Result<Arc<Mutex<SyncSubsystem>>, String>
The daemon-held multi-device-sync + execution-lease subsystem (B6),
lazily opened on first sync.*/lease.* contact and rooted at
<journal_dir>/sync/. Serialized first-open (see the Self::sync
field docs): the oplog journal’s exclusive advisory lock means only one
open may succeed, so the fallible init runs under the field’s
std::sync::Mutex while the subsystem itself is shared behind a
tokio::sync::Mutex.
Sourcepub fn subsystem_for_scope(
&self,
scope: &Scope,
) -> Result<Arc<Mutex<SyncSubsystem>>, String>
pub fn subsystem_for_scope( &self, scope: &Scope, ) -> Result<Arc<Mutex<SyncSubsystem>>, String>
Route a write to the right subsystem BY SCOPE: an opted-in
Scope::Shared{org} goes to the shared org delivery subsystem (relay scope
org:{orgId}, so it converges cross-member); everything else — Personal,
or an org with no opted-in delivery — goes to the personal subsystem, which
is byte-identical to before when org-scope is off. Used ONLY at the write
sites; every reader keeps calling Self::sync_subsystem. (8a)
Sourcepub fn org_subsystem(&self) -> Option<Arc<Mutex<SyncSubsystem>>>
pub fn org_subsystem(&self) -> Option<Arc<Mutex<SyncSubsystem>>>
The opted-in shared org delivery subsystem, if any — for pumping both scopes.
Sourcepub async fn persist_chat_goals(&self) -> Result<(), String>
pub async fn persist_chat_goals(&self) -> Result<(), String>
Persist the daemon-wide standing chat-goal registry. The snapshot is taken under the async mutex, then serialized and atomically rewritten on the blocking pool so chat/event dispatch does not block a Tokio worker.
Sourcepub fn embedded(
journal_dir: PathBuf,
shared_memgine: Arc<Mutex<MemgineEngine>>,
) -> Self
pub fn embedded( journal_dir: PathBuf, shared_memgine: Arc<Mutex<MemgineEngine>>, ) -> Self
Constructor for embedders (e.g. tokhn-daemon). The shared
memgine handle is required: every WS session created by
this state will reuse the same engine, preventing the
dual-memgine bug.
For embedders that also want to inject a pre-warmed inference
engine or other advanced wiring, build a ServerStateConfig
directly and call ServerState::with_config.
Sourcepub fn with_config(cfg: ServerStateConfig) -> Self
pub fn with_config(cfg: ServerStateConfig) -> Self
Build a ServerState from a ServerStateConfig — the path
embedders use when they need to inject a shared memgine and
a pre-warmed inference engine, or any other advanced wiring
the convenience constructors don’t cover.
Sourcepub fn try_with_config(cfg: ServerStateConfig) -> Result<Self, String>
pub fn try_with_config(cfg: ServerStateConfig) -> Result<Self, String>
Fallible startup surface used by the daemon before it begins listening. Durable lifecycle outbox replay is acknowledgement-bounded; any durability-unknown result stops construction and leaves its RunStore preimage/quarantine intact for an exact later retry.
Sourcepub fn set_harness_measurer(&self, m: Arc<dyn HarnessMeasurer>)
pub fn set_harness_measurer(&self, m: Arc<dyn HarnessMeasurer>)
Install the in-process harness evaluator evolution.run’s
harness_measure path grades candidates with. Called once at daemon
startup; a later call replaces it.
Sourcepub fn harness_measurer(&self) -> Option<Arc<dyn HarnessMeasurer>>
pub fn harness_measurer(&self) -> Option<Arc<dyn HarnessMeasurer>>
The installed harness evaluator, if any. Clones the Arc out and drops
the guard so callers never hold a std::sync lock across an await.
Sourcepub async fn spawn_durable_operation<F, E>(
&self,
operation_name: impl Into<String>,
operation: F,
) -> Receiver<Result<Value, E>> ⓘ
pub async fn spawn_durable_operation<F, E>( &self, operation_name: impl Into<String>, operation: F, ) -> Receiver<Result<Value, E>> ⓘ
Start one daemon-owned operation and return a connection-scoped result receiver. Dropping the receiver never cancels the operation.
The caller awaits the receiver exactly as it would have awaited the
operation, so response shape and timing are unchanged for a connection
that stays alive. What changes is the failure mode: when the connection
goes away, the per-connection JoinSet aborts only this waiter, and the
operation itself runs to completion on the daemon.
Sourcepub async fn spawn_durable_task<F>(
&self,
operation_name: impl Into<String>,
operation: F,
)
pub async fn spawn_durable_task<F>( &self, operation_name: impl Into<String>, operation: F, )
Start one fire-and-reconcile daemon-owned operation.
Sourcepub fn install_auth_token(&self, token: String) -> Result<(), String>
pub fn install_auth_token(&self, token: String) -> Result<(), String>
Enable the per-launch auth handshake. After this call, every
new WS connection must call session.auth with token as
the first frame; otherwise the connection is closed. Called
by car-server at startup unless --no-auth is set
(default flipped 2026-05); embedders supply their own token
if they want the same posture. Returns Err(token) when
auth was already installed.
Sourcepub fn install_host_token(&self, token: String) -> Result<(), String>
pub fn install_host_token(&self, token: String) -> Result<(), String>
Install the per-launch host token (Parslee-ai/car#254). A
session that later presents this via session.auth { host_token }
is granted the host-management role (ClientSession::is_host),
which authorize_run_access requires for cross-agent run-trace
reads. Set by car-server at startup (mints + writes the 0600
host-token file) unless --no-auth is set. Returns Err(token)
when a host token was already installed.
pub fn install_mobile_runtime_url(&self, url: String) -> Result<(), String>
pub fn install_mobile_registration_url(&self, url: String) -> Result<(), String>
pub fn install_parslee_session( &self, session: ParsleeSession, ) -> Result<(), ParsleeSession>
Sourcepub fn install_channel_supervisor(
&self,
supervisor: Arc<ChannelSupervisor>,
) -> Result<(), Arc<ChannelSupervisor>>
pub fn install_channel_supervisor( &self, supervisor: Arc<ChannelSupervisor>, ) -> Result<(), Arc<ChannelSupervisor>>
Install the runtime channel supervisor (Units 1/2/3). Called by
car-server at boot right after spawn_channel_pollers builds it, so
the host-gated messaging.config.set handler can reach it via
ServerState to spawn a channel’s watcher on enable. Idempotent on the
first call; a second install is rejected (returns the supervisor back).
Sourcepub fn install_mcp_url(&self, url: String) -> Result<(), String>
pub fn install_mcp_url(&self, url: String) -> Result<(), String>
Install the bound MCP URL after car-server’s listener is up.
Idempotent on the first call; subsequent calls are accepted
silently (matches the supervisor / a2a_dispatcher install
idiom). Returns Err(()) when an MCP URL was already
installed — embedders should treat this as “another
component beat us to it” and use whichever value is now set.
Sourcepub fn connectors(&self) -> Arc<ConnectorManager> ⓘ
pub fn connectors(&self) -> Arc<ConnectorManager> ⓘ
Lazy-initialize and return the remote MCP connector manager,
backed by ~/.car/connectors.json over the shared
mcp_executor. Construction is sync and
touches no network — call ensure_connectors_loaded to dial
persisted connectors.
Sourcepub async fn ensure_connectors_loaded(&self)
pub async fn ensure_connectors_loaded(&self)
Load persisted connectors and dial them, once per process.
Idempotent: the first caller does the work; later callers return
immediately. car-server calls this at boot so connector tools
are live before clients connect; the connectors.* handlers call
it too, so an embedder that skips the boot call still gets a
lazy load on first use. Enabled tools discovered here are
registered into any already-open sessions.
Sourcepub async fn register_connector_entries(&self, entries: &[ToolEntry])
pub async fn register_connector_entries(&self, entries: &[ToolEntry])
Register connector car_engine::ToolEntrys into every currently-open
session’s runtime registry, so a tool enabled mid-session
becomes visible without a reconnect. New sessions pick up the
enabled set in create_session.
Sourcepub async fn unregister_connector_tools(&self, canonical_names: &[String])
pub async fn unregister_connector_tools(&self, canonical_names: &[String])
Unregister connector tools (by canonical name) from every open session’s runtime, so a disabled or removed connector’s tools stop being visible to the model and accepted by the validator.
Sourcepub async fn any_host_connected(&self) -> bool
pub async fn any_host_connected(&self) -> bool
Whether any currently connected WS client holds the host-management
role — authenticated via session.auth { host_token }
(ClientSession::is_host), NOT the same thing as host.subscribe
membership (see the doc comment on the auth check this mirrors,
around handle_session_auth). In practice: is a CarHost / Command
Deck app connected right now, so its browser drawer is a real
visible surface?
Scans the (small) live session set fresh on every call — no caching,
so this is always current. The deciding signal for Task 7’s
headless/headed browser-launch default and its sign-in host-gone
fallback (assistant::browser_tools::HostConnectivity).
Sourcepub fn declagents(&self) -> Result<Arc<DeclRegistry>, String>
pub fn declagents(&self) -> Result<Arc<DeclRegistry>, String>
Lazy-initialize and return the agent supervisor. The first
call constructs a car_registry::supervisor::Supervisor backed by
~/.car/agents.json + ~/.car/logs/. Embedders that need a
non-default location should call
ServerState::install_supervisor before any handler runs.
In observer mode (set via install_observer_manifest),
returns a clear error mentioning the manifest path the
primary daemon owns. This prevents the second daemon from
re-attempting user_default() (which would also fail with
AlreadyRunning) on every WS call, and gives mutation
handlers a stable refusal path. Read-only handlers
(agents.list, agents.health) should call
Self::observer_manifest_path first and fall back to
car_registry::supervisor::Supervisor::list_from_manifest /
health_from_manifest when set. Closes
Parslee-ai/car-releases#44.
The declarative-agent registry, lazy-initialized on first use.
pub fn routing(&self) -> Result<Arc<RoutingStore>, String>
pub fn supervisor(&self) -> Result<Arc<Supervisor>, String>
Sourcepub fn install_supervisor(
&self,
supervisor: Arc<Supervisor>,
) -> Result<(), Arc<Supervisor>>
pub fn install_supervisor( &self, supervisor: Arc<Supervisor>, ) -> Result<(), Arc<Supervisor>>
Replace the lazy default with a caller-supplied supervisor.
Returns Err(()) when a supervisor was already installed.
Used by the standalone car-server binary to call
start_all() on a known-good handle without paying the
lazy-init lookup cost.
Sourcepub fn supervisor_if_installed(&self) -> Option<Arc<Supervisor>>
pub fn supervisor_if_installed(&self) -> Option<Arc<Supervisor>>
Non-acquiring read of the currently-installed supervisor.
Unlike supervisor, this does NOT lazy-
init via user_default() — it returns None instead of
constructing a fresh Supervisor and acquiring the
<manifest>.lock as a side effect. Use this from read-only
metadata paths (host.subscribe identity, status surfaces)
where causing lock acquisition on observation would be a
Heisenberg subscribe — the act of asking “do you own the
lock?” must not be the act of taking it.
Sourcepub fn install_observer_manifest(&self, path: PathBuf) -> Result<(), PathBuf>
pub fn install_observer_manifest(&self, path: PathBuf) -> Result<(), PathBuf>
Mark this daemon as observing a manifest owned by another
car-server process. After this call, supervisor() returns
an “observe-only” error and read-only handlers
(agents.list, agents.health) fall back to the static
Supervisor::list_from_manifest / health_from_manifest
paths. Idempotent — subsequent calls with the same path are
no-ops; a different path returns Err(()). Closes
Parslee-ai/car-releases#44.
Sourcepub fn observer_manifest_path(&self) -> Option<&PathBuf>
pub fn observer_manifest_path(&self) -> Option<&PathBuf>
Path of the manifest this daemon is observing but not
supervising. None when this daemon owns the supervisor
(the normal case) or when no manifest is configured at all
(no HOME, embedder didn’t install one).
Sourcepub async fn a2a_dispatcher(&self) -> Arc<A2aDispatcher> ⓘ
pub async fn a2a_dispatcher(&self) -> Arc<A2aDispatcher> ⓘ
Lazy-initialize and return the in-core A2A dispatcher. The
first call constructs an car_a2a::A2aDispatcher from
either the embedder’s overrides (set via
ServerStateConfig::with_a2a_runtime / with_a2a_store /
with_a2a_card_source) or sensible defaults: a fresh
Runtime with register_agent_basics registered, an
InMemoryTaskStore, and a card built from the runtime’s
tool schemas advertising ws://127.0.0.1:9100/ as the
public URL. Closes Parslee-ai/car-releases#28.
Sourcepub async fn reserve_run(&self, meta: RunMeta) -> Result<RunReservation, String>
pub async fn reserve_run(&self, meta: RunMeta) -> Result<RunReservation, String>
Atomically reserve a run id/idempotency key across every live session.
A per-run durability lock serializes the disk absence check with the
final process-wide map insert. The global runs lock is held only for
the two bounded map checks, never for filesystem enumeration/replay.
pub async fn persist_run_start( &self, run_id: &str, ) -> Result<RunStarted, String>
pub async fn commit_run_start(&self, run_id: &str) -> Result<(), String>
Sourcepub async fn start_run(&self, meta: RunMeta) -> Result<(), String>
pub async fn start_run(&self, meta: RunMeta) -> Result<(), String>
Compatibility/internal path for CAR-owned in-process callers that do
not expose runs.start. The public RPC uses reserve/persist/journal/
commit explicitly so its acknowledgement includes both durable stores.
Sourcepub async fn prepare_run_completion(
&self,
run_id: &str,
termination: RunTermination,
) -> Result<RunEnded, String>
pub async fn prepare_run_completion( &self, run_id: &str, termination: RunTermination, ) -> Result<RunEnded, String>
Make a run terminal with a harness-reported outcome
(runs.complete). Returns Err if the run_id is unknown or
already terminal — the handler maps that to a JSON-RPC error so
a double-complete or stale id is visible, not silently swallowed.
U3: also appends the terminal RunEnded line to the run’s JSONL
file so REPLAY (U5) reports the right status (Completed) after a
restart. The disk write happens after the lock is released; disk
failures are logged, never fatal.
pub async fn commit_run_completion( &self, ended: &RunEnded, ) -> Result<(), String>
pub async fn complete_run( &self, run_id: &str, termination: RunTermination, ) -> Result<RunEnded, String>
Sourcepub async fn persist_run_cancellation_requested(
&self,
owner_session: &ClientSession,
requested: &RunCancellationRequested,
) -> Result<(), String>
pub async fn persist_run_cancellation_requested( &self, owner_session: &ClientSession, requested: &RunCancellationRequested, ) -> Result<(), String>
Commit the body-free request to RunStore and the authenticated opening client’s journal under the same per-run durability lock used by terminal completion. Memory changes only after both receipts exist.
pub async fn persist_run_cancellation_result( &self, owner_session: &ClientSession, result: &RunCancelResponse, ) -> Result<(), String>
Sourcepub async fn persist_recovered_run_cancellation_result(
&self,
agent_id: &str,
result: &RunCancelResponse,
) -> Result<(), String>
pub async fn persist_recovered_run_cancellation_result( &self, agent_id: &str, result: &RunCancelResponse, ) -> Result<(), String>
Finish an unconfirmed cancellation after the opening socket has gone
away. The durable RunStarted client id is the journal authority: a
reconnect may request recovery, but it may not redirect the receipt to
its own journal. RunStore, that authenticated journal, live state, and
subscribers advance as one retryable transaction before success is
returned.
Sourcepub async fn prepare_run_incomplete(&self, run_id: &str) -> Option<RunEnded>
pub async fn prepare_run_incomplete(&self, run_id: &str) -> Option<RunEnded>
Mark a run Incomplete (R5) — used by disconnect cleanup when a
harness drops without runs.complete. No-op if the run is
already terminal (the common healthy-close case where
runs.complete won the race). Returns true if it actually
wrote the Incomplete marker.
U3: on the transition to Incomplete, appends the terminal
RunEnded { Incomplete } line to disk so an orphaned run reports
Incomplete (not InProgress) on REPLAY after a restart. Disk
failures are logged, never fatal.
pub async fn mark_run_incomplete(&self, run_id: &str) -> Option<RunEnded>
Sourcepub async fn run_meta(&self, run_id: &str) -> Option<RunMeta>
pub async fn run_meta(&self, run_id: &str) -> Option<RunMeta>
Non-acquiring read of a run’s current metadata (clone). Used by
tests and by the U2 recorder to learn a run’s owning agent_id.
NOTE: this clones the ENTIRE RunMeta, including its turns buffer
(full prompts + CLI output — up to the per-run caps). Hot callers
that need only the run’s header facts must use Self::run_header
instead; cloning the whole buffer per batch RPC is a ~512 MB
worst-case copy at this PR’s caps (ADV-2).
Sourcepub async fn run_header(
&self,
run_id: &str,
) -> Option<(String, bool, usize, Option<String>)>
pub async fn run_header( &self, run_id: &str, ) -> Option<(String, bool, usize, Option<String>)>
Lightweight, non-cloning read of a run’s header facts under the
runs lock: (agent_id, terminal, turn_count, trace_corruption).
The owning agent id and optional error are the only heap allocations;
This is what handle_runs_record_turns needs — owner for the write
authz check, terminality for the distinguishable run_terminal drop,
and the current turn count for the (fast-path) ceiling pre-check —
without the run_meta deep clone of every recorded turn’s payload.
Sourcepub async fn run_lifecycle_binding(
&self,
run_id: &str,
) -> Option<(String, bool)>
pub async fn run_lifecycle_binding( &self, run_id: &str, ) -> Option<(String, bool)>
Lightweight authenticated producer binding for run lifecycle handlers:
(client_id, terminal). Never falls back to disk because an on-disk row
has no live WebSocket producer to authorize new writes.
Sourcepub async fn run_lifecycle_state(
&self,
run_id: &str,
) -> Option<(String, bool, bool, bool)>
pub async fn run_lifecycle_state( &self, run_id: &str, ) -> Option<(String, bool, bool, bool)>
(client_id, terminal, start_committed, completion_pending) for the
authenticated producer path. Proposal writes require the exact owner,
a committed start, and no prepared terminal transaction.
pub async fn run_lifecycle_state_with_corruption( &self, run_id: &str, ) -> Option<(String, bool, bool, bool, Option<String>)>
Sourcepub async fn record_run_turns(
self: &Arc<Self>,
run_id: &str,
records: Vec<RunRecord>,
) -> RecordRunTurnsOutcome
pub async fn record_run_turns( self: &Arc<Self>, run_id: &str, records: Vec<RunRecord>, ) -> RecordRunTurnsOutcome
Append per-turn trace records to a run’s in-memory buffer (agent
run tracing, U2). The recorder calls this after every
proposal.submit on a session with a current run, passing the
RunRecord::Turns the recorder produced for that proposal.
Returns a RecordRunTurnsOutcome:
Appended { new_total }— the batch landed;new_totalis the run’s new total turn count (the caller computes its nextstart_indexfrom it).RefusedCeiling— the batch would take the run pastRECORD_TURNS_RUN_CEILINGand was refused WHOLE (the under-lock runaway backstop, ADV-1). Nothing was appended.UnknownOrTerminal— therun_idis unknown (the bracket was never opened) or already terminal (a turn arriving afterruns.completeis dropped — the run is closed). Nothing was appended.
U3 flushes the in-memory buffer to disk; U4 broadcasts it. Both read
via Self::run_turns.
U3: the same turn records are appended to the run’s JSONL file so REPLAY (U5) sees them after a restart. Persistence is the commit point: only after the exact JSONL batch is durable do memory and subscribers advance. The operation runs in a daemon-owned task, so dropping a request waiter cannot strand a persisted batch between disk and memory.
Sourcepub async fn ensure_proposal_run_turns(
self: &Arc<Self>,
pending: &PendingProposalFinalization,
) -> Result<(), String>
pub async fn ensure_proposal_run_turns( self: &Arc<Self>, pending: &PendingProposalFinalization, ) -> Result<(), String>
Make a proposal’s authenticated final trace durable exactly once and reflect a newly-appended batch into the live registry/subscription stream before proposal finalization can acknowledge success.
Sourcepub async fn run_turn_count(&self, run_id: &str) -> usize
pub async fn run_turn_count(&self, run_id: &str) -> usize
Number of turns already recorded for a run — the start_index the
recorder passes so per-turn index stays monotonic across the
run’s proposals. 0 for an unknown run.
Sourcepub async fn run_turns(&self, run_id: &str) -> Vec<RunRecord>
pub async fn run_turns(&self, run_id: &str) -> Vec<RunRecord>
Clone of a run’s ordered per-turn trace (agent run tracing, U2). This is the accessor U3 reads to flush turns to disk and U4 reads to broadcast them. Empty Vec for an unknown run or a run with no turns yet.
Sourcepub async fn unsubscribe_run(&self, run_id: &str, host_client_id: &str) -> bool
pub async fn unsubscribe_run(&self, run_id: &str, host_client_id: &str) -> bool
Remove a live run-trace subscriber for (run_id, host_client_id)
(agent run tracing, U4). Returns true if a subscription existed.
Dropping the crate::host::RunTraceSubscriber drops its channel
sender, which ends the drain task.
Sourcepub async fn drop_run_subscribers_for_client(&self, host_client_id: &str)
pub async fn drop_run_subscribers_for_client(&self, host_client_id: &str)
Drop every live run-trace subscription owned by host_client_id
(agent run tracing, U4 — R8 cleanup). Called from
remove_session on disconnect so a CarHost that drops doesn’t
leave dangling drain tasks. Reconnect-durability is client-side:
the run stays subscribable while it lives, and the CarHost re-issues
runs.subscribe {run_id} on its new connection — the server never
synthesizes a failure on subscriber drop.
Sourcepub async fn create_session(
&self,
client_id: &str,
channel: Arc<WsChannel>,
) -> Result<Arc<ClientSession>, String>
pub async fn create_session( &self, client_id: &str, channel: Arc<WsChannel>, ) -> Result<Arc<ClientSession>, String>
Fails only when <home>/.car/policies/ holds a malformed rule file —
see [apply_project_policies] for why that is fatal rather than a
warning.
Sourcepub async fn bind_substrate_to_connector(
&self,
session: &Arc<ClientSession>,
slug: &str,
) -> Result<String, String>
pub async fn bind_substrate_to_connector( &self, session: &Arc<ClientSession>, slug: &str, ) -> Result<String, String>
Bind a session’s runtime to the car_engine::McpSubstrate backed by an
already-connected MCP connector, making that session coherent:
its commodity built-ins (read_file/write_file/…/run_command)
act on the same environment as its mcp_{slug}_* connector tools,
instead of being split between the host/WS client and the remote MCP
server (docs/execution-substrate.md §3, phase 3).
This is opt-in and explicit — it is reached only via the
session.bindSubstrate JSON-RPC method, never inferred from a
connector merely being enabled. Sessions that never call it keep the
exact historic composition (share_with_fallback(ws_executor) over a
default LocalSubstrate), so GUI / A2A-over-WS / connectors-as-tools
/ coder behavior is unchanged.
On bind we:
- Wrap the live connector session (shared with the
mcp_{slug}_*routes — no second dial) as ancar_engine::McpSubstrateandset_substrateit on the session runtime. register_agent_basics()so the model sees portable names (read_file/run_command/…) that the engine routes to the bound substrate.- Swap the session executor to a
SubstrateShadowExecutorover the sameshare_with_fallback(ws_executor)composition, so the bare substrate-owned names fall through to the substrate while connector routes and host tool callbacks keep working unchanged (option (a)).
Returns the substrate’s environment name (the connector slug) on success, or an error string if the connector is not connected.
Sourcepub async fn remove_session(
&self,
client_id: &str,
) -> Option<Arc<ClientSession>>
pub async fn remove_session( &self, client_id: &str, ) -> Option<Arc<ClientSession>>
Remove a per-client session from the registry on disconnect.
Returns the removed session if present so callers can drop any
remaining strong refs (e.g. drain pending tool callbacks). Fix
for MULTI-4 / WS-3 — without this, state.sessions retains
Arc<ClientSession> for every connection that ever existed.
Auto Trait Implementations§
impl !Freeze for ServerState
impl !RefUnwindSafe for ServerState
impl !UnwindSafe for ServerState
impl Send for ServerState
impl Sync for ServerState
impl Unpin for ServerState
impl UnsafeUnpin for ServerState
Blanket Implementations§
Source§impl<T> BorrowMut<T> for Twhere
T: ?Sized,
impl<T> BorrowMut<T> for Twhere
T: ?Sized,
Source§fn borrow_mut(&mut self) -> &mut T
fn borrow_mut(&mut self) -> &mut T
impl<ST, DT> CastableFrom<ST, Initialized, Initialized> for DT
impl<ST, DT> CastableFrom<ST, Uninit, Uninit> for DT
impl<S, T> Duplex<S> for Twhere
T: FromSample<S> + ToSample<S>,
impl<T> ErasedDestructor for Twhere
T: 'static,
Source§impl<S> FromSample<S> for S
impl<S> FromSample<S> for S
fn from_sample_(s: S) -> S
Source§impl<T> Instrument for T
impl<T> Instrument for T
Source§fn instrument(self, span: Span) -> Instrumented<Self> ⓘ
fn instrument(self, span: Span) -> Instrumented<Self> ⓘ
Source§fn in_current_span(self) -> Instrumented<Self> ⓘ
fn in_current_span(self) -> Instrumented<Self> ⓘ
Source§impl<T> IntoEither for T
impl<T> IntoEither for T
Source§fn into_either(self, into_left: bool) -> Either<Self, Self> ⓘ
fn into_either(self, into_left: bool) -> Either<Self, Self> ⓘ
self into a Left variant of Either<Self, Self>
if into_left is true.
Converts self into a Right variant of Either<Self, Self>
otherwise. Read moreSource§fn into_either_with<F>(self, into_left: F) -> Either<Self, Self> ⓘ
fn into_either_with<F>(self, into_left: F) -> Either<Self, Self> ⓘ
self into a Left variant of Either<Self, Self>
if into_left(&self) returns true.
Converts self into a Right variant of Either<Self, Self>
otherwise. Read more