Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
8 changes: 8 additions & 0 deletions .changeset/sync-connections-metric.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,8 @@
---
'@powersync/service-core': minor
'@powersync/service-types': minor
'@powersync/service-rsocket-router': minor
'@powersync/service-module-core': minor
---

Add `powersync_sync_connections_total` to track sync stream close outcomes and setup rejections across HTTP and RSocket.
3 changes: 3 additions & 0 deletions modules/module-core/package.json
Original file line number Diff line number Diff line change
Expand Up @@ -24,5 +24,8 @@
"@powersync/service-rsocket-router": "workspace:*",
"@powersync/service-types": "workspace:*",
"fastify": "^5.8.5"
},
"devDependencies": {
"@opentelemetry/sdk-metrics": "^2.0.1"
}
}
17 changes: 13 additions & 4 deletions modules/module-core/src/CoreModule.ts
Original file line number Diff line number Diff line change
Expand Up @@ -19,13 +19,14 @@ export class CoreModule extends core.modules.AbstractModule {
this.registerAPIRoutes(context);
}

// Shutdown runs in reverse order: drain streams before the final metric export.
await this.configureMetrics(context);

// Configures a Fastify server and RSocket server
this.configureRouterImplementation(context);

// Configures health check probes based off configuration
this.configureHealthChecks(context);

await this.configureMetrics(context);
}

protected configureTags(context: core.ServiceContextContainer) {
Expand Down Expand Up @@ -75,7 +76,14 @@ export class CoreModule extends core.modules.AbstractModule {
});

const socketRouter = new ReactiveSocketRouter<core.routes.Context>({
max_concurrent_connections: context.configuration.api_parameters.max_concurrent_connections
max_concurrent_connections: context.configuration.api_parameters.max_concurrent_connections,
on_concurrency_limit_rejected: (error) => {
core.metrics.recordSyncConnection(context.metricsEngine, {
transport: core.metrics.SyncTransport.RSocket,
closeReason: core.metrics.SyncCloseReason.ConcurrencyLimit,
error
});
}
});

core.routes.configureRSocket(socketRouter, {
Expand All @@ -101,7 +109,8 @@ export class CoreModule extends core.modules.AbstractModule {
}
};
});
}
},
stop: (routerEngine) => routerEngine.shutDown()
});
}

Expand Down
58 changes: 58 additions & 0 deletions modules/module-core/test/src/metrics-shutdown.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,58 @@
import { CoreModule } from '@module/CoreModule.js';
import { MeterProvider } from '@opentelemetry/sdk-metrics';
import { container } from '@powersync/lib-services-framework';
import { ServiceContextContainer, ServiceContextMode, utils } from '@powersync/service-core';
import { setImmediate } from 'node:timers/promises';
import { expect, it, onTestFinished, vi } from 'vitest';

it.each([ServiceContextMode.API, ServiceContextMode.UNIFIED])(
'drains the router before shutting down metrics and storage in %s mode',
async (serviceMode) => {
container.registerDefaults();
onTestFinished(() => {
vi.restoreAllMocks();
});
const configuration = await new utils.CompoundConfigCollector().collectConfig({
config_base64: Buffer.from(
`
telemetry:
disable_telemetry_sharing: true
storage:
type: memory
`
).toString('base64')
});
const context = new ServiceContextContainer({ configuration, serviceMode });
await new CoreModule().initialize(context);

const drain = Promise.withResolvers<void>();
const routerDrained = vi.fn();
await context.routerEngine.start(async () => ({
onShutdown: async () => {
await drain.promise;
routerDrained();
}
}));
const routerShutdown = vi.spyOn(context.routerEngine, 'shutDown');
// MeterProvider.shutdown() performs the final export.
const metricsShutdown = vi.spyOn(MeterProvider.prototype, 'shutdown');
const storageShutdown = vi.spyOn(context.storageEngine, 'shutDown');
const stopping = context.lifeCycleEngine.stop();
try {
await expect.poll(() => routerShutdown.mock.calls.length).toBe(1);
// Let an incorrectly unawaited shutdown advance while the router is still draining.
await setImmediate();
expect(metricsShutdown).not.toHaveBeenCalled();
expect(storageShutdown).not.toHaveBeenCalled();
} finally {
drain.resolve();
await stopping;
}

expect(routerDrained).toHaveBeenCalledOnce();
expect(metricsShutdown).toHaveBeenCalledOnce();
expect(storageShutdown).toHaveBeenCalledOnce();
expect(routerDrained).toHaveBeenCalledBefore(metricsShutdown);
expect(metricsShutdown).toHaveBeenCalledBefore(storageShutdown);
}
);
5 changes: 3 additions & 2 deletions packages/rsocket-router/src/router/ReactiveSocketRouter.ts
Original file line number Diff line number Diff line change
Expand Up @@ -94,16 +94,17 @@ export class ReactiveSocketRouter<C extends SocketBaseContext> {
accept: async (payload, rsocket) => {
const connection = (rsocket as any).connection as WebsocketDuplexConnection;

const { max_concurrent_connections } = this.options ?? {};
const { max_concurrent_connections, on_concurrency_limit_rejected } = this.options ?? {};
logger.info(`Currently have ${wss.clients.size} active WebSocket connection(s)`);
// wss.clients.size includes this connection, so we check for greater than
// TODO: Share connection limit between this and http stream connections
if (max_concurrent_connections && wss.clients.size > max_concurrent_connections) {
const err = new errors.ServiceError({
status: 429,
code: errors.ErrorCode.PSYNC_S2304,
code: ErrorCode.PSYNC_S2304,
description: `Maximum active concurrent connections limit has been reached`
});
on_concurrency_limit_rejected?.(err);
logger.warn(err);
throw err;
}
Expand Down
7 changes: 6 additions & 1 deletion packages/rsocket-router/src/router/types.ts
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
import { router } from '@powersync/lib-services-framework';
import { errors, router } from '@powersync/lib-services-framework';
import * as t from 'ts-codec';

import { OnExtensionSubscriber, OnNextSubscriber, OnTerminalSubscriber } from 'rsocket-core';
Expand All @@ -19,6 +19,11 @@ export type RequestMeta = t.Decoded<typeof RSocketRequestMeta>;

export type ReactiveSocketRouterOptions<C> = {
max_concurrent_connections?: number;
/**
* Invoked with the error returned to the client when a connection SETUP is rejected because
* `max_concurrent_connections` was exceeded.
*/
on_concurrency_limit_rejected?: (error: errors.ServiceError) => void;
};

export type SocketResponder = OnTerminalSubscriber & OnNextSubscriber & OnExtensionSubscriber;
Expand Down
72 changes: 72 additions & 0 deletions packages/rsocket-router/test/src/concurrency-limit.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,72 @@
import { ErrorCode, logger } from '@powersync/lib-services-framework';
import { deserialize, serialize } from 'bson';
import * as http from 'http';
import { RSocketConnector } from 'rsocket-core';
import { WebsocketClientTransport } from 'rsocket-websocket-client';
import { afterEach, describe, expect, it, vi } from 'vitest';
import * as WebSocket from 'ws';
import { ReactiveSocketRouter, SocketBaseContext } from '../../src/router/ReactiveSocketRouter.js';

// Port range distinct from socket.test.ts (5433+) to avoid clashes between parallel workers
let nextPort = 5600;

describe('Concurrency limit', () => {
let cleanup: (() => Promise<void> | void)[] = [];

afterEach(async () => {
for (const fn of cleanup.reverse()) {
await fn();
}
cleanup = [];
});

function createConnector(address: string) {
return new RSocketConnector({
transport: new WebsocketClientTransport({
url: address,
wsCreator: (url) => new WebSocket.WebSocket(url) as any
}),
setup: {
dataMimeType: 'application/bson',
metadataMimeType: 'application/bson',
payload: {
data: null,
metadata: Buffer.from(serialize({ token: 'test-token' }))
}
}
});
}

it('rejects connections over max_concurrent_connections and reports each rejection', async () => {
const port = nextPort++;
const address = `ws://localhost:${port}`;

const httpServer = http.createServer();
await new Promise<void>((resolve) => httpServer.listen(port, resolve));
cleanup.push(() => new Promise<void>((resolve) => httpServer.close(() => resolve())));

const onLimitRejected = vi.fn();
const router = new ReactiveSocketRouter<SocketBaseContext>({
max_concurrent_connections: 1,
on_concurrency_limit_rejected: onLimitRejected
});

router.applyWebSocketEndpoints(httpServer, {
contextProvider: async () => ({ logger }),
endpoints: [],
metaDecoder: async (meta) => deserialize(meta.contents) as any,
payloadDecoder: async (rawData) => rawData && deserialize(rawData.contents)
});

const firstClient = await createConnector(address).connect();
cleanup.push(() => firstClient.close());
expect(onLimitRejected).not.toHaveBeenCalled();

// SETUP rejection reaches onClose, not connect().
const secondClient = await createConnector(address).connect();
const closeError = await new Promise<Error | undefined>((resolve) => secondClient.onClose(resolve));
expect(closeError?.message).toMatch(/Maximum active concurrent connections limit has been reached/);
expect(onLimitRejected).toHaveBeenCalledTimes(1);
expect(onLimitRejected.mock.calls[0][0].errorData.code).toBe(ErrorCode.PSYNC_S2304);
});
});
7 changes: 7 additions & 0 deletions packages/service-core/src/api/api-metrics.ts
Original file line number Diff line number Diff line change
@@ -1,5 +1,6 @@
import { APIMetric } from '@powersync/service-types';
import { MetricsEngine } from '../metrics/MetricsEngine.js';
import { initializeSyncConnectionMetrics } from '../metrics/connection-metrics.js';

/**
* Create and register the core API metrics.
Expand Down Expand Up @@ -27,6 +28,11 @@ export function createCoreAPIMetrics(engine: MetricsEngine): void {
name: APIMetric.CONCURRENT_CONNECTIONS,
description: 'Number of concurrent sync connections'
});

engine.createCounter({
name: APIMetric.SYNC_CONNECTIONS,
description: 'Sync stream closes and setup rejections, by outcome, close_reason, error_code and transport'
});
}

/**
Expand All @@ -38,4 +44,5 @@ export function initializeCoreAPIMetrics(engine: MetricsEngine): void {

// Initialize the metric, so that it reports a value before connections have been opened.
concurrent_connections.add(0);
initializeSyncConnectionMetrics(engine);
}
103 changes: 103 additions & 0 deletions packages/service-core/src/metrics/connection-metrics.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,103 @@
import { ErrorCode, errors } from '@powersync/lib-services-framework';
import { APIMetric } from '@powersync/service-types';
import { MetricsEngine } from './MetricsEngine.js';

export enum SyncCloseReason {
/** Routine client disconnect/reconnect. */
ClientClosed = 'client_closed',
/** Server-side close: auth expiry or switch to a new sync config. */
ServiceClosed = 'service_closed',
ProcessShutdown = 'process_shutdown',
StreamError = 'stream_error',
ServiceUnavailable = 'service_unavailable',
NoSyncConfig = 'no_sync_config',
StorageError = 'storage_error',
SyncConfigError = 'sync_config_error',
ConcurrencyLimit = 'concurrency_limit',
Unknown = 'unknown'
}

export enum SyncTransport {
HttpStream = 'http_stream',
RSocket = 'rsocket'
}

type SyncConnectionReasonPolicy = Readonly<{
outcome: 'success' | 'error' | 'rejected';
logText: string;
}>;

const SYNC_CONNECTION_REASON_POLICY = {
[SyncCloseReason.ClientClosed]: { outcome: 'success', logText: 'client closing stream' },
[SyncCloseReason.ServiceClosed]: { outcome: 'success', logText: 'service closing stream' },
[SyncCloseReason.ProcessShutdown]: { outcome: 'success', logText: 'process shutdown' },
[SyncCloseReason.StreamError]: { outcome: 'error', logText: 'stream error' },
[SyncCloseReason.ServiceUnavailable]: { outcome: 'rejected', logText: 'service unavailable' },
[SyncCloseReason.NoSyncConfig]: { outcome: 'rejected', logText: 'no sync config' },
[SyncCloseReason.StorageError]: { outcome: 'rejected', logText: 'storage error' },
[SyncCloseReason.SyncConfigError]: { outcome: 'rejected', logText: 'sync config error' },
[SyncCloseReason.ConcurrencyLimit]: { outcome: 'rejected', logText: 'concurrency limit' },
// Nothing was thrown or reported, so the stream ended cleanly without a specific reason.
[SyncCloseReason.Unknown]: { outcome: 'success', logText: 'unknown' }
} as const satisfies Readonly<Record<SyncCloseReason, SyncConnectionReasonPolicy>>;

/** Wording used by existing log-based dashboards for the `close_reason` field. */
export function syncConnectionCloseReasonLogText(closeReason?: SyncCloseReason): string {
return SYNC_CONNECTION_REASON_POLICY[closeReason ?? SyncCloseReason.Unknown].logText;
}

export interface SyncConnectionMetric {
transport: SyncTransport;
closeReason: SyncCloseReason;
/** The original failure, before any transport-specific wrapping. */
error?: unknown;
}

const ERROR_CODES: ReadonlySet<string> = new Set(Object.values(ErrorCode));

/** Seed valid series so a scrape before the first close supplies a zero baseline. */
export function initializeSyncConnectionMetrics(engine: MetricsEngine): void {
const counter = engine.getCounter(APIMetric.SYNC_CONNECTIONS);
const failureCodes = [...ERROR_CODES, 'other'];
for (const transport of Object.values(SyncTransport)) {
for (const closeReason of Object.values(SyncCloseReason)) {
const { outcome } = SYNC_CONNECTION_REASON_POLICY[closeReason];
let errorCodes: readonly string[];
switch (closeReason) {
case SyncCloseReason.ServiceUnavailable:
errorCodes = [ErrorCode.PSYNC_S2003];
break;
case SyncCloseReason.NoSyncConfig:
errorCodes = [ErrorCode.PSYNC_S2302];
break;
case SyncCloseReason.ConcurrencyLimit:
// HTTP keeps its existing empty 429 response, without a PowerSync error code.
errorCodes = [transport === SyncTransport.HttpStream ? 'other' : ErrorCode.PSYNC_S2304];
break;
default:
errorCodes = outcome === 'success' ? ['none'] : failureCodes;
}
for (const errorCode of errorCodes) {
counter.add(0, { outcome, close_reason: closeReason, error_code: errorCode, transport });
}
}
}
}

export function recordSyncConnection(engine: MetricsEngine, metric: SyncConnectionMetric): void {
const { outcome } = SYNC_CONNECTION_REASON_POLICY[metric.closeReason];

// Keep successful closes in one error_code series, even if an error was supplied.
let errorCode = 'none';
if (outcome !== 'success') {
const code = errors.ServiceError.isServiceError(metric.error) ? metric.error.errorData?.code : undefined;
errorCode = typeof code == 'string' && ERROR_CODES.has(code) ? code : 'other';
}

engine.getCounter(APIMetric.SYNC_CONNECTIONS).add(1, {
outcome,
close_reason: metric.closeReason,
error_code: errorCode,
transport: metric.transport
});
}
1 change: 1 addition & 0 deletions packages/service-core/src/metrics/metrics-index.ts
Original file line number Diff line number Diff line change
@@ -1,3 +1,4 @@
export * from './connection-metrics.js';
export * from './metrics-interfaces.js';
export * from './MetricsEngine.js';
export * from './open-telemetry/OpenTelemetryMetricsFactory.js';
Expand Down
9 changes: 8 additions & 1 deletion packages/service-core/src/metrics/metrics-interfaces.ts
Original file line number Diff line number Diff line change
@@ -1,9 +1,16 @@
/**
* Bounded-cardinality attributes (labels) attached to a metric data point.
* Only use low-cardinality enum-like values here — never user/client/request identifiers.
*/
export type MetricAttributes = Record<string, string | number | boolean>;

export interface Counter {
/**
* Increment the counter by the given value. Only positive numbers are valid.
* @param value
* @param attributes optional low-cardinality labels for this increment
*/
add(value: number): void;
add(value: number, attributes?: MetricAttributes): void;
}

export interface UpDownCounter {
Expand Down
Loading
Loading