From 2630089cee0abc67794a9ebc49f675a251da04f6 Mon Sep 17 00:00:00 2001 From: songzhendong <289505773@qq.com> Date: Sun, 16 Aug 2026 22:31:11 +0800 Subject: [PATCH] feat(grpc): native multi-backend failover via grpc-js --- .github/workflows/test.yaml | 33 ++- docs/en/advanced/troubleshooting.md | 51 +++- docs/en/setup/configuration.md | 6 +- src/agent/core/boot/ServiceManager.ts | 8 +- src/agent/core/meter/MeterSender.ts | 133 +++++----- .../core/remote/BackendAddressResolver.ts | 161 ++++++++++++ src/agent/core/remote/GRPCChannel.ts | 48 ++-- src/agent/core/remote/GRPCChannelManager.ts | 245 +++++++++++++----- .../core/remote/ServiceManagementClient.ts | 61 ++--- .../core/remote/TraceSegmentServiceClient.ts | 118 ++++----- src/agent/core/remote/coalesceReport.ts | 188 ++++++++++++++ src/index.ts | 27 +- src/logging/index.ts | 7 +- src/trace/context/SpanContext.ts | 47 +++- tests/plugins/common/Dockerfile.agent | 19 +- tests/plugins/common/base-compose.yml | 11 + tests/remote-e2e/common/server.ts | 61 +++++ .../static-failover/docker-compose.yml | 71 +++++ tests/remote-e2e/static-failover/test.ts | 81 ++++++ tests/remote-e2e/support/helpers.ts | 146 +++++++++++ tests/remote/BackendAddressResolver.test.ts | 63 +++++ tests/remote/GRPCChannelManager.test.ts | 159 ++++++++++-- tests/remote/MeterSender.test.ts | 57 +++- .../remote/TraceSegmentServiceClient.test.ts | 130 ++++++++++ tests/remote/coalesceReport.test.ts | 164 ++++++++++++ 25 files changed, 1782 insertions(+), 313 deletions(-) create mode 100644 src/agent/core/remote/BackendAddressResolver.ts create mode 100644 src/agent/core/remote/coalesceReport.ts create mode 100644 tests/remote-e2e/common/server.ts create mode 100644 tests/remote-e2e/static-failover/docker-compose.yml create mode 100644 tests/remote-e2e/static-failover/test.ts create mode 100644 tests/remote-e2e/support/helpers.ts create mode 100644 tests/remote/BackendAddressResolver.test.ts create mode 100644 tests/remote/TraceSegmentServiceClient.test.ts create mode 100644 tests/remote/coalesceReport.test.ts diff --git a/.github/workflows/test.yaml b/.github/workflows/test.yaml index 4b0fe787..78930dd0 100644 --- a/.github/workflows/test.yaml +++ b/.github/workflows/test.yaml @@ -73,7 +73,7 @@ jobs: - name: Unit tests on Node@${{ matrix.node-version }} run: | npm i - npx jest --testPathIgnorePatterns "/node_modules/" "/tests/plugins/" --runInBand + npx jest --testPathIgnorePatterns "/node_modules/" "/tests/plugins/" "/tests/remote-e2e/" --runInBand build-matrix: runs-on: ubuntu-latest @@ -119,6 +119,37 @@ jobs: npm i npm run test tests/plugins/${{ matrix.plugin }}/ + + TestRemoteE2E: + runs-on: ubuntu-latest + timeout-minutes: 45 + strategy: + matrix: + node-version: [ 20, 22, 24 ] + env: + SW_NODE_VERSION: ${{ matrix.node-version }} + steps: + - uses: actions/checkout@v4 + with: + submodules: true + + - uses: actions/cache@v4 + with: + path: ~/.npm + key: ${{ runner.os }}-node-${{ hashFiles('**/package-lock.json') }} + restore-keys: | + ${{ runner.os }}-node- + + - name: Set Up NodeJS ${{ matrix.node-version }} + uses: actions/setup-node@v4 + with: + node-version: ${{ matrix.node-version }} + + - name: Remote E2E On Node@${{ matrix.node-version }} (static-failover) + run: | + npm i + npm run test tests/remote-e2e/static-failover/ + TestLib: runs-on: ubuntu-latest timeout-minutes: 30 diff --git a/docs/en/advanced/troubleshooting.md b/docs/en/advanced/troubleshooting.md index 4218ee48..5636867a 100644 --- a/docs/en/advanced/troubleshooting.md +++ b/docs/en/advanced/troubleshooting.md @@ -15,7 +15,8 @@ logs. 6. Check that the path, suffix, and HTTP method are not filtered by trace settings. 7. Check the agent log for plugin version or OAP connection errors. -If the process is short-lived, call `await agent.flush()` before exit. +If the process is short-lived, call `await agent.flush()` before exit. `flush()` is a bounded +best-effort attempt, not a delivery guarantee. See [Data is missing when the process stops](#data-is-missing-when-the-process-stops). ## OAP connection errors @@ -25,8 +26,12 @@ The collector setting must use `host:port` form: export SW_AGENT_COLLECTOR_BACKEND_SERVICES=oap.example.com:11800 ``` -Check DNS, network access, firewall rules, and the OAP gRPC port. The current agent uses only the -first item if the value is a comma-separated list. +Check DNS, network access, firewall rules, and the OAP gRPC port. + +One `host:port` uses the grpc-js `dns:` resolver (all A/AAAA records become endpoints and are +re-resolved periodically). A comma-separated list uses a static resolver with `pick_first`: each +name is a literal endpoint only — no DNS expansion or re-resolution per name, so discovery is weaker +than a single DNS name (for example a headless Kubernetes service). An option passed to `agent.start()` replaces the environment value. This includes an empty string: @@ -39,11 +44,26 @@ For a TLS OAP endpoint, set `SW_AGENT_SECURE=true`. The agent uses the system tr configuration for a custom CA or client certificate. Set `SW_AGENT_AUTHENTICATION` if OAP requires an agent token. +Under TLS with multiple hostnames, certificate verification follows the channel authority (the first +list entry in the configured target). Endpoint pick order may be shuffled by grpc-js, but the target +string — and therefore authority / SNI — stays in config order. Every backend must present a +certificate that shares the needed SANs, or failover handshakes fail. Prefer one DNS name with +multiple A/AAAA records for TLS high availability. + +Channel disconnect lines are logged at `error` and recover lines at `warn` (throttled separately so a +recover line is not swallowed by the disconnect window). For per-address grpc-js detail, set +`GRPC_TRACE=pick_first,subchannel`. + ## Traces are missing during an OAP outage The agent keeps finished segments in a memory buffer. When the buffer reaches -`SW_AGENT_MAX_BUFFER_SIZE`, it removes the oldest finished segment. Restore the OAP connection; -increasing the buffer only delays data loss and uses more process memory. +`SW_AGENT_MAX_BUFFER_SIZE`, it removes the oldest finished segment and may log that the trace buffer +reached maximum size. Restore the OAP connection; increasing the buffer only delays data loss and +uses more process memory. + +If a report attempt fails, the agent may log that it discarded N trace segment(s) after report +failure. Those segments are not re-sent. Reporting is best-effort: failures discard data, and +`flush()` only waits briefly then tries once more. ## A library has no spans @@ -65,12 +85,28 @@ See [Node.js runtime metrics](../features/runtime-metrics.md) for meter names an ## Agent logs are hard to find -The default agent log level is `error`. +The default agent log level is `warn`. + +| Level | What you typically see | +| --- | --- | +| `error` | Auth rejection, channel disconnect (throttled), fatal boot failures | +| `warn` (default) | Channel recovered (throttled), trace buffer full, discarded segment batches, meter report failures (throttled) | +| `info` / `debug` | Lifecycle noise; span debug lines | + +Levels below the configured threshold are silent by design (`warn` / `info` / `debug` become no-ops +when the threshold is higher). - When `NODE_ENV` is not `production`, logs go to the console. - When `NODE_ENV=production`, logs go to `skywalking.log` in the process working directory. - Set `SW_LOGGING_TARGET=console` to use the console in production. +## Agent to OAP over an HTTP proxy + +Agent to OAP over an HTTP proxy is not supported. Every gRPC channel sets `grpc.enable_http_proxy=0`, +so host `http_proxy` / `https_proxy` never affect OAP uplink. This applies to single-address and +multi-address targets. After upgrade, traffic that previously relied on an HTTP CONNECT proxy to +reach OAP will no longer use that proxy. + ## Data is missing when the process stops `agent.destroy()` stops reporters but does not flush them. Use this order: @@ -80,6 +116,9 @@ await agent.flush(); agent.destroy(); ``` +`flush()` waits a short time for in-flight work and may start one more report attempt. It does not +guarantee delivery if OAP is slow or unreachable. + Do not use stop and restart as a normal agent update method inside one process. Module patches stay installed after `destroy()`. diff --git a/docs/en/setup/configuration.md b/docs/en/setup/configuration.md index 2660baa2..e4d77b3f 100644 --- a/docs/en/setup/configuration.md +++ b/docs/en/setup/configuration.md @@ -24,7 +24,7 @@ Node.js process. `SW_DISABLE` is checked when `start()` is called. | --- | --- | --- | --- | | `SW_AGENT_NAME` | `serviceName` | `your-nodejs-service` | Service name shown in SkyWalking. | | `SW_AGENT_INSTANCE` | `serviceInstance` | Host name | Service instance name shown in SkyWalking. | -| `SW_AGENT_COLLECTOR_BACKEND_SERVICES` | `collectorAddress` | `127.0.0.1:11800` | OAP gRPC address in `host:port` form. | +| `SW_AGENT_COLLECTOR_BACKEND_SERVICES` | `collectorAddress` | `127.0.0.1:11800` | OAP gRPC address(es). One `host:port` uses grpc-js `dns:` (all A/AAAA become endpoints; periodically re-resolved). A comma-separated list uses a static resolver with `pick_first` (literal endpoints only; no per-name DNS expansion or re-resolution). Endpoint pick order is shuffled by grpc-js (`shuffleAddressList`); the target string keeps config order so TLS authority / SNI stay on the first list entry. Under TLS, all backends must present certificates that share the needed SANs. Prefer one DNS name with multiple A/AAAA records for TLS high availability. | | `SW_AGENT_SECURE` | `secure` | `false` | Use TLS for the OAP gRPC connection. | | `SW_AGENT_AUTHENTICATION` | `authorization` | Not set | Authentication token sent to OAP. | | `SW_AGENT_TRACE_TIMEOUT` | `traceTimeout` | `10000` | gRPC deadline in milliseconds for trace and meter reports and service management requests. Must be a positive integer. | @@ -35,14 +35,12 @@ For token authentication, set the same token in OAP with `SW_AUTHENTICATION`. Se When `secure` is enabled, the agent uses the system trust store. It does not provide options for a custom CA, client certificate, or mutual TLS. -Use one OAP address. If a comma-separated list is set, the current agent uses only the first entry. - ## Agent control and logging | Environment variable | `agent.start()` option | Default | Description | | --- | --- | --- | --- | | `SW_DISABLE` | None | Not set | Set the exact value `true` to keep the agent stopped. | -| `SW_AGENT_LOGGING_LEVEL` | None | `error` | Agent log level: `error`, `warn`, `info`, or `debug`. | +| `SW_AGENT_LOGGING_LEVEL` | None | `warn` | Agent log level: `error`, `warn`, `info`, or `debug`. | | `SW_LOGGING_TARGET` | None | See below | Set to `console` to log to the console in production. | | `SW_AGENT_MAX_BUFFER_SIZE` | `maxBufferSize` | `1000` | Limit for active and buffered trace segments. Must be a positive integer. | | `SW_AGENT_DISABLE_PLUGINS` | `disablePlugins` | Empty | Comma-separated plugin file names without the `Plugin` suffix, such as `mysql,express`. | diff --git a/src/agent/core/boot/ServiceManager.ts b/src/agent/core/boot/ServiceManager.ts index e5dd18db..43b6daeb 100644 --- a/src/agent/core/boot/ServiceManager.ts +++ b/src/agent/core/boot/ServiceManager.ts @@ -51,7 +51,7 @@ class ServiceManager { try { service.prepare(); } catch (error) { - logger.error('ServiceManager prepare failed: %s', error); + logger.error(`ServiceManager prepare failed: ${error}`); } } @@ -59,7 +59,7 @@ class ServiceManager { try { service.boot(); } catch (error) { - logger.error('ServiceManager boot failed: %s', error); + logger.error(`ServiceManager boot failed: ${error}`); } } @@ -67,7 +67,7 @@ class ServiceManager { try { service.onComplete(); } catch (error) { - logger.error('ServiceManager onComplete failed: %s', error); + logger.error(`ServiceManager onComplete failed: ${error}`); } } @@ -79,7 +79,7 @@ class ServiceManager { try { service.shutdown(); } catch (error) { - logger.error('ServiceManager shutdown failed: %s', error); + logger.error(`ServiceManager shutdown failed: ${error}`); } } this.bootedServices.clear(); diff --git a/src/agent/core/meter/MeterSender.ts b/src/agent/core/meter/MeterSender.ts index 93b8d713..2101b423 100644 --- a/src/agent/core/meter/MeterSender.ts +++ b/src/agent/core/meter/MeterSender.ts @@ -28,9 +28,16 @@ import { RuntimeSnapshot } from './RuntimeSampler'; import GRPCChannelManager from '../remote/GRPCChannelManager'; import { GRPCChannelListener } from '../remote/GRPCChannelListener'; import { GRPCChannelStatus } from '../remote/GRPCChannelStatus'; +import { + coalesceReport, + flushCoalesced, + FLUSH_WAIT_MS, + ReportCoalesceState, + runCollectStream, +} from '../remote/coalesceReport'; const logger = createLogger(__filename); -const logReportError = throttled(logger, 'error', 30000); +const logReportError = throttled(logger, 'warn', 30000); /** Reports Node.js runtime metrics via gRPC MeterReportService (Go/Python-compatible pipeline). */ export default class MeterSender implements BootService, GRPCChannelListener { @@ -41,7 +48,7 @@ export default class MeterSender implements BootService, GRPCChannelListener { /** Latest gauge snapshot only — stale samples have no value after reconnect. */ private latestSnapshot?: RuntimeSnapshot; private timer?: NodeJS.Timeout; - private reporting?: Promise; + private readonly reportState: ReportCoalesceState = {}; private collector!: RuntimeMetricsCollector; @@ -75,12 +82,7 @@ export default class MeterSender implements BootService, GRPCChannelListener { if (!this.channelManager) { return undefined; } - - return new MeterReportServiceClient( - config.collectorAddress, - grpc.credentials.createInsecure(), - this.channelManager.getClientOptions(), - ); + return this.channelManager.createClient(MeterReportServiceClient); } private startTimer(): void { @@ -95,83 +97,61 @@ export default class MeterSender implements BootService, GRPCChannelListener { } private collectSample(): void { + // Always sample so RuntimeSampler CPU deltas stay on the report period (not the outage length). this.latestSnapshot = this.collector.sample(); } private reportBufferedMetrics(): Promise { + return coalesceReport( + this.reportState, + () => this.doReportBufferedMetrics(), + () => this.closed, + ); + } + + private doReportBufferedMetrics(): Promise { if (this.closed) { return Promise.resolve(); } - if (this.reporting) { - return this.reporting; + if (!config.serviceName || !config.serviceInstance) { + return Promise.resolve(); } - this.reporting = this.doReportBufferedMetrics().finally(() => { - this.reporting = undefined; - }); - - return this.reporting; - } - - private doReportBufferedMetrics(): Promise { - return new Promise((resolve) => { - try { - if (this.closed) { - resolve(); - return; - } - - const snapshot = this.latestSnapshot; - this.latestSnapshot = undefined; - if (!snapshot || this.status !== GRPCChannelStatus.CONNECTED || !this.reporterClient) { - resolve(); - return; - } + // Check connectivity before consuming the snapshot so disconnect-window samples are kept. + if (this.status !== GRPCChannelStatus.CONNECTED || !this.reporterClient) { + return Promise.resolve(); + } - if (!config.serviceName || !config.serviceInstance) { - resolve(); - return; - } + const snapshot = this.latestSnapshot; + if (!snapshot) { + return Promise.resolve(); + } + this.latestSnapshot = undefined; - const stream = this.reporterClient.collect( - new grpc.Metadata(), - { deadline: Date.now() + (config.traceTimeout || 10000) }, - (error: grpc.ServiceError | null) => { - if (error) { - logReportError('Failed to report runtime meter data', error); - this.reportGrpcError(error); - } - resolve(); - }, - ); - - try { - let metadataWritten = false; - for (const meterData of this.collector.toMeterData(snapshot)) { - // Meter.proto: service / instance / timestamp on the first stream element only. - if (!metadataWritten) { - meterData - .setService(config.serviceName) - .setServiceinstance(config.serviceInstance) - .setTimestamp(snapshot.collectedAt); - metadataWritten = true; - } - stream.write(meterData); - } - } finally { - try { - stream.end(); - } catch (error) { - logReportError('Failed to end meter collect stream', error); - resolve(); + const client = this.reporterClient; + const serviceName = config.serviceName; + const serviceInstance = config.serviceInstance; + return runCollectStream({ + open: (onStatus) => + client.collect(new grpc.Metadata(), { deadline: Date.now() + (config.traceTimeout || 10000) }, onStatus), + writeAll: (stream) => { + let metadataWritten = false; + for (const meterData of this.collector.toMeterData(snapshot)) { + // Meter.proto: service / instance / timestamp on the first stream element only. + if (!metadataWritten) { + meterData.setService(serviceName).setServiceinstance(serviceInstance).setTimestamp(snapshot.collectedAt); + metadataWritten = true; } + stream.write(meterData); } - } catch (error) { - logReportError('Failed to report runtime meter data', error); + }, + onFailure: (reason, error) => { + logReportError(reason, error); this.reportGrpcError(error); - resolve(); - } + }, + openFailureReason: 'Failed to report runtime meter data', + endFailureReason: 'Failed to end meter collect stream', }); } @@ -183,12 +163,21 @@ export default class MeterSender implements BootService, GRPCChannelListener { this.channelManager?.reportError(error); } + /** + * Best-effort: one shared FLUSH_WAIT_MS budget for forceReport of latest snapshot, then in-flight wait. + */ flush(): Promise | null { if (this.closed) { return null; } this.collectSample(); - return this.reportBufferedMetrics(); + return flushCoalesced( + this.reportState, + () => this.doReportBufferedMetrics(), + () => this.closed, + () => this.latestSnapshot != null, + FLUSH_WAIT_MS, + ); } shutdown(): void { @@ -197,7 +186,7 @@ export default class MeterSender implements BootService, GRPCChannelListener { clearInterval(this.timer); this.timer = undefined; } - this.reporting = undefined; + this.reportState.reporting = undefined; this.reporterClient = undefined; this.latestSnapshot = undefined; this.collector.destroy(); diff --git a/src/agent/core/remote/BackendAddressResolver.ts b/src/agent/core/remote/BackendAddressResolver.ts new file mode 100644 index 00000000..e01eefd6 --- /dev/null +++ b/src/agent/core/remote/BackendAddressResolver.ts @@ -0,0 +1,161 @@ +/*! + * + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + * + */ + +import * as grpc from '@grpc/grpc-js'; +import { createLogger } from '../../../logging'; + +const logger = createLogger(__filename); + +const SW_STATIC_SCHEME = 'sw-static'; + +let swStaticResolverRegistered = false; + +/** + * Parse one host:port via grpc.experimental.splitHostPort. + * Returns normalized host:port or null when invalid (logged at error). + */ +export function tryParseHostPort(entry: string): string | null { + const trimmed = entry.trim(); + if (!trimmed) { + return null; + } + const parsed = grpc.experimental.splitHostPort(trimmed); + if (!parsed?.host || parsed.port == null) { + logger.error(`Invalid collector address: ${entry}`); + return null; + } + if (parsed.port <= 0 || parsed.port > 65535) { + logger.error(`Invalid collector address (bad port): ${entry}`); + return null; + } + const host = parsed.host.includes(':') ? `[${parsed.host}]` : parsed.host; + return `${host}:${parsed.port}`; +} + +/** Parse comma-separated backend host:port entries. Invalid entries are logged and dropped. */ +export function parseStaticBackendAddresses(raw: string): string[] { + const result: string[] = []; + for (const part of raw.split(',')) { + const normalized = tryParseHostPort(part); + if (normalized) { + result.push(normalized); + } + } + return result; +} + +/** + * Build a grpc-js channel target. + * - One address: plain host:port (default dns: resolver — multi-IP + periodic re-resolve + natural TLS authority). + * - Multiple: sw-static:/// list for pick_first across explicit backends (literal endpoints only; + * no DNS expansion / re-resolution of each name — weaker discovery than a single dns: target). + */ +export function buildNativeGrpcTarget(addresses: string[]): string { + // Caller (openChannel) already requires a non-empty list. + if (addresses.length === 1) { + return addresses[0]!; + } + ensureSwStaticResolverRegistered(); + return `${SW_STATIC_SCHEME}:///${addresses.join(',')}`; +} + +function ensureSwStaticResolverRegistered(): void { + if (swStaticResolverRegistered) { + return; + } + const { registerResolver, statusOrFromValue, statusOrFromError, splitHostPort: grpcSplit } = grpc.experimental; + + /** + * Static multi-address resolver for comma-separated backends. + * Mirrors grpc-js resolver-ip.js (static endpoint list, no DNS expansion). + */ + class SwStaticResolver { + private readonly listener: grpc.experimental.ResolverListener; + private readonly endpoints: grpc.experimental.Endpoint[]; + private readonly error: { code: number; details: string; metadata: grpc.Metadata } | null; + private hasReturnedResult = false; + + constructor( + target: grpc.experimental.GrpcUri, + listener: grpc.experimental.ResolverListener, + _channelOptions: grpc.ChannelOptions, + ) { + this.listener = listener; + this.endpoints = []; + this.error = null; + + if (target.scheme !== SW_STATIC_SCHEME) { + this.error = { + code: grpc.status.UNAVAILABLE, + details: `Unrecognized scheme ${target.scheme} in sw-static resolver`, + metadata: new grpc.Metadata(), + }; + return; + } + + const pathList = target.path + .split(',') + .map((p) => p.trim()) + .filter(Boolean); + const endpoints: grpc.experimental.Endpoint[] = []; + for (const path of pathList) { + const hp = grpcSplit(path); + if (!hp?.host || hp.port == null) { + this.error = { + code: grpc.status.UNAVAILABLE, + details: `Failed to parse sw-static address ${path}`, + metadata: new grpc.Metadata(), + }; + return; + } + endpoints.push({ addresses: [{ host: hp.host, port: hp.port }] }); + } + this.endpoints = endpoints; + } + + updateResolution(): void { + if (this.hasReturnedResult) { + return; + } + this.hasReturnedResult = true; + process.nextTick(() => { + if (this.error) { + this.listener(statusOrFromError(this.error), {}, null, ''); + } else { + this.listener(statusOrFromValue(this.endpoints), {}, null, ''); + } + }); + } + + destroy(): void { + this.hasReturnedResult = false; + } + + static getDefaultAuthority(target: grpc.experimental.GrpcUri): string { + const first = target.path.split(',')[0]?.trim(); + if (!first) { + throw new Error('sw-static target path must contain at least one host:port'); + } + return first; + } + } + + registerResolver(SW_STATIC_SCHEME, SwStaticResolver); + swStaticResolverRegistered = true; +} diff --git a/src/agent/core/remote/GRPCChannel.ts b/src/agent/core/remote/GRPCChannel.ts index 4eae5588..9a315103 100644 --- a/src/agent/core/remote/GRPCChannel.ts +++ b/src/agent/core/remote/GRPCChannel.ts @@ -18,7 +18,7 @@ */ import * as grpc from '@grpc/grpc-js'; -import { ClientOptions, connectivityState } from '@grpc/grpc-js'; +import { ClientOptions, connectivityState, ChannelOptions } from '@grpc/grpc-js'; import ChannelBuilder, { ChannelBuildContext } from './ChannelBuilder'; import ChannelDecorator from './ChannelDecorator'; @@ -26,31 +26,41 @@ export default class GRPCChannel { private readonly originChannel: grpc.Channel; private readonly interceptors: grpc.Interceptor[]; - private constructor(host: string, port: number, channelBuilders: ChannelBuilder[], decorators: ChannelDecorator[]) { + /** + * Builders must spread `context.options` when they replace the options object + * so native channel options (keepalive / service_config) are preserved. + * `extraOptions` are merged exactly once into the initial context. + */ + private constructor( + target: string, + channelBuilders: ChannelBuilder[], + decorators: ChannelDecorator[], + extraOptions: ChannelOptions, + ) { let context: ChannelBuildContext = { credentials: grpc.credentials.createInsecure(), - options: {}, + options: { ...extraOptions }, }; for (const builder of channelBuilders) { context = builder.build(context); } - this.originChannel = new grpc.Channel(`${host}:${port}`, context.credentials, context.options); + this.originChannel = new grpc.Channel(target, context.credentials, context.options); this.interceptors = decorators.map((decorator) => decorator.build()); } static create( - host: string, - port: number, + target: string, channelBuilders: ChannelBuilder[], decorators: ChannelDecorator[], + extraOptions: ChannelOptions = {}, ): GRPCChannel { - return new GRPCChannel(host, port, channelBuilders, decorators); + return new GRPCChannel(target, channelBuilders, decorators, extraOptions); } - static newBuilder(host: string, port: number): GRPCChannelBuilder { - return new GRPCChannelBuilder(host, port); + static newBuilder(target: string): GRPCChannelBuilder { + return new GRPCChannelBuilder(target); } getChannel(): grpc.Channel { @@ -68,20 +78,23 @@ export default class GRPCChannel { return this.originChannel.getConnectivityState(requestConnection) === connectivityState.READY; } + getConnectivityState(requestConnection = false): connectivityState { + return this.originChannel.getConnectivityState(requestConnection); + } + shutdownNow(): void { this.originChannel.close(); } } class GRPCChannelBuilder { - private readonly host: string; - private readonly port: number; + private readonly target: string; private readonly channelBuilders: ChannelBuilder[] = []; private readonly decorators: ChannelDecorator[] = []; + private extraOptions: ChannelOptions = {}; - constructor(host: string, port: number) { - this.host = host; - this.port = port; + constructor(target: string) { + this.target = target; } addManagedChannelBuilder(builder: ChannelBuilder): this { @@ -94,7 +107,12 @@ class GRPCChannelBuilder { return this; } + withChannelOptions(options: ChannelOptions): this { + this.extraOptions = { ...this.extraOptions, ...options }; + return this; + } + build(): GRPCChannel { - return GRPCChannel.create(this.host, this.port, this.channelBuilders, this.decorators); + return GRPCChannel.create(this.target, this.channelBuilders, this.decorators, this.extraOptions); } } diff --git a/src/agent/core/remote/GRPCChannelManager.ts b/src/agent/core/remote/GRPCChannelManager.ts index 5eb2fedc..8fbf0163 100644 --- a/src/agent/core/remote/GRPCChannelManager.ts +++ b/src/agent/core/remote/GRPCChannelManager.ts @@ -18,11 +18,12 @@ */ import * as grpc from '@grpc/grpc-js'; -import { ClientOptions } from '@grpc/grpc-js'; +import { ClientOptions, ChannelOptions, ChannelCredentials } from '@grpc/grpc-js'; import config from '../../../config/AgentConfig'; -import { createLogger } from '../../../logging'; +import { createLogger, throttled } from '../../../logging'; import AgentIDDecorator from './AgentIDDecorator'; import AuthenticationDecorator from './AuthenticationDecorator'; +import { buildNativeGrpcTarget, parseStaticBackendAddresses } from './BackendAddressResolver'; import GRPCChannel from './GRPCChannel'; import { GRPCChannelListener } from './GRPCChannelListener'; import { GRPCChannelStatus } from './GRPCChannelStatus'; @@ -31,57 +32,95 @@ import StandardChannelBuilder from './StandardChannelBuilder'; import TLSChannelBuilder from './TLSChannelBuilder'; const logger = createLogger(__filename); +const logAuthRejected = throttled(logger, 'error', 30000); +const logChannelDisconnected = throttled(logger, 'error', 30000); +const logChannelRecovered = throttled(logger, 'warn', 30000); -function isGrpcNetworkError(error: unknown): boolean { +/** Placeholder authority — ignored when channelOverride is set (grpc-js client.js). */ +const STUB_AUTHORITY = 'skywalking-backend'; + +function isGrpcAuthError(error: unknown): boolean { const code = (error as grpc.ServiceError | undefined)?.code; + return code === grpc.status.PERMISSION_DENIED || code === grpc.status.UNAUTHENTICATED; +} +function isGrpcNetworkError(error: unknown): boolean { + const code = (error as grpc.ServiceError | undefined)?.code; + if (isGrpcAuthError(error)) { + return false; + } return ( code === grpc.status.UNAVAILABLE || - code === grpc.status.PERMISSION_DENIED || - code === grpc.status.UNAUTHENTICATED || code === grpc.status.RESOURCE_EXHAUSTED || - code === grpc.status.UNKNOWN + code === grpc.status.UNKNOWN || + code === grpc.status.DEADLINE_EXCEEDED ); } +function nativeChannelOptions(): ChannelOptions { + // No gRPC keepalive: stock OAP (grpc-java) rejects frequent idle pings with GOAWAY + // ENHANCE_YOUR_CALM. Agent traffic (trace/heartbeat/metrics) provides liveness. + // + // Agent→OAP via HTTP proxy is not supported. Disable grpc-js proxy interception + // uniformly (including single-address dns: targets that grpc-js could otherwise + // proxy) so host-app http_proxy/https_proxy cannot affect the agent channel — + // multi-address targets are also broken under HTTP CONNECT (proxy sees the + // unresolved comma list). + return { + 'grpc.enable_http_proxy': 0, + 'grpc.initial_reconnect_backoff_ms': 1_000, + 'grpc.max_reconnect_backoff_ms': 30_000, + 'grpc.service_config': JSON.stringify({ + // shuffleAddressList: distribute agents across backends without rewriting the + // target string (keeps channel authority / SNI stable under TLS). + loadBalancingConfig: [{ pick_first: { shuffleAddressList: true } }], + // Retry only unary idempotent reportInstanceProperties. Client-streaming + // collect (trace/meter) must not retry — grpc-js replays the write buffer and + // OAP does not dedupe segments. keepAlive is covered by the next 20s tick. + // FQDN matches Management.proto (skywalking.v3), not ManagementCompat. + methodConfig: [ + { + name: [{ service: 'skywalking.v3.ManagementService', method: 'reportInstanceProperties' }], + retryPolicy: { + maxAttempts: 3, + initialBackoff: '1s', + maxBackoff: '10s', + backoffMultiplier: 2, + retryableStatusCodes: ['UNAVAILABLE'], + }, + }, + ], + }), + }; +} + /** - * Shared gRPC channel manager (Java GRPCChannelManager skeleton). - * V1: single address; V2 reserved: multi-address failover via reportError(). + * Shared gRPC channel manager using grpc-js native multi-address failover + * (pick_first + reconnect backoff + service-config retry). Opens one channel at boot. */ export default class GRPCChannelManager implements BootService { private managedChannel: GRPCChannel | null = null; private readonly listeners: GRPCChannelListener[] = []; private lastStatus: GRPCChannelStatus | null = null; + private lastConnectivityState: grpc.connectivityState | null = null; private closed = false; - - /** V1: first address when comma-separated; V2: failover selection. */ - resolveAddress(): string { - const raw = config.collectorAddress ?? ''; - const first = raw.split(',')[0]?.trim(); - if (!first) { - throw new Error('collectorAddress is not configured'); - } - return first; - } - - getChannel(): grpc.Channel { - if (!this.managedChannel) { - throw new Error('gRPC channel is not available'); - } - - return this.managedChannel.getChannel(); - } + private grpcServers: string[] = []; getClientOptions(): ClientOptions { if (!this.managedChannel) { throw new Error('gRPC channel is not available'); } - return this.managedChannel.getClientOptions(); } - isConnected(): boolean { - return this.managedChannel?.isConnected(true) ?? false; + /** + * Construct a generated gRPC client bound to the shared channel. + * Address/credentials are placeholders — transport uses channelOverride. + */ + createClient( + ClientCtor: new (address: string, credentials: ChannelCredentials, options?: ClientOptions) => TClient, + ): TClient { + return new ClientCtor(STUB_AUTHORITY, grpc.credentials.createInsecure(), this.getClientOptions()); } addChannelListener(listener: GRPCChannelListener): void { @@ -95,26 +134,40 @@ export default class GRPCChannelManager implements BootService { return Number.MAX_SAFE_INTEGER; } - /** Align local status with grpc-js connectivity; avoid permanent DISCONNECT while channel stays READY. */ + /** + * Auth failures: throttled error only (same token across cluster backends). + * Network errors on READY: leave to grpc-js retry / pick_first. + * DEADLINE_EXCEEDED on a READY channel is intentionally not treated as failover — + * failover targets unreachable backends, not slow/overloaded RPCs on a live connection. + */ reportError(error: unknown): void { + if (this.closed) { + return; + } + + if (isGrpcAuthError(error)) { + logAuthRejected('gRPC authentication rejected by OAP; check SW_AGENT_AUTHENTICATION configuration', error); + return; + } + if (!isGrpcNetworkError(error)) { - logger.debug('gRPC report error (ignored): %s', error); + logger.debug(`gRPC report error (ignored): ${error}`); return; } const managed = this.managedChannel; - if (!managed || this.closed) { + if (!managed) { this.notify(GRPCChannelStatus.DISCONNECT); return; } if (managed.isConnected(false)) { - logger.debug('gRPC network error but channel still connected: %s', error); - this.notify(GRPCChannelStatus.CONNECTED); + // Includes DEADLINE_EXCEEDED while READY — do not tear down / rotate (see method doc). + logger.debug(`gRPC network error but channel still READY (native reconnect/retry): ${error}`); return; } - logger.debug('gRPC network error, notify DISCONNECT: %s', error); + logger.debug(`gRPC network error with non-READY channel: ${error}`); this.notify(GRPCChannelStatus.DISCONNECT); } @@ -122,23 +175,18 @@ export default class GRPCChannelManager implements BootService { boot(): void { this.closed = false; - const address = this.resolveAddress(); - const [host, portText] = address.split(':'); - const port = Number.parseInt(portText, 10); - - if (!host || Number.isNaN(port)) { - throw new Error(`Invalid collector address: ${address}`); + this.lastConnectivityState = null; + const parsed = parseStaticBackendAddresses(config.collectorAddress ?? ''); + if (parsed.length === 0) { + logger.error('Collector server addresses are not set.'); + logger.error('Agent will not uplink any data.'); + this.notify(GRPCChannelStatus.DISCONNECT); + return; } - - this.managedChannel = GRPCChannel.newBuilder(host, port) - .addManagedChannelBuilder(new StandardChannelBuilder()) - .addManagedChannelBuilder(new TLSChannelBuilder()) - .addChannelDecorator(new AgentIDDecorator()) - .addChannelDecorator(new AuthenticationDecorator()) - .build(); - - this.watchConnectivityState(); - this.notifyCurrentConnectivityState(true); + // Keep config order in the target so channel authority / SNI stay stable. + // Endpoint pick order is shuffled by pick_first.shuffleAddressList in service_config. + this.grpcServers = [...parsed]; + this.openChannel(this.grpcServers); } onComplete(): void {} @@ -150,6 +198,47 @@ export default class GRPCChannelManager implements BootService { managed?.shutdownNow(); this.notify(GRPCChannelStatus.DISCONNECT); this.listeners.length = 0; + this.grpcServers = []; + this.lastConnectivityState = null; + } + + private openChannel(addresses: string[]): void { + if (this.closed || addresses.length === 0) { + return; + } + + let target: string; + try { + target = buildNativeGrpcTarget(addresses); + } catch (error) { + logger.error(`Failed to build gRPC target: ${error}`); + this.notify(GRPCChannelStatus.DISCONNECT); + return; + } + + let built: GRPCChannel; + try { + built = GRPCChannel.newBuilder(target) + .withChannelOptions(nativeChannelOptions()) + .addManagedChannelBuilder(new StandardChannelBuilder()) + .addManagedChannelBuilder(new TLSChannelBuilder()) + .addChannelDecorator(new AgentIDDecorator()) + .addChannelDecorator(new AuthenticationDecorator()) + .build(); + } catch (error) { + logger.error(`Failed to build gRPC channel for target [${target}]: ${error}`); + this.notify(GRPCChannelStatus.DISCONNECT); + return; + } + + const previous = this.managedChannel; + this.managedChannel = built; + // Defensive: boot() opens once today (ServiceManager.booted), but close any prior + // channel if openChannel is ever invoked again. + previous?.shutdownNow(); + this.watchConnectivityState(); + // watchConnectivityState already requested a connection; do not request again. + this.notifyCurrentConnectivityState(false); } private watchConnectivityState(): void { @@ -157,20 +246,16 @@ export default class GRPCChannelManager implements BootService { if (this.closed || !managed) { return; } - const channel = managed.getChannel(); const currentState = channel.getConnectivityState(true); - channel.watchConnectivityState(currentState, Infinity, (error) => { if (this.closed || this.managedChannel !== managed) { return; } - if (error) { - logger.debug('Channel connectivity watch stopped: %s', error.message); + logger.debug(`Channel connectivity watch stopped: ${error.message}`); return; } - this.notifyCurrentConnectivityState(false); this.watchConnectivityState(); }); @@ -181,23 +266,63 @@ export default class GRPCChannelManager implements BootService { if (this.closed || !managed) { return; } + const state = managed.getConnectivityState(requestConnection); + const previousConnectivity = this.lastConnectivityState; + this.lastConnectivityState = state; - const channel = managed.getChannel(); - const ready = channel.getConnectivityState(requestConnection) === grpc.connectivityState.READY; - this.notify(ready ? GRPCChannelStatus.CONNECTED : GRPCChannelStatus.DISCONNECT); + if (state === grpc.connectivityState.READY) { + this.notify(GRPCChannelStatus.CONNECTED); + return; + } + // Handshake in progress — do not treat as disconnect. + if (state === grpc.connectivityState.CONNECTING) { + return; + } + // READY→IDLE is grpc-js's normal path after the active connection drops. + if (state === grpc.connectivityState.IDLE) { + if (previousConnectivity === grpc.connectivityState.READY) { + this.notify(GRPCChannelStatus.DISCONNECT); + } + return; + } + this.notify(GRPCChannelStatus.DISCONNECT); } private notify(status: GRPCChannelStatus): void { if (this.lastStatus === status) { return; } + const previous = this.lastStatus; this.lastStatus = status; + this.logStatusTransition(status, previous); + for (const listener of this.listeners) { try { listener.statusChanged(status); } catch (err) { - logger.error('GRPCChannelListener failed: %s', err); + logger.error(`GRPCChannelListener failed: ${err}`); + } + } + } + + /** Error-level connectivity logs; skipped on deliberate shutdown (closed). Throttled. */ + private logStatusTransition(status: GRPCChannelStatus, previous: GRPCChannelStatus | null): void { + if (this.closed) { + return; + } + const backends = this.grpcServers.join(',') || config.collectorAddress || ''; + if (status === GRPCChannelStatus.DISCONNECT) { + if (previous === GRPCChannelStatus.CONNECTED) { + logChannelDisconnected( + `gRPC channel disconnected from backends [${backends}]; reconnecting with exponential backoff`, + ); + } else { + logChannelDisconnected( + `gRPC channel not connected to backends [${backends}]; connecting with exponential backoff`, + ); } + } else if (status === GRPCChannelStatus.CONNECTED && previous === GRPCChannelStatus.DISCONNECT) { + logChannelRecovered(`gRPC channel recovered; connected to backends [${backends}]`); } } } diff --git a/src/agent/core/remote/ServiceManagementClient.ts b/src/agent/core/remote/ServiceManagementClient.ts index a558caf1..53c0e19f 100644 --- a/src/agent/core/remote/ServiceManagementClient.ts +++ b/src/agent/core/remote/ServiceManagementClient.ts @@ -32,7 +32,8 @@ import { GRPCChannelListener } from './GRPCChannelListener'; import { GRPCChannelStatus } from './GRPCChannelStatus'; const logger = createLogger(__filename); -const logHeartbeatError = throttled(logger, 'error', 30000); +const logKeepAliveError = throttled(logger, 'error', 30000); +const logPropertiesError = throttled(logger, 'error', 30000); export default class ServiceManagementClient implements BootService, GRPCChannelListener { private closed = false; @@ -115,39 +116,46 @@ export default class ServiceManagementClient implements BootService, GRPCChannel const reportProperties = Math.abs(this.sendPropertiesCounter++) % ServiceManagementClient.PROPERTIES_REPORT_PERIOD_FACTOR === 0; + // Separate try/catch per RPC. Properties ticks are mutually exclusive with keepAlive + // (same if/else as Java PROPERTIES_REPORT_PERIOD_FACTOR / master). if (reportProperties) { - this.managementServiceClient.reportInstanceProperties( - this.instanceProperties, - new grpc.Metadata(), - options, - (error) => { - if (error) { - logHeartbeatError('Failed to send instance properties', error); - this.reportGrpcError(error); - } - }, - ); + try { + this.managementServiceClient.reportInstanceProperties( + this.instanceProperties, + new grpc.Metadata(), + options, + (error) => { + if (error) { + logPropertiesError('Failed to send instance properties', error); + this.reportGrpcError(error); + } + }, + ); + } catch (error) { + logPropertiesError('Failed to send instance properties', error); + this.reportGrpcError(error); + } return; } - this.managementServiceClient.keepAlive(this.keepAlivePkg, new grpc.Metadata(), options, (error) => { - if (error) { - logHeartbeatError('Failed to send heartbeat', error); - this.reportGrpcError(error); - } - }); + try { + this.managementServiceClient.keepAlive(this.keepAlivePkg, new grpc.Metadata(), options, (error) => { + if (error) { + logKeepAliveError('Failed to send heartbeat', error); + this.reportGrpcError(error); + } + }); + } catch (error) { + logKeepAliveError('Failed to send heartbeat', error); + this.reportGrpcError(error); + } } private createManagementClient(): ManagementServiceClient | undefined { if (!this.channelManager) { return undefined; } - - return new ManagementServiceClient( - config.collectorAddress, - grpc.credentials.createInsecure(), - this.channelManager.getClientOptions(), - ); + return this.channelManager.createClient(ManagementServiceClient); } private reportGrpcError(error: unknown): void { if (this.closed) { @@ -156,9 +164,4 @@ export default class ServiceManagementClient implements BootService, GRPCChannel this.channelManager?.reportError(error); } - - flush(): Promise | null { - logger.warn('ServiceManagementClient does not need flush().'); - return null; - } } diff --git a/src/agent/core/remote/TraceSegmentServiceClient.ts b/src/agent/core/remote/TraceSegmentServiceClient.ts index 3f1812d7..c85e491f 100644 --- a/src/agent/core/remote/TraceSegmentServiceClient.ts +++ b/src/agent/core/remote/TraceSegmentServiceClient.ts @@ -28,10 +28,11 @@ import Segment from '../../../trace/context/Segment'; import GRPCChannelManager from './GRPCChannelManager'; import { GRPCChannelListener } from './GRPCChannelListener'; import { GRPCChannelStatus } from './GRPCChannelStatus'; +import { coalesceReport, flushCoalesced, FLUSH_WAIT_MS, ReportCoalesceState, runCollectStream } from './coalesceReport'; const logger = createLogger(__filename); -const logReportError = throttled(logger, 'error', 30000); const logBufferFull = throttled(logger, 'warn', 30000); +const logDiscardedBatch = throttled(logger, 'warn', 30000); export default class TraceSegmentServiceClient implements BootService, GRPCChannelListener { private closed = false; @@ -40,7 +41,9 @@ export default class TraceSegmentServiceClient implements BootService, GRPCChann private reporterClient?: TraceSegmentReportServiceClient; private readonly buffer: Segment[] = []; private timeout?: NodeJS.Timeout; - private reporting?: Promise; + private readonly reportState: ReportCoalesceState = {}; + /** Monotonic count of segments discarded after report failure (for throttled logs). */ + private discardedSegmentTotal = 0; private segmentFinishedListener?: (segment: Segment) => void; prepare(): void { @@ -58,7 +61,7 @@ export default class TraceSegmentServiceClient implements BootService, GRPCChann if (this.buffer.length >= config.maxBufferSize) { logBufferFull( - `Trace buffer reached maximum size (${config.maxBufferSize}); discarding oldest segments. The collector at ${config.collectorAddress} is likely unreachable.`, + `Trace buffer reached maximum size (${config.maxBufferSize}); discarding oldest segments. Configured backends [${config.collectorAddress}] are likely unreachable.`, ); this.buffer.shift(); } @@ -87,7 +90,7 @@ export default class TraceSegmentServiceClient implements BootService, GRPCChann this.timeout = undefined; } - this.reporting = undefined; + this.reportState.reporting = undefined; this.reporterClient = undefined; this.buffer.length = 0; this.channelManager = undefined; @@ -107,12 +110,7 @@ export default class TraceSegmentServiceClient implements BootService, GRPCChann if (!this.channelManager) { return undefined; } - - return new TraceSegmentReportServiceClient( - config.collectorAddress, - grpc.credentials.createInsecure(), - this.channelManager.getClientOptions(), - ); + return this.channelManager.createClient(TraceSegmentReportServiceClient); } private scheduleNextReport(): void { @@ -131,19 +129,11 @@ export default class TraceSegmentServiceClient implements BootService, GRPCChann } private reportOnce(): Promise { - if (this.closed) { - return Promise.resolve(); - } - - if (this.reporting) { - return this.reporting; - } - - this.reporting = this.doReport().finally(() => { - this.reporting = undefined; - }); - - return this.reporting; + return coalesceReport( + this.reportState, + () => this.doReport(), + () => this.closed, + ); } private doReport(): Promise { @@ -153,7 +143,12 @@ export default class TraceSegmentServiceClient implements BootService, GRPCChann return; } - emitter.emit('segments-sent'); + try { + emitter.emit('segments-sent'); + } catch (error) { + // Listener errors must not reject the report promise (host unhandledRejection). + logger.debug(`segments-sent listener failed: ${error}`); + } if (this.buffer.length === 0) { resolve(); @@ -165,41 +160,34 @@ export default class TraceSegmentServiceClient implements BootService, GRPCChann return; } - let stream: ReturnType | undefined; - try { - stream = this.reporterClient.collect( - new grpc.Metadata(), - { deadline: Date.now() + config.traceTimeout }, - (error) => { - if (error) { - logReportError('Failed to report trace data', error); - this.reportGrpcError(error); - } - resolve(); - }, - ); - - for (const segment of this.buffer) { - if (segment) { - if (logger._isDebugEnabled) { - logger.debug('Sending segment ', { segment }); + // Take ownership. On failure discard once (never re-send): disconnect-window + // data is already protected by READY→IDLE → DISCONNECT (status !== CONNECTED skips splice). + const batch = this.buffer.splice(0, this.buffer.length); + const client = this.reporterClient; + void runCollectStream({ + open: (onStatus) => + client.collect(new grpc.Metadata(), { deadline: Date.now() + config.traceTimeout }, onStatus), + writeAll: (stream) => { + for (const segment of batch) { + if (segment) { + if (logger._isDebugEnabled) { + logger.debug('Sending segment ', { segment }); + } + stream.write(segment.transform()); } - stream.write(segment.transform()); } - } - } catch (error) { - logReportError('Failed to report trace data', error); - this.reportGrpcError(error); - resolve(); - } finally { - this.buffer.length = 0; - try { - stream?.end(); - } catch (error) { - logReportError('Failed to end trace collect stream', error); - resolve(); - } - } + }, + onFailure: (reason, error) => { + this.discardedSegmentTotal += batch.length; + logDiscardedBatch( + `Discarded ${batch.length} trace segment(s) after report failure (${reason}) (total discarded: ${this.discardedSegmentTotal})`, + error, + ); + this.reportGrpcError(error); + }, + openFailureReason: 'Failed to report trace data', + endFailureReason: 'Failed to end trace collect stream', + }).then(resolve); }); } @@ -211,6 +199,9 @@ export default class TraceSegmentServiceClient implements BootService, GRPCChann this.channelManager?.reportError(error); } + /** + * Best-effort: one shared FLUSH_WAIT_MS budget for forceReport of remaining buffer, then in-flight wait. + */ flush(): Promise | null { if (this.closed) { return null; @@ -221,11 +212,12 @@ export default class TraceSegmentServiceClient implements BootService, GRPCChann this.timeout = undefined; } - if (this.buffer.length === 0) { - this.scheduleNextReport(); - return null; - } - - return this.reportOnce().finally(() => this.scheduleNextReport()); + return flushCoalesced( + this.reportState, + () => this.doReport(), + () => this.closed, + () => this.buffer.length > 0, + FLUSH_WAIT_MS, + ).finally(() => this.scheduleNextReport()); } } diff --git a/src/agent/core/remote/coalesceReport.ts b/src/agent/core/remote/coalesceReport.ts new file mode 100644 index 00000000..f64e67c5 --- /dev/null +++ b/src/agent/core/remote/coalesceReport.ts @@ -0,0 +1,188 @@ +/*! + * + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + * + */ + +/** Single-flight report guard: if a report is in flight, return it; otherwise start one. */ +export type ReportCoalesceState = { + reporting?: Promise; +}; + +/** + * Total budget for one flushCoalesced call (force pending + wait prior in-flight). + * agent.flush() uses the same constant per phase (SpanContext, then services) so the + * outer wait is never shorter than what the inner path can consume. + */ +export const FLUSH_WAIT_MS = 2000; + +/** Minimal client-streaming surface used by Trace / Meter collect. */ +export type CollectStream = { + write(message: unknown): void; + end(): void; + cancel(): void; +}; + +/** Resolve when `promise` settles or `ms` elapses — never rejects. */ +export function awaitWithTimeout(promise: Promise, ms: number): Promise { + return new Promise((resolve) => { + let settled = false; + const done = (): void => { + if (settled) { + return; + } + settled = true; + clearTimeout(timer); + resolve(); + }; + const timer = setTimeout(done, ms); + timer.unref?.(); + promise.then(done, done); + }); +} + +function beginFlight( + state: ReportCoalesceState, + doReport: () => Promise, + isClosed: () => boolean, +): Promise { + // Assign before invoking doReport so sync re-entry (e.g. segments-sent → flush) sees the guard. + let flight!: Promise; + state.reporting = flight = Promise.resolve() + .then(() => { + if (isClosed()) { + return; + } + return doReport(); + }) + .catch(() => undefined) + .finally(() => { + if (state.reporting === flight) { + state.reporting = undefined; + } + }); + return flight; +} + +/** + * Never rejects — callers (timers with `void`, flush) must not see unhandled rejections. + */ +export function coalesceReport( + state: ReportCoalesceState, + doReport: () => Promise, + isClosed: () => boolean, +): Promise { + if (isClosed()) { + return Promise.resolve(); + } + + if (state.reporting) { + return state.reporting.catch(() => undefined); + } + + return beginFlight(state, doReport, isClosed); +} + +/** + * Always start a new report flight (ignores an existing in-flight promise). + * Used by flush() after a bounded wait times out so remaining buffer/snapshot is attempted. + */ +export function forceReport( + state: ReportCoalesceState, + doReport: () => Promise, + isClosed: () => boolean, +): Promise { + if (isClosed()) { + return Promise.resolve(); + } + return beginFlight(state, doReport, isClosed); +} + +/** + * Shared Trace/Meter flush: one total `budgetMs` for forceReport of pending data, then + * waiting on any prior in-flight work (deadline shared — never 2× the constant). + * Pending has priority: an in-flight batch already had one attempt. + */ +export async function flushCoalesced( + state: ReportCoalesceState, + doReport: () => Promise, + isClosed: () => boolean, + hasPending: () => boolean, + budgetMs: number = FLUSH_WAIT_MS, +): Promise { + const deadline = Date.now() + budgetMs; + const remaining = (): number => Math.max(0, deadline - Date.now()); + const inflight = state.reporting; + + if (!isClosed() && hasPending()) { + await awaitWithTimeout(forceReport(state, doReport, isClosed), remaining()); + } + if (inflight) { + await awaitWithTimeout(inflight, remaining()); + } +} + +/** + * Shared client-streaming collect: open → write → end, with once-only failure + cancel. + */ +export function runCollectStream(options: { + open: (onStatus: (error: Error | null) => void) => CollectStream; + writeAll: (stream: CollectStream) => void; + onFailure: (reason: string, error: unknown) => void; + openFailureReason: string; + endFailureReason: string; +}): Promise { + return new Promise((resolve) => { + let failed = false; + const fail = (reason: string, error: unknown): void => { + if (failed) { + return; + } + failed = true; + options.onFailure(reason, error); + }; + + let stream: CollectStream | undefined; + try { + stream = options.open((error) => { + if (error) { + fail(options.openFailureReason, error); + } + resolve(); + }); + options.writeAll(stream); + try { + stream.end(); + } catch (error) { + fail(options.endFailureReason, error); + try { + stream.cancel(); + } catch { + /* ignore */ + } + resolve(); + } + } catch (error) { + fail(options.openFailureReason, error); + try { + stream?.cancel(); + } catch { + /* ignore */ + } + resolve(); + } + }); +} diff --git a/src/index.ts b/src/index.ts index 33d67a21..047569b0 100644 --- a/src/index.ts +++ b/src/index.ts @@ -22,6 +22,7 @@ import ServiceManager from './agent/core/boot/ServiceManager'; import { createLogger } from './logging'; import PluginInstaller from './core/PluginInstaller'; import SpanContext from './trace/context/SpanContext'; +import { awaitWithTimeout, FLUSH_WAIT_MS } from './agent/core/remote/coalesceReport'; const logger = createLogger(__filename); @@ -56,18 +57,20 @@ class Agent { return null; } - const spanContextFlush = SpanContext.flush(); - if (!spanContextFlush) { - return ServiceManager.INSTANCE.flush(); - } - - return new Promise((resolve) => { - spanContextFlush.then(() => { - const serviceFlush = ServiceManager.INSTANCE.flush(); - if (!serviceFlush) resolve(null); - else serviceFlush.then(() => resolve(null)); - }); - }); + // Bound waits so a stuck in-flight RPC (or SpanContext waiting on the flush + // request's own open span) cannot stall hosts (e.g. Lambda) for the gRPC deadline. + // Each phase uses FLUSH_WAIT_MS; Trace/Meter flushCoalesced shares one budget internally + // so the service wait is never shorter than force-pending + wait-in-flight. + return (async () => { + const spanContextFlush = SpanContext.flush(FLUSH_WAIT_MS); + if (spanContextFlush) { + await spanContextFlush; + } + const serviceFlush = ServiceManager.INSTANCE.flush(); + if (serviceFlush) { + await awaitWithTimeout(serviceFlush, FLUSH_WAIT_MS); + } + })(); } destroy(): void { diff --git a/src/logging/index.ts b/src/logging/index.ts index 32db1aac..e8c84dba 100644 --- a/src/logging/index.ts +++ b/src/logging/index.ts @@ -22,14 +22,15 @@ import { Logger } from 'winston'; type LoggerLevelAware = Logger & { _isDebugEnabled: boolean; - _isInfoEnabled: boolean; }; export function createLogger(name: string): LoggerLevelAware { - const loggingLevel = (process.env.SW_AGENT_LOGGING_LEVEL || 'error').toLowerCase(); + const loggingLevel = (process.env.SW_AGENT_LOGGING_LEVEL || 'warn').toLowerCase(); const logger = winston.createLogger({ level: loggingLevel, + // No format.splat(): user data (e.g. span operation URLs) must not be treated as + // printf format strings, and Error args must not expand full stacks into the buffer. format: winston.format.json(), defaultMeta: { file: name, @@ -52,11 +53,9 @@ export function createLogger(name: string): LoggerLevelAware { const loggerLevel = logger.levels[logger.level]; const _isDebugEnabled = loggerLevel >= logger.levels.debug; - const _isInfoEnabled = loggerLevel >= logger.levels.info; Object.assign(logger, { _isDebugEnabled, - _isInfoEnabled, }); const nop = (): void => { diff --git a/src/trace/context/SpanContext.ts b/src/trace/context/SpanContext.ts index bd199264..fba2e870 100644 --- a/src/trace/context/SpanContext.ts +++ b/src/trace/context/SpanContext.ts @@ -233,16 +233,43 @@ export default class SpanContext implements Context { ContextManager.restore(span); } - static flush(): Promise | null { - // This function explicitly returns null instead of a resolved Promise in case of nothing to flush so that in this - // case passing control back to the event loop can be avoided. Even a resolved Promise will run other things in - // the event loop when it is awaited and before it continues. - - return !SpanContext.nTotalSegments - ? null - : new Promise((resolve: (value: unknown) => void) => { - SpanContext.flushResolve.push(resolve); - }); + /** + * Wait until unfinished segments drain (`nTotalSegments === 0`). + * Returns null when there is nothing to wait for (avoids yielding to the event loop). + * When `timeoutMs` is set, the waiter removes itself on timeout so long-lived processes + * that call `agent.flush()` while traffic keeps `nTotalSegments > 0` do not leak resolvers. + */ + static flush(timeoutMs?: number): Promise | null { + if (!SpanContext.nTotalSegments) { + return null; + } + + return new Promise((resolve: (value: unknown) => void) => { + let settled = false; + let timer: NodeJS.Timeout | undefined; + + const entry = (value: unknown): void => { + if (settled) { + return; + } + settled = true; + if (timer) { + clearTimeout(timer); + } + const idx = SpanContext.flushResolve.indexOf(entry); + if (idx >= 0) { + SpanContext.flushResolve.splice(idx, 1); + } + resolve(value); + }; + + SpanContext.flushResolve.push(entry); + + if (timeoutMs != null && timeoutMs >= 0) { + timer = setTimeout(() => entry(null), timeoutMs); + timer.unref?.(); + } + }); } traceId(): string { diff --git a/tests/plugins/common/Dockerfile.agent b/tests/plugins/common/Dockerfile.agent index 5e907d4f..7db6a1c8 100644 --- a/tests/plugins/common/Dockerfile.agent +++ b/tests/plugins/common/Dockerfile.agent @@ -14,13 +14,18 @@ # limitations under the License. ARG SW_NODE_VERSION - FROM node:${SW_NODE_VERSION} -ARG ROOT=. - WORKDIR /app - -ADD $ROOT /app - -RUN npm install request && npm install +# Install deps before copying sources so package-lock / Dockerfile edits alone +# bust the npm ci layer; source edits reuse the cached install. +# Plugin e2e needs the full tree in the image (no src mount). +# Remote-e2e uses the image-built src/ (including prepared src/proto); only test +# scripts are bind-mounted at runtime. +# +# --ignore-scripts: `prepare` needs scripts/ + protos (copied next). Then rebuild +# packages whose install scripts are required for natives/protoc in this image. +COPY package.json package-lock.json ./ +RUN npm ci --ignore-scripts +COPY . . +RUN npm rebuild grpc-tools protobufjs && npm run prepare diff --git a/tests/plugins/common/base-compose.yml b/tests/plugins/common/base-compose.yml index 504dfa91..10eb7d5b 100644 --- a/tests/plugins/common/base-compose.yml +++ b/tests/plugins/common/base-compose.yml @@ -30,6 +30,17 @@ services: timeout: 60s retries: 120 + # Same image/healthcheck as collector but no host port — for remote-e2e multi-collector cases. + collector-no-host-port: + image: ghcr.io/apache/skywalking-agent-test-tool/mock-collector:fa81b1b6d9caef484a65b5019efa28cac4e3d21d + networks: + - traveling-light + healthcheck: + test: [ "CMD", "bash", "-c", "cat < /dev/null > /dev/tcp/127.0.0.1/12800" ] + interval: 5s + timeout: 60s + retries: 120 + agent: build: context: ../../../ diff --git a/tests/remote-e2e/common/server.ts b/tests/remote-e2e/common/server.ts new file mode 100644 index 00000000..1fca36d9 --- /dev/null +++ b/tests/remote-e2e/common/server.ts @@ -0,0 +1,61 @@ +/*! + * + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + * + */ + +import agent from '../../../src'; +import * as http from 'http'; + +agent.start({ serviceName: 'server', maxBufferSize: 1000 }); + +const server = http.createServer((req, res) => { + const url = req.url || '/'; + if (url === '/ping' || url.startsWith('/probe/')) { + res.statusCode = 200; + res.end('ok'); + return; + } + if (url === '/flush') { + const deadlineMs = 8000; + let completed = false; + const timer = setTimeout(() => { + if (completed) { + return; + } + completed = true; + // Distinguish hung flush from success (helpers assert status + body). + res.statusCode = 504; + res.end('flush deadline'); + }, deadlineMs); + void Promise.resolve(agent.flush()) + .catch(() => undefined) + .finally(() => { + if (completed) { + return; + } + completed = true; + clearTimeout(timer); + res.statusCode = 200; + res.end('flushed'); + }); + return; + } + res.statusCode = 404; + res.end('not found'); +}); + +server.listen(5000, () => console.info('Listening on port 5000...')); diff --git a/tests/remote-e2e/static-failover/docker-compose.yml b/tests/remote-e2e/static-failover/docker-compose.yml new file mode 100644 index 00000000..aba99e33 --- /dev/null +++ b/tests/remote-e2e/static-failover/docker-compose.yml @@ -0,0 +1,71 @@ +# +# Licensed to the Apache Software Foundation (ASF) under one or more +# contributor license agreements. See the NOTICE file distributed with +# this work for additional information regarding copyright ownership. +# The ASF licenses this file to You under the Apache License, Version 2.0 +# (the "License"); you may not use this file except in compliance with +# the License. You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, +# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +# See the License for the specific language governing permissions and +# limitations under the License. +# + +version: '2.1' + +services: + collector-a: + extends: + file: ../../plugins/common/base-compose.yml + service: collector-no-host-port + ports: + - '12810:12800' + networks: + - traveling-light + + collector-b: + extends: + file: ../../plugins/common/base-compose.yml + service: collector-no-host-port + ports: + - '12811:12800' + networks: + - traveling-light + + server: + build: + context: ../../.. + dockerfile: tests/plugins/common/Dockerfile.agent + args: + - SW_NODE_VERSION=${SW_NODE_VERSION:-22} + ports: + - '5010:5000' + environment: + SW_AGENT_COLLECTOR_BACKEND_SERVICES: collector-a:19876,collector-b:19876 + SW_AGENT_RUNTIME_METRICS_REPORTER_ACTIVE: 'false' + volumes: + # Image already ran `npm run prepare` (generated src/proto). Do not bind-mount host + # src/ over it — a checkout without protocol/ submodule would hide stubs and only + # fail as an opaque server health-check timeout. + - ../common:/app/tests/remote-e2e/common + - .:/app/tests/remote-e2e/static-failover + entrypoint: ['bash', '-c', 'npx ts-node /app/tests/remote-e2e/common/server.ts'] + depends_on: + collector-a: + condition: service_healthy + collector-b: + condition: service_healthy + healthcheck: + test: ['CMD', 'bash', '-c', 'cat < /dev/null > /dev/tcp/127.0.0.1/5000'] + interval: 5s + timeout: 60s + retries: 120 + networks: + - traveling-light + +networks: + traveling-light: diff --git a/tests/remote-e2e/static-failover/test.ts b/tests/remote-e2e/static-failover/test.ts new file mode 100644 index 00000000..5230691d --- /dev/null +++ b/tests/remote-e2e/static-failover/test.ts @@ -0,0 +1,81 @@ +/*! + * + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + * + */ + +/* eslint-env jest */ + +import waitForExpect from 'wait-for-expect'; +import { StartedDockerComposeEnvironment } from 'testcontainers'; +import { createRemoteE2eContext } from '../support/helpers'; + +const e2e = createRemoteE2eContext(__dirname, { + serverPort: 5010, + collectorAHttpPort: 12810, + collectorBHttpPort: 12811, +}); + +describe('remote-e2e static failover', () => { + let compose: StartedDockerComposeEnvironment | undefined; + + beforeAll(async () => { + compose = await e2e.upCompose(); + }); + + afterAll(async () => { + if (compose) { + await compose.down(); + } + }); + + it('fails over to the standby collector after the active one stops', async () => { + if (!compose) { + throw new Error('Docker Compose environment failed to start'); + } + + await waitForExpect(async () => e2e.pingServer()); + await e2e.flushServer(); + + const activeInfo = { activeService: '' as 'collector-a' | 'collector-b', standbyHttpPort: 0 }; + let ambiguousProbe: string | undefined; + await waitForExpect(async () => { + const resolved = await e2e.resolveActiveCollector(); + if (resolved.kind === 'ambiguous') { + ambiguousProbe = resolved.probePath; + return; + } + if (resolved.kind === 'pending') { + throw new Error(`Probe ${resolved.probePath} not received by either collector yet`); + } + activeInfo.activeService = resolved.activeService; + activeInfo.standbyHttpPort = resolved.standbyHttpPort; + }); + expect(ambiguousProbe).toBeUndefined(); + expect(activeInfo.activeService).toMatch(/^collector-[ab]$/); + await e2e.stopComposeService(compose, activeInfo.activeService); + + await waitForExpect( + async () => { + await e2e.pingServer(); + await e2e.flushServer(); + await e2e.assertCollectorReceivedPing(activeInfo.standbyHttpPort); + }, + 120000, + 3000, + ); + }, 300000); +}); diff --git a/tests/remote-e2e/support/helpers.ts b/tests/remote-e2e/support/helpers.ts new file mode 100644 index 00000000..5646fe22 --- /dev/null +++ b/tests/remote-e2e/support/helpers.ts @@ -0,0 +1,146 @@ +/*! + * + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + * + */ + +/* eslint-env jest */ + +import * as path from 'path'; +import { DockerComposeEnvironment, StartedDockerComposeEnvironment, Wait } from 'testcontainers'; +import axios from 'axios'; + +export type RemoteE2ePorts = { + serverPort: number; + collectorAHttpPort: number; + collectorBHttpPort: number; +}; + +export function createRemoteE2eContext(caseDir: string, ports: RemoteE2ePorts) { + const rootDir = path.resolve(caseDir); + const { serverPort, collectorAHttpPort, collectorBHttpPort } = ports; + + async function pingServer(): Promise { + const response = await axios.get(`http://localhost:${serverPort}/ping`); + expect(response.status).toBe(200); + } + + async function flushServer(): Promise { + const response = await axios.get(`http://localhost:${serverPort}/flush`, { + validateStatus: () => true, + }); + expect(response.status).toBe(200); + expect(response.data).toBe('flushed'); + } + + async function collectorReceiveData(port: number): Promise { + return String((await axios.get(`http://localhost:${port}/receiveData`)).data); + } + + async function collectorHasPing(port: number): Promise { + const data = await collectorReceiveData(port); + return data.includes('operationName: GET:/ping'); + } + + async function assertCollectorReceivedPing(port: number): Promise { + const data = await collectorReceiveData(port); + expect(data).toContain('serviceName: server'); + expect(data).toContain('operationName: GET:/ping'); + expect(data).toContain("http.status_code, value: '200'"); + } + + /** + * Identify which mock collector the agent is currently reporting to via a + * unique per-probe path. + * - pending: neither has the probe yet (caller may retry) + * - ambiguous: both received it (fail outside waitForExpect — do not retry) + * - ready: exactly one collector has the probe + */ + async function resolveActiveCollector(): Promise< + | { + kind: 'ready'; + activeService: 'collector-a' | 'collector-b'; + standbyService: 'collector-a' | 'collector-b'; + activeHttpPort: number; + standbyHttpPort: number; + } + | { kind: 'pending'; probePath: string } + | { kind: 'ambiguous'; probePath: string } + > { + const token = `p${Date.now().toString(36)}${Math.random().toString(36).slice(2, 8)}`; + const probePath = `/probe/${token}`; + const needle = `operationName: GET:${probePath}`; + + const probeRes = await axios.get(`http://localhost:${serverPort}${probePath}`); + expect(probeRes.status).toBe(200); + await flushServer(); + + const aData = await collectorReceiveData(collectorAHttpPort); + const bData = await collectorReceiveData(collectorBHttpPort); + const aHas = aData.includes(needle); + const bHas = bData.includes(needle); + + if (!aHas && !bHas) { + return { kind: 'pending', probePath }; + } + if (aHas && bHas) { + return { kind: 'ambiguous', probePath }; + } + if (aHas) { + return { + kind: 'ready', + activeService: 'collector-a', + standbyService: 'collector-b', + activeHttpPort: collectorAHttpPort, + standbyHttpPort: collectorBHttpPort, + }; + } + return { + kind: 'ready', + activeService: 'collector-b', + standbyService: 'collector-a', + activeHttpPort: collectorBHttpPort, + standbyHttpPort: collectorAHttpPort, + }; + } + + async function upCompose(): Promise { + // Server depends_on collectors with service_healthy — waiting on server alone is enough. + // compose `build:` for the agent image — Docker layer cache invalidates on lockfile/Dockerfile change. + return new DockerComposeEnvironment(rootDir, 'docker-compose.yml') + .withWaitStrategy('server-1', Wait.forHealthCheck()) + .up(); + } + + /** Stop one compose service (Compose v2 / testcontainers: `-1`). */ + async function stopComposeService( + compose: StartedDockerComposeEnvironment, + service: 'collector-a' | 'collector-b', + ): Promise { + await compose.getContainer(`${service}-1`).stop(); + } + + return { + rootDir, + pingServer, + flushServer, + collectorHasPing, + assertCollectorReceivedPing, + resolveActiveCollector, + upCompose, + stopComposeService, + }; +} diff --git a/tests/remote/BackendAddressResolver.test.ts b/tests/remote/BackendAddressResolver.test.ts new file mode 100644 index 00000000..2cdca8d1 --- /dev/null +++ b/tests/remote/BackendAddressResolver.test.ts @@ -0,0 +1,63 @@ +/*! + * + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + * + */ + +/* eslint-env jest */ + +import { + buildNativeGrpcTarget, + parseStaticBackendAddresses, + tryParseHostPort, +} from '../../src/agent/core/remote/BackendAddressResolver'; + +describe('BackendAddressResolver (comma-separated static backends)', () => { + it('parses comma-separated host:port entries', () => { + expect(parseStaticBackendAddresses('a:11800, b:11800 ,c:11800')).toEqual(['a:11800', 'b:11800', 'c:11800']); + }); + + it('drops invalid entries', () => { + expect(parseStaticBackendAddresses('good:11800,bad,also-bad:xyz,:9')).toEqual(['good:11800']); + }); + + it('supports bracketed IPv6 literals', () => { + expect(parseStaticBackendAddresses('[::1]:11800, [2001:db8::1]:11800')).toEqual([ + '[::1]:11800', + '[2001:db8::1]:11800', + ]); + expect(tryParseHostPort('[::1]:11800')).toBe('[::1]:11800'); + }); + + it('rejects unbracketed IPv6 host:port', () => { + expect(tryParseHostPort('2001:db8::1:11800')).toBeNull(); + expect(parseStaticBackendAddresses('2001:db8::1:11800,good:11800')).toEqual(['good:11800']); + }); + + it('buildNativeGrpcTarget uses plain host:port for a single address', () => { + expect(buildNativeGrpcTarget(['oap.example.com:11800'])).toBe('oap.example.com:11800'); + expect(buildNativeGrpcTarget(['10.0.0.1:11800'])).toBe('10.0.0.1:11800'); + }); + + it('buildNativeGrpcTarget uses sw-static for multiple addresses', () => { + expect(buildNativeGrpcTarget(['10.0.0.1:11800', '10.0.0.2:11800'])).toBe( + 'sw-static:///10.0.0.1:11800,10.0.0.2:11800', + ); + expect(buildNativeGrpcTarget(['collector-a:19876', 'collector-b:19876'])).toBe( + 'sw-static:///collector-a:19876,collector-b:19876', + ); + }); +}); diff --git a/tests/remote/GRPCChannelManager.test.ts b/tests/remote/GRPCChannelManager.test.ts index ad6f16e3..cf17f37a 100644 --- a/tests/remote/GRPCChannelManager.test.ts +++ b/tests/remote/GRPCChannelManager.test.ts @@ -22,64 +22,173 @@ import * as grpc from '@grpc/grpc-js'; import GRPCChannelManager from '../../src/agent/core/remote/GRPCChannelManager'; import { GRPCChannelStatus } from '../../src/agent/core/remote/GRPCChannelStatus'; +import config from '../../src/config/AgentConfig'; const mockShutdownNow = jest.fn(); const mockGetConnectivityState = jest.fn(); const mockWatchConnectivityState = jest.fn(); +const mockIsConnected = jest.fn(() => true); +const mockNewBuilder = jest.fn(); +const mockWithChannelOptions = jest.fn().mockReturnThis(); jest.mock('../../src/agent/core/remote/GRPCChannel', () => ({ __esModule: true, default: { - newBuilder: jest.fn(() => ({ - addManagedChannelBuilder: jest.fn().mockReturnThis(), - addChannelDecorator: jest.fn().mockReturnThis(), - build: jest.fn(() => ({ - getChannel: () => ({ - getConnectivityState: mockGetConnectivityState, - watchConnectivityState: mockWatchConnectivityState, - }), - getClientOptions: () => ({ channelOverride: {} }), - isConnected: jest.fn(() => true), - shutdownNow: mockShutdownNow, - })), - })), + newBuilder: (...args: unknown[]) => mockNewBuilder(...args), }, })); -jest.mock('../../src/config/AgentConfig', () => ({ - __esModule: true, - default: { - collectorAddress: '127.0.0.1:11800', - }, -})); +function installChannelMock(): void { + mockNewBuilder.mockImplementation(() => ({ + withChannelOptions: mockWithChannelOptions, + addManagedChannelBuilder: jest.fn().mockReturnThis(), + addChannelDecorator: jest.fn().mockReturnThis(), + build: jest.fn(() => ({ + getChannel: () => ({ + getConnectivityState: mockGetConnectivityState, + watchConnectivityState: mockWatchConnectivityState, + }), + getClientOptions: () => ({ channelOverride: {} }), + isConnected: mockIsConnected, + getConnectivityState: mockGetConnectivityState, + shutdownNow: mockShutdownNow, + })), + })); +} + +describe('GRPCChannelManager (native grpc-js multi-backend failover)', () => { + const originalCollector = config.collectorAddress; + const originalSecure = config.secure; -describe('GRPCChannelManager initial connectivity', () => { beforeEach(() => { jest.clearAllMocks(); + installChannelMock(); mockWatchConnectivityState.mockImplementation(() => undefined); + mockGetConnectivityState.mockReturnValue(grpc.connectivityState.READY); + mockIsConnected.mockReturnValue(true); + config.collectorAddress = '127.0.0.1:11800'; + config.secure = false; + jest.spyOn(Math, 'random').mockReturnValue(0); }); - it('notifies CONNECTED when channel is already READY at boot', () => { - mockGetConnectivityState.mockReturnValue(grpc.connectivityState.READY); + afterEach(() => { + config.collectorAddress = originalCollector; + config.secure = originalSecure; + jest.restoreAllMocks(); + }); + it('notifies CONNECTED when channel is READY after boot', () => { const listener = { statusChanged: jest.fn() }; const manager = new GRPCChannelManager(); + manager.addChannelListener(listener); + manager.boot(); + expect(mockNewBuilder).toHaveBeenCalled(); + expect(mockWatchConnectivityState).toHaveBeenCalled(); + expect(listener.statusChanged).toHaveBeenCalledWith(GRPCChannelStatus.CONNECTED); + expect(mockWithChannelOptions).toHaveBeenCalled(); + const options = mockWithChannelOptions.mock.calls[0][0]; + expect(options['grpc.enable_http_proxy']).toBe(0); + expect(options['grpc.keepalive_time_ms']).toBeUndefined(); + manager.shutdown(); + }); + + it('uses plain host:port target for a single backend', () => { + config.collectorAddress = 'oap.example.com:11800'; + const manager = new GRPCChannelManager(); + manager.boot(); + expect(mockNewBuilder.mock.calls[0][0]).toBe('oap.example.com:11800'); + manager.shutdown(); + }); + + it('builds sw-static multi-address target for multiple backends', () => { + config.collectorAddress = '10.0.0.1:11800,10.0.0.2:11800'; + const manager = new GRPCChannelManager(); + manager.boot(); + const target = mockNewBuilder.mock.calls[0][0] as string; + expect(target.startsWith('sw-static:///')).toBe(true); + expect(target).toContain('10.0.0.1:11800'); + expect(target).toContain('10.0.0.2:11800'); + manager.shutdown(); + }); + it('preserves config address order in the channel target (LB shuffles endpoints)', () => { + config.collectorAddress = 'a:11800,b:11800'; + jest.spyOn(Math, 'random').mockReturnValue(0.99); + const manager = new GRPCChannelManager(); + manager.boot(); + expect(mockNewBuilder.mock.calls[0][0]).toBe('sw-static:///a:11800,b:11800'); + manager.shutdown(); + }); + + it('preserves address order under TLS for stable authority', () => { + config.secure = true; + config.collectorAddress = 'a:11800,b:11800'; + jest.spyOn(Math, 'random').mockReturnValue(0.99); + const manager = new GRPCChannelManager(); + manager.boot(); + expect(mockNewBuilder.mock.calls[0][0]).toBe('sw-static:///a:11800,b:11800'); + manager.shutdown(); + }); + + it('notifies DISCONNECT when collector addresses are empty', () => { + config.collectorAddress = ''; + const listener = { statusChanged: jest.fn() }; + const manager = new GRPCChannelManager(); manager.addChannelListener(listener); manager.boot(); + expect(mockNewBuilder).not.toHaveBeenCalled(); + expect(listener.statusChanged).toHaveBeenCalledWith(GRPCChannelStatus.DISCONNECT); + manager.shutdown(); + }); - expect(listener.statusChanged).toHaveBeenCalledWith(GRPCChannelStatus.CONNECTED); + it('does not rebuild on network error while READY', () => { + const manager = new GRPCChannelManager(); + manager.boot(); + mockNewBuilder.mockClear(); + mockShutdownNow.mockClear(); + manager.reportError({ code: grpc.status.UNAVAILABLE, message: 'transient' }); + expect(mockNewBuilder).not.toHaveBeenCalled(); + expect(mockShutdownNow).not.toHaveBeenCalled(); + manager.shutdown(); }); - it('notifies DISCONNECT when channel is not READY at boot', () => { - mockGetConnectivityState.mockReturnValue(grpc.connectivityState.CONNECTING); + it('does not rebuild on UNAUTHENTICATED (auth is not fixed by rotating backends)', () => { + config.collectorAddress = 'a:11800,b:11800'; + const manager = new GRPCChannelManager(); + manager.boot(); + const targetBefore = mockNewBuilder.mock.calls[0][0]; + mockNewBuilder.mockClear(); + manager.reportError({ code: grpc.status.UNAUTHENTICATED, message: 'bad token' }); + expect(mockNewBuilder).not.toHaveBeenCalled(); + expect(targetBefore).toContain('sw-static:///'); + manager.shutdown(); + }); + it('does not treat CONNECTING as DISCONNECT', () => { + mockGetConnectivityState.mockReturnValue(grpc.connectivityState.CONNECTING); const listener = { statusChanged: jest.fn() }; const manager = new GRPCChannelManager(); + manager.addChannelListener(listener); + manager.boot(); + expect(mockNewBuilder).toHaveBeenCalled(); + expect(mockWatchConnectivityState).toHaveBeenCalled(); + expect(listener.statusChanged).not.toHaveBeenCalledWith(GRPCChannelStatus.DISCONNECT); + expect(listener.statusChanged).not.toHaveBeenCalledWith(GRPCChannelStatus.CONNECTED); + manager.shutdown(); + }); + it('treats READY then IDLE as DISCONNECT', () => { + mockGetConnectivityState.mockReturnValue(grpc.connectivityState.READY); + const listener = { statusChanged: jest.fn() }; + const manager = new GRPCChannelManager(); manager.addChannelListener(listener); manager.boot(); + expect(listener.statusChanged).toHaveBeenCalledWith(GRPCChannelStatus.CONNECTED); + mockGetConnectivityState.mockReturnValue(grpc.connectivityState.IDLE); + const watchCb = mockWatchConnectivityState.mock.calls[0][2] as (err?: Error) => void; + watchCb(); expect(listener.statusChanged).toHaveBeenCalledWith(GRPCChannelStatus.DISCONNECT); + manager.shutdown(); }); }); diff --git a/tests/remote/MeterSender.test.ts b/tests/remote/MeterSender.test.ts index c1a3018d..8fde7aa8 100644 --- a/tests/remote/MeterSender.test.ts +++ b/tests/remote/MeterSender.test.ts @@ -26,7 +26,7 @@ let sampleSequence = 0; const mockChannelManager = { addChannelListener: jest.fn(), - getClientOptions: jest.fn(() => ({})), + createClient: jest.fn((ClientCtor: new (...args: unknown[]) => unknown) => new ClientCtor()), reportError: jest.fn(), }; @@ -39,6 +39,7 @@ const createMeterData = () => ({ const mockStream = { write: jest.fn(), end: jest.fn(), + cancel: jest.fn(), }; const mockSnapshot = () => ({ collectedAt: 1_000_000 + sampleSequence++ * 500, cpu: 1 }); @@ -101,6 +102,7 @@ describe('MeterSender', () => { mockChannelManager.reportError.mockClear(); mockStream.write.mockReset(); mockStream.end.mockReset(); + mockStream.cancel.mockReset(); pendingCollectCallback = undefined; sender = new MeterSender(); sender.prepare(); @@ -130,6 +132,7 @@ describe('MeterSender', () => { const reportPromise = (sender as unknown as { reportBufferedMetrics: () => Promise }).reportBufferedMetrics(); await Promise.resolve(); + await Promise.resolve(); pendingCollectCallback?.(null); await reportPromise; @@ -149,6 +152,7 @@ describe('MeterSender', () => { const reportPromise = (sender as unknown as { reportBufferedMetrics: () => Promise }).reportBufferedMetrics(); await Promise.resolve(); + await Promise.resolve(); pendingCollectCallback?.(null); await reportPromise; @@ -178,10 +182,61 @@ describe('MeterSender', () => { jest.advanceTimersByTime(1000); await Promise.resolve(); + await Promise.resolve(); pendingCollectCallback?.(null); await Promise.resolve(); expect(collector.sample).toHaveBeenCalled(); expect(mockStream.write).toHaveBeenCalled(); }); + + it('flush drains a new snapshot after an in-flight report', async () => { + (sender as unknown as { latestSnapshot?: { collectedAt: number } }).latestSnapshot = { + collectedAt: 1, + }; + const first = (sender as unknown as { reportBufferedMetrics: () => Promise }).reportBufferedMetrics(); + await Promise.resolve(); + await Promise.resolve(); + + (sender as unknown as { latestSnapshot?: { collectedAt: number } }).latestSnapshot = { + collectedAt: 2, + }; + const second = (sender as unknown as { flush: () => Promise | null }).flush(); + + pendingCollectCallback?.(null); + await first; + await Promise.resolve(); + await Promise.resolve(); + pendingCollectCallback?.(null); + await second; + + expect(mockStream.write.mock.calls.length).toBeGreaterThanOrEqual(2); + }); + + it('keeps snapshot when disconnected instead of consuming it', async () => { + sender.statusChanged(GRPCChannelStatus.DISCONNECT); + (sender as unknown as { latestSnapshot?: { collectedAt: number } }).latestSnapshot = { + collectedAt: 99, + }; + await (sender as unknown as { reportBufferedMetrics: () => Promise }).reportBufferedMetrics(); + expect((sender as unknown as { latestSnapshot?: { collectedAt: number } }).latestSnapshot?.collectedAt).toBe(99); + }); + + it('reports meter failure once when sync end fails and callback also errors', async () => { + (sender as unknown as { latestSnapshot?: { collectedAt: number } }).latestSnapshot = { + collectedAt: 7, + }; + mockStream.end.mockImplementation(() => { + throw new Error('end failed'); + }); + + const report = (sender as unknown as { reportBufferedMetrics: () => Promise }).reportBufferedMetrics(); + await Promise.resolve(); + await Promise.resolve(); + pendingCollectCallback?.(new Error('callback error')); + await report; + + expect(mockChannelManager.reportError).toHaveBeenCalledTimes(1); + expect(mockStream.cancel).toHaveBeenCalled(); + }); }); diff --git a/tests/remote/TraceSegmentServiceClient.test.ts b/tests/remote/TraceSegmentServiceClient.test.ts new file mode 100644 index 00000000..7921efc3 --- /dev/null +++ b/tests/remote/TraceSegmentServiceClient.test.ts @@ -0,0 +1,130 @@ +/*! + * + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + * + */ + +/* eslint-env jest */ + +import TraceSegmentServiceClient from '../../src/agent/core/remote/TraceSegmentServiceClient'; +import { GRPCChannelStatus } from '../../src/agent/core/remote/GRPCChannelStatus'; +import ServiceManager from '../../src/agent/core/boot/ServiceManager'; +import Segment from '../../src/trace/context/Segment'; + +const mockStream = { + write: jest.fn(), + end: jest.fn(), + cancel: jest.fn(), +}; + +let pendingCollectCallback: ((error: Error | null) => void) | undefined; +const mockCollect = jest.fn((_meta: unknown, _opts: unknown, cb: (error: Error | null) => void) => { + pendingCollectCallback = cb; + return mockStream; +}); +const mockCreateClient = jest.fn(() => ({ collect: mockCollect })); +const mockChannelManager = { + addChannelListener: jest.fn(), + createClient: mockCreateClient, + reportError: jest.fn(), +}; + +jest.mock('../../src/agent/core/boot/ServiceManager', () => ({ + __esModule: true, + default: { + INSTANCE: { + findService: jest.fn(() => mockChannelManager), + }, + }, +})); + +function fakeSegment(): Segment { + return { + transform: () => ({ fake: true }), + } as unknown as Segment; +} + +describe('TraceSegmentServiceClient flush / coalesce', () => { + let client: TraceSegmentServiceClient; + + beforeEach(() => { + jest.clearAllMocks(); + pendingCollectCallback = undefined; + mockStream.write.mockReset(); + mockStream.end.mockReset(); + mockStream.cancel.mockReset(); + mockCollect.mockImplementation((_meta, _opts, cb) => { + pendingCollectCallback = cb; + return mockStream; + }); + (ServiceManager.INSTANCE.findService as jest.Mock).mockReturnValue(mockChannelManager); + + client = new TraceSegmentServiceClient(); + client.prepare(); + client.statusChanged(GRPCChannelStatus.CONNECTED); + client.boot(); + }); + + afterEach(() => { + client.shutdown(); + }); + + it('flush while a report is in flight runs a second report and drains the buffer', async () => { + const buffer = (client as unknown as { buffer: Segment[] }).buffer; + buffer.push(fakeSegment()); + + const reportOnce = (client as unknown as { reportOnce: () => Promise }).reportOnce.bind(client); + const first = reportOnce(); + await Promise.resolve(); + expect(mockCollect).toHaveBeenCalledTimes(1); + expect(buffer.length).toBe(0); + + // Segment arrives while first report is still in flight. + buffer.push(fakeSegment()); + const flushPromise = client.flush(); + expect(flushPromise).not.toBeNull(); + + pendingCollectCallback?.(null); + await first; + await Promise.resolve(); + await Promise.resolve(); + + // Second report should have started for the buffered segment (forceReport). + expect(mockCollect).toHaveBeenCalledTimes(2); + pendingCollectCallback?.(null); + await flushPromise; + + expect(buffer.length).toBe(0); + expect(mockStream.write.mock.calls.length).toBeGreaterThanOrEqual(2); + }); + + it('discardBatch runs at most once when sync end fails and callback also errors', async () => { + const buffer = (client as unknown as { buffer: Segment[] }).buffer; + buffer.push(fakeSegment()); + mockStream.end.mockImplementation(() => { + throw new Error('end failed'); + }); + + const reportOnce = (client as unknown as { reportOnce: () => Promise }).reportOnce.bind(client); + const p = reportOnce(); + await Promise.resolve(); + pendingCollectCallback?.(new Error('callback error')); + await p; + + expect(mockChannelManager.reportError).toHaveBeenCalledTimes(1); + expect(mockStream.cancel).toHaveBeenCalled(); + }); +}); diff --git a/tests/remote/coalesceReport.test.ts b/tests/remote/coalesceReport.test.ts new file mode 100644 index 00000000..a4558c85 --- /dev/null +++ b/tests/remote/coalesceReport.test.ts @@ -0,0 +1,164 @@ +/*! + * + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + * + */ + +/* eslint-env jest */ + +import { + awaitWithTimeout, + coalesceReport, + flushCoalesced, + forceReport, + ReportCoalesceState, +} from '../../src/agent/core/remote/coalesceReport'; + +describe('coalesceReport', () => { + it('is single-flight: concurrent callers share the in-flight promise', async () => { + const state: ReportCoalesceState = {}; + let starts = 0; + let finish!: () => void; + const gate = new Promise((r) => { + finish = r; + }); + const doReport = jest.fn(async () => { + starts += 1; + await gate; + }); + + const a = coalesceReport(state, doReport, () => false); + const b = coalesceReport(state, doReport, () => false); + await Promise.resolve(); + expect(starts).toBe(1); + finish(); + await Promise.all([a, b]); + expect(starts).toBe(1); + }); + + it('never rejects when doReport rejects', async () => { + const state: ReportCoalesceState = {}; + await expect( + coalesceReport( + state, + async () => { + throw new Error('boom'); + }, + () => false, + ), + ).resolves.toBeUndefined(); + }); + + it('forceReport starts a new flight even when one is in progress', async () => { + const state: ReportCoalesceState = {}; + let starts = 0; + let finishFirst!: () => void; + const firstGate = new Promise((r) => { + finishFirst = r; + }); + const doReport = jest.fn(async () => { + starts += 1; + if (starts === 1) { + await firstGate; + } + }); + + const first = coalesceReport(state, doReport, () => false); + await Promise.resolve(); + expect(starts).toBe(1); + + const second = forceReport(state, doReport, () => false); + await Promise.resolve(); + expect(starts).toBe(2); + + finishFirst(); + await Promise.all([first, second]); + }); + + it('flushCoalesced forces pending before waiting on in-flight', async () => { + const state: ReportCoalesceState = {}; + let starts = 0; + let finishFirst!: () => void; + let finishSecond!: () => void; + const firstGate = new Promise((r) => { + finishFirst = r; + }); + const secondGate = new Promise((r) => { + finishSecond = r; + }); + const doReport = jest.fn(async () => { + starts += 1; + if (starts === 1) { + await firstGate; + } else { + await secondGate; + } + }); + + void coalesceReport(state, doReport, () => false); + await Promise.resolve(); + expect(starts).toBe(1); + + const flushPromise = flushCoalesced( + state, + doReport, + () => false, + () => true, + 80, + ); + + await Promise.resolve(); + await Promise.resolve(); + // Pending is forced while the first flight is still blocked (not after a full wait). + expect(starts).toBe(2); + + finishSecond(); + await new Promise((r) => setTimeout(r, 30)); + finishFirst(); + await flushPromise; + expect(starts).toBe(2); + }); + + it('flushCoalesced still awaits forceReport when in-flight would exhaust the budget', async () => { + const state: ReportCoalesceState = {}; + let forceDone = false; + const never = new Promise(() => undefined); + state.reporting = never; + + const doReport = jest.fn(async () => { + await new Promise((r) => setTimeout(r, 40)); + forceDone = true; + }); + + await flushCoalesced( + state, + doReport, + () => false, + () => true, + 100, + ); + // Under wait-then-force, budget is spent on `never` and flush returns before force finishes. + expect(forceDone).toBe(true); + expect(doReport).toHaveBeenCalledTimes(1); + }); + + it('awaitWithTimeout resolves when the timer wins', async () => { + const never = new Promise(() => undefined); + const started = Date.now(); + await awaitWithTimeout(never, 50); + expect(Date.now() - started).toBeLessThan(500); + }); +});