Skip to content

Member-scale stress testing: sam-bench join, member-journey harness, soak fleet; fix one-caller stdio bridge and the per-IP router cap - #612

Merged
aojea merged 5 commits into
google:mainfrom
aojea:member-journey
Oct 9, 2026
Merged

aojea merged 5 commits into
google:mainfrom
aojea:member-journey

Conversation

@aojea

@aojea aojea commented Oct 9, 2026

Copy link
Copy Markdown
Collaborator

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

  • A command-backed MCP service served exactly one caller when its backend was written with the Go or TypeScript SDK. The mesh ingress for /sam/{peer}/mcp/{service} lands on the shared StdioBridge, one process for every caller, and a strict stdio server refuses a second initialize on its one session. The Python SDK tolerates repeated initialize, which is why nothing had noticed. Found by running real servers from both SDKs as providers.
  • Routers refused connections with eight members behind one address. libp2p's per-source-IP cap of 8 encodes "one IP is one peer" (its own comment ties it to the dialer's concurrency); members of a mesh share addresses (pods behind a node's SNAT, a cluster behind a NAT, an office) and hold two connections per router while enrolling. Measured on bananas: per_ip_limit: 3 refused 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 later initialize from the result under each caller's own id; concurrent initializes wait rather than race; a refused handshake is not cached; one notifications/initialized reaches 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-watermark unless 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 — N sam-node processes from one bootstrap token, each timed start → /readyz → first MCP call through the mesh, every provider record tried and counted; --hold samples 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 under tests/scale. TestMemberBurstJoinsConcurrently (10 members, ~3 s) in make test; a member-journey workflow at 30 members on push to main.
  • 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; SamCanaryMemberDisconnected on sam_node_mesh_connected == 0 for 5m.

Tested on bananas before this PR

  • Soak fleet applied as the deploy step does: 20/20 ready in 31 s, 0 restarts, routers' authenticated peers 10 → 30, 0 refusals, the GMP collector scraping every pod; the Rules object accepted.
  • 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: the per_ip_limit refusals 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 lint 0 issues, every commit builds and vets on its own. tests/e2e/canary_manifests.bats picks the new template up by glob; it passed a server-side dry run on a kind cluster.

Follow-ups

  • The ~3 s gap in the node's cold join.
  • The 500-member run from one minion against bananas, once the router default has deployed (member-journey.sh --env bananas --count 500 --late 20 --hold 24h --rollout).

aojea added 5 commits October 9, 2026 11:09
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.
@gemini-code-assist

Copy link
Copy Markdown
Contributor

Warning

Gemini encountered an error creating the review. You can try again by commenting /gemini review.

@aojea

aojea commented Oct 9, 2026

Copy link
Copy Markdown
Collaborator Author

/gemini review

@gemini-code-assist gemini-code-assist Bot left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

Comment thread internal/bench/join.go
Comment on lines +574 to +579
peers := make([]string, 0, len(providers))
for _, p := range providers {
if p.GetPeerId() != "" {
peers = append(peers, p.GetPeerId())
}
}

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

high

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
  1. Peer IDs must be canonicalized at the boundary to prevent security bugs where different string encodings of the same peer ID bypass checks. (link)

Comment on lines +125 to +129
for _, n := range status.GetEnrolledNodes() {
if n.GetPeerId() != provider.peerID.String() && !strings.EqualFold(n.GetPeerId(), srv.PeerID()) {
enrolled++
}
}

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

high

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
  1. Raw string compared for identity. Decode both sides and compare peer.ID values, or compare .String() of two decoded IDs. (link)

Comment thread internal/bench/join.go
Comment on lines +33 to +35

"github.com/google/sam/api"
)

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

medium

Import "github.com/libp2p/go-libp2p/core/peer" to support canonicalizing peer IDs decoded from the wire.

Suggested change
"github.com/google/sam/api"
)
"github.com/google/sam/api"
"github.com/libp2p/go-libp2p/core/peer"
)

Comment thread internal/bench/join.go
Comment on lines +694 to +743
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
}
}

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

medium

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
		}
	}

Comment thread internal/bench/join.go
Comment on lines +766 to +771
for i, m := range members {
if m.cmd == nil || m.hasExited() {
continue
}
wg.Add(1)
sem <- struct{}{}

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

medium

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{}{}

Comment on lines +29 to +31
"github.com/google/sam/internal/bench"
"github.com/google/sam/internal/standalone"
)

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

medium

Import "github.com/libp2p/go-libp2p/core/peer" to support decoding and comparing peer IDs in the test.

Suggested change
"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"
)

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant