#Messaging and room state
On this page
Where it livesHow it worksThe two WebSocket feedsREST and the shared HTTP boundThe inbound pipelineRouting to host sessionsThe unrouted queueLocal noticesLifecycle of a queued messageOne inbound message to an owned runtime and backOutbound send and mention resolutionThe human room feed and room message cacheLifecycle and system rowsRoom connectivityWorking-state lease publicationEvent fan-out and the client event streamState and ownershipContractsTraits this subsystem providesControl methodsBand protocol consumedSemantic guarantees enforced by codeInvariantsFailure and recoveryExtension pointsAdding a coding-agent harnessAdding a messaging featureRefactor notesThis subsystem moves conversation between Band rooms and the coding agents and people that Jam hosts locally. Inbound messages for an agent arrive on that agent's Band WebSocket, are written to durable local storage before anything else happens, are routed by room to exactly one host session, and are handed to that session's provider adapter. They stay queued until the agent explicitly replies or acknowledges. In parallel, one Human-API WebSocket per signed-in account feeds the desktop's room timelines, rosters, presence, and activity through a best-effort event stream that clients re-hydrate after any gap.
The subsystem does not decide what an agent does with a message, does not own provider processes, and does not poll Band for room history on every event. It never routes a message to a guessed session: a room that no session owns keeps its messages in an explicit unrouted bucket. It treats the event stream to clients as a change feed, not as a source of truth.
#Where it lives
| Area | Crate and module | Key types and functions |
|---|---|---|
| Band surface as traits | crates/jam-transport/src/lib.rs, dto.rs |
UserApi (39 methods), AgentApi (22), Subscriber (1), HumanSubscriber (11), Transport (blanket union), RoomEvent, Subscription, HumanSubscription, InboundMessage, RoomInbound, RoomMutationInbound, ConnState |
| Band implementation | crates/jam-transport-band/src/lib.rs |
BandClient, BandHttpClient, MAX_HTTP_IN_FLIGHT, MAX_HTTP_IDLE_PER_HOST, HumanRoomWatch, MAX_HUMAN_ROOM_WATCHES |
| Agent WebSocket | crates/jam-transport-band/src/socket.rs |
impl Subscriber for BandClient, drive, run_session, handle_text, catch_up_pending_messages, ReconnectBackoff, Heartbeat, connect_with_shared_tls |
| Human WebSocket | crates/jam-transport-band/src/human_socket.rs |
impl HumanSubscriber for BandClient, drive, run_session, reconcile_room_channels, room_inbound, SessionEnd::Superseded |
| Inbound pipeline and outbound send | crates/jam-core/src/engine.rs |
Engine, Engine::run, handle_inbound, RoomSpool, Bridge, Deliverer, AttachmentStore, claim_message_processing, Engine::send_with_attachments, Engine::reply_parts, Engine::ack, Engine::enqueue_local_notice |
| Mentions and peer cache | crates/jam-core/src/mentions.rs, peers.rs |
extract_handles, resolve_mentions, rewrite_inbound_mentions, prepare_outbound_mentions, project_outbound_mentions, mentions_for_recipients, PeerCache |
| Queue and routing | crates/jam-core/src/queue.rs, session.rs |
Queue, FileQueue, Session, SessionRegistry |
| Queue backend port | crates/jam-manager/src/queue.rs |
QueueBackend, SqliteQueueBackend, FileQueueBackend, SqliteQueue, import_jsonl_tree |
| SQLite queue tables | crates/jam-store/src/sqlite.rs |
queue_append, queue_read, queue_rehome, queue_remove, run_maintenance; table queue_messages |
| Worker glue | crates/jam-manager/src/worker.rs |
spawn_worker, supervise, redeliver, reconcile_queue_dir, commit_turn_disposition, post_assistant_reply, QuestionGate, InboundAttachmentStore |
| Human room feed | crates/jam-manager/src/roomfeed.rs |
run_account_feed, run_account_feed_inner, RoomProjection, apply_projection_mutation_once, merge_snapshot_with_live_mutations, run_room_persistence_worker, schedule_membership_join_fallback |
| Room connectivity | crates/jam-manager/src/room_connectivity.rs |
ObservedRoom, Manager::accept_room_connectivity, run_connectivity_room, observed_room_participants |
| Lifecycle rows and latch | crates/jam-manager/src/lifecycle_events.rs |
OnlineCause, OfflineCause, LifecycleLatch, membership_room_message |
| Working-state lease | crates/jam-manager/src/working_state.rs |
WorkingStatePublisher, run_publisher, PLATFORM_WORKING_TTL, WORKING_REFRESH_INTERVAL |
| Event fan-out | crates/jam-manager/src/fanout.rs, logfile.rs |
Fanout::publish, Fanout::subscribe_with_room_snapshot, emit_locked, format_event_line, PeerLog |
| Client event stream | crates/jam-manager/src/manager.rs Manager::event_stream; crates/jam-daemon/src/lib.rs events; crates/jam-client/src/lib.rs event_stream; crates/jam-ipc/src/events.rs relay |
EventKind::StreamResync, EventKind::RoomLiveStateReset, MAX_EVENT_LINE_BYTES |
| Desktop consumer | apps/desktop/src-tauri/src/events.rs TauriSink; apps/desktop/src/stores/daemon.ts applyEvent; apps/desktop/src/App.tsx |
DaemonEvent, ResyncEvent, subscribeDaemonEvents, subscribeResync |
The crate edges follow the intended direction: jam-transport depends only on jam-domain; jam-transport-band adds jam-tls; jam-core depends on jam-domain, jam-transport, and jam-contract (for one constant, see Refactor notes). jam-manager is the only crate that binds these together, and bins/jam/src/jamd.rs chooses the concrete queue backend and constructs the single BandHttpClient. The wider process picture is in Runtime topology.
#How it works
#The two WebSocket feeds
Jam holds two kinds of Band WebSocket, and they serve different consumers.
The agent socket is one per running agent identity. Engine::run in jam-core calls Subscriber::subscribe, which BandClient implements in socket.rs. It joins agent_rooms:{agent_id} and agent_contacts:{agent_id}, lists the agent's chats over REST (list_chats, bounded by CATCH_UP_DEADLINE of 20 seconds), joins chat_room:{room}, room_participants:{room}, and, when the platform advertises the WorkingActivity API contract, room_activity:{room} for each room. After each room join it calls catch_up_pending_messages, which pages GET /api/v1/agent/chats/{room}/messages. Band's default query returns every message that is not yet processed, so reconnect replays all work the agent has not settled. Only after catch-up does the socket publish ConnState::Up. A failed chat listing is treated as a disconnect rather than going Up with no room channels (test reconnect.rs::catchup_list_failure_reconnects_instead_of_going_up_deaf).
The human socket is one per signed-in account. roomfeed::run_account_feed calls HumanSubscriber::subscribe_human, implemented in human_socket.rs. It joins user_rooms:{user_uuid} and user_contacts:{user_uuid} and then a bounded working set of rooms, up to three channels per room (room_participants and room_activity are skipped when the platform denied or does not support them), chosen by HumanRoomWatch. The set is capped at MAX_HUMAN_ROOM_WATCHES (48). Up to 8 slots go first to explicitly opened or locally changed rooms (MAX_PINNED_HUMAN_ROOM_WATCHES); the remaining slots are filled in order by leased watches, rooms with recently observed traffic, and the server directory baseline ranked by last message activity (ranked_human_room_ids). message_created, message_updated, and event_created frames decode to RoomMutationInbound::Upsert, and message_deleted to RoomMutationInbound::Delete. The human socket has no message catch-up of its own; after a reconnect the manager reconciles room heads over REST (see Failure and recovery).
Both drivers share frame plumbing from socket.rs (send_join, join_room_channels, Heartbeat, ReconnectBackoff, reconnect_jitter, connect_with_shared_tls). Each has its own drive and run_session loop. A heartbeat is sent every 25 seconds (HEARTBEAT); a heartbeat that is still unanswered at the next tick ends the session. Reconnect backoff starts at 2 seconds, doubles per consecutive failure up to 60 seconds, resets after a session that lasted at least 60 seconds, and adds up to 1 second of jitter. Each reconnect rebuilds the handshake request so a refreshed OAuth token is used (test bearer.rs::ws_reconnect_re_reads_the_token_provider).
Only the human socket handles the Band supersede push. It ends the session with SessionEnd::Superseded, publishes ConnState::Superseded, and does not reconnect (test supersede.rs::supersede_is_terminal_and_does_not_reconnect). The agent socket has no supersede branch; jam-core::engine::map_conn_state maps Superseded to PeerState::Stopped only to keep the match total.
#REST and the shared HTTP bound
Every REST call goes through BandClient::execute, which clones the current credential out of a watch channel, sends the request, and on AuthExpired asks the owner of the Bearer credential to refresh once and retries. jamd constructs exactly one BandHttpClient with BandClient::http_client and clones it into every user and agent BandClient (bins/jam/src/jamd.rs). The clone shares the reqwest pool, a Semaphore of MAX_HTTP_IN_FLIGHT (16) active requests, and pool_max_idle_per_host(MAX_HTTP_IDLE_PER_HOST) (8). A BandResponse holds its permit until the body is consumed or dropped, so a slow body cannot escape the bound (test jam-transport-band::lib::shared_http_client_bounds_parallel_requests). WebSockets do not take permits. The default request timeout is DEFAULT_HTTP_TIMEOUT (15 seconds), and redirects are refused.
#The inbound pipeline
The components one inbound agent message passes through, and where it is stored.
Engine::run owns the subscription receivers. For each InboundMessage it discards messages with an empty id, skips ids already present in the routes index, and writes the message into that room's RoomSpool with push_wait. A spool is a private directory under <app-dir>/inbound-spool/<sha256(agent_id)>/<sha256(peer, room)>/, one JSON file per message named by a monotonic sequence, written to a .tmp file, fsynced, renamed, and followed by a directory fsync. Directories are mode 0700 and files 0600. A spool is bounded to 4,096 entries and 64 MiB across all rooms under one agent root (INBOUND_SPOOL_MAX_ENTRIES, INBOUND_SPOOL_MAX_BYTES). When the bound is reached, push_wait waits for capacity rather than dropping a message, which stalls the agent's whole Engine::run loop until a room worker frees space. On startup RoomSpool::discover finds existing spool directories, promotes valid .tmp files, and starts a worker for each (test engine::run_recovers_private_spooled_messages_without_a_new_room_event).
Each room has one inbound worker task (spawn_room_inbound_worker). It peeks the oldest spool entry, acquires one of eight engine-wide inbound_slots, and awaits handle_inbound. It removes the spool entry only when the message was consumed: either it was a system-marker message, or Engine::resolve now finds it in the routing index. Otherwise it waits 5 seconds and retries the same entry. This gives one ordered lane per room while allowing different rooms to progress concurrently (test lifecycle.rs::long_running_room_turn_does_not_block_unrelated_room_delivery_or_daemon_status).
handle_inbound performs these steps in order:
- Replay guard. An in-memory
inbound_attemptsset prevents two concurrent attempts for the same message id. A message for a room inretired_roomsis dropped with a warning. A message already inroutesis ignored (testengine::duplicate_inbound_is_not_delivered_or_marked_processing_twice). - System markers. Content that parses as a
[system:…]marker (jam_domain::parse_system_marker) is claimed and marked processed immediately and never queued (testengine::system_marker_inbound_is_auto_processed_never_enqueued). - Sender and mentions. The sender handle comes from
PeerCache::lookup_by_id, andrewrite_inbound_mentionsreplaces each@[[uuid]]token whose id the cache knows with@owner/handle. Unknown ids stay as tokens.PeerCacherefreshes from/agent/peersat most once per 5 seconds (MIN_REFRESH). - Attachments. Up to
jam_contract::MAX_MESSAGE_ATTACHMENTSattachments are downloaded with the agent's own credential under one per-message byte budget (max_artifact_bytes), verified against the platform's SHA-256 when one is present, and stored atomically throughAttachmentStore::store_batch. Transient download or storage failures retry within the room lane every 5 seconds (up to three attempts each); too many attachments, permanent download failures, digest mismatches, exhausted retries, or a budget that is already spent before the next download deliver the message with an "Attachments unavailable" note instead. One exception does not deliver: if a single download returns more bytes than the remaining budget,handle_inboundemits a warning ("the message remains pending") and returns before enqueue, so the spool entry is not consumed. Local paths, not remote filenames, are appended to the content, because filenames are participant-controlled. - Addressing flag.
addressed_to_useris computed from the raw content (@[[<owner user id>]]) before the rewrite. It is recorded for the UI's manual-ack affordance and gates nothing. - Durable enqueue. Under the
RoutingStatewrite lock, the engine re-checks retirement and duplicates, routes the room throughSessionRegistry::route, appends to the owning session's queue or the unrouted queue, and inserts the message intoroutes. Retirement takes the same lock, so a room cannot be retired between route selection and append (testengine::concurrent_retirement_reply_and_ack_share_one_routing_lock). - Processing claim. After the append,
claim_message_processingcallsPOST /api/v1/agent/chats/{room}/messages/{id}/processing. Band requires this claim before it accepts settlement. Claiming after the local write means a refused or lost response cannot create remote work with no local record. - Events. The engine stamps activity and emits
EventKind::Inbound. If no session owns the room it also emitsEventKind::QueuedUnroutedand stops. - Delivery. If the claim failed, the engine spawns
retry_processing_claim_and_deliver(bounded to 32 concurrent retry tasks, backoff 1 to 30 seconds, six attempts) and returns. Otherwise it stamps the session's activity atomic and awaitsDeliverer::deliver. A delivery error is reported asEventKind::RuntimeIssuewith sourceDeliveryand the message stays queued (testengine::inbound_enqueues_even_when_host_delivery_fails).
The engine never marks a message processed after delivery. Delivered to an inbox is not the same as handled, so the only transitions to processed are an explicit reply or acknowledgement (test engine::inbound_delivers_but_does_not_mark_processed).
#Routing to host sessions
A jam_core::Session binds exactly one room to one queue and one Deliverer. SessionRegistry keeps sessions by SessionId and a by_room index. insert evicts any previous owner of the room (the steal path) and releases a session's old room when the same id is re-inserted on a different room. Routing is exact-match only; there is no default session. These in-memory sessions are the engine's view of the durable HostSession bindings that Execution and lifecycle owns.
Physical queue location must equal routing, so reply and ack can remove a message from the queue that actually holds it. QueueBackend::reconcile enforces that before a worker is built (spawn_worker) and at 11 more call sites in manager.rs (attach, detach, steal, restore, window reconcile). On SQLite, SqliteStore::queue_rehome sets every row of the peer to the unrouted bucket and then claims each owned room for its session, in one transaction. On the legacy file backend, worker::reconcile_queue_dir reads every JSONL file, recomputes the desired layout, and rewrites only files whose contents changed. Both first purge obsolete Claude Code bind-prompt advisories (jam-bind-prompt: ids) for rooms that are now owned.
When a session is bound live, Engine::bind_session_reconciled runs the reconcile closure under the same write lock that handle_inbound holds across route and append, then inserts the session. A concurrent inbound either lands before the reconcile and is re-homed, or routes directly to the new session afterwards (test engine::live_bind_reconcile_excludes_concurrent_routing). The manager then calls Engine::reindex_queue for the session and unrouted queues and worker::redeliver, which reasserts the processing claim and redelivers each queued message (tests lifecycle.rs::unrouted_message_is_delivered_on_attach_and_ackable_without_ghost, lifecycle.rs::stealing_a_room_carries_its_queued_messages_to_the_new_owner).
#The unrouted queue
The unrouted bucket is the queue row set whose session column is the empty string (queue::UNROUTED), or _unrouted.jsonl on the file backend. A message lands there when its room has no bound session. The worker's supervise loop sees the Inbound event first and calls the manager's on_inbound_wake hook; if the room belongs to a dormant, idle-reaped owned runtime, the hook starts a wake and the following QueuedUnrouted event is swallowed. Otherwise queued_unrouted_warning turns it into a peer Warning ("no host session owns room …"). Unrouted rows are aged out at daemon start by run_maintenance(now, 90 days, 1000, 30 days) in jamd.rs: rows older than 30 days are deleted, then the bucket is capped to the newest 1,000 rows. Routed queues shrink only on settlement.
#Local notices
A local notice is a jam_domain::Message whose id starts with LOCAL_NOTICE_ID_PREFIX (jam-notice:). Engine::enqueue_local_notice queues and delivers it exactly like an inbound Band message, but nothing is posted, claimed, or settled remotely, and a queued notice with the same id is a no-op. claim_message_processing and Engine::ack skip Band for notice ids, and Engine::reply_parts refuses them with CoreError::LocalNoticeReply. The one producer on main is the shared-board bridge in crates/jam-manager/src/manager/board_bridge.rs, which tells an agent that a task was assigned to it (test engine::local_notice_is_delivered_and_acked_without_band).
#Lifecycle of a queued message
This state machine answers: which states can one inbound agent message be in, and which event moves it between them?
QueuedRouted includes a message whose processing claim failed; it moves to Delivered when the bounded claim retry or a later redelivery succeeds. Settled means Band recorded processed and the row left the local queue. If Band settles but local removal fails, the row stays in Delivered and is redelivered. The state names are this chapter's labels; only QueuedUnrouted is also a code identifier, as an engine event.
#One inbound message to an owned runtime and back
The durable commit points and acknowledgements when a Band message reaches an owned Codex runtime and the reply goes back.
Owned runtimes settle through the worker, not through the CLI. In worker::project_runtime_event, a TurnDispositionStaged event is buffered per turn key (at most TurnDisposition::MAX_PER_TURN per turn and 256 open turns). On TurnComplete with outcome Complete, commit_turn_disposition commits each staged disposition in order: Reply goes through post_assistant_reply (which splits text at chunk::BAND_CONTENT_CAP, 16,000 bytes, and calls Engine::reply_parts), and NoReply calls Engine::ack. Without a staged disposition, a buffered final assistant text is posted as a compatibility fallback unless the source was an agent message and explicit dispositions are required. A failed or cancelled turn consumes nothing (tests lifecycle.rs::explicit_reply_posts_tool_text_and_consumes_agent_inbound, lifecycle.rs::staged_disposition_on_failed_or_cancelled_turn_does_not_consume_agent_inbound, lifecycle.rs::runtime_final_reply_failure_keeps_source_queued_for_retry). Provider details are in Providers and attached agents.
Attached agents settle through Control::reply and Control::ack (the band reply and band ack CLI, or the desktop). Manager::reply and Manager::ack call the same Engine methods, then settle onboarding bookkeeping and close any Copilot bridge job for the message. Only Manager::reply also records outbound analytics (record_outbound).
#Outbound send and mention resolution
How outbound text becomes a Band message with structured mentions.
Band requires at least one mention on every message, and a mention is the only way to wake an agent, so the engine refuses to send text that resolves to nobody. Explicit recipient_ids are resolved against the current room roster. @handle text is resolved against the agent's peer index, retried once after MENTION_RETRY (500 ms) only when nothing resolved, then against room participants for members that the peer index has not yet propagated. An ambiguous handle is never guessed; it stays literal and produces a "skipped" warning in SendResult.warnings (test engine::send_falls_back_to_room_participant_when_not_in_peer_index). The agent's own id and handle are always excluded. prepare_outbound_mentions rewrites a short human handle that is a prefix of a longer owner/name handle in the same text so Band's substitution cannot hit the wrong one. project_outbound_mentions returns the canonical @[[id]] body so an optimistic local row matches Band's later echo (test engine::projected_send_returns_the_canonical_mention_token_body).
Engine::send_to_participants is a separate exact-recipient path used for directed questions. It does not parse @handle text, so a question that names an agent does not wake that agent (test engine::directed_question_does_not_wake_an_agent_named_in_its_text).
Engine::reply_parts adds the reply rules. It re-claims processing first (a queue row can outlive a refused claim), reads the room roster once, drops mentions of anyone no longer in the room, and auto-mentions the original sender only if the sender is still present. If the sender left, it posts to the sole remaining participant or to explicitly mentioned current members, and otherwise fails with NoResolvableMention (tests engine::reply_to_removed_sender_uses_the_sole_remaining_participant, engine::reply_to_removed_sender_without_a_unique_remaining_recipient_is_rejected). All parts are posted before mark_processed; the source is removed from the queue only after Band accepts settlement. A settlement or local-cleanup failure after posting is reported as a warning on the last SendResult, and the source stays queued (test engine::reply_reports_failed_remote_settlement_after_sending).
The human's own messages use a different path. Manager::room_send resolves the human's mentions, refreshes attachment uploads when the ff_file_transfer platform flag is on, and calls UserApi::send_room_message with the account's user credential. Its echo arrives through the human socket like any other message.
#The human room feed and room message cache
Manager::spawn_room_feeds starts one run_account_feed task per account (through spawn_account_room_feed). Before the feed starts, that task lists the account's rooms and seeds BandClient::set_human_room_directory and the manager's human_rooms set (which the outbound analytics funnel uses to tell agent-only rooms from human rooms). run_account_feed then seeds the in-memory RoomProjections from the encrypted store (room_message_mutation_sequences, list_cached_room_messages) before it subscribes.
Every room event from the human feed is published under the sentinel peer human/<profile>. For each RoomMutationInbound, run_account_feed_inner:
- Maps it to a
RoomMessageMutationthroughroom_message_from_inbound, which lifts tool name, tool call id, system fields, mention names, and thought-stream ids from Band metadata. - Calls
apply_projection_mutation_once. An upsert identical to the projected row is ignored; otherwise the room'ssequenceincrements, the retained window is trimmed toROOM_MESSAGE_CACHE_PER_ROOM_LIMIT(100) newest messages, and the mutation is appended to a journal of the last 512 mutations. - Publishes
EventKind::RoomMessageMutation { account_instance, sequence, mutation }before persistence (testroomfeed::live_message_is_published_before_persistence_starts). - Queues stats buckets and a
RoomPersistenceOpon bounded channels of 256. If the persistence channel is full, the operation is coalesced into a bounded dirty checkpoint that the worker flushes every 100 ms (testroomfeed::full_queue_keeps_only_a_bounded_checkpoint_for_ordinary_messages).
run_room_persistence_worker writes through Store::apply_room_message_mutation_with_attention, which rejects a sequence at or below the room's stored watermark and, for a direct mention of the account human, stages an attention item in the same transaction. The retained cache lives in room_message_cache_state and room_message_cache_items, keyed by (account_instance, chat_id).
Control::room_messages, room_message_snapshot, and refresh_room_messages read that cache and Band history. A refresh captures the projection's sequence, fetches the newest page, and merges it with live mutations that arrived after the capture through merge_snapshot_with_live_mutations. If the snapshot predates the journal floor, the merge returns None and the caller must refetch. Concurrent refreshes for the same (account, generation, room) share one flight (room_message_refreshes).
Messages an owned runtime posts are published immediately as an unsequenced EventKind::RoomMessage through RuntimeRoomMessageSink and Manager::record_runtime_room_message, which also stages a direct-mention attention item when the text mentions the human (test lifecycle.rs::local_agent_send_persists_a_direct_human_mention_without_a_feed_echo). The desktop inserts it by id and replaces it when Band's sequenced echo arrives.
Room add, remove, title, participant, presence, and activity events go through publish_fenced_membership_event and publish_room_event_inner. Presence and activity events carry an observation generation from HumanRoomWatch and a Phoenix join ref; Fanout rejects an event whose generation no longer matches the watched room (test roomfeed::stale_membership_and_delete_events_cannot_cross_a_room_reobservation). A room removal clears the room projection and queues RoomPersistenceOp::DeleteRoom, which retires the room's cache and its local system rows. Marker events (ContactsChanged, AgentRoomsChanged, BoardChanged) carry no payload; the feed calls a manager callback that refetches the authoritative data over REST.
#Lifecycle and system rows
lifecycle_events.rs has two jobs. First, membership_room_message builds a local row in Band's own participant message shape. roomfeed::schedule_membership_join_fallback waits PLATFORM_MEMBERSHIP_GRACE (2 seconds) after a participant_added event and, only if no platform membership row for that participant arrived within 15 seconds of the event, publishes the local row as an EventKind::RoomMessage and stores it with Store::append_room_system_message (table room_system_messages, 200 rows per room). This covers Band deployments that send the socket event but no durable row. Manager::merge_room_system_messages merges these rows into history reads, and a platform row supersedes the local one.
Second, LifecycleLatch gates online and offline transitions per (peer, room). Every live-bind and spawn path passes an explicit OnlineCause (25 OnlineCause:: variant references in manager.rs); only FreshSpawn and FirstLiveBind may announce. On main, Manager::emit_agent_online and emit_agent_offline update the latch only and create no room message. The offline result still decides whether the liveness sweep publishes a "stopped reporting" peer Warning. JAM_SYSTEM_EVENT_DEBOUNCE_SECS overrides the 300-second online floor.
#Room connectivity
External-agent connectivity in a room is observed through the human socket. When a watched room's channels join, the socket emits RoomEvent::Connectivity(RoomConnectivityEvent::Observed). Manager::accept_room_connectivity creates one ObservedRoom per (profile, account_instance, room) and spawns run_connectivity_room, which hydrates the roster over REST within RECONCILIATION_TIMEOUT (20 seconds). Later Changed and MembershipChanged events apply only if their revision is newer than applied_revision. A membership removal or block suspends the projection until an authorized snapshot arrives, and local agents owned by this account are excluded. The result is published as EventKind::RoomConnectivityChanged. The module documents its lock order at the top: account generation, then local ownership, then transport observation, then room state, then fan-out. No lock spans an await. with_participant_mutation fences observed rosters before and after a local add or remove, so a roster read captured before the change cannot satisfy a later read (test room_connectivity::connectivity_observed_read_rejects_a_cycle_that_predates_a_local_membership_change).
#Working-state lease publication
Band shows other participants which agents are working in a room through a content-free lease: POST /api/v1/agent/chats/{room}/activity with {"working": true|false}. WorkingStatePublisher publishes it per room execution without blocking the runtime path. begin_turn and finish_turn only update a watch value; a detached task (run_publisher) serializes the network calls. It refreshes working: true every 3 seconds while a turn is open (WORKING_REFRESH_INTERVAL), bounds each report to 2 seconds (WORKING_REPORT_TIMEOUT), and relies on the platform's 10-second TTL (PLATFORM_WORKING_TTL) as the failure backstop, so there is no per-request retry. Compile-time assertions keep the timeout below the interval and the interval below half the TTL. A completion from an older turn cannot clear a newer one, and a terminal clear is always attempted before the next true. An Unsupported response pauses publication until the execution epoch changes. Shutdown waits at most WORKING_SHUTDOWN_DRAIN (2.25 seconds). JAM_BAND_ACTIVITY_DISABLED disables publication. Publishers are created in worker.rs for owned-runtime turns and in manager.rs for attached-agent activity reports (AttachedWorkingEntry).
#Event fan-out and the client event stream
Fanout is the single point through which every peer event flows. Engine events reach it through worker::supervise, which adjusts some kinds before publishing: it degrades Connected to Degraded when a host reports a readiness warning, routes ChatAdded and ChatRemoved to the room-membership sink, swallows the engine's empty ContactsChanged placeholder after scheduling the real refresh, and rewrites QueuedUnrouted as described in The unrouted queue. The human feed, room connectivity, questions, tasks, and many manager paths publish directly (61 fan.publish call sites in manager.rs alone).
Fanout::publish holds one state mutex while it updates the live projections (peer states, current warnings, runtime issues, participant caches, presence, remote activity, human removal tombstones), appends one line to the peer's rotating log, and sends on a tokio::sync::broadcast channel of capacity 256. Broadcast delivery is best-effort: a subscriber that falls behind by more than 256 events receives RecvError::Lagged. Appending and broadcasting under the same lock (emit_locked) guarantees that a newer presence observation cannot be emitted before an accepted older one. Peer logs live at <app-dir>/logs/<profile>/<scope>.log, rotate at 5 MiB to a single .1 backup, and at most 16 are open at once (LRU).
How a human-socket message reaches the desktop store, and how the chain recovers from a lagging subscriber and from a dropped connection.
Manager::event_stream(cancel, replay_rooms) wraps a broadcast receiver in an mpsc channel of 256. For Control::events (IPC clients) replay_rooms is true: the stream starts with RoomLiveStateReset followed by the current presence and activity snapshots from subscribe_with_room_snapshot, taken under the fan-out lock. On Lagged, it resubscribes, takes a fresh snapshot, and pushes StreamResync in front (test lifecycle.rs::lagging_desktop_subscription_replaces_live_room_state). Manager::events uses replay_rooms false: on lag it forwards a bare StreamResync (through forward_event_stream_result) without a snapshot. Manager::wait_inbox does not use either stream; it subscribes to Fanout directly and, on RecvError::Lagged, re-reads the room's durable queue.
The daemon route GET /v1/events (jam-daemon::events) subscribes first, then calls Control::list and prepends one PeerAdded per peer, then streams NDJSON. A failed snapshot does not stop the live feed (test roundtrip.rs::events_continue_when_peer_snapshot_fails). The stream is tied to the daemon shutdown token so graceful shutdown cannot hang on it (test resilience.rs::shutdown_completes_while_an_events_stream_is_open), and it bypasses the unary IPC budget (test resilience.rs::event_stream_bypasses_a_full_unary_budget).
jam-client's event_stream reads lines through LineReader with a 1 MiB per-line cap (MAX_EVENT_LINE_BYTES). An over-cap line is dropped and replaced by an in-band StreamResync. jam_ipc::events::relay pumps the stream into a Sink. It translates StreamResync into Sink::resync and does not forward it. When the stream ends, it arms a gap, waits RELAY_RETRY (2 seconds, apps/desktop/src-tauri/src/conn.rs), resubscribes, and fires resync once before the first resumed event (tests ipc.rs::relay_signals_resync_once_after_a_gap_when_events_resume, ipc.rs::relay_does_not_resync_while_daemon_stays_down). The relay never emits connection status; the separate monitor owns that.
In the desktop, TauriSink::event emits DaemonEvent to the WebView first and then raises notifications. TauriSink::resync resets notification state and emits ResyncEvent. App.tsx subscribes to both; on ResyncEvent it resets lane revisions, inbox generations, and participant snapshots and calls rehydrate("reconnect"). applyEvent in stores/daemon.ts handles room_live_state_reset by clearing all platform-derived live room state, and room_message_mutation by ignoring any sequence not greater than the last one seen for ${accountInstance}:${chatId}. The full kind list and desktop coverage are in Daemon event kinds.
#State and ownership
| Aggregate | Key and scope | Authority | Mutation coordinator | Durable representation | Projections | Freshness or fence | Recovery | Deletion authority |
|---|---|---|---|---|---|---|---|---|
| Band message processing state | Band message id in a room | Band | Engine via Bridge::mark_processing and mark_processed |
Band platform | Pending list returned to catch-up | Only processed excludes a message from catch-up |
Catch-up replays every unprocessed message on each connect | Explicit reply or ack only; a retired room's local work is dropped because Band returns 404 |
| Inbound room spool entry | (agent, peer, room) directory, sequence file |
Engine::run for that agent |
One room worker per room | JSON files under <app-dir>/inbound-spool/, mode 0600 |
None | Message id dedup on push; removed only when consumed | RoomSpool::discover on engine start; .tmp promotion |
Room worker after consumption; RoomSpool::cleanup for a stopped worker |
| Durable inbox queue row | (profile, scope, message_id); bucket column session |
QueueBackend chosen by jamd |
Engine appends under the routing lock; reconcile re-homes |
queue_messages (SQLCipher) or legacy queue/<profile>/<scope>/*.jsonl |
Engine.routes index, jam inbox, desktop inbox |
Primary key dedup; first write wins | Reindexed at engine start; redelivered on spawn and live bind | Reply or ack; retired-room settlement; peer purge or archive; unrouted maintenance |
| Routing state | One engine per agent identity | Manager binding decisions |
Engine RoutingState write lock |
Memory; rebuilt from HostSession rows and queues |
inbox_room, delivery target |
Live bind reconcile runs under the same lock as route plus append | Rebuilt by spawn_worker |
unbind_session, retire_room |
| Room retirement obligation | room-retirement: setting per (profile, scope, room) |
Manager |
Manager::admit_room_retirement under room_retirement_transition_gate |
Store setting with phase Retiring or Retired |
Engine retired_rooms set |
Persisted before the engine fences the room; binding generation captured | Loaded at rebuild; retried when persistence failed | Manager after rebind or completed retirement |
| Human room projection | (account_instance, chat_id) |
Band human socket | roomfeed::apply_projection_mutation_once |
Memory, seeded from the room cache | RoomMessageMutation events |
Monotonic per-room sequence; 512-entry journal floor |
Seeded from store at feed start | Room removal clears it |
| Room message cache | (account_instance, chat_id), 100 messages |
Band | Room persistence worker; refresh flights | room_message_cache_state, room_message_cache_items |
room_messages, room_message_snapshot |
Stored last_mutation_sequence rejects older writes |
Reconnect and open-room reconciliation over REST | retire_room_message_cache on room removal |
| Local system rows | (profile, chat_id, message_id), 200 per room |
Jam, superseded by Band rows | schedule_membership_join_fallback |
room_system_messages |
Merged into history reads | Platform row within 15 seconds wins | None needed | delete_room_system_messages on room removal |
| Fan-out live state | Peer key and room | Latest accepted event | Fanout state mutex |
Memory only | status, participants, IPC snapshot |
Observation and join generations on presence and activity | Reset on stream subscribe; transport resets | Fanout::forget, PeerRemoved |
| Peer event log | (profile, scope) file |
Fanout |
Fanout::append_line |
<app-dir>/logs/<profile>/<scope>.log plus .1 |
Control::logs |
Best-effort append | None | Rotation at 5 MiB |
| Human room watch set | One per BandClient (account) |
Manager directory refresh plus local pins | HumanRoomWatch |
Memory | Human socket channel joins | selection_generation, per-room observation generation |
Rebuilt from the room directory | unwatch_human_room, departed rooms |
| Observed room connectivity | (profile, account_instance, room) |
Band roster and push | Manager::accept_room_connectivity |
Memory | RoomConnectivityChanged |
applied_revision, membership_generation |
Re-hydrated per observation | Observation cancellation |
| Working-state lease | Room execution | Band (TTL) | WorkingStatePublisher actor |
Band only | Other participants' RemoteActivity |
Generation per turn, execution epoch | TTL expiry | Terminal clear or TTL |
| Lifecycle latch | (peer, room) |
Manager |
lifecycle_latch mutex |
Memory; cleared by restart | Liveness sweep warning | Transition latch plus 300-second floor | Restart re-arms | Not applicable |
#Contracts
#Traits this subsystem provides
jam_core::Bridge(14 methods, blanket-implemented for everyAgentApi). The engine's consumer-side slice of Band: peers, send, attachments, events, working lease, processing claims, chats, history, participants. Optional methods default to typedUnsupportedorTransporterrors.jam_core::Deliverer. One method,deliver(peer, msg). The engine awaits it inside the room's ordered lane. An error leaves the message queued. A no-op implementation (jam_host::generic::Generic) is the pull baseline: the message is already in the queue.jam_core::Queue.append,read,find_by_message_id,remove_by_message_ids.appendmust deduplicate by message id andreadmust return arrival order (FileQueueandSqliteStore::queue_read, which orders byenqueued_at, rowid).jam_core::AttachmentStoreandRoomRetirementAdmission. Injected by the worker so the engine can store verified attachment bytes and persist a room-removal obligation without depending on the store.jam_manager::QueueBackend.session_queue,unrouted_queue,all_queues,reconcile,purge,archive_into,restore_from. Re-import is safe because rows deduplicate by(profile, scope, message_id)(testlifecycle.rs::sqlite_archive_then_restore_preserves_the_queue).jam_manager::ManagerTransport. The union of the four transport roles withinto_bridge,into_subscriber,into_user_api, andinto_human_subscriber, blanket-implemented.
#Control methods
The messaging methods on jam_contract::Control are events, inbox, inbox_unrouted, inbox_room, wait_inbox, send, send_with_artifacts, reply, ack, chat_history, participants, add_participant, remove_participant, room_messages, room_message_snapshot, refresh_room_messages, room_participants, and room_send. Their routes, Tauri commands, and CLI use are listed in the Control API crosswalk. send has no route of its own: jam-client posts it to the same POST /v1/send route, whose daemon handler calls send_with_artifacts. events (GET /v1/events) and wait_inbox (POST /v1/waitInbox) are held-open streaming routes without the unary timeout.
#Band protocol consumed
| Surface | Endpoint or topic | Used by |
|---|---|---|
| Agent send | POST /api/v1/agent/chats/{room}/messages |
Engine send and reply |
| Processing claim | POST /api/v1/agent/chats/{room}/messages/{id}/processing |
claim_message_processing |
| Settlement | POST /api/v1/agent/chats/{room}/messages/{id}/processed |
reply_parts, ack, system markers |
| Catch-up | GET /api/v1/agent/chats/{room}/messages (cursor pages, no status filter) |
catch_up_pending_messages |
| Agent history | GET /api/v1/agent/chats/{room}/context |
Engine::chat_history |
| Working lease | POST /api/v1/agent/chats/{room}/activity |
WorkingStatePublisher |
| Agent channels | agent_rooms:{id}, agent_contacts:{id}, chat_room:{room}, room_participants:{room}, room_activity:{room} |
Agent socket |
| Human channels | user_rooms:{uuid}, user_contacts:{uuid}, plus the three room channels for up to 48 rooms |
Human socket |
#Semantic guarantees enforced by code
- At-least-once delivery to the host. A message is redelivered after a worker respawn, a live bind, or a daemon restart until it is settled (tests
lifecycle.rs::spawn_redelivers_unacked_queue_and_is_idempotent,lifecycle.rs::redelivery_failure_does_not_abort_the_rest). Hosts must tolerate duplicates; the attached Claude Code mailbox and the Copilot bridge deduplicate by Band message id. - Deduplication by message id at four layers: the engine
routesindex, theinbound_attemptsset, the spool push, and the queue primary key. - Settlement only on explicit intent. Nothing marks a message processed except
reply_parts,ack, and the system-marker guard. - Per-room order for the ordinary path: one spool and one worker per room. See the exceptions in Invariants.
- Best-effort event stream.
Control::eventsdocuments that events can drop and that query methods remain the truth. Code guarantees aStreamResyncor relay resync after every observed gap, aRoomLiveStateResetplus snapshots on every subscription, and a monotonicsequenceper room for human-feed mutations. - Bounded Band usage. At most 16 concurrent REST requests per daemon, 8 idle sockets per host, 15-second request timeout, 2-second working-report timeout.
#Invariants
| Rule | Enforced by | Known exceptions |
|---|---|---|
| An inbound message is durable locally before Band is told it is being processed. | handle_inbound appends before claim_message_processing; spool write precedes both. Test engine::inbound_stays_queued_and_does_not_reach_the_host_when_processing_claim_fails. |
System-marker messages are claimed without a queue row by design. |
| Only an explicit reply or ack marks a message processed. | No mark_processed call after deliver in handle_inbound. Test engine::inbound_does_not_mark_processed_even_when_user_addressed. |
System markers; retired rooms drop local work without a remote write. |
| A message is removed from the local queue only after Band accepts settlement. | reply_parts and ack call mark_processed before remove_by_message_ids. Test engine::ack_propagates_failed_remote_settlement_and_keeps_inbound_queued. |
Local notices have no remote settlement. |
| No message is routed to a guessed session. | SessionRegistry::route is exact-match; unowned rooms go to the unrouted bucket. Test engine::inbound_to_unowned_room_is_queued_unrouted_for_worker_classification. |
None. |
| Queue location equals routing before ack or reply. | QueueBackend::reconcile at spawn and bind; bind_session_reconciled runs it under the routing write lock. Tests engine::live_bind_reconcile_excludes_concurrent_routing, jam-store::queue_rehome_moves_messages_to_owning_session_else_unrouted. |
The file backend reconcile is best-effort per file and logs failures. |
| Messages in one room reach the host in arrival order. | One RoomSpool and one worker task per room. Test engine::inbound_appends_multiple_in_order. |
A message whose processing claim failed is delivered later by a detached retry task, after later messages. Messages drained by try_recv when a room removal arrives bypass the spool. |
| A retired room admits no new work. | retired_rooms checked before and under the routing lock; admission persisted before the fence. Tests engine::retired_room_rejects_late_inbound_until_explicit_rebind, lifecycle.rs::retired_room_fails_closed_for_late_inbound_without_an_unowned_queue_item. |
Engine::reopen_room after an authoritative room-added replacement. |
| Outbound text never addresses an ambiguous or unresolved handle. | match_unique_handle returns Ambiguous; unresolved handles become warnings. Test engine::send_skips_unresolvable_but_proceeds_with_real_mention. |
None. |
| Every IPC subscription starts from a clean live room state. | subscribe_with_room_snapshot prepends RoomLiveStateReset. Test lifecycle.rs::lagging_desktop_subscription_replaces_live_room_state. |
Internal waiters (replay_rooms false). |
| A lost event is always followed by a re-hydrate signal. | forward_event_stream_result, LineReader over-cap handling, relay gap tracking. Tests ipc.rs::relay_translates_in_band_stream_resync_into_a_sink_resync, ipc.rs::relay_resubscribes_after_stream_end. |
None found. |
| Human-feed mutations are applied at most once and in order per room. | Projection sequence; store watermark; desktop roomMutationSequenceByChat. Test roomfeed::identical_local_send_echo_does_not_advance_the_room_projection. |
Unsequenced local RoomMessage rows. |
| One process-wide Band REST budget. | Single BandHttpClient constructed in jamd.rs; Semaphore of 16. Test shared_http_client_bounds_parallel_requests. |
Public avatar fetches use a separate credential-free pool. |
| A superseded human connection does not reconnect. | SessionEnd::Superseded returns from drive; feed returns after on_auth_changed. Tests supersede.rs::supersede_is_terminal_and_does_not_reconnect, roomfeed::feed_reports_superseded_and_stops_without_reconnecting. |
Agent sockets do not handle supersede. |
| Working-lease cadence stays inside the platform TTL. | Compile-time assert! in working_state.rs. Test working_state::policy_values_preserve_the_platform_ttl_margin. |
None. |
| Room connectivity locks are taken in one order and never across an await. | Module comment and structure in room_connectivity.rs. |
Not enforced by a type or test; review only. |
#Failure and recovery
Daemon crash or restart. Spool files survive and are rediscovered by RoomSpool::discover. Queue rows survive in SQLCipher. On rebuild, each worker reconciles queues, the engine rebuilds its routes index from every session queue plus the unrouted queue, spawn_worker redelivers each session's backlog through redeliver (which reasserts the processing claim first), and the agent socket's catch-up replays every message Band still considers unprocessed. Duplicates collapse on message id. Room projections are reseeded from the room cache and its stored sequences so sequences stay monotonic across restarts. Retirement obligations are reloaded from settings. The lifecycle latch, fan-out live state, and working leases are memory-only and restart empty; leases expire by TTL.
Crash between Band settlement and local removal. If mark_processed succeeded but remove_by_message_ids failed, the reply's warning says so and the row stays queued and indexed. On the next spawn worker::redeliver reasserts the processing claim before delivering, and skips the message if the claim fails. Inferred: the outcome depends on whether Band accepts a processing claim for a message it already records as processed, which no Jam code or test establishes. If Band accepts it, the agent sees the message again and can reply twice, because nothing in the engine prevents the second reply. If Band refuses it, the row is never redelivered and stays queued until an explicit ack, room retirement, or peer purge.
Crash after posting a reply but before settlement. The reply exists on Band and the source is still processing. On restart the source is redelivered and the agent may reply again. The engine posts all parts before settling so a partial multi-part reply is retried as a whole (test engine::reply_parts_posts_ordered_parts_and_processes_the_source_once).
Processing claim failures. A retryable claim failure (HTTP 408, 409, 425, 429, 5xx, expired authentication, unavailable, or transport errors; retryable_processing_claim_error) starts a bounded retry that re-resolves the current room owner before delivering. A permanent failure or exhausted retries publishes a peer Warning; the message stays queued and is reclaimed on the next redelivery (test engine::a_transient_processing_claim_failure_retries_into_the_live_session).
Host delivery failure. The engine reports a RuntimeIssue for the exact session and keeps the message. Redelivery happens on respawn or rebind. The engine does not retry delivery on a timer.
Enqueue failure. handle_inbound emits EventKind::Error and returns. The spool entry is not consumed, so the room worker retries it every 5 seconds. Inferred: a persistent store failure blocks that room's lane indefinitely while other rooms continue.
Agent socket disconnect. drive publishes Reconnecting, emits RemoteActivityReset, backs off, reconnects with a fresh credential, rejoins every room, and replays pending messages. The engine maps socket states to PeerState events. A catch-up that times out for one room logs a warning and continues live; the remote queue is preserved for the next connect.
Human socket disconnect. drive emits ParticipantPresenceReset and RemoteActivityReset (ordered after every event from the old connection), backs off, reconnects, and reconciles channels. When the socket reports Up again, the feed calls on_reconnected: the manager publishes PeerState::Connected for human/<profile>, replays pending shared-board operations, and runs rooms_at_with_cancel, which refreshes the room directory and reconciles the cached head of every room with activity in the last 30 days whose last_message_at differs from the cache. If the whole subscription ends, run_account_feed_inner re-subscribes with its own backoff of 3 seconds doubling to 5 minutes; a subscription that dropped within 60 seconds counts as a failure.
Supersede. The human feed marks the account AccountAuthState::Superseded and stops. It stays stopped until the user signs in again.
Room removed while messages are in flight. On RoomEvent::Removed or HumanRoomDeleted, Engine::run first handles every message already in the transport channel, then marks the room's worker to retire when empty and waits for it, then takes the retirement gate, persists the obligation through RoomRetirementAdmission, and fences the room with retire_room. If persistence fails, the route stays active and the manager retries. Manager::settle_retired_room_queue then removes the room's rows from every queue bucket without a remote write, because Band returns 404 for a deleted room. A failed local cleanup keeps the binding and is retried (test lifecycle.rs::deleted_room_queue_cleanup_failure_defers_and_retries_binding_retirement). Inferred: the wait for the room worker spins with yield_now inside the main select arm, so a spool entry that can never be consumed (for example an attachment over the byte limit, which returns before enqueue) would keep Engine::run from processing any other event for that agent until cancellation.
Late message for a retired room. handle_inbound drops it without a queue row or claim. Inferred: because the spool entry is not consumed, it stays in the room spool and is re-evaluated every 5 seconds, counting against the agent's spool budget, until the spool is cleaned by a later stopped-worker path or the worker is cancelled.
Event stream lag or drop. Covered in Event fan-out and the client event stream. The desktop re-hydrates through the same rehydrate path used at startup and on reconnect. Room message views re-read room_message_snapshot and refresh on open.
Persistence backpressure in the room feed. A full persistence channel coalesces operations into a bounded dirty set flushed every 100 ms. Direct-mention attention retries are bounded per account (test roomfeed::direct_mention_retry_state_is_account_bounded_and_compact). The live event is never delayed by persistence.
#Extension points
#Adding a coding-agent harness
A harness plugs into this subsystem only through jam_core::Deliverer, which jam_host::Host extends. There are seven production implementations: CodexAppServer, CopilotSdk, CopilotCli, ClaudeCode, ClaudeCodeOwned, AcpHost, and Generic, plus the manager's QuestionGate wrapper, which intercepts answered-question messages before they reach the inner host. To add one:
- Implement
Deliverer::deliverin the new adapter undercrates/jam-host/src/. It must return quickly: the room lane and one of the engine's eightinbound_slotsare held until it returns. The owned adapters only send a command on an internal channel. - Tolerate duplicate delivery of the same
Message.message_idafter respawn or rebind. - Decide how the harness settles. Owned runtimes emit
AgentEventKind::TurnDispositionStagedwith asource_message_idand thenTurnComplete;worker::project_runtime_eventandcommit_turn_dispositiondo the rest. Attached harnesses must callband replyorband ack, which reachControl::replyandControl::ack. - If the harness can report turn start and end, the worker's
RuntimeWorkingPublisherspublishes the working lease automatically fromTurnStartedandTurnComplete. An attached harness reports through activity telemetry, whichmanager.rsmaps to anAttachedWorkingEntry. - If the harness needs synthetic advisories in the queue, give them a reserved id prefix. The only existing one besides
jam-notice:isjam_host::claudecode::BIND_PROMPT_ID_PREFIX, whichjam-manager/src/queue.rsand eight code references inmanager.rsmatch by string prefix.
Host construction, session binding, and teardown are owned by Execution and lifecycle, not by this subsystem.
#Adding a messaging feature
- A new Band room event. Add a
RoomEventvariant injam-transport/src/lib.rs(20 today). Produce it insocket.rs::handle_textand orhuman_socket.rs. Handle it inEngine::handle_room_event(exhaustive match) androomfeed::publish_room_event_inner(exhaustive match).run_account_feed_innerhas explicit arms for seven variants and a catch-all for the rest, so a new variant silently takes the membership path unless you add an arm. - A new daemon event kind. Add the variant to
jam_domain::EventKind(59 today), add its line tofanout::format_event_line(exhaustive), decide whetherFanout::publishmust cache it, runjust bindings, add a case toapplyEventinapps/desktop/src/stores/daemon.ts, and regenerate Daemon event kinds.TauriSink::eventmatches with a wildcard and needs no change unless the event should notify. - A new outbound send shape. Add a method on
Enginethat ends inBridge::send_message_with_attachments, reuseprepare_outbound_mentionsandproject_outbound_mentions, and callmark_activity. Then follow the nine-step vertical inAGENTS.md("Add a Tauri command") for theControlmethod. - A new transport. Implement
UserApi(39 methods, most with defaults),AgentApi(22),Subscriber(1), andHumanSubscriber(11).ManagerTransportis then blanket-implemented. Replace thenew_transportclosure inbins/jam/src/jamd.rs. The fakes incrates/jam-manager/tests/lifecycle.rs,discovery_band.rs, andoauth_refresh.rsare complete examples. - A new queue backend. Implement
jam_core::QueueandQueueBackend, includingreconcileso physical location equals routing, and select it injamd.rsnext toSqliteQueueBackend.
#Refactor notes
File size and mixing of concerns. crates/jam-manager/src/manager.rs is 58,145 lines and contains messaging entry points (send_with_artifacts, reply, ack, room_send, room_messages, refresh_room_messages, event_stream, wait_inbox), feed spawning, lifecycle latch calls, room retirement, and unrelated subsystems. worker.rs is 10,609 lines and mixes the engine event pump, runtime event projection, disposition commit, usage ingestion, and managed-sandbox context. jam-core/src/engine.rs has about 2,900 production lines, of which handle_inbound alone contains the attachment retry loop with seven copies of the "Attachments unavailable" note. roomfeed.rs has about 3,000 production lines and human_socket.rs about 2,000.
Two queue reconcile implementations. SQLite re-homes in SqliteStore::queue_rehome; the legacy file backend re-homes in worker::reconcile_queue_dir, which lives in the worker module rather than in queue.rs. FileQueue::append reads the whole file to deduplicate. The default store is SQLite (JAM_STORE unset), so the file path is legacy but still reachable and still used directly by tests.
Spool cost and at-rest protection. RoomSpool::push scans every file under the agent's spool root to compute usage and reads every file in the room directory to deduplicate, under a process-wide static lock (INBOUND_SPOOL_BUDGET_LOCK). Spool files are plain JSON protected only by file mode, while the queue they feed is in the encrypted store. Message content therefore exists unencrypted on disk between receipt and enqueue.
Ordering exceptions. The processing-claim retry path delivers from a detached task after later messages in the same room, and the removal drain calls handle_inbound directly from the main loop, outside the room lane and the inbound_slots bound.
Boundary exceptions. jam-core imports jam_contract::MAX_MESSAGE_ATTACHMENTS, although its crate comment describes dependencies on Bridge, Deliverer, and Queue only. jam-manager/src/queue.rs depends on the provider-specific constant jam_host::claudecode::BIND_PROMPT_ID_PREFIX, and manager.rs references it eight times in code, so a Claude Code advisory convention is part of generic queue reconciliation. jam-host depends on jam-core (optional, behind runtime-providers) for Deliverer.
Events used as control signals. The engine emits an empty EventKind::ContactsChanged placeholder that the worker swallows to trigger a refresh, and QueuedUnrouted is rewritten into a Warning or dropped by the worker. RoomEvent mixes agent and human-only variants; Engine::handle_room_event returns without emitting for five of them. jam_transport::Transport is defined and blanket-implemented but has no users outside its definition; ManagerTransport is the union actually used.
Wide callback surfaces. run_account_feed takes 17 parameters, including six callbacks into the manager. WorkerHooks carries 28 fields. Both are the seams a refactor must preserve or replace.
Lifecycle latch. 25 OnlineCause references feed a latch whose online result is discarded; only the offline result gates a warning. The engine's system-marker guard comment still refers to lifecycle rows emitted locally.
Single broadcast channel. All peers and all human feeds share one broadcast channel of 256 events. One burst (for example many presence changes) can lag every slow subscriber at once. Correctness then depends on every consumer handling StreamResync and re-hydrating.
Stale comments that disagree with code. The crate docs of jam-transport/src/lib.rs and jam-transport-band/src/lib.rs still say the Subscriber WebSocket "lands in P2b"; it is implemented in socket.rs. The jam-core crate doc lists only Bridge, Deliverer, and Queue as dependencies, but the crate also depends on jam-contract. SessionRegistry::iter says reply and ack use it to find a message's session; they use the routes index through Engine::resolve. The FileQueue module doc says one engine is the only writer per peer, but reconcile_queue_dir, archive, restore, and retirement settlement also write queue files.
Test concentration. crates/jam-manager/tests/lifecycle.rs is 68,007 lines and holds most end-to-end messaging regressions named in this chapter. Its fake transport implements every Band trait and is the practical contract test for a transport refactor.