Correlator
NAME
MCP::Client::Correlator - the pending-request table for one MCP connection
DESCRIPTION
A stdio MCP server is one pipe carrying every conversation at once: requests go
down it in whatever order the caller made them, and answers come back in
whatever order the server finished them. The correlator is the bookkeeping that
makes that usable ā it hands out ids, hands back a Promise per request, and
matches each inbound answer to the promise that is waiting for it.
Everything it does is thread-safe, and it is careful about how. All table
mutation happens under a lock, but no Promise is ever kept or broken while
that lock is held: the vows are collected inside the lock and settled outside
it. A continuation that runs on the settling thread therefore cannot re-enter
the correlator and deadlock ā the same defer-mutations-during-a-walk discipline
used elsewhere in this tree.
Time is injectable
&.now and &.schedule-after exist so timeouts can be tested without
sleeping. Override both with a virtual clock and a timer you fire by hand, and
a test for "this request timed out after sixty seconds" runs in microseconds
and never flakes on a loaded CI box.
Late answers
A request that times out, is cancelled, or is failed by fail-all is removed
from the table there and then. If the server answers it afterwards, resolve
returns False and the answer is dropped ā a promise is only settled once,
and the caller has already been told what happened.
EXAMPLES
The shape a transport uses it in:
use MCP::Client::Correlator;
my $correlator = MCP::Client::Correlator.new(default-timeout => 30);
# Sending side
my $id = $correlator.next-id;
my $answer = $correlator.register($id, method => 'tools/call', timeout => 10);
$proc.print(format-message(build-request($id, 'tools/call', %args, :era<modern>)));
# Receiving side, on the read loop
given parse-inbound($line) -> %in {
if %in<kind> eq 'response' {
%in<error>:exists
?? $correlator.reject(%in<id>, X::MCP::Client::Protocol.new(
detail => %in<error><message> // 'server error',
code => %in<error><code>,
data => %in<error><data>,
))
!! $correlator.resolve(%in<id>, normalize-result(%in<result>));
}
}
my %result = await $answer; # or throws whichever failure got there first
When the child dies, everyone waiting hears about it at once, and anything registered afterwards fails immediately instead of hanging forever:
$correlator.fail-all(X::MCP::Client::ServerGone.new(
exit-code => $proc-exit, stderr-tail => $tail,
));
say $correlator.closed; # True
await $correlator.register(99); # throws ServerGone straight away
Timeouts under a virtual clock:
my @timers;
my $c = MCP::Client::Correlator.new(
now => { 0 },
schedule-after => -> Real $s, &cb { @timers.push($s => &cb); Promise.kept },
);
my $p = $c.register(1, timeout => 5);
@timers[0].value.(); # fire the timer by hand
say $p.status; # Broken
say $p.cause ~~ X::MCP::Client::Timeout; # True