Skip to content

feat: multi-backend gRPC addresses via static pick_first resolver - #256

Merged
wu-sheng merged 14 commits into
apache:mainfrom
songzhendong:feature/multi-backend-grpc
Sep 29, 2026
Merged

wu-sheng merged 14 commits into
apache:mainfrom
songzhendong:feature/multi-backend-grpc

Conversation

@songzhendong

@songzhendong songzhendong commented Sep 26, 2026 •

Copy link
Copy Markdown
Contributor

Summary

Support comma-separated reporter.grpc.backend_service addresses on one shared gRPC channel with native pick_first failover. Dial is the only multi-address branch; reporters reuse the existing long-lived Collect loops for both single- and multi-address modes.

Behavior

  • Normalize and deduplicate addresses, validate numeric ports, and skip invalid entries in comma-separated lists with a warning once during initialization. One surviving endpoint uses the existing single-address path; no valid endpoints disables the gRPC reporter with a warning and lets the application continue. A single non-host:port target (for example dns:///... or unix:///...) is passed through to grpc.Dial.
  • Shuffle the static resolver list once per channel. Resolver refreshes retain that order. Each resolver address sets TLS ServerName from that host (explicit credential overrides still win).
  • Failover stays on the shared ClientConn: a failed Send reopens Collect on the same channel without replaying telemetry. Only idempotent reportInstanceProperties retries UNAVAILABLE.
  • Connecting / Idle / Ready report Connected so registration is not stalled on check_interval. Client keepalive and per-address dial timeouts help leave dead or silent peers so pick_first can move on.
  • Bound send watchdogs and cancel-on-unready protect Collect against half-open peers. Graceful shutdown closes the stream before canceling the context.
  • Refresh instance properties after disconnect and periodically so a standby receives metadata.
  • An all-invalid gRPC backend list does not disable the Kafka reporter.

Validation

  • Fork CI: unit/race/lint, plugin tests, Windows plugin tests, and multi-backend E2E.

Support comma-separated reporter.grpc.backend_service (>=2 after normalize)
with a static endpoint list and pick_first, aligned with Node sw-static
failover. Single-address paths keep historical behavior.

Multi-backend only: dial/RPC timeouts, proxy bypass, auth-failure throttled
logs, Collect Send/CloseAndRecv bounds, in-place RecreateConnection on
half-open peers, and UNAVAILABLE retries limited to reportInstanceProperties
so streaming Collect cannot replay segments.

GetConnection shares the connection-manager mutex with Recreate/Release/Peek.
Multi-backend status watcher skips Connect on Shutdown ClientConns.
MultiBackendSend cancels the stream context on deadline and uses a buffered
done channel so a timed-out Send does not need a second drain goroutine;
Recreate closes the prior conn to unblock it.

gRPC service stubs and PprofTaskClient are published as immutable bundles via
atomic.Pointer so bind/recreate cannot race with send, heartbeat, profile, or
pprof poll/upload goroutines; multi-backend pprof upload retries once after
refreshing the stub.

Tests: mid-stream kill->standby, concurrent Get/Recreate/Release, TestRace*
(including client-bundle and pprof-client swap), service-config retry scope,
hung-send unblock, reporter auto-failover for trace/metrics/log without
manual RecreateConnection. CI runs go test -race on reporter packages. E2E
resolves active via unique /sw-failover-probe/{token} and requires standby
POST:/info + toolkit log growth, with bounded waits/retries so the case
stays under the GHA job timeout.
@wu-sheng wu-sheng added this to the 0.8.0 milestone Sep 26, 2026
@wu-sheng wu-sheng added the enhancement New feature or request label Sep 26, 2026
@mrproliu
mrproliu requested a lite review from Copilot September 26, 2026 00:24

Copilot AI left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

Copilot review overview

🟡 Changes recommended

Unresolved critical and moderate findings remain in connection lifecycle, stream retry/recovery, pprof handling, and E2E configuration.

Get a fresh assessment by requesting another Copilot review.

Review effort: Lite
Findings: 4 High severity · 1 Medium severity

Open (5)
What changed in this PR

Adds multi-address gRPC backend support with static resolution and pick_first failover while preserving single-address behavior.

Changes:

  • Adds backend normalization, static resolution, reconnection, and bounded operations.
  • Updates reporter, pprof, CDS, tooling, CI, and configuration support.
  • Adds unit, race, and end-to-end failover coverage and documentation.
File Summary
tools/​go-agent/​tools/​dst.go Extends package-reference rewriting.
tools/​go-agent/​tools/​dst_test.go Adds AST regression tests.
tools/​go-agent/​config/​agent.default.yaml Documents multi-address configuration.
test/​e2e/​case/​grpc-multi-backend/​verify-failover.sh Verifies backend failover.
test/​e2e/​case/​grpc-multi-backend/​expected/​failover-ok.yml Defines expected E2E output.
test/​e2e/​case/​grpc-multi-backend/​e2e.yaml Registers the E2E scenario.
test/​e2e/​case/​grpc-multi-backend/​docker-compose.yml Defines the multi-backend topology.
test/​e2e/​base/​consumer/​main.go Adds failover probe endpoints.
plugins/​core/​reporter/​static_backend_resolver.go Implements static gRPC resolution.
plugins/​core/​reporter/​pprof_manager.go Adds swappable pprof clients and retries.
plugins/​core/​reporter/​multi_backend_stream_test.go Tests stream failover and concurrency.
plugins/​core/​reporter/​grpc/​multi_backend_failover_test.go Tests reporter failover.
plugins/​core/​reporter/​grpc/​grpc.go Integrates multi-backend reporter logic.
plugins/​core/​reporter/​grpc/​client_bundle_race_test.go Tests atomic client swaps.
plugins/​core/​reporter/​conn_manager.go Manages multi-backend connections and recreation.
plugins/​core/​reporter/​cds_manager.go Adds bounded multi-backend CDS calls.
plugins/​core/​reporter/​backend_addresses.go Normalizes backend address lists.
plugins/​core/​reporter/​backend_addresses_test.go Tests backend address parsing.
plugins/​core/​reporter/​auth_failure_log.go Throttles authentication failure logs.
plugins/​core/​reporter/​auth_failure_log_test.go Tests authentication logging.
docs/​menu.yml Adds troubleshooting navigation.
docs/​en/​advanced-features/​grpc-backend-troubleshooting.md Documents backend failover and diagnostics.
agent/​reporter/​imports.go Adds generated reporter imports.
.github/​workflows/​skywalking-go.yaml Adds reporter race testing.
.github/​workflows/​e2e.yaml Registers multi-backend E2E testing.

💡 Add a code-review agent skill or configure MCP servers for context-aware, tailored reviews. Learn more in the docs.

Comment thread plugins/core/reporter/conn_manager.go Outdated
Comment thread plugins/core/reporter/conn_manager.go Outdated
Comment thread plugins/core/reporter/grpc/grpc.go Outdated
cases:
- name: fail over to standby after active collector stops
query: bash test/e2e/case/grpc-multi-backend/verify-failover.sh
expected: expected/failover-ok.yml
Comment thread plugins/core/reporter/conn_manager.go Outdated
@wu-sheng
wu-sheng requested a review from mrproliu September 26, 2026 00:42
@wu-sheng

wu-sheng commented Sep 26, 2026 •

Copy link
Copy Markdown
Member

Cross-agent alignment follow-up, updated for dc4aa2c (compared with merged Python 5666826 and Node.js d7f266f): shared-channel/pick_first recovery, shuffled selection, no replay, and skipping invalid entries in comma-separated lists are aligned. The DNS/authority distinction below remains.

  1. Invalid-entry mismatch addressed in dc4aa2c. Go now skips invalid entries and emits a warning per skipped entry during initialization. For example, oap-a:11800,invalid,oap-b:11800 keeps both valid endpoints. Parsing is cached so channel acquisition/reconnects do not repeat these warnings. One surviving endpoint uses the existing single-address path; if none survive, reporting is disabled with a warning and the application continues without an initialization error or panic. Standalone targets retain their historical resolver behavior.

    This implements the requested warning-only policy. The skip behavior matches Python and Node.js; diagnostic severity and all-invalid handling intentionally follow the requested Go behavior. Reporter/tooling race tests and the pinned reporter linter passed. An instrumented application with an all-invalid list emitted warnings, reached application code, and exited 0 without panic or reporter initialization errors.

  2. Multi-hostname DNS and authority fallback still differ from Python. Go resolves hostname endpoints when dialing, following the Node.js static-endpoint model. Python expands multi-address hostnames once at startup, skips names that fail resolution, and chooses the first usable original endpoint as its authority. Consequently, startup DNS failures can produce different TLS identities, and later DNS changes can affect reconnects differently. Go's fixed first-valid-configured authority matches Node.js, but Python's first-DNS-usable fallback is an edge-case difference.

    Recommendation: explicitly accept and document the Go/Node.js DNS behavior unless identical Python startup-DNS and authority-fallback semantics are required. The common requirement that certificates cover the selected channel identity is not itself an outstanding Node.js mismatch.

The invalid-entry policy is now implemented. The remaining decision is whether to accept the Go/Node.js DNS and authority behavior or require Python's startup-DNS and first-usable-authority semantics.

@songzhendong

Copy link
Copy Markdown
Contributor Author

Thanks for the professional follow‑up and related handling

@mrproliu mrproliu 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.

Thanks for the PR.

My main concern is that this adds a second implementation (dialing, connection status, send loops, heartbeat) that only runs when there are two or more addresses. I don't think the fork is needed.

gRPC's ClientConn already handles failover across addresses: pick_first is the default LB policy, and when the active backend dies the reporter just sees a failed Send, then re-opens the stream on the same ClientConn. That is exactly what the existing single-address loops already do.

Proposed shape

  1. Always parse backend_service into a list. When it has 2+ entries, dial through the static resolver; otherwise dial as today. That should be the only branch.
  2. Keep the existing connection status checker and the existing long-lived stream loops. Remove multiBackendTraceSendLoop, the IsMultiBackend() branches in the metrics/log/profile loops, and the multi-only heartbeat logic.
  3. The Send timeout and "cancel stream when the channel leaves READY" protections are useful, but they fix a real problem for single-address too (a half-open connection blocks Send forever). Enable them for both modes.
  4. serviceClients, pprofClientBundle, CDS stub re-binding and shutdownCtx only exist to support connection re-creation, which this PR says does not happen. They can go.

This should make the change much smaller and keep one code path to maintain.

A few things to fix regardless

  • In multi-backend mode the first ReportInstanceProperties waits a full check_interval (20s) because Connecting is treated as disconnected. Single-address registers immediately.
  • The per-batch trace stream (CloseAndRecv after every batch) makes throughput RTT-bound. The existing long-lived stream already has the same "no replay" behavior.
  • TLS verifies every backend against the first address's host. Setting resolver.Address.ServerName per address avoids this instead of documenting it.
  • An all-invalid list also disables the Kafka reporter, since initManager runs on that path too.
  • Please remove the "Codex P2" / "Codex: ..." comments and the Node.js / Python references from code and CI.

@@ -0,0 +1,81 @@
# gRPC Backend Connectivity Troubleshooting

## Multi-address failover

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.

Address review feedback: dial is the only multi-address branch; reuse
long-lived Collect loops with BoundSend (including CloseAndRecv) and
cancel-on-unready for both modes; unify connection status under lock
(Connecting is Connected; closed conn reports Shutdown immediately);
refresh instance properties after disconnect and periodically without
blocking heartbeats; drop recreate support structures; set per-address
TLS ServerName; keep Kafka usable when gRPC backend_service is
all-invalid; merge troubleshooting docs into grpc-tls.md and strip
Codex/cross-agent noise from code/CI.
@songzhendong

Copy link
Copy Markdown
Contributor Author

Thanks for the review.

We converged the multi-backend work onto a single reporter path, as requested.

Structural changes

  1. Dial is the only mode branch. �ackend_service is always normalized to a list; with 2+ addresses we dial through the static resolver + pick_first, otherwise we keep the existing single-address dial.
  2. One set of long-lived Collect loops for trace / metrics / log / profile. The multi-only send loops, IsMultiBackend() branches in those pipelines, and the multi-only heartbeat path are removed. Failover is handled by the shared ClientConn: a failed Send reopens the stream on the same channel; segments are not replayed.
  3. BoundSend and cancel-on-unready apply in both single- and multi-address modes, including CloseAndRecv, so a half-open peer cannot block Send or stream close and stall failover.
  4. Connection recreation support is gone (serviceClients, pprofClientBundle, CDS stub re-binding, shutdownCtx, etc.). Clients are created once on the shared connection.

Other fixes

  • Connecting / Idle / Ready all report Connected, so the first ReportInstanceProperties no longer waits a full check_interval in multi-backend mode.
  • Connection status updates compare and write managed.status under the same lock, after re-checking the map entry, to avoid the status data race.
  • Trace reporting stays on the long-lived Collect stream (no per-batch CloseAndRecv).
  • Instance properties are re-reported after disconnect and refreshed periodically; a properties failure does not block heartbeats.
  • Each resolver address sets ServerName from that host for TLS SNI / verification.
  • An all-invalid �ackend_service list disables the gRPC reporter only; the Kafka reporter continues without CDS.
  • Merged the troubleshooting notes into docs/en/advanced-features/grpc-tls.md and removed the separate doc / Codex and cross-agent noise from code and CI.

Status

These changes are on a single commit and have passed fork CI (unit/race/lint, plugin tests, Windows plugin tests, and multi-backend E2E).

@mrproliu mrproliu 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.

BoundSend: goroutine + timer per message on the hot path

After the convergence, BoundSend runs for every segment, metric batch and log line for all users. Each call spawns a goroutine, allocates a result channel and creates a timer.

grpc-go's SendMsg already returns as soon as the stream context is cancelled (the write-quota wait is interrupted by stream done), so the only thing we actually need is "cancel the context if Send takes too long". That can be done without a goroutine:

watchdog := time.AfterFunc(timeout, cancel)
err := stream.Send(s)
watchdog.Stop()

with one timer reused per stream via Reset. The "send ignores cancel, wait for grace" branch guards against a situation grpc-go does not produce on a real stream; TestMultiBackendSendReturnsWhenSendIgnoresCancel only reproduces it with a time.Sleep.

Remaining multi-address-only branches

Two places still behave differently depending on the address count. Since everything else is now one path, I'd suggest removing both:

  1. checkConnectionStatus polls every checkInterval (20s by default) for multi-address but every 5s for single-address. That makes Disconnect → Connected detection four times slower in multi mode for no reason. The conn.Connect() nudge can also go: on grpc-go 1.55 Connect() is a no-op while in TransientFailure. With those two gone, the multi variable in this function is unused.

  2. directTCPContextDialer is only installed for multi-address. Its effects are bypassing HTTP_PROXY / HTTPS_PROXY and changing TCP keepalive from 15s to 10s. Dial time is already bounded by MinConnectTimeout: 5s. Dropping the custom dialer removes the last behavioral difference between the two modes, and the proxy paragraph in the docs can be removed with it.

Comment thread .github/workflows/e2e.yaml Outdated
@@ -47,6 +47,8 @@ jobs:
case:
- name: gRPC

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.

Suggested change
- name: gRPC
- name: gRPC single-backend

}
}
r.closeTracingStream(stream)
cancel()

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.

This runs on the normal shutdown path (channel closed by Close()), and it cancels the stream context before CloseAndRecv. Two consequences:

  • CloseAndRecv returns Canceled immediately, so we no longer wait for the OAP ack. Segments still sitting in the transport write buffer at shutdown can be lost. The previous code waited for the ack here.
  • Every graceful shutdown now logs send closing error context canceled.

Suggest: on this path call closeTracingStream first (BoundSend already cancels on timeout, so it can't hang), then cancel(). The error path at L266 is fine to keep as is, since the stream is already broken there.

Same applies to the metrics (L322), log (L368) and profile (L439) loops.

Copilot AI left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

Comment thread plugins/core/reporter/grpc/grpc.go Outdated
Comment on lines +271 to +272
cancel()
r.closeTracingStream(cancel, stream)
Comment on lines +100 to +111
// configuredAddressesAsResolverState builds resolver addresses with ServerName
// set from each endpoint's host so TLS SNI / certificate verification follows
// the dialed backend rather than a fixed channel authority.
func configuredAddressesAsResolverState(backends []string) []resolver.Address {
addresses := make([]resolver.Address, 0, len(backends))
for _, cfg := range backends {
host, _, err := net.SplitHostPort(cfg)
if err != nil {
addresses = append(addresses, resolver.Address{Addr: cfg})
continue
}
addresses = append(addresses, resolver.Address{Addr: cfg, ServerName: host})
Comment on lines +411 to +412
aGS.Stop()
_ = aLis.Close()
Comment on lines +47 to +50
_, err := stream.Recv()
if err == io.EOF {
return status.Error(s.code, "collector failed after receiving the segment")
}
Close Collect streams before canceling on graceful drain; simplify
BoundSend to a context watchdog without a per-message goroutine; use one
status poll interval with no Connect nudge, and drop the multi-only
custom TCP dialer; tighten failover/RPC-error tests and rename the
single-backend E2E case.
@songzhendong

songzhendong commented Sep 28, 2026 •

Copy link
Copy Markdown
Contributor Author

Shutdown and BoundSend

  1. Graceful drain closes before cancel. On normal channel close, each Collect loop (trace / metrics / log / profile) runs CloseAndRecv under BoundSend first, then cancels the stream context, so OAP can still ack in-flight data. The error path keeps cancel-first, since the stream is already broken there and a hanging CloseAndRecv must not block reconnect.
  2. BoundSend is a context watchdog only. It uses time.AfterFunc to cancel the stream context on timeout. There is no per-send helper goroutine and no post-cancel grace wait. Panics from send still propagate to sendWithRecover on the caller goroutine.

Single-path cleanup

  1. One status poll for both modes. checkConnectionStatus always sleeps 5s. The multi-only faster poll and the conn.Connect() nudge are removed; the next RPC wakes Idle, and Connect() is a no-op on TransientFailure in grpc-go 1.55.
  2. Custom multi-address dialer removed. directTCPContextDialer is gone, so multi- and single-address dial share the same gRPC dialer behavior. The HTTP_PROXY / HTTPS_PROXY bypass note was dropped from docs/en/advanced-features/grpc-tls.md with it.

Docs, CI naming, and tests

  1. Renamed the existing E2E matrix case to gRPC single-backend (alongside gRPC multi-backend).
  2. TLS ServerName remains per address, matching the code, tests, and grpc-tls.md (SNI / verification follow the dialed host; an explicit credential override still wins).
  3. Failover unit coverage stops the active peer (ResolvedBackendAddresses()[0]), not a hardcoded first config address.
  4. The RPC-error policy test fails while the client stream is still active (after the first Recv), then asserts no Collect replay and that the shared ClientConn is kept.

Status

These changes are on a single commit and have passed fork CI (unit/race/lint, plugin tests, Windows plugin tests, and multi-backend E2E).

@mrproliu mrproliu 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.

You haven't changed the CHANGES.md? Please added a new feature about this one.

Comment thread .github/workflows/skywalking-go.yaml Outdated
run: make test
- name: Test Race
run: make test-race
- name: Test Race (reporter packages)

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.

What’s the difference with the existing make test-race?

Add the multi-address gRPC backend feature to CHANGES.md, and remove the
extra reporter race step that duplicates make test-race (TestRace*).
@songzhendong

Copy link
Copy Markdown
Contributor Author

CHANGES.md and CI

  1. Documented the feature in CHANGES.md. Added under 0.8.0 Features: support for multi-address gRPC backends with pick_first failover.
  2. Removed the redundant reporter race CI step. make test-race already runs -run '^TestRace' ./..., which covers the reporter race regression (TestRaceConnManagerGetRelease). The extra go test -race over the full reporter packages duplicated that coverage.

Status

These changes are on a single commit and have passed fork CI (unit/race/lint, plugin tests, Windows plugin tests, and multi-backend E2E).

mrproliu
mrproliu previously approved these changes Sep 28, 2026

@mrproliu mrproliu 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.

Generally LGTM. Once the CI all passed, it can be merge.

@songzhendong

Copy link
Copy Markdown
Contributor Author

Thanks for the review

@wu-sheng wu-sheng left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

Thanks for converging onto a single path, it is much easier to follow now. I re-reviewed at aab0e69. Three issues should be fixed before merge (1–3); the rest can be follow-ups.

Blocking

1. TLS enabled + no valid backend_service panics at startup (instrument.go#L229)

In the generated initManager, tc, err := generateTLSCredential(...) declares a new err inside the TLS branch, so connManager, err = NewConnectionManager(...) assigns to that inner variable. The outer err stays nil when NewConnectionManager returns (nil, errNoValidBackendService), NewCDSManager then calls GetConnection on a nil *ConnectionManager, and the application crashes. That contradicts the "reporting is disabled, the application continues" behavior from dc4aa2c. The shadowing was harmless before because NewConnectionManager never returned an error.

Fix:

-		tc, err := generateTLSCredential({{.Config.Reporter.GRPC.TLS.CAPath.ToGoStringValue}}, 
+		tc, tlsErr := generateTLSCredential({{.Config.Reporter.GRPC.TLS.CAPath.ToGoStringValue}}, 
 			{{.Config.Reporter.GRPC.TLS.ClientKeyPath.ToGoStringValue}},
 			{{.Config.Reporter.GRPC.TLS.ClientCertChainPath.ToGoStringValue}},
 			{{.Config.Reporter.GRPC.TLS.InsecureSkipVerify.ToGoBoolValue}})
-		if err != nil {
-			panic(fmt.Sprintf("generate go agent tls credential error: %v", err))
+		if tlsErr != nil {
+			panic(fmt.Sprintf("generate go agent tls credential error: %v", tlsErr))
 		}
Regression test: tools/go-agent/instrument/reporter/instrument_test.go (fails at aab0e69, passes with the fix)

Lint cannot see variable scopes inside a template string, so this renders initManagerFunc, stubs the constructors, and runs the emitted Go for TLS off and on.

package reporter

import (
	"html"
	"os"
	"os/exec"
	"path/filepath"
	"strings"
	"testing"

	"github.com/apache/skywalking-go/tools/go-agent/config"
	"github.com/apache/skywalking-go/tools/go-agent/tools"
)

func TestGeneratedInitManagerPropagatesConnectionError(t *testing.T) {
	if err := config.LoadConfig(""); err != nil {
		t.Fatal(err)
	}
	generated := html.UnescapeString(tools.ExecuteTemplate(initManagerFunc, struct {
		Config *config.Config
	}{Config: config.GetConfig()}))
	// Exercise the emitted Go code: lint cannot inspect variable scopes inside
	// a template string. Stub external constructors to isolate error propagation.
	generated = strings.ReplaceAll(generated, "operator.LogOperator", "interface{}")
	source := "package main\nimport (\"fmt\"; \"os\"; \"strconv\"; \"strings\"; \"time\")\n" +
		generated + initManagerErrorHarness
	sourcePath := filepath.Join(t.TempDir(), "main.go")
	if err := os.WriteFile(sourcePath, []byte(source), 0o600); err != nil {
		t.Fatal(err)
	}
	cmd := exec.Command("go", "run", sourcePath)
	if output, err := cmd.CombinedOutput(); err != nil {
		t.Fatalf("generated initManager failed to propagate connection error: %v\n%s", err, output)
	}
}

const initManagerErrorHarness = `
type ConnectionManager struct{}
type CDSManager struct{}
type PprofTaskManager struct{}

var connectionError = fmt.Errorf("no valid backend service addresses")

func generateTLSCredential(string, string, string, bool) (interface{}, error) {
	return struct{}{}, nil
}

func NewConnectionManager(interface{}, time.Duration, string, string, interface{}) (*ConnectionManager, error) {
	return nil, connectionError
}

func NewCDSManager(interface{}, string, time.Duration, *ConnectionManager) (*CDSManager, error) {
	panic("CDS must not start after connection initialization fails")
}

func NewPprofTaskManager(interface{}, string, time.Duration, *ConnectionManager, string) (*PprofTaskManager, error) {
	panic("pprof must not start after connection initialization fails")
}

func main() {
	for _, tlsEnabled := range []string{"false", "true"} {
		if err := os.Setenv("SW_AGENT_REPORTER_GRPC_TLS_ENABLE", tlsEnabled); err != nil {
			panic(err)
		}
		conn, cds, pprof, err := initManager(nil, time.Second)
		if err != connectionError || conn != nil || cds != nil || pprof != nil {
			panic(fmt.Sprintf("TLS=%s: connection error was not propagated: %v", tlsEnabled, err))
		}
	}
}
`

2. A first address that silently drops connections prevents failover (conn_manager.go#L206)

In grpc-go v1.55.0, pick_first keeps all addresses in one addrConn, and resetTransport computes a single connectDeadline that tryAllAddrs shares across every address (clientconn.go#L1150-L1220). With MinConnectTimeout: 5s, if the first address in the shuffled list silently drops the connection attempt (host powered off or network-partitioned, no RST), each attempt spends the whole deadline on it. The next address starts with the deadline already passed and fails immediately. The channel goes TransientFailure, backs off, and repeats. The shuffle is fixed per channel, so the healthy standby is never used for the life of the process. The unit and e2e tests stop servers gracefully, so the connection is refused immediately and this case is never exercised.

Suggestion: give each address its own dial timeout, e.g. a grpc.WithContextDialer that caps each attempt at a slice of the remaining deadline, plus a test with a non-responding listener as the first address.

3. Regression: gRPC target URIs are now rejected (backend_addresses.go#L70)

On main, backend_service is passed as-is to grpc.Dial, so dns:///oap-headless:11800, unix:///path.sock, etc. work. dc4aa2c preserved that by parsing only comma-separated values. Since fd76a1a, every entry goes through net.SplitHostPort, which rejects these targets ("too many colons"). parseBackendServiceList then returns errNoValidBackendService and the reporter becomes a discard reporter with only a warning. Please restore passthrough for a single entry that is not host:port.

Should fix

4. No keepalive, so a backend that stops responding without closing the connection never triggers failover (conn_manager.go#L300)

No grpc.WithKeepaliveParams is set. BoundSend and WatchConnCancelOnUnready cancel the stream, but the transport stays READY, so pick_first keeps choosing the dead backend. Each reopened stream buffers about 64KB of Sends that return nil and are lost, then blocks for 8s, and this repeats until the OS gives up on the TCP connection (about 15 minutes). Client keepalive would close the transport and let pick_first move on. For the same reason, WatchConnCancelOnUnready (L374) rarely fires: when a READY transport is lost, pick_first goes to IDLE, not TransientFailure, and a backend that stops responding without closing the connection stays READY.

5. The 8s BoundSend also applies to single-backend setups (grpc.go#L221)

I understand this was requested for both modes, but note the tradeoff: an OAP that stops reading for more than 8s (load, long GC) now gets the Collect stream reset, and everything buffered in it is dropped, where it previously just applied backpressure. With keepalive from (4) handling dead peers, this bound could be much longer or configurable.

6. One timer is created and stopped for every message sent (conn_manager.go#L307)

time.AfterFunc + Stop per Send. One timer per stream reused via Reset gives the same bound without per-message timers.

Pre-existing, touched by this PR (fine as follow-ups)

  1. CDS / profile-task / pprof calls attach no auth metadata (cds_manager.go#L82, also GetProfileTaskCommands, GetPprofTaskCommands, pprof Collect). Same on main, but in multi-backend mode the new authFailureLogger now logs "check reporter.grpc.authentication" every 30s even when the token is correct.
  2. case ConnectionStatusShutdown: break in InitCDS (cds_manager.go#L74) only exits the switch, so the loop keeps calling a closed connection forever. This PR fixed the same pattern in grpc.go with return.

Cleanup

  1. Test-only / unused exports ship in every instrumented binary (conn_manager.go#L289): the no-op "deprecated" *CancelGrace*ForTest shims, the MultiBackend*ForTest aliases, MultiBackendSend, PeekConnection, ResolvedBackendAddresses, isIPLiteralHost. The comment says they keep existing tests compiling, but those tests are new in this PR; please call BoundSend / SetBoundSendTimeoutForTest directly and drop the rest.
  2. BackendRPCContext / BackendStreamContext take serverAddr and discard it with _ = serverAddr (L319, L334); please remove the parameter.
  3. The four close*Stream functions (grpc.go#L450-L497) and the four open/watch blocks are near-identical copies; one helper each would keep future fixes in one place.
  4. E2E: verify-failover.sh#L95 counts POST:/info on the standby, but the provider also produces a POST:/info entry span and reports to its own, independently shuffled backend. If the provider picked the standby from the start, the trace half passes without the consumer failing over. Only the log check is specific to the consumer.
  5. The PR description is out of date: it still describes batched trace sends, waiting for READY, and the first endpoint as the fixed TLS identity, and says single-address behavior is unchanged. Please update it to match the current code.

Propagate connection errors when TLS is enabled; give each multi-backend
address its own dial timeout; restore single-target gRPC URI passthrough;
add client keepalive and a longer reusable BoundSend watchdog; attach auth
metadata to CDS/profile/pprof RPCs; fix CDS shutdown; harden failover E2E;
and drop unused test-only exports.
Move resolver-address inspection into a test-only helper so instrumented
binaries no longer ship the diagnostics accessor called out in review.
@songzhendong

Copy link
Copy Markdown
Contributor Author

Thanks for the review.

Blocking

  1. TLS + no valid backend_service no longer panics. The generated initManager uses tlsErr for credential setup so NewConnectionManager errors propagate on the outer err. A regression test covers the emitted template.
  2. Per-address dial timeout. Multi-backend dials use a context dialer that caps each attempt, so a silent first peer cannot consume the whole MinConnectTimeout. A blackhole-first unit test covers failover to the healthy standby.
  3. Single-target gRPC URIs restored. A lone non-host:port value (for example dns:///... or unix:///...) is passed through to grpc.Dial again. Comma-separated lists still normalize host:port entries and skip invalid ones.

Should fix

  1. Client keepalive is enabled on both single- and multi-address dial paths so half-open peers can leave READY and pick_first can move on.
  2. BoundSend default is 60s (was 8s). With keepalive handling dead peers, Collect streams are less likely to reset under temporary OAP backpressure.
  3. Reusable send watchdog. Pipeline Sends share one BoundSendWatchdog per stream (timer.Reset) instead of allocating an AfterFunc per message.
  4. Auth metadata is attached for CDS FetchConfigurations, profile GetProfileTaskCommands, and pprof GetPprofTaskCommands / Collect.
  5. CDS shutdown uses return on ConnectionStatusShutdown so the loop does not keep calling a closed connection.

Cleanup

  1. Test-only exports removed from production (CancelGrace* / MultiBackend* aliases, MultiBackendSend, PeekConnection, ResolvedBackendAddresses, isIPLiteralHost). Resolver inspection for tests lives in _test.go. BoundSendTimeoutForTest remains for unit tests.
  2. BackendRPCContext / BackendStreamContext no longer take an unused serverAddr.
  3. Shared helpers in the gRPC reporter: openBackendStream and closeStream.
  4. Multi-backend E2E asserts consumer-specific toolkit log + unique probe growth after killing the active collector (not POST:/info, which the provider may already send to standby).
  5. PR description updated to match the current single-path behavior (long-lived Collect, no wait-for-READY registration stall, per-address TLS ServerName).

WatchConnCancelOnUnready also cancels after Ready when the channel enters Idle / TransientFailure / Shutdown.

Status

These changes are on this PR and have passed fork CI (unit/race/lint, plugin tests, Windows plugin tests, and multi-backend E2E).

@wu-sheng

Copy link
Copy Markdown
Member

Thanks for the quick turnaround. I re-reviewed at 09305d3; the unit tests pass under -race.

Blocking items: verified

  1. TLS startup panic: fixed; the template now uses tlsErr (instrument.go#L229). Note that the new test only checks the template text for tc, tlsErr :=, so it guards this exact spelling rather than the behavior. Rendering the template and running it with a failing NewConnectionManager stub would catch any future shadowing.
  2. Silent first address: the per-address dialer fixes it, and TestMultiBackendPerAddrDialTimeoutFailsOverPastBlackhole shows failover past a listener that never accepts.
  3. gRPC target URIs: a single non-host:port value passes through to grpc.Dial again (backend_addresses.go#L44).

Keepalive vs. the server ping policy

The new client keepalive (conn_manager.go#L76-L80) pings every 30s. OAP's GRPCServer does not configure keepalive enforcement (GRPCServer.java#L183), so grpc-java's default applies: each ping less than 5 minutes after the previous one is a strike, and on the third strike the server sends GOAWAY ENHANCE_YOUR_CALM / too_many_pings and closes the connection.

I reproduced this with this PR's ConnectionManager against a stock grpc-go server, which has the same 5-minute default, and an idle Collect stream:

ERROR: [transport] Client received GoAway with error code ENHANCE_YOUR_CALM and debug data equal to ASCII "too_many_pings".
t=120.0s state READY -> IDLE
original stream after idle: rpc error: code = Unavailable ... received prior goaway: code: ENHANCE_YOUR_CALM, debug data: "too_many_pings"

With the default check_interval (20s) this does not happen in practice. grpc-go only pings after 30s with no reads, and the management heartbeat response arrives every 20s. It does happen when a user sets check_interval to about 2 minutes or more: every Collect stream on the connection is dropped on each GOAWAY, and this repeats a few times until grpc-go has doubled its ping interval past 5 minutes. For reference, the Python agent intentionally omits keepalive for this reason (grpc_channel.py#L34).

Suggestion: derive the keepalive time from the heartbeat, e.g. Time: max(30s, checkInterval + 10s). While the backend is healthy, the heartbeat response then always comes first and pings never fire. When it goes dark, the ping still detects it within Time + Timeout.

Three or more addresses share the connect budget

Each address gets 2s (conn_manager.go#L74), but grpc-go still shares one MinConnectTimeout (5s, L223) across the whole address list. If the first two shuffled addresses are silent, the third gets about 1s on the first attempt, and a fourth gets nothing. It recovers within about 20s as the backoff grows the connect deadline toward check_interval, so this is minor. Setting MinConnectTimeout to max(5s, perAddrTimeout * len(backends)) would remove the gap.

Leftover duplication

closeStream was added, but the same 7-line CloseAndRecv / io.EOF closure is still pasted 8 times (grpc.go#L295, L308, L362, L373, L422, L433, L492, L518). A generic helper taking the stream would remove it:

func closeStream[R any](r *gRPCReporter, cancel context.CancelFunc,
	stream interface{ CloseAndRecv() (R, error) }, errLog string) {
	if err := reporter.BoundSend(cancel, func() error {
		if _, err := stream.CloseAndRecv(); err != io.EOF {
			return err
		}
		return nil
	}, 0); err != nil {
		r.logger.Errorf("%s %v", errLog, err)
	}
}

Also, openBackendStream returns a ctx that every caller discards with _ = ctx (L284, L347, L411, L472); the return value can be dropped.

Derive client keepalive Time from check_interval to avoid too_many_pings
against default server policy; scale multi-backend MinConnectTimeout by
address count; collapse CloseAndRecv into a generic helper; and run the
generated initManager under stubs for TLS on/off.
Satisfy gocritic typeDefFirst so lint passes in CI.
@songzhendong

Copy link
Copy Markdown
Contributor Author

Thanks for the review.

Keepalive and connect budget

  1. Client keepalive Time is now max(30s, check_interval + 10s). While the backend is healthy, management heartbeats suppress pings and avoid too_many_pings against the default server policy. When the peer goes dark, keepalive still detects it within Time + Timeout.
  2. Multi-backend MinConnectTimeout is max(5s, per-address timeout × address count), so later peers still get dial budget when earlier shuffled addresses are silent.

Cleanup and tests

  1. CloseAndRecv / io.EOF handling is one generic closeStream helper; the eight pasted closures are gone. openBackendStream no longer returns an unused ctx.
  2. The TLS initManager regression test now renders the template and go runs it with stubbed constructors for TLS off and on (SW_AGENT_REPORTER_GRPC_TLS_ENABLE), so error propagation is exercised rather than only string-matching tlsErr.

Status

These changes are on this PR and have passed fork CI (unit/race/lint, plugin tests, Windows plugin tests, and multi-backend E2E).

@wu-sheng

Copy link
Copy Markdown
Member

Two additional issues remain at 09305d3:

  1. TLS handshake traffic clears the per-address timeout before HTTP/2 is ready. In conn_manager.go#L278-L281, any successful raw socket read clears the deadline. With TLS, that first read is handshake traffic. If the peer completes TLS but never sends HTTP/2 SETTINGS, it consumes the entire shared connection deadline; later addresses start with that deadline already expired. The fixed resolver order repeats this on subsequent attempts, preventing failover even with a healthy second endpoint.

    Reproduced with a TLS listener that completes the handshake and then remains silent, followed by a healthy TLS gRPC server: the channel never reached READY during a 13-second test. The existing never-accept test receives no bytes and does not cover this case. The per-address timeout needs to remain effective through TLS and HTTP/2 negotiation, rather than being cleared after the first raw read.

  2. Canceling streams when the channel becomes Idle breaks graceful GOAWAY draining. conn_manager.go#L447-L450 treats Idle as evidence that an existing stream has failed. However, graceful GOAWAY can move the channel to Idle while already accepted streams remain usable. The watcher cancels those streams, and the reporter discards the next dequeued message when Send fails. This affects single-backend configurations too and is independent of the too_many_pings issue.

    Reproduced with a gRPC server configured with MaxConnectionAge=300ms and MaxConnectionAgeGrace=5s: without the watcher, the channel was Idle but the second Send and CloseAndRecv succeeded; with the watcher, the context was canceled and the second Send returned EOF. Avoid treating channel Idle alone as a reason to cancel established streams; stream errors and the send watchdog can detect failed operations while accepted streams are allowed to drain.

Clear the per-address dial deadline only after the first post-handshake
read (HTTP/2 SETTINGS) so a TLS-only silent peer cannot starve failover.
Stop treating channel Idle as stream failure so graceful GOAWAY can drain.
Rely on embedded TransportCredentials for OverrideServerName so
staticcheck SA1019 no longer fails CI lint.
@songzhendong

Copy link
Copy Markdown
Contributor Author

Thanks for the review.

1. TLS handshake clearing the per-address dial deadline — Fixed. The context dialer still sets a per-address deadline for the full TCP/TLS/HTTP/2 window, but no longer clears it on the first raw socket read. A credentials wrapper clears the deadline only after ClientHandshake returns and the first post-handshake application read succeeds (typically the peer's HTTP/2 SETTINGS). A peer that completes TLS and then stays silent still hits the deadline so pick_first can move on. Added a regression test with a TLS listener that handshakes and then remains silent ahead of a healthy TLS gRPC backend.

2. Canceling streams on channel Idle (graceful GOAWAY) — Fixed. WatchConnCancelOnUnready no longer treats Idle alone as failure; it cancels only on TransientFailure or Shutdown. Graceful GOAWAY can leave accepted streams usable while the channel is Idle; stream errors and BoundSend still cover failed operations. Added a regression test with MaxConnectionAge / MaxConnectionAgeGrace that asserts the second Send and CloseAndRecv succeed during Idle drain with the watcher running.

Please take another look when convenient.

wu-sheng and others added 2 commits September 29, 2026 16:11
…cher

Keep the per-address dial deadline until the first complete HTTP/2 frame
(the server SETTINGS preface) has been read. Clearing it on the first
post-handshake read let a peer that sends only the 9-byte frame header and
withholds the payload consume the shared connect deadline, so pick_first
never reached the healthy standby.

Remove WatchConnCancelOnUnready. Channel connectivity state describes the
channel, not the transport a stream runs on: after a graceful GOAWAY, a
refused reconnect moved the channel to TRANSIENT_FAILURE and the watcher
canceled a stream that was still draining. A stream on a dead transport
already fails on its own, and keepalive plus BoundSend cover dead peers
and stuck sends.

Add regression tests for a partial SETTINGS first peer and for a draining
stream after GracefulStop with a refused reconnect.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01Q5YiM91L3aNmFHTZWTWTMV
The connection-state watcher was removed in the previous commit. Describe
the remaining mechanisms instead: keepalive closes half-open transports and
the bounded send timeout unblocks a stuck Send. Also update the dial-deadline
comments to say the deadline clears after the server's first complete HTTP/2
frame, not the first post-handshake read.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01Q5YiM91L3aNmFHTZWTWTMV
@wu-sheng
wu-sheng merged commit 6d75af8 into apache:main Sep 29, 2026
50 checks passed
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

enhancement New feature or request

Projects

None yet

Development

Successfully merging this pull request may close these issues.

4 participants