Repository navigation
Member-scale stress testing: sam-bench join, member-journey harness, soak fleet; fix one-caller stdio bridge and the per-IP router cap - #612
Conversation
A command-backed MCP service is one process shared by every mesh caller,
and the bridge already owns the JSON-RPC id space for that reason. Stdio
MCP is one session per process, so a server built with the Go or
TypeScript SDK refuses a second initialize on it with "duplicate
initialize received"; only the first caller of such a service ever got
in, and the Python SDK's tolerance of repeated initialize hid it.
The bridge now performs the handshake once, on behalf of the first
caller, and answers every later initialize from the backend's result
under the caller's own id. Callers arriving while the handshake is in
flight wait for its outcome instead of racing it, a refused handshake
is not cached so the next caller retries, and the backend receives
exactly one notifications/initialized. What is lost is per-caller
capability negotiation; all callers see what the first one negotiated.
Also corrects the comment in service.go: mesh callers of
/sam/{peer}/mcp/{service} do land on the shared bridge; only probes and
tool listings get a subprocess of their own.
libp2p caps inbound connections at 8 per source IP, a number its own code explains as matching the dialer's concurrency: it assumes one IP is one peer, as on a public DHT, and the matching rate limit (0.2 conns/s, burst 16) assumes the same. Members of a mesh share addresses: pods leave a node through its address, a cluster sits behind a NAT, so does an office, and each member holds one connection per router steadily and two while enrolling. On the testnet, eight members behind one address were enough for a router to refuse connections. The handshake, not the address, is what keeps strangers out. The cap's job is to bound how much of a router one address can occupy, so it is now a share of the connection budget: a quarter of --high-watermark unless --conns-per-source-ip is set, 1000 with the stock 4000. Filling a router takes at least four addresses, one address carries about 500 members, and the rate limit scales with it as before. The resource manager with the configured cap is now always installed; the branch that fell back to libp2p's limits is gone, and the exported sam_router_conns_per_source_ip_limit reports the derived value. sam-one keeps pinning the whole budget for its single proxied address.
The density and fleet runs say how many agents one member can carry. This measures members: N sam-node processes are started from one bootstrap token and each is timed from start to an authenticated router connection (/readyz) to its first MCP call of a service through the mesh, with every provider record tried and the attempts counted, so a record naming a peer that no longer answers shows up as a number rather than as a slower mean. A member that never gets there is reported by the stage it stopped at and contributes no latency. With --hold the fleet stays resident and every member's readiness is sampled, so a router rollout or a key rotation meanwhile appears as outage windows per member with their lengths; --resident-marker and the join report written at that moment let a script run its own phase against the resident fleet. Cancelling ends the hold with the report intact, and the observation carries metrics scraped before and after, in the same shape sam-bench run writes. The members are the binary the run is pointed at; the tests re-enter the test binary as a node that answers /readyz and the two sidecar routes the journey uses.
tests/scale/member-journey.sh takes a fleet through the phases a testnet sees: a burst of N members joining at once, more joining the populated mesh, a router rollout under the resident fleet with more joining behind it, and a hold. Each phase is a sam-bench join observation; the summary renders one table and a verdict against thresholds the environment can override, and prints every call attempt that failed. With --env it mints a bootstrap token from the control plane's admin secret and revokes it on exit, and snapshots the control plane and routers per pod through the API server, so the verdict also covers 5xx, request p99, and connections the routers refused. The members call an MCP service, so two real servers ship for a local mesh: tests/scale/stdio-mcp.py on the Python SDK and tests/scale/stdio-mcp on the Go SDK. Running both found the stdio bridge accepting one caller per command-backed service. Two tiers run in CI. TestMemberBurstJoinsConcurrently starts ten members at once against a standalone server and a provider, the concurrent enrollment path nothing else in the suite exercises, in about three seconds. The member-journey workflow boots sam-one, one provider per SDK, and runs the harness at thirty members on every push to main, with loose latency bounds since a runner's latency says nothing about the mesh, and the table in the step summary.
…mber loses the mesh
The canaries are two pods each and restart when their own service
misbehaves, so between rollouts nobody watches what a plain member
experiences over days. sam-soak-${ENV_NAME} is twenty socket-only
members that join and stay, publishing and calling nothing. Readiness
is the node's /readyz, 200 only while it holds an authenticated router
connection, and there is deliberately no liveness probe: a restart
would hide the reconnection the fleet exists to observe. The rollout
in deploy.yaml completes only once every member holds a router, and
every later deploy restarts the routers under the fleet.
Each pod carries sam-canary=true, so the canary PodMonitoring scrapes
sam_node_mesh_connected from it, and SamCanaryMemberDisconnected fires
when a member reads 0 for five minutes: seconds is a rollout, five
minutes is a member that cannot get back, which after a key rotation
or a refresh is how a fleet falls off the mesh one pod at a time.
Applied to bananas before this commit: 20/20 ready in 31 s, no
connection refused, the collector scraping every pod.
|
Warning Gemini encountered an error creating the review. You can try again by commenting |
|
/gemini review |
There was a problem hiding this comment.
Code Review
This pull request introduces a new sam-bench join command and associated scale testing tools to measure the cost of starting a fleet of members. It also updates the router's default conns-per-source-ip to scale with the connection budget, and modifies StdioBridge to perform the MCP handshake once and cache the result for concurrent callers. The reviewer feedback highlights critical style guide violations regarding peer ID canonicalization at boundaries in both the benchmark tool and integration tests. Additionally, the reviewer suggests optimizing the long-running hold loop by replacing time.After with a time.Ticker to prevent memory leaks, and exiting early in sampleReady if the context is cancelled.
| peers := make([]string, 0, len(providers)) | ||
| for _, p := range providers { | ||
| if p.GetPeerId() != "" { | ||
| peers = append(peers, p.GetPeerId()) | ||
| } | ||
| } |
There was a problem hiding this comment.
According to the SAM Code Review Style Guide (Section 1: Peer IDs must be canonicalized at the boundary), wire fields like p.GetPeerId() from api.DiscoveredProvider are untrusted and non-canonical. They must be decoded using peer.Decode and canonicalized using .String() at the boundary (here, inside the discover function) before being used or compared. This prevents silent misses in bans, admission, or dedup checks due to different valid string encodings of the same peer ID.
peers := make([]string, 0, len(providers))
for _, p := range providers {
if p.GetPeerId() != "" {
pID, err := peer.Decode(p.GetPeerId())
if err != nil {
return nil, fmt.Errorf("discover: invalid peer ID %q: %w", p.GetPeerId(), err)
}
peers = append(peers, pID.String())
}
}References
- Peer IDs must be canonicalized at the boundary to prevent security bugs where different string encodings of the same peer ID bypass checks. (link)
| for _, n := range status.GetEnrolledNodes() { | ||
| if n.GetPeerId() != provider.peerID.String() && !strings.EqualFold(n.GetPeerId(), srv.PeerID()) { | ||
| enrolled++ | ||
| } | ||
| } |
There was a problem hiding this comment.
According to the SAM Code Review Style Guide (Section 1, Rule 5), raw peer ID strings from wire fields (like n.GetPeerId()) must not be compared directly for identity with canonical peer IDs. Both sides should be decoded and compared as peer.ID values to avoid silent mismatches due to different valid string encodings.
for _, n := range status.GetEnrolledNodes() {
nID, err := peer.Decode(n.GetPeerId())
if err != nil {
t.Fatalf("invalid peer ID: %v", err)
}
srvID, err := peer.Decode(srv.PeerID())
if err != nil {
t.Fatalf("invalid srv peer ID: %v", err)
}
if nID != provider.peerID && nID != srvID {
enrolled++
}
}References
- Raw string compared for identity. Decode both sides and compare peer.ID values, or compare .String() of two decoded IDs. (link)
|
|
||
| "github.com/google/sam/api" | ||
| ) |
| client := &http.Client{Timeout: 2 * time.Second} | ||
| start := time.Now() | ||
| end := start.Add(opts.Hold) | ||
| for { | ||
| now := time.Now() | ||
| ready := sampleReady(ctx, client, members) | ||
| count := 0 | ||
| for i, m := range members { | ||
| if ready[i] { | ||
| count++ | ||
| } | ||
| s := &states[i] | ||
| switch { | ||
| case report.Samples == 0: | ||
| s.ready = ready[i] | ||
| if !ready[i] { | ||
| s.downSince = now | ||
| } | ||
| case s.ready && !ready[i]: | ||
| s.ready = false | ||
| s.downSince = now | ||
| s.flapped = true | ||
| report.Flaps++ | ||
| opts.logf("member %04d lost its router", m.index) | ||
| case !s.ready && ready[i]: | ||
| s.ready = true | ||
| outages = append(outages, now.Sub(s.downSince)) | ||
| opts.logf("member %04d back after %.1fs", m.index, now.Sub(s.downSince).Seconds()) | ||
| } | ||
| } | ||
| report.Series = append(report.Series, ReadySample{At: now, Ready: count}) | ||
| if report.Samples == 0 { | ||
| report.Resident = count | ||
| report.MinReady = count | ||
| } else if count < report.MinReady { | ||
| report.MinReady = count | ||
| } | ||
| report.Samples++ | ||
|
|
||
| if now.After(end) || ctx.Err() != nil { | ||
| break | ||
| } | ||
| select { | ||
| case <-ctx.Done(): | ||
| case <-time.After(opts.SampleInterval): | ||
| } | ||
| if ctx.Err() != nil { | ||
| break | ||
| } | ||
| } |
There was a problem hiding this comment.
Using time.After inside a long-running loop (like the hold loop, which can run for up to 24 hours) creates a new timer on every iteration. These timers will not be garbage collected until they expire, which can lead to unnecessary memory overhead. It is highly recommended to use a single time.Ticker defined outside the loop to control the sampling interval.
client := &http.Client{Timeout: 2 * time.Second}
start := time.Now()
end := start.Add(opts.Hold)
ticker := time.NewTicker(opts.SampleInterval)
defer ticker.Stop()
for {
now := time.Now()
ready := sampleReady(ctx, client, members)
count := 0
for i, m := range members {
if ready[i] {
count++
}
s := &states[i]
switch {
case report.Samples == 0:
s.ready = ready[i]
if !ready[i] {
s.downSince = now
}
case s.ready && !ready[i]:
s.ready = false
s.downSince = now
s.flapped = true
report.Flaps++
opts.logf("member %04d lost its router", m.index)
case !s.ready && ready[i]:
s.ready = true
outages = append(outages, now.Sub(s.downSince))
opts.logf("member %04d back after %.1fs", m.index, now.Sub(s.downSince).Seconds())
}
}
report.Series = append(report.Series, ReadySample{At: now, Ready: count})
if report.Samples == 0 {
report.Resident = count
report.MinReady = count
} else if count < report.MinReady {
report.MinReady = count
}
report.Samples++
if now.After(end) || ctx.Err() != nil {
break
}
select {
case <-ctx.Done():
case <-ticker.C:
}
if ctx.Err() != nil {
break
}
}| for i, m := range members { | ||
| if m.cmd == nil || m.hasExited() { | ||
| continue | ||
| } | ||
| wg.Add(1) | ||
| sem <- struct{}{} |
There was a problem hiding this comment.
When the context is cancelled, the loop in sampleReady will continue to push tasks to the semaphore channel, which can block if the semaphore is full. Checking ctx.Err() != nil at the beginning of the loop allows exiting early and avoiding spawning unnecessary goroutines once the context is cancelled.
for i, m := range members {
if ctx.Err() != nil {
break
}
if m.cmd == nil || m.hasExited() {
continue
}
wg.Add(1)
sem <- struct{}{}| "github.com/google/sam/internal/bench" | ||
| "github.com/google/sam/internal/standalone" | ||
| ) |
There was a problem hiding this comment.
Import "github.com/libp2p/go-libp2p/core/peer" to support decoding and comparing peer IDs in the test.
| "github.com/google/sam/internal/bench" | |
| "github.com/google/sam/internal/standalone" | |
| ) | |
| "github.com/google/sam/internal/bench" | |
| "github.com/google/sam/internal/standalone" | |
| "github.com/libp2p/go-libp2p/core/peer" | |
| ) |
Stress coverage for the one dimension the testnets had never been measured on: many members. The datapath and agent density are measured (boundary overhead below the noise floor, 1000 agents on one host); what a new member experiences while N-1 others are present was not. This adds the instrument, runs it, and fixes the two things it found.
What was found
/sam/{peer}/mcp/{service}lands on the sharedStdioBridge, one process for every caller, and a strict stdio server refuses a secondinitializeon its one session. The Python SDK tolerates repeatedinitialize, which is why nothing had noticed. Found by running real servers from both SDKs as providers.per_ip_limit: 3refused with 8 members from one workstation.Commits
node: perform the MCP handshake once in the stdio bridge— the bridge handshakes once on behalf of the first caller and answers laterinitializefrom the result under each caller's own id; concurrent initializes wait rather than race; a refused handshake is not cached; onenotifications/initializedreaches the backend. Through the mesh with Go + Python providers: 1.60 → 1.00 calls per member, 12 → 0 failed attempts.router: derive the per-source-IP connection cap from the high watermark— a quarter of--high-watermarkunless set (1000 with the stock 4000): filling a router takes at least four addresses, one address carries ~500 members, the rate limit scales with it. The handshake is the trust boundary, not the address.bench: add sam-bench join— Nsam-nodeprocesses from one bootstrap token, each timed start →/readyz→ first MCP call through the mesh, every provider record tried and counted;--holdsamples readiness and turns flaps into outage windows per member.scale: member-journey harness, with an integration test and a CI job— phases: burst, late joiners, router rollout under the resident fleet, hold; per-pod control-plane and router metric snapshots; a table and a verdict. Real stdio MCP servers on both SDKs undertests/scale.TestMemberBurstJoinsConcurrently(10 members, ~3 s) inmake test; amember-journeyworkflow at 30 members on push tomain.deploy: a resident soak fleet on each testnet, and an alert when a member loses the mesh—sam-soak-${ENV_NAME}, 20 socket-only members that join and stay, readiness =/readyz, no liveness probe by design;SamCanaryMemberDisconnectedonsam_node_mesh_connected == 0for 5m.Tested on bananas before this PR
member-journey.sh --env bananas(6 + 2 members, 90 s hold): all completed, 0 flaps, 0 control-plane 5xx, control-plane p99 ≤ 50 ms, enrolled nodes 13485 → 13496; the minted bootstrap token revoked on exit. Two findings: theper_ip_limitrefusals above, and a cold join of 6.05 s for every burst member (late joiners 3.5–4 s), ~3 s of it between "Successfully enrolled" and the router re-handshake — a fixed cost in the join path, left for a follow-up; the harness's 5 s p50 check is red on bananas until then.Validation
go test ./...green,make lint0 issues, every commit builds and vets on its own.tests/e2e/canary_manifests.batspicks the new template up by glob; it passed a server-side dry run on a kind cluster.Follow-ups
member-journey.sh --env bananas --count 500 --late 20 --hold 24h --rollout).