Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
56 commits
Select commit Hold shift + click to select a range
87a59c0
RED: no code in busbar-plane-mcp routes a call to another plane or im…
MattJackson Oct 7, 2026
c7e4875
busbar-plane-mcp: delete the unserved Plane/SessionPlane impl and the…
MattJackson Oct 7, 2026
6bc5958
RED: an mcp tool result reaches its caller as the upstream sent it (f…
MattJackson Oct 7, 2026
34df85d
mcp: a tool result is relayed as the upstream sent it; the markup str…
MattJackson Oct 7, 2026
67e0b43
RED: the door's holds and gathers are bounded (plane-mcp findings 7, …
MattJackson Oct 7, 2026
1372aa0
plane-mcp: bound an upstream answer, the uri dedupe, the refused unit…
MattJackson Oct 7, 2026
f8af9d2
RED: an argument URL is judged by the kernel's dest.judge (plane-mcp …
MattJackson Oct 7, 2026
381e4c3
RED: an upstream's JSON-RPC error message reaches the mcp caller unch…
MattJackson Oct 7, 2026
5bee00b
mcp: an upstream's JSON-RPC error message reaches the caller unchange…
MattJackson Oct 7, 2026
717c0c5
mcp: BUSBAR-7065 mcp-output-schema-violation is retired (finding 3)
MattJackson Oct 7, 2026
64a2244
RED: a full unit table refuses the next arrival and evicts nothing (p…
MattJackson Oct 7, 2026
9d1044c
plane-mcp: a full unit table refuses the new arrival (429) instead of…
MattJackson Oct 7, 2026
a8d7849
mcp argument guard asks the kernel's dest.judge; the plane-local host…
MattJackson Oct 7, 2026
5bc1d1e
bounds test: name the first refused arrival
MattJackson Oct 7, 2026
7c81fbb
RED: the task store is host records: a live task answered from anothe…
MattJackson Oct 7, 2026
4c4b186
mcp tasks: the task store is host records; leftover live handles are …
MattJackson Oct 7, 2026
362c554
plane-mcp: the unit-table refusal is stated before the arrival's rout…
MattJackson Oct 7, 2026
11b3ca7
RED: two callers on one stdio child each get only their own ask; a mi…
MattJackson Oct 7, 2026
098c48c
Merge remote-tracking branch 'origin/predev' into lane-mcpfix-passthr…
MattJackson Oct 8, 2026
a80cdd8
Merge remote-tracking branch 'origin/predev' into lane-mcpfix-destjudge
MattJackson Oct 8, 2026
1bd92d4
Merge branch 'lane-mcpfix-passthrough' into lane-mcpfix-destjudge
MattJackson Oct 8, 2026
e75b9dd
Merge origin/predev into lane-mcpfix-deadplane: tool_plane.rs and its…
MattJackson Oct 8, 2026
5759cf7
Merge branch 'lane-mcpfix-passthrough' into lane-mcpfix-deadplane
MattJackson Oct 8, 2026
7150aa7
RED: the session revisions, the legacy event stream and the client ne…
MattJackson Oct 7, 2026
e6b1246
mcp: serve the session revisions and the 2024-11-05 event stream in b…
MattJackson Oct 7, 2026
1046b4e
Merge remote-tracking branch 'origin/predev' into lane-mcpfix-session…
MattJackson Oct 8, 2026
2e6d071
Merge branch 'lane-mcpfix-destjudge' into lane-mcpfix-sessions-v2
MattJackson Oct 8, 2026
5e4a407
Merge lane-mcpfix-line into lane-mcpfix-sessions: the line carrier's …
MattJackson Oct 8, 2026
bf4b93a
mcp stdio: calls whose asks are relayed reach a child one at a time; …
MattJackson Oct 8, 2026
14ce150
Merge remote-tracking branch 'origin/predev' into lane-mcpfix-stdio-v2
MattJackson Oct 8, 2026
9678c43
Merge lane-mcpfix-deadplane into lane-mcpfix-stdio-v2: tests/alloc_ga…
MattJackson Oct 8, 2026
e066934
Merge branch 'lane-mcpfix-sessions-v2' into lane-mcpfix-stdio-v2
MattJackson Oct 8, 2026
b695fd9
mcp tasks: cite the design in words, not with the section sign (tests…
MattJackson Oct 8, 2026
d0479cc
Merge remote-tracking branch 'origin/predev' into lane-mcpfix-tasks
MattJackson Oct 8, 2026
27cec2c
Merge lane-mcpfix-stdio-v2 into lane-mcpfix-tasks: the run's call tak…
MattJackson Oct 8, 2026
52ff323
Merge remote-tracking branch 'origin/predev' into lane-mcpfix-bounds
MattJackson Oct 8, 2026
ce62f73
Merge lane-mcpfix-stdio-v2 into lane-mcpfix-tasks: the run's call tak…
MattJackson Oct 8, 2026
74db9ba
Merge lane-mcpfix-tasks into lane-mcpfix-bounds: the upstream answer'…
MattJackson Oct 8, 2026
1967da2
Merge branch 'lane-mcpfix-tasks' into lane-mcpfix-bounds
MattJackson Oct 8, 2026
3784ba7
bounds: the over-bound upstream answer is failed at the gather (upstr…
MattJackson Oct 8, 2026
56082a2
bounds: the over-bound reset in one statement, so the reply ceiling, …
MattJackson Oct 8, 2026
56fb4aa
Merge remote-tracking branch 'origin/lane-mcpfix-ask' into lane-mcpfi…
MattJackson Oct 8, 2026
db7b802
Merge origin/lane-mcpfix-sampling (#615, queued ahead) into lane-mcpf…
MattJackson Oct 8, 2026
fe802bb
Merge branch 'lane-mcpfix-passthrough' into lane-mcpfix-destjudge
MattJackson Oct 8, 2026
5bcd536
Merge lane-mcpfix-passthrough (with #615 sampling and #616 ask, queue…
MattJackson Oct 8, 2026
13fa14c
Merge branch 'lane-mcpfix-destjudge' into lane-mcpfix-sessions-v2
MattJackson Oct 8, 2026
5345deb
Merge branch 'lane-mcpfix-deadplane' into lane-mcpfix-stdio-v2
MattJackson Oct 8, 2026
1459852
Merge branch 'lane-mcpfix-sessions-v2' into lane-mcpfix-stdio-v2
MattJackson Oct 8, 2026
d89d9c0
Merge lane-mcpfix-stdio-v2 (with #616 ask, queued ahead) into lane-mc…
MattJackson Oct 8, 2026
49adf93
Merge branch 'lane-mcpfix-tasks' into lane-mcpfix-bounds
MattJackson Oct 8, 2026
3c57455
hop reds from the mcpfix stack (#660 run 37720341630): mcp_stdio_serv…
MattJackson Oct 8, 2026
3df8338
Merge origin/predev into lane-mcpfix-bounds: keep dest.judge in arggu…
MattJackson Oct 10, 2026
c5c4e88
Merge remote-tracking branch 'origin/predev' into lane-mcpfix-bounds
MattJackson Oct 10, 2026
853b964
plane-mcp: drop serde_json/raw_value, keep the loader in the conforma…
MattJackson Oct 10, 2026
164ff06
Merge remote-tracking branch 'origin/predev' into lane-mcpfix-bounds
MattJackson Oct 10, 2026
a5e2f69
pre-flight: fmt / lock / abi header
MattJackson Oct 10, 2026
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
28 changes: 22 additions & 6 deletions crates/busbar-core-connector/src/program.rs
Original file line number Diff line number Diff line change
Expand Up @@ -13,8 +13,10 @@
//! spawns from `1`. Two leases reading one generation reach one running program; a plugin that
//! must greet each program once (a handshake) greets each generation once.
//! * Every frame the program writes is handed to EVERY lease of its generation that is open when
//! it is read, in order; a lease reads only what arrived after it opened. Whichever lease reads
//! drives the program's pipes for all, and a frame read wakes every lease waiting on one.
//! it is read, in order; a lease reads only what arrived after it opened, and only WHOLE frames:
//! one opened while a frame is part way through skips that frame's rest (another lease's
//! message, never the start of its own). Whichever lease reads drives the program's pipes for
//! all, and a frame read wakes every lease waiting on one.
//! * A write is ONE WHOLE MESSAGE ([`Connection::write_whole`]): leases sharing the program never
//! interleave part of one message with another's.
//! * Closing a lease leaves the program running. A program that ends (its output closed, a failed
Expand Down Expand Up @@ -136,6 +138,8 @@ struct Inbox {
end_read: bool,
/// The message its open carried, while the program's queue had no room for it.
unsent: Option<Vec<u8>>,
/// It opened part way through a frame: that frame's rest is skipped.
mid_frame: bool,
}

/// The program as it runs.
Expand All @@ -147,6 +151,8 @@ struct Live {
struct State {
live: Option<Live>,
generation: u64,
/// The last piece handed to the leases did not end its frame.
mid_frame: bool,
retired: bool,
next_lease: u64,
leases: HashMap<u64, Inbox>,
Expand Down Expand Up @@ -246,6 +252,7 @@ impl Member {
state: Mutex::new(State {
live: None,
generation: 0,
mid_frame: false,
retired: false,
next_lease: 1,
leases: HashMap::new(),
Expand Down Expand Up @@ -305,6 +312,7 @@ impl Member {
st.generation += 1;
let generation = st.generation;
st.live = Some(Live { conn, generation });
st.mid_frame = false;
}
Err(f) => {
// A spawn that fails repeats on every attempt: it is counted like any end.
Expand All @@ -314,6 +322,7 @@ impl Member {
}
}
let generation = st.generation;
let mid_frame = st.mid_frame;
let lease = st.next_lease;
st.next_lease += 1;
let mut head = Vec::new();
Expand All @@ -330,6 +339,7 @@ impl Member {
ended: None,
end_read: false,
unsent: (!first.is_empty()).then(|| first.to_vec()),
mid_frame,
},
);
let waker = self.fan_waker();
Expand Down Expand Up @@ -404,12 +414,18 @@ impl Member {
continue;
}
for inbox in st.leases.values_mut() {
if inbox.generation == generation && inbox.ended.is_none() {
inbox
.pieces
.push_back((got.bytes.clone(), got.end_of_frame));
if inbox.generation != generation || inbox.ended.is_some() {
continue;
}
if inbox.mid_frame {
inbox.mid_frame = !got.end_of_frame;
continue;
}
inbox
.pieces
.push_back((got.bytes.clone(), got.end_of_frame));
}
st.mid_frame = !got.end_of_frame;
self.fan.wake_all();
}
Poll::Ready(Ok(None)) => st.ended(generation, None, &self.fan),
Expand Down
55 changes: 53 additions & 2 deletions crates/busbar-core-connector/src/tests/program_tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -16,15 +16,20 @@ use busbar_contract::conn::{
};

use crate::registry::{Entry, Transports};
use crate::support::{worker, TestDoor};
use crate::support::{worker, Knobs, TestDoor};
use crate::Connector;

const OWNER: InstanceId = InstanceId(1);
const NEED: NeedId = NeedId(1);

fn connector() -> Connector {
connector_over(TestDoor::identity("bytes"))
}

/// A connector whose `bytes` entry is `door`.
fn connector_over(door: TestDoor) -> Connector {
let view = Transports::new(vec![Entry {
door: Arc::new(TestDoor::identity("bytes")),
door: Arc::new(door),
alpn: Vec::new(),
}])
.unwrap();
Expand Down Expand Up @@ -234,6 +239,52 @@ fn a_frame_reaches_every_open_lease_and_an_opens_body_is_its_first_message() {
});
}

/// RED (finding 13): a lease opened while the program is part way through one frame does not read
/// that frame's rest as though it began there (on a program its leases share, another exchange's
/// answer taken as its own): its first body bytes begin a frame. The lease already reading the frame reads
/// it whole.
#[test]
fn a_lease_opened_mid_frame_reads_from_the_next_frame() {
worker().block_on(async {
let c = connector_over(TestDoor::new(
"bytes",
&["bytes"],
&[],
Knobs {
piece: Some(16),
..Knobs::default()
},
));
let long = "A".repeat(100);
let script = format!("read x; echo {long}; read y; echo second; cat >/dev/null");
declare(&c, &[("one", sh(&script, &[]))]).unwrap();
let a = open(&c, "one", b"go\n").unwrap();
assert_eq!(generation(&c, a).await, 1);
let mut buf = [0_u8; 64];
let first = read(&c, a, &mut buf)
.await
.expect("the long frame's first piece");
assert_eq!(first.kind, PieceKind::Body);
assert!(!first.end, "the long frame arrives in pieces");
let mut ha = buf[..first.len].to_vec();
let b = open(&c, "one", b"next\n").unwrap();
assert_eq!(generation(&c, b).await, 1);
let mut hb = Vec::new();
assert_eq!(
line(&c, b, &mut hb).await.as_deref(),
Some("second"),
"the lease opened mid-frame took the rest of a frame it never saw begin"
);
assert_eq!(
line(&c, a, &mut ha).await,
Some(long),
"a reads its frame whole"
);
c.close(OWNER, a).unwrap();
c.close(OWNER, b).unwrap();
});
}

/// RED: a program that ends ends every lease of its generation; the next open spawns it anew — the
/// next generation, another process — once the restart backoff has passed.
#[test]
Expand Down
11 changes: 10 additions & 1 deletion crates/busbar-core-connector/src/tests/support.rs
Original file line number Diff line number Diff line change
Expand Up @@ -36,6 +36,9 @@ pub struct Knobs {
pub text: bool,
/// `locate` answers this protocol offer (ProtocolNameList bytes).
pub offer: Option<&'static [u8]>,
/// Every frame piece the framing yields holds at most this many bytes, so one frame arrives
/// in several pieces (the last ends it), as a line framer cuts a line its sink cannot hold.
pub piece: Option<usize>,
}

#[derive(Default)]
Expand All @@ -46,6 +49,7 @@ struct State {
heard: bool,
deadline_ns: u64,
text: bool,
piece: Option<usize>,
}

/// One `begin` crossing's opening head fields, name and value.
Expand Down Expand Up @@ -148,7 +152,11 @@ fn answer(st: &mut State, sink: &FramerSink, o: &mut FramerOut, silence: Option<
put(sink.wire, &mut st.outbound, w);
y.wire_len = w as u64;
if !st.inbound.is_empty() && sink.pieces_cap > 0 {
let n = st.inbound.len().min(sink.frame_cap);
let n = st
.inbound
.len()
.min(sink.frame_cap)
.min(st.piece.unwrap_or(usize::MAX));
put(sink.frame, &mut st.inbound, n);
let flags = if st.inbound.is_empty() {
PIECE_END_OF_FRAME
Expand Down Expand Up @@ -267,6 +275,7 @@ impl FramerDoor for TestDoor {
let token = self.next.fetch_add(1, Ordering::Relaxed);
let mut st = State {
text: self.knobs.text,
piece: self.knobs.piece,
..State::default()
};
answer(&mut st, &i.sink, o, silence);
Expand Down
6 changes: 3 additions & 3 deletions crates/busbar-plane-mcp/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -16,9 +16,9 @@
#
# THE DEPENDENCY EDGE, AND THE THING THAT CHANGED. There is no longer a codec crate to name. The
# MCP wire dialect — the mount path, the revision, the method names, the `_meta` keys, the error
# code table, the notification pair, the durable record vocabulary, the content sanitizer and the
# structured-output check — folded IN here, because definition 16–20 puts everything one protocol
# needs and no other protocol may name in the protocol's own plane crate (#39: no `busbar-*-codec`).
# code table, the notification pair, the durable record vocabulary and the content sanitizer —
# folded IN here, because definition 16–20 puts everything one protocol needs and no other protocol
# may name in the protocol's own plane crate (#39: no `busbar-*-codec`).
# The engine's registry row and its two operation cells did NOT come with it: they name the kernel's
# handler matrix, which is kernel-side machinery, and they now read this crate's vocabulary rather
# than restating it. The adapter still names no runtime, opens nothing, holds no connection and
Expand Down
26 changes: 23 additions & 3 deletions crates/busbar-plane-mcp/src/adapt.rs
Original file line number Diff line number Diff line change
Expand Up @@ -13,7 +13,9 @@
//! What a session client is allowed to reach is decided here too ([`session_method`]): the methods
//! the session revisions define. The stateless revision's own methods (`server/discover`,
//! `subscriptions/listen`, the tasks extension) are never offered to a session client, whether it
//! asks for them by name or reads the capabilities `initialize` answered with.
//! asks for them by name or reads the capabilities `initialize` answered with. The session's own
//! verbs (`resources/subscribe`, `resources/unsubscribe`, `logging/setLevel`) stay for the old
//! revisions (THE DESIGN section 2, the mcp bullet) and are answered from the session's state.
//!
//! The `2024-11-05` event-stream framing lives here as values ([`frame`], [`endpoint_event`]);
//! writing them to a connection is the caller's.
Expand Down Expand Up @@ -96,6 +98,8 @@ pub enum SessionMethod {
Dispatch,
/// `ping`, answered here with an empty result.
Ping,
/// A verb of the session's own state: answered from it, never dispatched.
Session,
/// A notification or a response from the client: accepted (202) and not answered.
Accept,
/// Not a method of the session revisions: `-32601`.
Expand All @@ -114,6 +118,13 @@ const SESSION_DISPATCHED: &[&str] = &[
"completion/complete",
];

/// The verbs a session's own state answers (subscribe stays for the old revisions).
const SESSION_STATE_VERBS: &[&str] = &[
"resources/subscribe",
"resources/unsubscribe",
"logging/setLevel",
];

/// Classifies one session message. `method` is `None` for a client's response to a request.
#[must_use]
pub fn session_method(method: Option<&str>, has_id: bool) -> SessionMethod {
Expand All @@ -122,6 +133,7 @@ pub fn session_method(method: Option<&str>, has_id: bool) -> SessionMethod {
Some(_) if !has_id => SessionMethod::Accept,
Some(METHOD_PING) => SessionMethod::Ping,
Some(m) if SESSION_DISPATCHED.contains(&m) => SessionMethod::Dispatch,
Some(m) if SESSION_STATE_VERBS.contains(&m) => SessionMethod::Session,
Some(_) => SessionMethod::NotFound,
}
}
Expand Down Expand Up @@ -232,8 +244,10 @@ pub fn lower_result(result: &mut Value) -> Result<(), NotExpressible> {
/// Only the capability groups the session revisions define are carried, and only when discovery
/// declared them for this caller: `tools`, `prompts`, `resources` and (from `2025-06-18`)
/// `completions`. `listChanged` is `list_changed` for all three lists, the caller's statement of
/// whether it delivers those notifications on the session's stream. Resource subscription is never
/// declared: the session revisions reach it by methods this plane does not offer a session.
/// whether it delivers those notifications on the session's stream. Resource subscription and
/// `logging` are declared on every session revision: subscribe stays for the old revisions, and the
/// session's own state answers both (an upstream's `notifications/resources/updated` reaches the
/// session's stream, and its `notifications/message` past the session's floor).
#[must_use]
pub fn initialize_result(discovery: &Value, revision: Revision, list_changed: bool) -> Value {
let declared = discovery.get("capabilities");
Expand All @@ -247,6 +261,12 @@ pub fn initialize_result(discovery: &Value, revision: Revision, list_changed: bo
);
}
}
if let Some(resources) = caps.get_mut("resources").and_then(Value::as_object_mut) {
resources.insert("subscribe".into(), Value::Bool(true));
}
if has("logging") {
caps.insert("logging".into(), Value::Object(Map::new()));
}
if has("completions") && revision != Revision::R2024_11_05 {
caps.insert("completions".into(), Value::Object(Map::new()));
}
Expand Down
Loading
Loading