Compare commits

..

6 Commits

Author SHA1 Message Date
syuilo a98ac7216d New translations ja-jp.yml (Persian)
[ci skip]
2026-08-01 21:20:20 +09:00
syuilo 491fea52a5 New translations ja-jp.yml (Persian)
[ci skip]
2026-08-01 20:01:25 +09:00
syuilo c507ae9988 New translations ja-jp.yml (Persian)
[ci skip]
2026-08-01 18:39:21 +09:00
syuilo daaa0dac88 New translations ja-jp.yml (Persian)
[ci skip]
2026-08-01 17:42:59 +09:00
syuilo 8e467aec09 New translations ja-jp.yml (Chinese Traditional)
[ci skip]
2026-08-01 02:12:15 +09:00
4ster1sk 4db251696d fix(test): setup.sh で $1.config.json の生成条件が $1.default.yml を参照していたのを修正 (#17826)
* fix(test): setup.sh で $1.config.json の生成条件が $1.default.yml を参照していたのを修正

* update .gitignore
2026-07-31 19:34:17 +09:00
39 changed files with 162 additions and 3288 deletions
-57
View File
@@ -329,63 +329,6 @@ id: 'aidx'
# #tracePropagationTargets:
# # - 'internal-service.example'
# ┌─────────┐
#───┘ Tracing └────────────────────────────────────────────────
# OpenTelemetry distributed tracing (Traces only, opt-in).
# Existing API and Queue spans are exported to an OTLP http/protobuf endpoint.
#
# Standard OTEL_* environment variables are honored for values omitted below:
# OTEL_EXPORTER_OTLP_TRACES_ENDPOINT, OTEL_EXPORTER_OTLP_ENDPOINT,
# OTEL_EXPORTER_OTLP_HEADERS, OTEL_TRACES_SAMPLER, OTEL_TRACES_SAMPLER_ARG, etc.
#
# When this is enabled together with sentryForBackend, Misskey shares Sentry's
# TracerProvider and adds an OTLP SpanProcessor to it. In that mode:
# - sampleRate below is ignored; sampling follows sentryForBackend.options
# tracesSampleRate / tracesSampler.
# - resourceAttributes below is ignored; set OTEL_SERVICE_NAME and
# OTEL_RESOURCE_ATTRIBUTES in the environment if you need to override them.
# - spans produced by Sentry integrations are exported to both Sentry and OTLP.
# - sentry-trace / baggage propagation to outbound remote requests is disabled
# by default while OTel is enabled. Set propagateTraceToRemote: true only if
# you intentionally want Sentry trace propagation for your deployment.
#otelForBackend:
# endpoint: 'http://localhost:4318/v1/traces'
# #headers:
# # authorization: 'Bearer xxxxx'
# #sampleRate: 1.0
# # Record PostgreSQL query spans. Disabled by default to avoid instrumentation
# # overhead when database-level diagnostics are not needed.
# #capturePgSpans: false
# # Include raw SQL text in PostgreSQL spans. Requires capturePgSpans and is
# # disabled by default because non-parameterized queries can expose values to
# # the OTLP Collector.
# #capturePgStatement: false
# # Include PostgreSQL connection-pool spans with an existing parent span.
# # Requires capturePgSpans and is disabled by default to avoid noisy spans
# # from connection churn.
# #capturePgConnectionSpans: false
# # Record Redis command spans. Disabled by default because subscribing to the
# # ioredis diagnostics channel adds work to every Redis command, including
# # commands without a parent span. Enable for development diagnostics or when
# # the added overhead is acceptable in production.
# #captureRedisCommandSpans: false
# # Include Redis startup/reconnect spans. Disabled by default to avoid noisy
# # connection churn; these spans are recorded even without an HTTP or job
# # parent span.
# #captureRedisConnectionSpans: false
# # Record Redis commands without an HTTP or job parent span. Disabled by
# # default because queue polling and background work can produce many root
# # traces; only applies when captureRedisCommandSpans is enabled.
# #captureRedisRootSpans: false
# #resourceAttributes:
# # deployment.environment: 'production'
# #propagateTraceToRemote: false
# # Queue worker spans are independent traces linked to their enqueuer by default.
# # Set to 'parent' to make them child spans instead. This can create very large
# # traces for high-fan-out queues such as federation delivery.
# #jobTraceContextMode: 'link'
#sentryForFrontend:
# vueIntegration:
# tracingOptions:
+2 -9
View File
@@ -7,15 +7,8 @@
-
### Server
- Feat: OpenTelemetryサポート
- 詳細な設定はconfigファイルを参照してください。
- Sentryとの併用も可能です。Sentry併用時は、PostgreSQL Query と Redis command は Sentry で計装されます。
- 以下の自動計装をサポートしています。(計装対象にする項目は設定可能)
- PostgreSQL query
- Redis command
- 全ての受信HTTPリクエスト
- 全ての送信HTTPリクエスト
- ジョブキュー(エンキュー元のトレースを含む)
-
## 2026.7.0
+81
View File
@@ -0,0 +1,81 @@
---
_lang_: "فارسی"
headlineMisskey: "یک شبکه از طریق یادداشت ها متصل گردید."
introMisskey: "خوش آمدید! Misskey یک سرویس میکروبلاگینگ متن باز و غیرمتمرکز می باشد.\nیک \"یادداشت\" درست کن و در مورد اتفاق هایی که داره می افته و یا به دیگران درمورد خودت بگو📡\nشما می توانید به سرعت با کمک قابلیت \"واکنش ها\" به یادداشت های دیگران واکنش نشون بدید👍\nیک دنیای جدید رو کاوش کنید🚀"
poweredByMisskeyDescription: "{name} یکی از سرور های متن باز پلتفرم <b>Misskey</b> می باشد."
monthAndDay: "{day}/{month}"
search: "جستجو"
reset: "بازنشانی"
notifications: "اعلان ها"
username: "نام کاربری"
password: "گذرواژه"
initialPasswordForSetup: "گذرواژه برای راه‌اندازی تنظیمات اولیه"
initialPasswordIsIncorrect: "گذرواژه برای راه‌اندازی تنظیمات اولیه نادرست می باشد."
initialPasswordForSetupDescription: "اگر خودتان Misskey را نصب کرده اید، لطفا از گذرواژه ای که در فایل تنظیمات وارد کرده اید استفاده کنید.\nاگه از سرویس میزبانی برای Misskey استفاده می کنید، لطفا از گذرواژه ای که توسط آنها ارائه شده استفاده کنید.\nاگر گذرواژه ای تنظیم نکردید، بخش را خالی بگذارید و ادامه دهید."
forgotPassword: "گذرواژه ام را فراموش کردم"
fetchingAsApObject: "در حال دریافت از فدیوِرس"
ok: "باشه"
gotIt: "متوجه شدم"
cancel: "لغو"
noThankYou: "الان نه"
enterUsername: "نام کاربری خود را وارد کنید"
renotedBy: "بازنشر شده توسط {user}"
noNotes: "هیچ یادداشتی موجود نیست"
noNotifications: "اعلانی موجود نیست"
instance: "سرور"
settings: "تنظیمات"
notificationSettings: "تنظیمات اعلان"
basicSettings: "تنظیمات اصلی"
otherSettings: "تنظیمات بیشتر"
openInWindow: "در پنجره باز کن"
profile: "نمایه"
timeline: "نوار زمانی"
noAccountDescription: "این کاربر هنوز درون شرح حال خودش چیزی ننوشته."
login: "ورود"
loggingIn: "وارد شده"
logout: "خروج"
signup: "ثبت‌نام"
uploading: "بارگذاری..."
save: "ذخیره"
users: "کاربران"
addUser: "افزودن کاربر"
favorite: "افزودن به برگزیده ها"
favorites: "برگزیده ها"
unfavorite: "حذف از برگزیده ها"
favorited: "به برگزیده ها اضافه شد."
alreadyFavorited: "قبلا به برگزیده ها اضافه شده."
cantFavorite: "نمی توان به برگزیده ها اضافه کرد."
pin: "سنجاق کردن"
unpin: "حذف سنجاق"
copyContent: "رونوشت از محتوا"
copyLink: "رونوشت از پیوند"
copyRemoteLink: "رونوشت از پیوند خارجی"
copyLinkRenote: "رونوشت از پیوند بازنشر"
delete: "حذف"
deleteAndEdit: "حذف و ویرایش"
deleteAndEditConfirm: "آیا می خواهید که این یادداشت را حذف و دوباره ویرایش کنید؟ در این صورت همچنین تمام واکنش ها، بازنشر ها و پاسخ ها به این یادداشت پاک خواهد شد."
addToList: "افزودن به فهرست"
pinned: "سنجاق کردن"
instances: "سرور"
remove: "حذف"
smtpUser: "نام کاربری"
smtpPass: "گذرواژه"
user: "کاربران"
searchByGoogle: "جستجو"
_sfx:
notification: "اعلان ها"
_2fa:
renewTOTPCancel: "الان نه"
_widgets:
profile: "نمایه"
notifications: "اعلان ها"
timeline: "نوار زمانی"
_profile:
username: "نام کاربری"
_notification:
_types:
login: "ورود"
_deck:
_columns:
notifications: "اعلان ها"
tl: "نوار زمانی"
+15
View File
@@ -619,6 +619,8 @@ output: "輸出"
script: "腳本"
disablePagesScript: "停用頁面的 AiScript 腳本"
updateRemoteUser: "更新遠端使用者資訊"
unsetMfa: "解除雙重驗證"
unsetMfaConfirm: "要解除雙重驗證嗎?"
unsetUserAvatar: "移除使用者的大頭貼"
unsetUserAvatarConfirm: "確定要移除使用者的大頭貼嗎?"
unsetUserBanner: "移除使用者的橫幅圖像"
@@ -1417,6 +1419,8 @@ addToEmojiPalette: "增加表情符號調色盤"
emojiPaletteAlreadyAddedConfirm: "此表情符號在這個表情符號調色盤裡已經有了。確定要增加嗎?"
append: "加在最後"
prepend: "加在前面"
urlPreviewSensitiveList: "限制縮圖顯示的 URL"
urlPreviewSensitiveListDescription: "以空格指定為 AND,以換行指定為 OR。若以斜線(/)包圍則視為正規表達式。符合條件時,將不再顯示縮圖。"
_imageEditing:
_vars:
caption: "檔案標題"
@@ -2149,6 +2153,15 @@ _sensitiveMediaDetection:
setSensitiveFlagAutomaticallyDescription: "即使將此設定關閉,判定結果也會保留在內部。"
analyzeVideos: "啟用影片分析"
analyzeVideosDescription: "除了靜止影像以外,也分析影片。伺服器的負荷會稍微增加。"
externalServiceInfo: "敏感媒體的判定已分離至外部服務(sensitive-detector)。若要使用此功能,必須另外設定 Sidecar 服務,並設定下方的連接資訊。若未設定連接資訊,則不會執行判定(視為非敏感內容)。"
apiUrl: "判定服務的連接資訊 URL"
apiUrlDescription: "sensitive-detector 服務的基礎網址(例如:http://localhost:3009)。若要連接位於私有網路上的服務,請在設定檔的 allowedPrivateNetworks 中允許對應的連接網路。若使用代理伺服器(Proxy),也請一併設定 proxyBypassHosts。若留白,則不會執行敏感內容判定。"
apiKey: "API 金鑰"
apiKeyDescription: "若判定服務端已設定認證(Bearer Token),請輸入此資訊。若未設定,請保持空白。"
timeout: "連線逾時(微秒)"
timeoutDescription: "單次判定請求的逾時時間。"
maxImagesPerRequest: "每筆請求的最大圖片數量"
maxImagesPerRequestDescription: "當判定包含多個影格的內容(如影片等)時,為每筆請求所允許的單次最大圖片數量。超過此數量的部分將會分割並依序發送。請確保此數值不超過 sensitive-detector 端的 maxParts 設定(預設值:10)。若超過該數值,該區塊(Chunk)的所有項目將被視為非敏感內容。"
_emailUnavailable:
used: "已被使用"
format: "格式無效"
@@ -2461,6 +2474,7 @@ _permissions:
"read:admin:show-moderation-log": "查看審查紀錄"
"read:admin:show-user": "查看使用者的私密資訊"
"write:admin:suspend-user": "凍結使用者"
"write:admin:unset-mfa": "解除使用者的雙重驗證"
"write:admin:unset-user-avatar": "刪除使用者的頭像"
"write:admin:unset-user-banner": "刪除使用者的橫幅"
"write:admin:unsuspend-user": "解除凍結使用者"
@@ -2976,6 +2990,7 @@ _moderationLogTypes:
createAvatarDecoration: "建立頭像裝飾"
updateAvatarDecoration: "更新頭像裝飾"
deleteAvatarDecoration: "刪除頭像裝飾"
unsetMfa: "解除使用者的雙重驗證"
unsetUserAvatar: "移除使用者的大頭貼"
unsetUserBanner: "移除使用者的橫幅圖像"
createSystemWebhook: "建立 SystemWebhook"
-10
View File
@@ -57,7 +57,6 @@
"@fastify/cors": "11.3.0",
"@fastify/http-proxy": "11.5.0",
"@fastify/multipart": "10.1.0",
"@fastify/otel": "0.20.1",
"@fastify/static": "10.1.2",
"@kitajs/html": "4.2.13",
"@misskey-dev/emoji-assets": "17.0.3",
@@ -68,15 +67,6 @@
"@nestjs/common": "11.1.28",
"@nestjs/core": "11.1.28",
"@nestjs/testing": "11.1.28",
"@opentelemetry/api": "1.9.1",
"@opentelemetry/core": "2.9.0",
"@opentelemetry/exporter-trace-otlp-proto": "0.220.0",
"@opentelemetry/instrumentation": "0.220.0",
"@opentelemetry/instrumentation-pg": "0.72.0",
"@opentelemetry/resources": "2.9.0",
"@opentelemetry/sdk-trace-base": "2.9.0",
"@opentelemetry/sdk-trace-node": "2.9.0",
"@opentelemetry/semantic-conventions": "1.43.0",
"@oxc-project/runtime": "0.139.0",
"@peertube/http-signature": "1.7.0",
"@sentry/node": "10.65.0",
-1
View File
@@ -87,7 +87,6 @@ export default defineConfig((args) => {
'class-validator',
/^@sentry\/.*/,
/^@sentry-internal\/.*/,
/^@opentelemetry\/.*/,
'@nestjs/websockets/socket-module',
'@nestjs/microservices/microservices-module',
'@nestjs/microservices',
+2 -31
View File
@@ -27,25 +27,6 @@ type SentryBackendConfig = {
disabledIntegrations?: string[];
};
type SentryBackendConfigSource = Omit<SentryBackendConfig, 'options'> & {
options?: Partial<Sentry.NodeOptions>;
};
type OtelBackendConfig = {
endpoint?: string;
headers?: Record<string, string>;
sampleRate?: number;
capturePgSpans?: boolean;
capturePgStatement?: boolean;
capturePgConnectionSpans?: boolean;
captureRedisCommandSpans?: boolean;
captureRedisConnectionSpans?: boolean;
captureRedisRootSpans?: boolean;
resourceAttributes?: Record<string, string>;
propagateTraceToRemote?: boolean;
jobTraceContextMode?: 'link' | 'parent';
};
/**
* 設定ファイルの型
*/
@@ -90,8 +71,7 @@ type Source = {
index: string;
scope?: 'local' | 'global' | string[];
};
sentryForBackend?: SentryBackendConfigSource;
otelForBackend?: OtelBackendConfig;
sentryForBackend?: SentryBackendConfig;
sentryForFrontend?: {
options: Partial<SentryVue.BrowserOptions> & { dsn: string };
vueIntegration?: SentryVue.VueIntegrationOptions | null;
@@ -237,7 +217,6 @@ export type Config = {
redisForTimelines: RedisOptions & RedisOptionsSource;
redisForReactions: RedisOptions & RedisOptionsSource;
sentryForBackend: SentryBackendConfig | undefined;
otelForBackend: OtelBackendConfig | undefined;
sentryForFrontend: {
options: Partial<SentryVue.BrowserOptions> & { dsn: string };
vueIntegration?: SentryVue.VueIntegrationOptions | null;
@@ -342,8 +321,7 @@ export function loadConfig(): Config {
redisForJobQueue: config.redisForJobQueue ? convertRedisOptions(config.redisForJobQueue, host) : redis,
redisForTimelines: config.redisForTimelines ? convertRedisOptions(config.redisForTimelines, host) : redis,
redisForReactions: config.redisForReactions ? convertRedisOptions(config.redisForReactions, host) : redis,
sentryForBackend: config.sentryForBackend == null ? undefined : normalizeSentryBackendConfig(config.sentryForBackend),
otelForBackend: config.otelForBackend,
sentryForBackend: config.sentryForBackend,
sentryForFrontend: config.sentryForFrontend,
id: config.id,
proxy: config.proxy,
@@ -380,13 +358,6 @@ export function loadConfig(): Config {
};
}
export function normalizeSentryBackendConfig(config: SentryBackendConfigSource): SentryBackendConfig {
return {
...config,
options: config.options ?? {},
};
}
function tryCreateUrl(url: string) {
try {
return new URL(url);
+1 -7
View File
@@ -9,7 +9,6 @@ import { DI } from '@/di-symbols.js';
import type { Config } from '@/config.js';
import { baseQueueOptions, QUEUE } from '@/queue/const.js';
import { allSettled } from '@/misc/promise-tracker.js';
import { instrumentQueue } from '@/core/telemetry/queue-instrumentation.js';
import {
DeliverJobData,
EndedPollNotificationJobData,
@@ -33,12 +32,7 @@ export type UserWebhookDeliverQueue = Bull.Queue<UserWebhookDeliverJobData>;
export type SystemWebhookDeliverQueue = Bull.Queue<SystemWebhookDeliverJobData>;
function createQueue<T extends object>(queueName: string, config: Config): Bull.Queue<T> {
const queue = new Bull.Queue<T>(queueName, baseQueueOptions(config, queueName));
// Queue のラップは、enqueue 時に OTel context をジョブデータへ埋め込むためのもの。
// Sentry 単独ではジョブ間の context 伝播を使わないので、OTel 未設定時は元の Queue を返す。
if (config.otelForBackend == null) return queue;
return instrumentQueue(queue);
return new Bull.Queue<T>(queueName, baseQueueOptions(config, queueName));
}
const $system: Provider = {
@@ -5,7 +5,7 @@
import { Injectable } from '@nestjs/common';
import { bindThis } from '@/decorators.js';
import { captureMessage, shutdownTelemetry, startSpan, startSpanWithTraceContext } from './telemetry-registry.js';
import { captureMessage, shutdownTelemetry, startSpan } from './telemetry-registry.js';
import type { OnApplicationShutdown } from '@nestjs/common';
import type { TelemetryCaptureMessageOptions } from './adapters/TelemetryAdapter.js';
@@ -21,13 +21,6 @@ export class TelemetryService implements OnApplicationShutdown {
return startSpan(name, fn);
}
@bindThis
public startSpanWithTraceContext<T>(name: string, jobData: object, fn: () => T): T {
// jobData に enqueue 元の context があれば worker span へ復元する。
// context の無い既存ジョブは、通常の startSpan と同じ扱いになる。
return startSpanWithTraceContext(name, jobData, fn);
}
@bindThis
public async onApplicationShutdown(_signal?: string): Promise<void> {
await shutdownTelemetry();
@@ -1,296 +0,0 @@
/*
* SPDX-FileCopyrightText: syuilo and misskey-project
* SPDX-License-Identifier: AGPL-3.0-only
*/
import * as os from 'node:os';
import cluster from 'node:cluster';
import { envOption } from '@/env.js';
import { registerDiagLogger } from '@/core/telemetry/telemetry-diag.js';
import { installHttpClientInstrumentation } from '@/core/telemetry/http-client-instrumentation.js';
import { installDatabaseInstrumentation } from '@/core/telemetry/database-instrumentation.js';
import { installRedisInstrumentation } from '@/core/telemetry/redis-instrumentation.js';
import { executeSpan, getQueueTraceContextMode, injectActiveTraceContext, recordSpanError, startSpanWithQueueTraceContext } from '@/core/telemetry/queue-trace-context.js';
import type { LogTraceContext } from '@/logging/types.js';
import type { Span, SpanStatusCode, Tracer } from '@opentelemetry/api';
import type { Resource, ResourceDetector } from '@opentelemetry/resources';
import type { ParentBasedSampler, Sampler } from '@opentelemetry/sdk-trace-base';
import type { OtelBackendRuntimeConfig, TelemetryAdapter, TelemetryCaptureMessageOptions } from './TelemetryAdapter.js';
import type { QueueTraceContextCarrier, QueueTraceContextDeps } from '../queue-trace-context.js';
const DEFAULT_SHUTDOWN_TIMEOUT = 5000;
type OpenTelemetryAdapterDeps = {
tracer: Pick<Tracer, 'startActiveSpan'>;
provider: {
shutdown(): Promise<void>;
};
getActiveSpan: () => Span | undefined;
spanStatusCodeError: SpanStatusCode;
shutdownTimeout: number;
shutdownHttpClientInstrumentation?: () => void;
shutdownDatabaseInstrumentation?: () => void;
shutdownRedisInstrumentation?: () => void;
queueTraceContext?: QueueTraceContextDeps;
};
type CreateSamplerDeps = {
ParentBasedSampler: new (config: { root: Sampler }) => ParentBasedSampler;
TraceIdRatioBasedSampler: new (sampleRate: number) => Sampler;
};
type CreateResourceDeps = {
defaultResource: () => Resource;
resourceFromAttributes: (attributes: Record<string, string>) => Resource;
detectResources: (config: { detectors: ResourceDetector[] }) => Resource;
envDetector: ResourceDetector;
serviceNameAttribute: string;
serviceInstanceIdAttribute: string;
serviceVersionAttribute: string;
serviceVersion: string;
};
type InstrumentationInstaller = () => (() => void) | Promise<() => void>;
export class OpenTelemetryAdapter implements TelemetryAdapter {
public constructor(
private readonly deps: OpenTelemetryAdapterDeps,
) {
}
public static async create(config: OtelBackendRuntimeConfig): Promise<OpenTelemetryAdapter> {
const [
{ context, diag, DiagLogLevel, propagation, ROOT_CONTEXT, SpanKind, SpanStatusCode, trace },
{ W3CTraceContextPropagator },
{ OTLPTraceExporter },
{ defaultResource, detectResources, envDetector, resourceFromAttributes },
{ BatchSpanProcessor, ParentBasedSampler, TraceIdRatioBasedSampler },
{ NodeTracerProvider },
{ ATTR_SERVICE_INSTANCE_ID, ATTR_SERVICE_NAME, ATTR_SERVICE_VERSION },
] = await Promise.all([
import('@opentelemetry/api'),
import('@opentelemetry/core'),
import('@opentelemetry/exporter-trace-otlp-proto'),
import('@opentelemetry/resources'),
import('@opentelemetry/sdk-trace-base'),
import('@opentelemetry/sdk-trace-node'),
import('@opentelemetry/semantic-conventions'),
]);
// OTel SDK内部のexport失敗は既定だと見えにくいため、Misskeyのloggerへ橋渡しする。
registerDiagLogger(diag, DiagLogLevel.WARN);
// endpoint/headersを未指定にしておくと、OTEL_EXPORTER_OTLP_* 環境変数の標準fallbackが効く。
const exporter = new OTLPTraceExporter({
...(config.endpoint != null ? { url: config.endpoint } : {}),
...(config.headers != null ? { headers: config.headers } : {}),
});
const spanProcessor = new BatchSpanProcessor(exporter);
// SDK 2.xではSpanProcessorをprovider生成時に渡す。ここでOTel単体用のproviderを作る。
const provider = new NodeTracerProvider({
resource: createResource(config, {
defaultResource,
resourceFromAttributes,
detectResources,
envDetector,
serviceNameAttribute: ATTR_SERVICE_NAME,
serviceInstanceIdAttribute: ATTR_SERVICE_INSTANCE_ID,
serviceVersionAttribute: ATTR_SERVICE_VERSION,
serviceVersion: config.serviceVersion,
}),
...(config.sampleRate != null ? { sampler: createSampler(config.sampleRate, {
ParentBasedSampler,
TraceIdRatioBasedSampler,
}) } : {}),
spanProcessors: [spanProcessor],
});
// HTTP送信には注入しないが、将来のQueue連結でpropagation APIを使える状態にする。
provider.register({
propagator: new W3CTraceContextPropagator(),
});
// provider操作をdepsに閉じ込め、span wrapper本体をユニットテストしやすくする。
const tracer = provider.getTracer('misskey-backend');
return installInstrumentationsWithCleanup([
() => installHttpClientInstrumentation({
tracer,
spanKindClient: SpanKind.CLIENT,
spanStatusCodeError: SpanStatusCode.ERROR,
}),
// pg のrequire hookとioredis diagnostics channelは、Nest moduleの動的importより前に有効化する。
() => installDatabaseInstrumentation(provider, {
capturePgSpans: config.capturePgSpans === true,
capturePgStatement: config.capturePgStatement === true,
capturePgConnectionSpans: config.capturePgConnectionSpans === true,
}),
() => installRedisInstrumentation(tracer, SpanKind.CLIENT, SpanStatusCode.ERROR, {
captureConnectionSpans: config.captureRedisConnectionSpans === true,
captureCommandSpans: config.captureRedisCommandSpans === true,
requireParentSpan: config.captureRedisRootSpans !== true,
}),
], ([
shutdownHttpClientInstrumentation,
shutdownDatabaseInstrumentation,
shutdownRedisInstrumentation,
]) => new OpenTelemetryAdapter({
tracer,
provider,
getActiveSpan: () => trace.getActiveSpan(),
spanStatusCodeError: SpanStatusCode.ERROR,
shutdownTimeout: DEFAULT_SHUTDOWN_TIMEOUT,
shutdownHttpClientInstrumentation,
shutdownDatabaseInstrumentation,
shutdownRedisInstrumentation,
queueTraceContext: {
tracer,
propagation,
trace,
getActiveContext: () => context.active(),
rootContext: ROOT_CONTEXT,
mode: getQueueTraceContextMode(config.jobTraceContextMode),
spanStatusCodeError: SpanStatusCode.ERROR,
},
}));
}
public captureMessage(message: string, _opts: TelemetryCaptureMessageOptions): void {
// captureMessageは例外通知APIなので、OTelでは対象spanにエラー状態を付ける。
// アクティブspanが無い場合(例: BullMQのjob処理が既に完了しspanが閉じた後の'failed'イベント)でも
// 通知を握り潰さないよう、報告専用の短命spanを作ってそこに記録する。
const span = this.deps.getActiveSpan();
if (span != null) {
recordSpanError(span, new Error(message), this.deps.spanStatusCodeError);
return;
}
this.deps.tracer.startActiveSpan('captureMessage', reportSpan => {
recordSpanError(reportSpan, new Error(message), this.deps.spanStatusCodeError);
reportSpan.end();
});
}
/** activeなSpanの識別子を、Logging基盤で扱える形式へ変換します。 */
public getActiveTraceContext(): LogTraceContext | undefined {
const activeSpan = this.deps.getActiveSpan();
if (activeSpan == null) return undefined;
const { traceId, spanId, traceFlags } = activeSpan.spanContext();
return { traceId, spanId, traceFlags };
}
public startSpan<T>(name: string, fn: () => T): T {
// 既存のTelemetryAdapter契約に合わせ、同期/非同期どちらでも同じspan lifetimeを保証する。
return this.deps.tracer.startActiveSpan(name, span => executeSpan(span, fn, this.deps.spanStatusCodeError));
}
public injectTraceContext(carrier: QueueTraceContextCarrier): void {
const queueTraceContext = this.deps.queueTraceContext;
// Queue context 用の依存は任意なので、無い場合はジョブデータを変更しない。
if (queueTraceContext == null) return;
injectActiveTraceContext(queueTraceContext, carrier);
}
public startSpanWithTraceContext<T>(name: string, jobData: object, fn: () => T): T {
const queueTraceContext = this.deps.queueTraceContext;
// Queue context 用の依存が無い場合は、従来の span 作成経路と同じ動作を保つ。
if (queueTraceContext == null) return this.startSpan(name, fn);
return startSpanWithQueueTraceContext(queueTraceContext, name, jobData, fn, () => this.startSpan(name, fn));
}
public async shutdown(): Promise<void> {
shutdownInstrumentations([
this.deps.shutdownHttpClientInstrumentation,
this.deps.shutdownDatabaseInstrumentation,
this.deps.shutdownRedisInstrumentation,
]);
// BatchSpanProcessorのflushが詰まってもプロセス終了を妨げないよう、上限時間を設ける。
// タイムアウト側のtimerは、flushが先に終わった場合にイベントループを無駄に引き留めないようclearする。
let timer: NodeJS.Timeout | undefined;
await Promise.race([
this.deps.provider.shutdown(),
new Promise<void>(resolve => {
timer = setTimeout(resolve, this.deps.shutdownTimeout);
}),
]).finally(() => {
if (timer != null) clearTimeout(timer);
});
}
}
export async function installInstrumentationsWithCleanup<T>(
installers: readonly InstrumentationInstaller[],
create: (shutdowns: ReadonlyArray<() => void>) => T,
): Promise<T> {
const shutdowns: Array<() => void> = [];
try {
for (const install of installers) {
shutdowns.push(await install());
}
return create(shutdowns);
} catch (error) {
shutdownInstrumentations(shutdowns.reverse());
throw error;
}
}
function shutdownInstrumentations(shutdowns: ReadonlyArray<(() => void) | undefined>): void {
for (const shutdown of shutdowns) {
try {
shutdown?.();
} catch {
// instrumentation の停止失敗で、残りの停止処理や provider の flush を妨げない。
}
}
}
export function createResource(config: OtelBackendRuntimeConfig, deps: CreateResourceDeps): Resource {
// resourceを明示指定するとSDKのdefaultResource()は自動付与されなくなる(マージではなく上書き)ため、
// telemetry.sdk.*等の標準属性を失わないよう明示的にmergeする。
const misskeyDefaultResource = deps.resourceFromAttributes({
[deps.serviceNameAttribute]: 'misskey-backend',
[deps.serviceInstanceIdAttribute]: `${os.hostname()}:${process.pid}`,
[deps.serviceVersionAttribute]: deps.serviceVersion,
'misskey.process.role': getMisskeyProcessRole(),
});
// OTel標準の OTEL_SERVICE_NAME / OTEL_RESOURCE_ATTRIBUTES を尊重する。
// mergeは右辺が優先されるため、config.resourceAttributesを最優先にする。
return deps.defaultResource()
.merge(misskeyDefaultResource)
.merge(deps.detectResources({ detectors: [deps.envDetector] }))
.merge(deps.resourceFromAttributes(config.resourceAttributes ?? {}));
}
export function createSampler(sampleRate: number, deps: CreateSamplerDeps): ParentBasedSampler {
// 設定ミスを無言でAlwaysOn/AlwaysOffに倒さず、起動時に明確に失敗させる。
// (YAMLでクォートされた数値文字列などnumber型の保証が無い値が来てもここで弾く)
if (typeof sampleRate !== 'number' || !Number.isFinite(sampleRate) || sampleRate < 0 || sampleRate > 1) {
throw new Error('otelForBackend.sampleRate must be a number between 0.0 and 1.0.');
}
return new deps.ParentBasedSampler({
root: new deps.TraceIdRatioBasedSampler(sampleRate),
});
}
export function getMisskeyProcessRole(): string {
// Trace backend上でserver/queue/workerを見分けられるよう、Misskey固有の役割をresourceに載せる。
if (envOption.disableClustering) {
if (envOption.onlyServer) return 'primary-server';
if (envOption.onlyQueue) return 'primary-queue';
return 'primary-server+queue';
}
if (cluster.isPrimary) {
if (envOption.onlyServer) return 'fork-only';
if (envOption.onlyQueue) return 'primary-queue';
return 'primary-server';
}
// worker.tsのworkerMainに合わせる: onlyServerならserver()、それ以外はjobQueue()を実行する。
if (envOption.onlyServer) return 'worker-server';
return 'worker-queue';
}
@@ -3,18 +3,13 @@
* SPDX-License-Identifier: AGPL-3.0-only
*/
import Logger from '@/logger.js';
import { registerDiagLogger } from '@/core/telemetry/telemetry-diag.js';
import { getQueueTraceContextMode, injectActiveTraceContext, startSpanWithQueueTraceContext } from '@/core/telemetry/queue-trace-context.js';
import type { LogTraceContext } from '@/logging/types.js';
import type * as SentryNode from '@sentry/node';
import type { NodeOptions } from '@sentry/node';
import type { OtelBackendRuntimeConfig, SentryBackendConfig, TelemetryAdapter, TelemetryCaptureMessageOptions } from './TelemetryAdapter.js';
import type { QueueTraceContextCarrier, QueueTraceContextDeps } from '../queue-trace-context.js';
import type { SentryBackendConfig, TelemetryAdapter, TelemetryCaptureMessageOptions } from './TelemetryAdapter.js';
// OpenTelemetryAdapterのDEFAULT_SHUTDOWN_TIMEOUTと揃え、Sentryのtransportが詰まってもプロセス終了を妨げないようにする。
// Sentryのtransportが詰まってもプロセス終了を妨げないようにする。
const DEFAULT_SHUTDOWN_TIMEOUT = 5000;
const logger = new Logger('telemetry', 'green');
type SentryIntegrationsOption = NonNullable<NodeOptions['integrations']>;
// eslint-disable-next-line @typescript-eslint/no-explicit-any
@@ -73,50 +68,9 @@ export function buildSentryNodeOptions(
};
}
type BuildSentryOtlpInitOptions = {
sentryConfig: SentryBackendConfig;
otelConfig: OtelBackendRuntimeConfig;
otlpProcessor: unknown;
nodeProfilingIntegration?: () => SentryIntegration;
warn?: (message: string) => void;
};
export function buildSentryOtlpInitOptions(options: BuildSentryOtlpInitOptions): SentryNodeOptions {
// OTel併存時も、remoteへtrace headerを漏らさないデフォルトはSentry単体時と揃える。
// propagateTraceToRemote: true か、options.tracePropagationTargets の明示指定がある場合のみ既定を上書きする。
const { tracePropagationTargets, ...sentryOptions } = options.sentryConfig.options;
const propagateTraceToRemote = options.otelConfig.propagateTraceToRemote === true || tracePropagationTargets != null;
const warn = options.warn ?? ((message: string) => logger.warn(message));
if (options.otelConfig.sampleRate != null) {
warn('otelForBackend.sampleRate is ignored when sentryForBackend is also configured; configure sentryForBackend.options.tracesSampleRate or tracesSampler instead.');
}
if (options.otelConfig.resourceAttributes != null) {
warn('otelForBackend.resourceAttributes is ignored when sentryForBackend is also configured; configure OTEL_RESOURCE_ATTRIBUTES instead.');
}
return {
...buildSentryNodeOptions({
...options.sentryConfig,
options: {
...sentryOptions,
...(propagateTraceToRemote ? { tracePropagationTargets } : {}),
},
}, options.nodeProfilingIntegration),
// Sentryの単一TracerProviderにOTLP processorを追加し、親欠損や二重providerを避ける。
openTelemetrySpanProcessors: [
...(options.sentryConfig.options.openTelemetrySpanProcessors ?? []),
options.otlpProcessor as NonNullable<SentryNodeOptions['openTelemetrySpanProcessors']>[number],
],
};
}
export class SentryTelemetryAdapter implements TelemetryAdapter {
private constructor(
private readonly Sentry: typeof SentryNode,
private readonly queueTraceContext?: QueueTraceContextDeps,
) {
}
@@ -129,45 +83,6 @@ export class SentryTelemetryAdapter implements TelemetryAdapter {
return new SentryTelemetryAdapter(Sentry);
}
public static async createWithOtlpExport(
sentryConfig: SentryBackendConfig,
otelConfig: OtelBackendRuntimeConfig,
): Promise<SentryTelemetryAdapter> {
const Sentry = await import('@sentry/node');
const { nodeProfilingIntegration } = await import('@sentry/profiling-node');
const { context, diag, DiagLogLevel, propagation, ROOT_CONTEXT, SpanStatusCode, trace } = await import('@opentelemetry/api');
const { BatchSpanProcessor } = await import('@opentelemetry/sdk-trace-base');
const { OTLPTraceExporter } = await import('@opentelemetry/exporter-trace-otlp-proto');
registerDiagLogger(diag, DiagLogLevel.WARN);
// OTLP送信だけを担うprocessorを作り、provider生成はSentry.init側に任せる。
const otlpProcessor = new BatchSpanProcessor(new OTLPTraceExporter({
...(otelConfig.endpoint != null ? { url: otelConfig.endpoint } : {}),
...(otelConfig.headers != null ? { headers: otelConfig.headers } : {}),
}));
// SentryとOTLPを同一providerに集約することで、どちらの宛先にも同じspan実体を流す。
Sentry.init(buildSentryOtlpInitOptions({
sentryConfig,
otelConfig,
otlpProcessor,
nodeProfilingIntegration,
}));
// Sentry が初期化した同じ OTel provider から tracer/context API を受け取り、
// Queue を跨ぐ context 伝播も Sentry と OTLP の両方へ同一 span として出力する。
return new SentryTelemetryAdapter(Sentry, {
tracer: trace.getTracer('misskey-backend'),
propagation,
trace,
getActiveContext: () => context.active(),
rootContext: ROOT_CONTEXT,
mode: getQueueTraceContextMode(otelConfig.jobTraceContextMode),
spanStatusCodeError: SpanStatusCode.ERROR,
});
}
public captureMessage(message: string, opts: TelemetryCaptureMessageOptions): void {
this.Sentry.captureMessage(message, {
level: opts.level,
@@ -189,19 +104,6 @@ export class SentryTelemetryAdapter implements TelemetryAdapter {
return this.Sentry.startSpan({ name }, fn);
}
public injectTraceContext(carrier: QueueTraceContextCarrier): void {
// Sentry 単体構成では queueTraceContext を持たず、従来どおりジョブデータを変更しない。
if (this.queueTraceContext == null) return;
injectActiveTraceContext(this.queueTraceContext, carrier);
}
public startSpanWithTraceContext<T>(name: string, jobData: object, fn: () => T): T {
// Sentry 単体構成では Sentry 既存の span 作成経路を使う。
if (this.queueTraceContext == null) return this.startSpan(name, fn);
return startSpanWithQueueTraceContext(this.queueTraceContext, name, jobData, fn, () => this.startSpan(name, fn));
}
public async shutdown(): Promise<void> {
// timeout未指定だとtransportのflushが詰まった際にプロセス終了を妨げるため、上限時間を設ける。
await this.Sentry.close(DEFAULT_SHUTDOWN_TIMEOUT);
@@ -5,19 +5,14 @@
import type { Config } from '@/config.js';
import type { LogTraceContext } from '@/logging/types.js';
import type { QueueTraceContextCarrier } from '../queue-trace-context.js';
export type SentryBackendConfig = NonNullable<Config['sentryForBackend']>;
export type OtelBackendConfig = NonNullable<Config['otelForBackend']>;
export type OtelBackendRuntimeConfig = OtelBackendConfig & {
serviceVersion: string;
};
export interface TelemetryCaptureMessageOptions {
/** 現在はエラー通知用途だけに絞る。追加する場合は各adapterでの扱いを揃えること。 */
level: 'error';
/** Sentryではuser.idへ渡す。OTel adapterは現在span属性へ付与していないため、必要ならadapter側で拡張する。 */
/** Sentryではuser.idへ渡す補助情報です。 */
userId?: string;
/** queue名やendpoint名など、通知先で調査に使う補助情報。 */
@@ -25,14 +20,14 @@ export interface TelemetryCaptureMessageOptions {
}
/**
* Sentry・OpenTelemetryなど、エラートラッキング/APMサービスごとの実装差異を隠蔽するための抽象。
* エラートラッキング/APMサービスごとの実装差異を隠蔽するための抽象。
* 新しいサービスを追加する場合はこのインターフェースを実装するアダプタをこのディレクトリに追加し、
* telemetry-registry.tsのinitTelemetry内で登録する。
*/
export interface TelemetryAdapter {
/**
* 実行中の処理で起きたエラー相当の事象を記録する。
* Sentryはmessage通知、OTelはactive spanまたは短命spanへの例外記録として扱う。
* Sentryはmessage通知として扱う。
*/
captureMessage(message: string, opts: TelemetryCaptureMessageOptions): void;
@@ -45,18 +40,6 @@ export interface TelemetryAdapter {
*/
startSpan<T>(name: string, fn: () => T): T;
/**
* BullMQ のジョブデータへ保存する carrier に、active trace context を注入する。
* OTel を使わない adapter は実装しない。
*/
injectTraceContext?(carrier: QueueTraceContextCarrier): void;
/**
* ジョブに保存された enqueue 元の context を、worker span の Link または parent として復元する。
* context を持たないジョブの互換性は adapter 側で保つ。
*/
startSpanWithTraceContext?<T>(name: string, jobData: object, fn: () => T): T;
/**
* プロセス終了時にtelemetry backendへ残りのデータをflushする。
* 実装側ではtransport停止に引きずられないよう、待機時間に上限を設ける。
@@ -1,86 +0,0 @@
/*
* SPDX-FileCopyrightText: syuilo and misskey-project
* SPDX-License-Identifier: AGPL-3.0-only
*/
import type { Span, TracerProvider } from '@opentelemetry/api';
type Instrumentation = {
setTracerProvider(provider: TracerProvider): void;
enable(): void;
disable(): void;
};
type PgInstrumentationConfig = {
enhancedDatabaseReporting: boolean;
requireParentSpan: boolean;
ignoreConnectSpans: boolean;
requestHook?(span: Span): void;
};
type InstrumentationConstructor = new (config: PgInstrumentationConfig) => Instrumentation;
type InstrumentationDeps = {
PgInstrumentation: InstrumentationConstructor;
};
type DatabaseInstrumentationOptions = {
capturePgStatement: boolean;
capturePgConnectionSpans: boolean;
};
type InstallDatabaseInstrumentationOptions = DatabaseInstrumentationOptions & {
capturePgSpans: boolean;
};
/**
* pg はアプリケーションが import する前に有効化しないと、require hook 型の
* 自動計装がモジュールを patch できない。そのため、
* telemetry provider の登録直後にこの関数を呼び出す。
*/
export async function installDatabaseInstrumentation(provider: TracerProvider, options: InstallDatabaseInstrumentationOptions): Promise<() => void> {
if (!options.capturePgSpans) return () => {};
const { PgInstrumentation } = await import('@opentelemetry/instrumentation-pg');
return installInstrumentation(provider, { PgInstrumentation }, options);
}
export function installInstrumentation(provider: TracerProvider, deps: InstrumentationDeps, options: DatabaseInstrumentationOptions = {
capturePgStatement: false,
capturePgConnectionSpans: false,
}): () => void {
const instrumentations = [
new deps.PgInstrumentation({
// SQLパラメータには投稿内容・認証情報などが含まれ得るため、常に記録しない。
enhancedDatabaseReporting: false,
requireParentSpan: true,
ignoreConnectSpans: !options.capturePgConnectionSpans,
// instrumentation-pgはSQL本文を無加工で属性へ追加する。明示opt-in時だけ残す。
...(options.capturePgStatement ? {} : {
requestHook: (span: Span) => {
span.setAttribute('db.statement', '[REDACTED]');
span.setAttribute('db.query.text', '[REDACTED]');
},
}),
}),
];
try {
for (const instrumentation of instrumentations) {
instrumentation.setTracerProvider(provider);
instrumentation.enable();
}
} catch (error) {
for (const instrumentation of instrumentations) {
instrumentation.disable();
}
throw error;
}
return () => {
for (const instrumentation of instrumentations) {
instrumentation.disable();
}
};
}
@@ -1,150 +0,0 @@
/*
* SPDX-FileCopyrightText: syuilo and misskey-project
* SPDX-License-Identifier: AGPL-3.0-only
*/
import { channel } from 'node:diagnostics_channel';
import { diag } from '@opentelemetry/api';
import type { ClientRequest, IncomingMessage } from 'node:http';
import type { Span, SpanOptions, SpanStatusCode, Tracer } from '@opentelemetry/api';
const HTTP_CLIENT_REQUEST_CREATED = 'http.client.request.created';
const HTTP_CLIENT_RESPONSE_FINISH = 'http.client.response.finish';
const HTTP_CLIENT_REQUEST_ERROR = 'http.client.request.error';
type HttpClientSpan = Pick<Span, 'end' | 'recordException' | 'setAttribute' | 'setStatus'>;
type HttpClientInstrumentationDeps = {
tracer: Pick<Tracer, 'startSpan'>;
spanKindClient: SpanOptions['kind'];
spanStatusCodeError: SpanStatusCode;
subscribe: (name: string, listener: (message: unknown) => void) => () => void;
reportError?: (error: unknown) => void;
};
type RequestCreatedMessage = { request: ClientRequest };
type ResponseFinishMessage = { request: ClientRequest; response: IncomingMessage };
type RequestErrorMessage = { request: ClientRequest; error: Error };
/**
* require フックを使わず、Node.js 組み込み HTTP クライアントの diagnostics channel を計装する。
* telemetry 初期化前に読み込まれたモジュールも対象になる。
*/
export function createHttpClientInstrumentation(deps: HttpClientInstrumentationDeps): () => void {
const spans = new WeakMap<ClientRequest, HttpClientSpan>();
const unsubscribeCreated = deps.subscribe(HTTP_CLIENT_REQUEST_CREATED, guardDiagnosticsListener(deps, (message: unknown) => {
const { request } = message as RequestCreatedMessage;
const { url, host, port } = getRequestDetails(request);
const method = request.method;
const span = deps.tracer.startSpan(method, {
kind: deps.spanKindClient,
attributes: {
'http.request.method': method,
'url.full': url,
'server.address': host,
'server.port': port,
},
});
spans.set(request, span);
}));
const unsubscribeResponseFinish = deps.subscribe(HTTP_CLIENT_RESPONSE_FINISH, guardDiagnosticsListener(deps, (message: unknown) => {
const { request, response } = message as ResponseFinishMessage;
const span = spans.get(request);
if (span == null) return;
const statusCode = response.statusCode;
if (statusCode != null) {
span.setAttribute('http.response.status_code', statusCode);
}
span.setAttribute('network.protocol.version', response.httpVersion);
if (statusCode != null && statusCode >= 400) {
span.setAttribute('error.type', String(statusCode));
span.setStatus({ code: deps.spanStatusCodeError });
}
span.end();
spans.delete(request);
}));
const unsubscribeRequestError = deps.subscribe(HTTP_CLIENT_REQUEST_ERROR, guardDiagnosticsListener(deps, (message: unknown) => {
const { request, error } = message as RequestErrorMessage;
const span = spans.get(request);
if (span == null) return;
span.recordException(error);
span.setAttribute('error.type', getErrorType(error));
span.setStatus({ code: deps.spanStatusCodeError });
span.end();
spans.delete(request);
}));
return () => {
unsubscribeCreated();
unsubscribeResponseFinish();
unsubscribeRequestError();
};
}
export function installHttpClientInstrumentation(deps: Omit<HttpClientInstrumentationDeps, 'subscribe'>): () => void {
return createHttpClientInstrumentation({
...deps,
reportError: deps.reportError ?? (error => diag.error('Failed to process HTTP client diagnostics event.', error)),
subscribe: (name, listener) => {
const diagnosticChannel = channel(name);
diagnosticChannel.subscribe(listener);
return () => diagnosticChannel.unsubscribe(listener);
},
});
}
function getRequestDetails(request: ClientRequest): { url: string; host: string; port: number } {
let protocol: 'http:' | 'https:' = 'http:';
try {
protocol = request.protocol === 'https:' ? 'https:' : 'http:';
const host = request.getHeader('host')?.toString() ?? request.host;
const url = new URL(request.path || '/', `${protocol}//${host}`);
// URL 属性には認証情報やクエリ文字列を含めない。
url.username = '';
url.password = '';
url.search = '';
url.hash = '';
return {
url: url.toString(),
host: url.hostname,
// URL.port は既定ポートでは空文字列になるため、スキームから補う。
port: url.port === '' ? (url.protocol === 'https:' ? 443 : 80) : Number(url.port),
};
} catch {
return {
url: `${protocol}//localhost/`,
host: 'localhost',
port: protocol === 'https:' ? 443 : 80,
};
}
}
function guardDiagnosticsListener(
deps: Pick<HttpClientInstrumentationDeps, 'reportError'>,
listener: (message: unknown) => void,
): (message: unknown) => void {
return (message) => {
try {
listener(message);
} catch (error) {
try {
deps.reportError?.(error);
} catch {
// diagnostics listener の失敗をアプリケーションへ伝播させない。
}
}
};
}
function getErrorType(error: Error): string {
// Node.js の system error code は安定した低カーディナリティの識別子になる。
const code = (error as NodeJS.ErrnoException).code;
return code ?? error.name;
}
@@ -1,33 +0,0 @@
/*
* SPDX-FileCopyrightText: syuilo and misskey-project
* SPDX-License-Identifier: AGPL-3.0-only
*/
import { injectTraceContext } from './telemetry-registry.js';
import { injectQueueTraceContext } from './queue-trace-context.js';
import type * as Bull from 'bullmq';
/**
* Queue の add/addBulk をラップし、全ての BullMQ enqueue 経路を一箇所で捕捉する。
* QueueService を通さず直接 add/addBulk する呼び出し元もあるため、それぞれで注入すると漏れやすい。
*/
export function instrumentQueue<T extends object>(queue: Bull.Queue<T>): Bull.Queue<T> {
// BullMQ のメソッドは Queue インスタンスを this として使うため、差し替え前に bind して保持する。
const add = queue.add.bind(queue);
queue.add = ((name, data, opts) => {
// BullMQ が data を Redis 用にシリアライズする前に、enqueue 元の context を内部フィールドへ追加する。
injectQueueTraceContext(data, injectTraceContext);
return add(name, data, opts);
}) as typeof queue.add;
// addBulk は複数ジョブを一度にシリアライズするので、各 data へ同じ context を注入する。
const addBulk = queue.addBulk.bind(queue);
queue.addBulk = ((jobs) => {
for (const job of jobs) {
injectQueueTraceContext(job.data, injectTraceContext);
}
return addBulk(jobs);
}) as typeof queue.addBulk;
return queue;
}
@@ -1,182 +0,0 @@
/*
* SPDX-FileCopyrightText: syuilo and misskey-project
* SPDX-License-Identifier: AGPL-3.0-only
*/
import type { Context, PropagationAPI, Span, SpanContext, SpanOptions, SpanStatusCode, Tracer } from '@opentelemetry/api';
/*
* enqueue (push) から worker による取得・処理 (pop) までをテレメトリ上で関連付け、
* Queue を挟む非同期処理を一連の流れとして追跡できるようにする。
* そのため、trace context を次の流れで伝播する:
*
* 1. producer 側で active context を W3C Trace Context 形式の carrier へ注入する。
* 2. carrier をジョブデータの内部フィールドに保存し、BullMQ/Redis 経由で worker へ渡す。
* 3. worker 側で carrier を context へ戻し、設定に応じて Link または parent に使う。
*
* OTel API への依存を引数にしているのは、OTel 単体構成と Sentry + OTLP 構成で
* 同じ伝播処理を使うため。
*/
/** Redis 上のジョブデータにだけ保存する、OpenTelemetry propagator 用の carrier。 */
export type QueueTraceContextCarrier = Record<string, string>;
export type QueueTraceContextMode = 'link' | 'parent';
// 通常のジョブプロセッサが参照しない Misskey 内部フィールドとして、ユーザー定義の data と区別する。
const QUEUE_TRACE_CONTEXT_KEY = '__misskeyTraceContext';
export type QueueTraceContextDeps = {
tracer: Pick<Tracer, 'startActiveSpan'>;
propagation: Pick<PropagationAPI, 'extract' | 'inject'>;
trace: {
getSpanContext(context: Context): SpanContext | undefined;
};
getActiveContext: () => Context;
rootContext: Context;
mode: QueueTraceContextMode;
spanStatusCodeError: SpanStatusCode;
};
export type QueueSpanContext = {
options: SpanOptions;
parentContext: Context;
};
type QueueSpanContextDeps = Pick<QueueTraceContextDeps, 'propagation' | 'trace' | 'rootContext' | 'mode'>;
/**
* enqueue 元の active context を、ジョブ本来のデータを壊さない内部フィールドとして保持する。
* propagator が何も注入しなかった場合は、Redis に不要な空オブジェクトを残さない。
*/
export function injectQueueTraceContext(data: unknown, inject: (carrier: QueueTraceContextCarrier) => void): void {
// BullMQ の型定義外から渡される値も考慮し、書き込めない値は無視する。
if (data == null || typeof data !== 'object') return;
const carrier: QueueTraceContextCarrier = {};
inject(carrier);
// active span が無い場合、propagator は何も注入しない。空の内部フィールドは Redis に保存しない。
if (Object.keys(carrier).length === 0) return;
Object.assign(data, { [QUEUE_TRACE_CONTEXT_KEY]: carrier });
}
/** 現在実行中の span を propagator の標準形式で carrier へ書き出す。 */
export function injectActiveTraceContext(deps: QueueTraceContextDeps, carrier: QueueTraceContextCarrier): void {
deps.propagation.inject(deps.getActiveContext(), carrier);
}
/**
* ジョブに保存された carrier から、worker span の開始に必要な context と options を組み立てる。
*
* - parent: enqueue span の子として同じ trace を継続する。
* - link: worker span を別の root trace にし、enqueue span への関連だけを Link に残す。
*
* link モードで trace を分けることで、worker 側の sampling 判定を enqueue 側から独立させられる。
*/
export function getQueueSpanContext(data: unknown, deps: QueueSpanContextDeps): QueueSpanContext | undefined {
const carrier = getQueueTraceContextCarrier(data);
if (carrier == null) return undefined;
// Redis から復元した carrier は別プロセス由来なので、現在の context ではなく ROOT_CONTEXT から展開する。
const extractedContext = deps.propagation.extract(deps.rootContext, carrier);
if (deps.mode === 'parent') {
return {
options: {},
parentContext: extractedContext,
};
}
// Link が受け取るのは Context ではなく SpanContext なので、extract 後に取り出す。
const spanContext = deps.trace.getSpanContext(extractedContext);
return {
options: {
root: true,
...(spanContext != null ? { links: [{ context: spanContext }] } : {}),
},
parentContext: deps.rootContext,
};
}
/**
* context を持つジョブは Link/parent の規則で span を開始する。
* デプロイ前に enqueue されたジョブなど、context を持たない場合は既存の adapter 固有実装へ委ねる。
*/
export function startSpanWithQueueTraceContext<T>(
deps: QueueTraceContextDeps,
name: string,
jobData: object,
fn: () => T,
fallback: () => T,
): T {
const spanContext = getQueueSpanContext(jobData, deps);
if (spanContext == null) return fallback();
return deps.tracer.startActiveSpan(name, spanContext.options, spanContext.parentContext, span => executeSpan(span, fn, deps.spanStatusCodeError));
}
/**
* 既存の TelemetryAdapter 契約に合わせ、同期値と Promise のどちらを返す処理も span で包む。
* 成功・失敗のどちらでも処理完了まで span を開き、必ず一度だけ閉じる。
*/
export function executeSpan<T>(span: Span, fn: () => T, spanStatusCodeError: SpanStatusCode): T {
try {
const result = fn();
if (isPromiseLike(result)) {
// fn() の戻り値は T のまま保ちつつ、Promise の settle 時に span を閉じる。
return result.then(
value => {
span.end();
return value;
},
error => {
recordSpanError(span, error, spanStatusCodeError);
span.end();
throw error;
},
) as T;
}
span.end();
return result;
} catch (error) {
// fn() が同期的に throw した場合も、Promise の reject と同じ形で記録して呼び出し元へ戻す。
recordSpanError(span, error, spanStatusCodeError);
span.end();
throw error;
}
}
/** Error 以外の throw 値も OTel exporter が扱える例外に正規化し、span をエラー状態にする。 */
export function recordSpanError(span: Span, error: unknown, spanStatusCodeError: SpanStatusCode): void {
const exception = error instanceof Error ? error : new Error(String(error));
span.recordException(exception);
span.setStatus({
code: spanStatusCodeError,
message: exception.message,
});
}
/**
* 未設定時は、enqueue と worker の sampling を分離できる link を使う。
* 設定ミスを無言でフォールバックさせないため、未知の値は起動時にエラーにする。
*/
export function getQueueTraceContextMode(mode: unknown): QueueTraceContextMode {
if (mode == null || mode === 'link') return 'link';
if (mode === 'parent') return 'parent';
throw new Error('otelForBackend.jobTraceContextMode must be either \'link\' or \'parent\'.');
}
function getQueueTraceContextCarrier(data: unknown): QueueTraceContextCarrier | undefined {
if (data == null || typeof data !== 'object') return undefined;
const carrier = (data as Record<string, unknown>)[QUEUE_TRACE_CONTEXT_KEY];
if (carrier == null || typeof carrier !== 'object' || Array.isArray(carrier)) return undefined;
// Redis 上の旧いジョブや壊れたデータで worker を落とさないよう、propagator に渡せる string map だけを受け入れる。
const entries = Object.entries(carrier);
if (entries.length === 0 || entries.some(([, value]) => typeof value !== 'string')) return undefined;
return Object.fromEntries(entries) as QueueTraceContextCarrier;
}
function isPromiseLike<T>(value: T): value is T & PromiseLike<Awaited<T>> {
// native Promise に限らず thenable も span の完了を待てるよう、instanceof ではなく then の有無で判定する。
return value != null && typeof (value as { then?: unknown }).then === 'function';
}
@@ -1,197 +0,0 @@
/*
* SPDX-FileCopyrightText: syuilo and misskey-project
* SPDX-License-Identifier: AGPL-3.0-only
*/
import { tracingChannel } from 'node:diagnostics_channel';
import { context, diag, trace } from '@opentelemetry/api';
import type { Span, SpanKind, SpanStatusCode, Tracer } from '@opentelemetry/api';
type IORedisCommandContext = {
command: string;
args: string[];
database: number | undefined;
serverAddress: string;
serverPort: number | undefined;
};
type IORedisConnectContext = {
serverAddress: string;
serverPort: number | undefined;
};
type TracingChannelSubscribers<T extends object> = {
start(message: T): void;
end(message: T & { error?: unknown }): void;
asyncStart(message: T & { error?: unknown }): void;
asyncEnd(message: T & { error?: unknown }): void;
error(message: T & { error: unknown }): void;
};
type TracingChannel<T extends object> = {
subscribe(subscribers: TracingChannelSubscribers<T>): void;
unsubscribe(subscribers: TracingChannelSubscribers<T>): void;
};
type RedisInstrumentationDeps = {
tracingChannel<T extends object>(name: string): TracingChannel<T>;
tracer: Pick<Tracer, 'startSpan'>;
getActiveSpan(): Span | undefined;
spanKindClient: SpanKind;
spanStatusCodeError: SpanStatusCode;
reportError?: (error: unknown) => void;
};
type RedisInstrumentationOptions = {
captureCommandSpans?: boolean;
/** requireParentSpan is the official ioredis instrumentation's default. */
requireParentSpan?: boolean;
captureConnectionSpans?: boolean;
};
/**
* ioredis 5.11以降が公開する native diagnostics channel を購読する。
* バンドル後もioredis自身が発行するイベントを使うため、require hook に
* 依存せず、ESM import 経路でもコマンド span を取得できる。
*/
export function installRedisInstrumentation(
tracer: Pick<Tracer, 'startSpan'>,
spanKindClient: SpanKind,
spanStatusCodeError: SpanStatusCode,
options: RedisInstrumentationOptions = {},
): () => void {
return createRedisInstrumentation({
tracingChannel,
tracer,
getActiveSpan: () => trace.getSpan(context.active()),
spanKindClient,
spanStatusCodeError,
reportError: error => diag.error('Failed to process Redis diagnostics event.', error),
}, options);
}
export function createRedisInstrumentation(deps: RedisInstrumentationDeps, options: RedisInstrumentationOptions = {}): () => void {
const requireParentSpan = options.requireParentSpan ?? true;
const cleanup: Array<() => void> = [];
if (options.captureCommandSpans === true) {
const commandChannel = deps.tracingChannel<IORedisCommandContext>('ioredis:command');
const commandSubscribers = createTracingChannelSubscribers(commandChannel, deps, requireParentSpan, message => ({
name: message.command,
attributes: {
'db.system.name': 'redis',
'db.operation.name': message.command,
'server.address': message.serverAddress,
...(message.database != null ? { 'db.namespace': message.database.toString(10) } : {}),
...(message.serverPort != null ? { 'server.port': message.serverPort } : {}),
},
}));
cleanup.push(() => commandChannel.unsubscribe(commandSubscribers));
}
if (options.captureConnectionSpans === true) {
const connectChannel = deps.tracingChannel<IORedisConnectContext>('ioredis:connect');
// Connection spans are explicitly opt-in and should include startup and reconnect attempts.
const connectSubscribers = createTracingChannelSubscribers(connectChannel, deps, false, message => ({
name: 'connect',
attributes: {
'db.system.name': 'redis',
'db.operation.name': 'connect',
'server.address': message.serverAddress,
...(message.serverPort != null ? { 'server.port': message.serverPort } : {}),
},
}));
cleanup.push(() => connectChannel.unsubscribe(connectSubscribers));
}
return () => {
for (const unsubscribe of cleanup) {
unsubscribe();
}
};
}
function createTracingChannelSubscribers<T extends object>(
channel: TracingChannel<T>,
deps: RedisInstrumentationDeps,
requireParentSpan: boolean,
getSpanOptions: (message: T) => { name: string; attributes: Record<string, string | number> },
): TracingChannelSubscribers<T> {
const spans = new WeakMap<object, { span: Span; recordedError: boolean }>();
const recordError = (state: { span: Span; recordedError: boolean }, value: unknown): void => {
if (state.recordedError) return;
const error = toError(value);
state.span.recordException(error);
state.span.setStatus({ code: deps.spanStatusCodeError, message: error.message });
state.span.setAttribute('error.type', getErrorType(error));
const statusCode = getRedisErrorStatusCode(error.message);
if (statusCode != null) state.span.setAttribute('db.response.status_code', statusCode);
state.recordedError = true;
};
const finish = (message: T & { error?: unknown }): void => {
const state = spans.get(message);
if (state == null) return;
if (message.error != null) {
recordError(state, message.error);
}
state.span.end();
spans.delete(message);
};
const subscribers: TracingChannelSubscribers<T> = {
start: guardTracingSubscriber(deps, (message: T) => {
if (requireParentSpan && deps.getActiveSpan() == null) return;
const options = getSpanOptions(message);
const span = deps.tracer.startSpan(options.name, {
kind: deps.spanKindClient,
attributes: options.attributes,
});
spans.set(message, { span, recordedError: false });
}),
// Promiseを返すコマンドでは完了前にもendイベントが来るため、同期例外だけここで閉じる。
end: guardTracingSubscriber(deps, (message: T & { error?: unknown }) => {
if (message.error != null) finish(message);
}),
asyncStart: () => {},
asyncEnd: guardTracingSubscriber(deps, finish),
error: guardTracingSubscriber(deps, (message: T & { error: unknown }) => {
const state = spans.get(message);
if (state == null) return;
recordError(state, message.error);
}),
};
channel.subscribe(subscribers);
return subscribers;
}
function guardTracingSubscriber<T>(
deps: Pick<RedisInstrumentationDeps, 'reportError'>,
subscriber: (message: T) => void,
): (message: T) => void {
return (message) => {
try {
subscriber(message);
} catch (error) {
try {
deps.reportError?.(error);
} catch {
// diagnostics subscriber の失敗をアプリケーションへ伝播させない。
}
}
};
}
function toError(value: unknown): Error {
return value instanceof Error ? value : new Error(String(value));
}
function getRedisErrorStatusCode(message: string): string | undefined {
return message.match(/^([A-Z][A-Z0-9_]*)\b/)?.[1];
}
function getErrorType(error: Error): string {
return (error as NodeJS.ErrnoException).code ?? error.name;
}
@@ -1,40 +0,0 @@
/*
* SPDX-FileCopyrightText: syuilo and misskey-project
* SPDX-License-Identifier: AGPL-3.0-only
*/
import Logger from '@/logger.js';
import type { DiagAPI, DiagLogger, DiagLogLevel } from '@opentelemetry/api';
export function registerDiagLogger(
diagApi: DiagAPI,
diagLogLevelWarn: DiagLogLevel,
): void {
// diagはプロセスグローバルなので、通常運用で必要なWARN以上だけをMisskeyのログに流す。
const logger = new Logger('otel', 'green');
const diagLogger: DiagLogger = {
error: (message, ...args) => logger.error(formatDiagMessage(message, args)),
warn: (message, ...args) => logger.warn(formatDiagMessage(message, args)),
info: (message, ...args) => logger.info(formatDiagMessage(message, args)),
debug: (message, ...args) => logger.debug(formatDiagMessage(message, args)),
verbose: (message, ...args) => logger.debug(formatDiagMessage(message, args)),
};
diagApi.setLogger(diagLogger, {
logLevel: diagLogLevelWarn,
suppressOverrideMessage: true,
});
}
function formatDiagMessage(message: string, args: unknown[]): string {
if (args.length === 0) return message;
return `${message} ${args.map(arg => {
if (arg instanceof Error) return arg.stack ?? arg.message;
if (typeof arg === 'string') return arg;
try {
return JSON.stringify(arg);
} catch {
return String(arg);
}
}).join(' ')}`;
}
@@ -5,10 +5,8 @@
import type { Config } from '@/config.js';
import { setLogTraceContextProvider } from '@/logging/logging-runtime.js';
import { OpenTelemetryAdapter } from './adapters/OpenTelemetryAdapter.js';
import { SentryTelemetryAdapter } from './adapters/SentryTelemetryAdapter.js';
import type { OtelBackendRuntimeConfig, TelemetryAdapter, TelemetryCaptureMessageOptions } from './adapters/TelemetryAdapter.js';
import type { QueueTraceContextCarrier } from './queue-trace-context.js';
import type { TelemetryAdapter, TelemetryCaptureMessageOptions } from './adapters/TelemetryAdapter.js';
/**
* NestのDIコンテナが構築される前(boot処理内)で初期化する必要があるため、
@@ -18,24 +16,8 @@ import type { QueueTraceContextCarrier } from './queue-trace-context.js';
const adapters: TelemetryAdapter[] = [];
export async function initTelemetry(config: Config): Promise<void> {
const otelForBackend: OtelBackendRuntimeConfig | undefined = config.otelForBackend == null ? undefined : {
...config.otelForBackend,
serviceVersion: config.version,
};
// SentryとOTelを同時に使う場合はproviderを分けず、Sentry側へOTLP processorを追加する。
let adapter: TelemetryAdapter | undefined;
if (config.sentryForBackend && otelForBackend) {
adapter = await SentryTelemetryAdapter.createWithOtlpExport(config.sentryForBackend, otelForBackend);
} else if (config.sentryForBackend) {
// Sentry単体時は既存のSentry adapterだけを登録する。
adapter = await SentryTelemetryAdapter.create(config.sentryForBackend);
} else if (otelForBackend) {
// OTel単体時だけMisskey自身でNodeTracerProviderを立てる。
adapter = await OpenTelemetryAdapter.create(otelForBackend);
}
if (adapter != null) {
if (config.sentryForBackend) {
const adapter = await SentryTelemetryAdapter.create(config.sentryForBackend);
adapters.push(adapter);
// Telemetryの初期化後に登録し、初期化前のBootstrapログは従来どおり出力する。
setLogTraceContextProvider(() => adapter.getActiveTraceContext?.());
@@ -63,22 +45,6 @@ export function startSpan<T>(name: string, fn: () => T): T {
return wrapped();
}
export function injectTraceContext(carrier: QueueTraceContextCarrier): void {
// Queue の carrier は共有データなので、通知と異なり全 adapter にブロードキャストしない。
// OTel provider は現在 1 つだけなので、同じ header を上書きしないよう最初の対応 adapter だけを使う。
adapters.find(adapter => adapter.injectTraceContext != null)?.injectTraceContext?.(carrier);
}
export function startSpanWithTraceContext<T>(name: string, jobData: object, fn: () => T): T {
// Queue context を解釈できる adapter に span 作成を任せる。
// 対応 adapter が無い構成では通常の startSpan へフォールバックする。
const adapter = adapters.find(adapter => adapter.startSpanWithTraceContext != null);
if (adapter?.startSpanWithTraceContext != null) {
return adapter.startSpanWithTraceContext(name, jobData, fn);
}
return startSpan(name, fn);
}
export async function shutdownTelemetry(): Promise<void> {
// 終了時は登録済みadapterを並列にflush/shutdownする。
await Promise.allSettled(adapters.map(adapter => adapter.shutdown()));
@@ -1,11 +0,0 @@
/*
* SPDX-FileCopyrightText: syuilo and misskey-project
* SPDX-License-Identifier: AGPL-3.0-only
*/
/**
* Telemetry-specific shutdown is implemented by telemetry-registry.ts.
* Signal coordination belongs to boot/shutdown-handler.ts so telemetry and
* logging remain independent domains.
*/
export { shutdownTelemetry } from './telemetry-registry.js';
@@ -11,7 +11,7 @@ import type Logger from '@/logger.js';
import { bindThis } from '@/decorators.js';
import { TelemetryService } from '@/core/telemetry/TelemetryService.js';
import { CheckModeratorsActivityProcessorService } from '@/queue/processors/CheckModeratorsActivityProcessorService.js';
import { runQueueJobWithTraceContext } from './queue-job-runner.js';
import { runQueueJob } from './queue-job-runner.js';
import { UserWebhookDeliverProcessorService } from './processors/UserWebhookDeliverProcessorService.js';
import { SystemWebhookDeliverProcessorService } from './processors/SystemWebhookDeliverProcessorService.js';
import { EndedPollNotificationProcessorService } from './processors/EndedPollNotificationProcessorService.js';
@@ -159,8 +159,7 @@ export class QueueProcessorService implements OnApplicationShutdown {
};
}
// 以下の各 Worker は job.data に保存された enqueue 元の trace context を復元し、
// ジョブの実処理全体を Link または parent の worker span で囲む。
// 以下の各 Worker はジョブの実処理全体を worker span で囲む。
//#region system
{
const processer = (job: Bull.Job) => {
@@ -180,10 +179,9 @@ export class QueueProcessorService implements OnApplicationShutdown {
const logger = this.logger.createSubLogger('system');
this.systemQueueWorker = new Bull.Worker(QUEUE.SYSTEM, (job) => {
return runQueueJobWithTraceContext(
return runQueueJob(
this.telemetryService,
'Queue: System: ' + job.name,
job.data,
() => processer(job) as Promise<void>,
err => {
logger.error(`failed(${err.name}: ${err.message}) id=${job.id}`, { job: renderJob(job), e: renderError(err) });
@@ -235,10 +233,9 @@ export class QueueProcessorService implements OnApplicationShutdown {
const logger = this.logger.createSubLogger('db');
this.dbQueueWorker = new Bull.Worker(QUEUE.DB, (job) => {
return runQueueJobWithTraceContext(
return runQueueJob(
this.telemetryService,
'Queue: DB: ' + job.name,
job.data,
() => processer(job),
err => {
logger.error(`failed(${err.name}: ${err.message}) id=${job.id}`, { job: renderJob(job), e: renderError(err) });
@@ -266,10 +263,9 @@ export class QueueProcessorService implements OnApplicationShutdown {
const logger = this.logger.createSubLogger('deliver');
this.deliverQueueWorker = new Bull.Worker(QUEUE.DELIVER, (job) => {
return runQueueJobWithTraceContext(
return runQueueJob(
this.telemetryService,
'Queue: Deliver',
job.data,
() => this.deliverProcessorService.process(job),
err => {
logger.error(`failed(${err.name}: ${err.message}) ${getJobInfo(job)} to=${job.data.to}`, { e: renderError(err) });
@@ -305,10 +301,9 @@ export class QueueProcessorService implements OnApplicationShutdown {
const logger = this.logger.createSubLogger('inbox');
this.inboxQueueWorker = new Bull.Worker(QUEUE.INBOX, (job) => {
return runQueueJobWithTraceContext(
return runQueueJob(
this.telemetryService,
'Queue: Inbox',
job.data,
() => this.inboxProcessorService.process(job),
err => {
const activityId = job.data.activity ? job.data.activity.id : 'none';
@@ -345,10 +340,9 @@ export class QueueProcessorService implements OnApplicationShutdown {
const logger = this.logger.createSubLogger('user-webhook');
this.userWebhookDeliverQueueWorker = new Bull.Worker(QUEUE.USER_WEBHOOK_DELIVER, (job) => {
return runQueueJobWithTraceContext(
return runQueueJob(
this.telemetryService,
'Queue: UserWebhookDeliver',
job.data,
() => this.userWebhookDeliverProcessorService.process(job),
err => {
logger.error(`failed(${err.name}: ${err.message}) ${getJobInfo(job)} to=${job.data.to}`, { e: renderError(err) });
@@ -384,10 +378,9 @@ export class QueueProcessorService implements OnApplicationShutdown {
const logger = this.logger.createSubLogger('system-webhook');
this.systemWebhookDeliverQueueWorker = new Bull.Worker(QUEUE.SYSTEM_WEBHOOK_DELIVER, (job) => {
return runQueueJobWithTraceContext(
return runQueueJob(
this.telemetryService,
'Queue: SystemWebhookDeliver',
job.data,
() => this.systemWebhookDeliverProcessorService.process(job),
err => {
logger.error(`failed(${err.name}: ${err.message}) ${getJobInfo(job)} to=${job.data.to}`, { e: renderError(err) });
@@ -432,10 +425,9 @@ export class QueueProcessorService implements OnApplicationShutdown {
const logger = this.logger.createSubLogger('relationship');
this.relationshipQueueWorker = new Bull.Worker(QUEUE.RELATIONSHIP, (job) => {
return runQueueJobWithTraceContext(
return runQueueJob(
this.telemetryService,
'Queue: Relationship: ' + job.name,
job.data,
() => processer(job),
err => {
logger.error(`failed(${err.name}: ${err.message}) id=${job.id}`, { job: renderJob(job), e: renderError(err) });
@@ -475,10 +467,9 @@ export class QueueProcessorService implements OnApplicationShutdown {
const logger = this.logger.createSubLogger('objectStorage');
this.objectStorageQueueWorker = new Bull.Worker(QUEUE.OBJECT_STORAGE, (job) => {
return runQueueJobWithTraceContext(
return runQueueJob(
this.telemetryService,
'Queue: ObjectStorage: ' + job.name,
job.data,
() => processer(job) as Promise<void>,
err => {
logger.error(`failed(${err.name}: ${err.message}) id=${job.id}`, { job: renderJob(job), e: renderError(err) });
@@ -507,10 +498,9 @@ export class QueueProcessorService implements OnApplicationShutdown {
const logger = this.logger.createSubLogger('ended-poll-notification');
this.endedPollNotificationQueueWorker = new Bull.Worker(QUEUE.ENDED_POLL_NOTIFICATION, (job) => {
return runQueueJobWithTraceContext(
return runQueueJob(
this.telemetryService,
'Queue: EndedPollNotification',
job.data,
() => this.endedPollNotificationProcessorService.process(job),
err => {
logger.error(`failed(${err.name}: ${err.message}) id=${job.id}`, { job: renderJob(job), e: renderError(err) });
@@ -532,10 +522,9 @@ export class QueueProcessorService implements OnApplicationShutdown {
const logger = this.logger.createSubLogger('post-scheduled-note');
this.postScheduledNoteQueueWorker = new Bull.Worker(QUEUE.POST_SCHEDULED_NOTE, (job) => {
return runQueueJobWithTraceContext(
return runQueueJob(
this.telemetryService,
'Queue: PostScheduledNote',
job.data,
() => this.postScheduledNoteProcessorService.process(job),
err => {
logger.error(`failed(${err.name}: ${err.message}) id=${job.id}`, { job: renderJob(job), e: renderError(err) });
@@ -5,17 +5,16 @@
import type { TelemetryService } from '@/core/telemetry/TelemetryService.js';
type QueueTelemetryService = Pick<TelemetryService, 'startSpanWithTraceContext'>;
type QueueTelemetryService = Pick<TelemetryService, 'startSpan'>;
/** QueueのprocessorをTrace Context付きで実行し、失敗処理をSpan内で行います。 */
export function runQueueJobWithTraceContext<T>(
/** Queueのprocessorを実行し、失敗処理をSpan内で行います。 */
export function runQueueJob<T>(
telemetryService: QueueTelemetryService,
spanName: string,
jobData: object,
processJob: () => T | Promise<T>,
onError: (error: Error) => void,
): Promise<T> {
return telemetryService.startSpanWithTraceContext(spanName, jobData, async (): Promise<T> => {
return telemetryService.startSpan(spanName, async (): Promise<T> => {
try {
return await processJob();
} catch (error) {
@@ -32,7 +32,6 @@ import { HealthServerService } from './HealthServerService.js';
import { ClientServerService } from './web/ClientServerService.js';
import { OpenApiServerService } from './api/openapi/OpenApiServerService.js';
import { OAuth2ProviderService } from './oauth/OAuth2ProviderService.js';
import { registerHttpServerInstrumentation } from './http-server-instrumentation.js';
import { registerHttpAccessLog } from './http-access-log.js';
const _dirname = fileURLToPath(new URL('.', import.meta.url));
@@ -82,7 +81,6 @@ export class ServerService implements OnApplicationShutdown {
logger: false,
});
this.#fastify = fastify;
await registerHttpServerInstrumentation(fastify, this.config);
registerHttpAccessLog(fastify);
// HSTS
@@ -123,7 +123,7 @@ export function registerHttpAccessLog(fastify: FastifyInstance, manager: LogMana
const states = new WeakMap<object, AccessRequestState>();
// HTTP計装の後に登録し、リクエスト開始時のactiveなTrace Contextを保存します。
// リクエスト開始時のactiveなTrace Contextを保存します。
fastify.addHook('onRequest', (request, _reply, done) => {
states.set(request, {
traceContext: manager.getActiveTraceContext(),
@@ -1,35 +0,0 @@
/*
* SPDX-FileCopyrightText: syuilo and misskey-project
* SPDX-License-Identifier: AGPL-3.0-only
*/
import type { Config } from '@/config.js';
import type { FastifyInstance } from 'fastify';
type TelemetryConfig = Pick<Config, 'otelForBackend' | 'sentryForBackend'>;
export function shouldRegisterHttpServerInstrumentation(config: TelemetryConfig): boolean {
// Sentryもリクエストspanを作成するため、両方を登録すると重複して出力される。
return config.otelForBackend != null && config.sentryForBackend == null;
}
/**
* ActivityPubや
* well-knownを含む全HTTP受信経路を1つのroot spanとして計測する
*/
export async function registerHttpServerInstrumentation(fastify: FastifyInstance, config: TelemetryConfig): Promise<void> {
if (!shouldRegisterHttpServerInstrumentation(config)) return;
const { FastifyOtelInstrumentation } = await import('@fastify/otel');
const instrumentation = new FastifyOtelInstrumentation({
requestHook: (span, request) => {
const route = request.routeOptions.url;
if (route != null) {
// デフォルトだとトレース名が「request」で固定されてしまうため、判別がつかなくなる。
// ルート名をspan名に設定することで、トレースビューでルートごとの処理時間を確認できるようになる。
span.updateName(`${request.method} ${route}`);
}
},
});
await fastify.register(instrumentation.plugin());
}
@@ -4,3 +4,4 @@ volumes
docker.env
*.test.conf
*.test.default.yml
*.test.config.json
+1 -1
View File
@@ -28,7 +28,7 @@ function generate {
-days 500
if [ ! -f .config/docker.env ]; then cp .config/example.docker.env .config/docker.env; fi
if [ ! -f .config/$1.conf ]; then sed "s/\${HOST}/$1/g" .config/example.conf > .config/$1.conf; fi
if [ ! -f .config/$1.default.yml ]; then sed "s/\${HOST}/$1/g" .config/example.config.json > .config/$1.config.json; fi
if [ ! -f .config/$1.config.json ]; then sed "s/\${HOST}/$1/g" .config/example.config.json > .config/$1.config.json; fi
}
generate a.test
@@ -1,89 +0,0 @@
/*
* SPDX-FileCopyrightText: syuilo and misskey-project
* SPDX-License-Identifier: AGPL-3.0-only
*/
import { randomUUID } from 'node:crypto';
import { describe, expect, test } from 'vitest';
import { Queue, Worker } from 'bullmq';
import { SpanKind, SpanStatusCode } from '@opentelemetry/api';
import { InMemorySpanExporter, SimpleSpanProcessor } from '@opentelemetry/sdk-trace-base';
import { NodeTracerProvider } from '@opentelemetry/sdk-trace-node';
import { loadConfig } from '@/config.js';
import { installRedisInstrumentation } from '@/core/telemetry/redis-instrumentation.js';
const config = loadConfig();
describe('Redis telemetry instrumentation', () => {
test('records Redis spans below HTTP and BullMQ worker spans without Redis arguments', async () => {
const exporter = new InMemorySpanExporter();
const provider = new NodeTracerProvider({
spanProcessors: [new SimpleSpanProcessor(exporter)],
});
provider.register();
const tracer = provider.getTracer('telemetry-redis-instrumentation-test');
const uninstall = installRedisInstrumentation(tracer, SpanKind.CLIENT, SpanStatusCode.ERROR, {
captureCommandSpans: true,
});
const queueName = `telemetry-${randomUUID()}`;
const prefix = `telemetry-${randomUUID()}`;
const connection = {
host: config.redis.host,
port: config.redis.port,
...(config.redis.password != null ? { password: config.redis.password } : {}),
};
const queue = new Queue(queueName, { connection, prefix });
let worker: Worker | undefined;
let httpSpanId: string | undefined;
let jobSpanId: string | undefined;
const secret = `secret-${randomUUID()}`;
try {
const processed = new Promise<void>((resolve, reject) => {
worker = new Worker(queueName, async job => {
return await tracer.startActiveSpan('Queue: telemetry test', async jobSpan => {
jobSpanId = jobSpan.spanContext().spanId;
try {
// updateData uses BullMQ's worker-side ioredis client.
await job.updateData({ secret });
return 'ok';
} finally {
jobSpan.end();
}
});
}, { connection, prefix });
worker.once('completed', () => resolve());
worker.once('failed', (_job, error) => reject(error));
worker.on('error', reject);
});
await tracer.startActiveSpan('HTTP POST /telemetry-test', async httpSpan => {
httpSpanId = httpSpan.spanContext().spanId;
try {
// Queue#add uses BullMQ's producer-side ioredis client.
await queue.add('probe', { secret });
} finally {
httpSpan.end();
}
});
await processed;
await provider.forceFlush();
const redisSpans = exporter.getFinishedSpans().filter(span => span.attributes['db.system.name'] === 'redis');
expect(redisSpans.some(span => span.parentSpanContext?.spanId === httpSpanId)).toBe(true);
expect(redisSpans.some(span => span.parentSpanContext?.spanId === jobSpanId)).toBe(true);
for (const span of redisSpans) {
expect(span.attributes).not.toHaveProperty('db.statement');
expect(span.attributes).not.toHaveProperty('db.query.text');
expect(Object.values(span.attributes)).not.toContain(secret);
}
} finally {
await worker?.close(true);
await queue.obliterate({ force: true });
await queue.close();
uninstall();
await provider.shutdown();
}
}, 30000);
});
@@ -1,437 +0,0 @@
/*
* SPDX-FileCopyrightText: syuilo and misskey-project
* SPDX-License-Identifier: AGPL-3.0-only
*/
import { afterEach, beforeEach, describe, expect, test, vi } from 'vitest';
import { SpanStatusCode } from '@opentelemetry/api';
import { defaultResource, detectResources, envDetector, resourceFromAttributes } from '@opentelemetry/resources';
import { ParentBasedSampler, TraceIdRatioBasedSampler } from '@opentelemetry/sdk-trace-base';
import { ATTR_SERVICE_INSTANCE_ID, ATTR_SERVICE_NAME, ATTR_SERVICE_VERSION } from '@opentelemetry/semantic-conventions';
import type { Context, SpanContext } from '@opentelemetry/api';
import { OpenTelemetryAdapter, createResource, createSampler, getMisskeyProcessRole, installInstrumentationsWithCleanup } from '@/core/telemetry/adapters/OpenTelemetryAdapter.js';
const mocks = vi.hoisted(() => {
return {
envOption: {
disableClustering: false,
onlyServer: false,
onlyQueue: false,
},
isPrimary: false,
};
});
vi.mock('@/env.js', () => ({
envOption: mocks.envOption,
}));
vi.mock('node:cluster', () => ({
default: {
get isPrimary() {
return mocks.isPrimary;
},
},
}));
const samplerDeps = {
ParentBasedSampler,
TraceIdRatioBasedSampler,
};
describe('OpenTelemetryAdapter', () => {
afterEach(() => {
vi.useRealTimers();
});
test('wraps async work in an active span and ends it after success', async () => {
const span = {
end: vi.fn(),
recordException: vi.fn(),
setStatus: vi.fn(),
};
const tracer = {
startActiveSpan: vi.fn(async (_name: string, fn: (spanArg: any) => Promise<string>) => fn(span)),
} as any;
const provider = {
shutdown: vi.fn(),
};
const adapter = new OpenTelemetryAdapter({
tracer,
provider,
getActiveSpan: () => undefined,
spanStatusCodeError: SpanStatusCode.ERROR,
shutdownTimeout: 10,
});
await expect(adapter.startSpan('API: test', async () => 'ok')).resolves.toBe('ok');
expect(tracer.startActiveSpan).toHaveBeenCalledWith('API: test', expect.any(Function));
expect(span.recordException).not.toHaveBeenCalled();
expect(span.setStatus).not.toHaveBeenCalled();
expect(span.end).toHaveBeenCalledTimes(1);
});
test('records thrown errors on the active span before rethrowing', async () => {
const error = new Error('boom');
const span = {
end: vi.fn(),
recordException: vi.fn(),
setStatus: vi.fn(),
};
const tracer = {
startActiveSpan: vi.fn(async (_name: string, fn: (spanArg: any) => Promise<void>) => fn(span)),
} as any;
const adapter = new OpenTelemetryAdapter({
tracer,
provider: { shutdown: vi.fn() },
getActiveSpan: () => undefined,
spanStatusCodeError: SpanStatusCode.ERROR,
shutdownTimeout: 10,
});
await expect(adapter.startSpan('Queue: test', async () => {
throw error;
})).rejects.toThrow(error);
expect(span.recordException).toHaveBeenCalledWith(error);
expect(span.setStatus).toHaveBeenCalledWith({
code: SpanStatusCode.ERROR,
message: error.message,
});
expect(span.end).toHaveBeenCalledTimes(1);
});
test('creates a root worker span linked to the enqueue span by default', () => {
const span = {
end: vi.fn(),
recordException: vi.fn(),
setStatus: vi.fn(),
};
const rootContext = {} as Context;
const extractedContext = {} as Context;
const sourceSpanContext = {
traceId: '0123456789abcdef0123456789abcdef',
spanId: '0123456789abcdef',
traceFlags: 1,
isRemote: true,
} as SpanContext;
const propagation = {
isPropagationApi: true,
inject: vi.fn(),
extract(this: { isPropagationApi: boolean }) {
if (!this.isPropagationApi) throw new Error('lost propagation API receiver');
return extractedContext;
},
};
const tracer = {
startActiveSpan: vi.fn((_name: string, _options: unknown, _context: unknown, fn: (spanArg: typeof span) => string) => fn(span)),
} as any;
const adapter = new OpenTelemetryAdapter({
tracer,
provider: { shutdown: vi.fn() },
getActiveSpan: () => undefined,
spanStatusCodeError: SpanStatusCode.ERROR,
shutdownTimeout: 10,
queueTraceContext: {
tracer,
propagation: propagation as any,
trace: { getSpanContext: () => sourceSpanContext },
getActiveContext: () => rootContext,
rootContext,
mode: 'link',
spanStatusCodeError: SpanStatusCode.ERROR,
},
});
expect(adapter.startSpanWithTraceContext('Queue: Deliver', {
__misskeyTraceContext: {
traceparent: '00-0123456789abcdef0123456789abcdef-0123456789abcdef-01',
},
}, () => 'ok')).toBe('ok');
expect(tracer.startActiveSpan).toHaveBeenCalledWith('Queue: Deliver', {
root: true,
links: [{ context: sourceSpanContext }],
}, rootContext, expect.any(Function));
expect(span.end).toHaveBeenCalledTimes(1);
});
test('returns the active span context for log enrichment', () => {
const adapter = new OpenTelemetryAdapter({
tracer: { startActiveSpan: vi.fn() },
provider: { shutdown: vi.fn() },
getActiveSpan: () => ({
spanContext: () => ({
traceId: '0123456789abcdef0123456789abcdef',
spanId: '0123456789abcdef',
traceFlags: 0,
}),
} as any),
spanStatusCodeError: SpanStatusCode.ERROR,
shutdownTimeout: 10,
});
expect(adapter.getActiveTraceContext()).toEqual({
traceId: '0123456789abcdef0123456789abcdef',
spanId: '0123456789abcdef',
traceFlags: 0,
});
});
test('bridges captureMessage to the active span when one exists', () => {
const activeSpan = {
recordException: vi.fn(),
setStatus: vi.fn(),
};
const adapter = new OpenTelemetryAdapter({
tracer: { startActiveSpan: vi.fn() },
provider: { shutdown: vi.fn() },
getActiveSpan: () => activeSpan as any,
spanStatusCodeError: SpanStatusCode.ERROR,
shutdownTimeout: 10,
});
adapter.captureMessage('Queue failed', {
level: 'error',
extra: { queue: 'deliver' },
});
expect(activeSpan.recordException).toHaveBeenCalledWith(expect.objectContaining({
message: 'Queue failed',
}));
expect(activeSpan.setStatus).toHaveBeenCalledWith({
code: SpanStatusCode.ERROR,
message: 'Queue failed',
});
});
test('times out shutdown instead of waiting forever', async () => {
vi.useFakeTimers();
const adapter = new OpenTelemetryAdapter({
tracer: { startActiveSpan: vi.fn() },
provider: { shutdown: vi.fn(() => new Promise<void>(() => {})) },
getActiveSpan: () => undefined,
spanStatusCodeError: SpanStatusCode.ERROR,
shutdownTimeout: 50,
});
const shutdown = adapter.shutdown();
await vi.advanceTimersByTimeAsync(50);
await expect(shutdown).resolves.toBeUndefined();
});
test('clears the shutdown timeout timer once provider.shutdown() resolves first', async () => {
vi.useFakeTimers();
const clearTimeoutSpy = vi.spyOn(global, 'clearTimeout');
const adapter = new OpenTelemetryAdapter({
tracer: { startActiveSpan: vi.fn() },
provider: { shutdown: vi.fn().mockResolvedValue(undefined) },
getActiveSpan: () => undefined,
spanStatusCodeError: SpanStatusCode.ERROR,
shutdownTimeout: 5000,
});
await adapter.shutdown();
expect(clearTimeoutSpy).toHaveBeenCalled();
clearTimeoutSpy.mockRestore();
});
test('continues shutdown and flushes the provider when instrumentation cleanup fails', async () => {
const shutdownHttpClientInstrumentation = vi.fn(() => { throw new Error('failed'); });
const shutdownDatabaseInstrumentation = vi.fn();
const shutdownRedisInstrumentation = vi.fn();
const provider = { shutdown: vi.fn().mockResolvedValue(undefined) };
const adapter = new OpenTelemetryAdapter({
tracer: { startActiveSpan: vi.fn() },
provider,
getActiveSpan: () => undefined,
spanStatusCodeError: SpanStatusCode.ERROR,
shutdownTimeout: 10,
shutdownHttpClientInstrumentation,
shutdownDatabaseInstrumentation,
shutdownRedisInstrumentation,
});
await expect(adapter.shutdown()).resolves.toBeUndefined();
expect(shutdownHttpClientInstrumentation).toHaveBeenCalledOnce();
expect(shutdownDatabaseInstrumentation).toHaveBeenCalledOnce();
expect(shutdownRedisInstrumentation).toHaveBeenCalledOnce();
expect(provider.shutdown).toHaveBeenCalledOnce();
});
test('captureMessage starts a standalone span to report the error when there is no active span', () => {
const reportSpan = {
end: vi.fn(),
recordException: vi.fn(),
setStatus: vi.fn(),
};
const tracer = {
startActiveSpan: vi.fn((_name: string, fn: (spanArg: typeof reportSpan) => void) => fn(reportSpan)),
};
const adapter = new OpenTelemetryAdapter({
tracer: tracer as any,
provider: { shutdown: vi.fn() },
getActiveSpan: () => undefined,
spanStatusCodeError: SpanStatusCode.ERROR,
shutdownTimeout: 10,
});
adapter.captureMessage('Queue: Deliver failed', {
level: 'error',
extra: { queue: 'deliver' },
});
expect(tracer.startActiveSpan).toHaveBeenCalledWith('captureMessage', expect.any(Function));
expect(reportSpan.recordException).toHaveBeenCalledWith(expect.objectContaining({
message: 'Queue: Deliver failed',
}));
expect(reportSpan.setStatus).toHaveBeenCalledWith({
code: SpanStatusCode.ERROR,
message: 'Queue: Deliver failed',
});
expect(reportSpan.end).toHaveBeenCalledTimes(1);
});
});
describe('installInstrumentationsWithCleanup', () => {
test('cleans up acquired instrumentations in reverse order and preserves the installation error', async () => {
const installationError = new Error('failed to install');
const shutdownHttpClientInstrumentation = vi.fn();
const shutdownDatabaseInstrumentation = vi.fn(() => { throw new Error('failed to clean up'); });
await expect(installInstrumentationsWithCleanup([
() => shutdownHttpClientInstrumentation,
() => shutdownDatabaseInstrumentation,
() => { throw installationError; },
], () => { throw new Error('should not create'); })).rejects.toBe(installationError);
expect(shutdownDatabaseInstrumentation).toHaveBeenCalledOnce();
expect(shutdownHttpClientInstrumentation).toHaveBeenCalledOnce();
expect(shutdownDatabaseInstrumentation.mock.invocationCallOrder[0]).toBeLessThan(shutdownHttpClientInstrumentation.mock.invocationCallOrder[0]);
});
test('cleans up all instrumentations when adapter construction fails', async () => {
const constructionError = new Error('failed to create adapter');
const shutdownHttpClientInstrumentation = vi.fn();
const shutdownDatabaseInstrumentation = vi.fn();
await expect(installInstrumentationsWithCleanup([
() => shutdownHttpClientInstrumentation,
() => shutdownDatabaseInstrumentation,
], () => { throw constructionError; })).rejects.toBe(constructionError);
expect(shutdownDatabaseInstrumentation).toHaveBeenCalledOnce();
expect(shutdownHttpClientInstrumentation).toHaveBeenCalledOnce();
});
});
describe('createSampler', () => {
test('accepts sample rates within [0, 1]', () => {
expect(() => createSampler(0, samplerDeps)).not.toThrow();
expect(() => createSampler(0.5, samplerDeps)).not.toThrow();
expect(() => createSampler(1, samplerDeps)).not.toThrow();
});
test('rejects sample rates outside [0, 1]', () => {
expect(() => createSampler(-0.1, samplerDeps)).toThrow();
expect(() => createSampler(1.1, samplerDeps)).toThrow();
});
test('rejects NaN instead of silently disabling sampling', () => {
expect(() => createSampler(Number.NaN, samplerDeps)).toThrow();
});
test('rejects non-number values that pass through YAML as strings', () => {
expect(() => createSampler('0.5' as unknown as number, samplerDeps)).toThrow();
});
});
describe('createResource', () => {
test('lets explicit config override OTEL resource env, and env override Misskey defaults', () => {
const previousServiceName = process.env['OTEL_SERVICE_NAME'];
const previousResourceAttributes = process.env['OTEL_RESOURCE_ATTRIBUTES'];
process.env['OTEL_SERVICE_NAME'] = 'env-service';
process.env['OTEL_RESOURCE_ATTRIBUTES'] = [
'deployment.environment=staging',
'misskey.process.role=env-role',
'service.instance.id=env-instance',
'env.only=value',
].join(',');
try {
const resource = createResource({
serviceVersion: '2026.1.0',
resourceAttributes: {
[ATTR_SERVICE_NAME]: 'config-service',
'deployment.environment': 'production',
'config.only': 'value',
},
}, {
defaultResource,
resourceFromAttributes,
detectResources,
envDetector,
serviceNameAttribute: ATTR_SERVICE_NAME,
serviceInstanceIdAttribute: ATTR_SERVICE_INSTANCE_ID,
serviceVersionAttribute: ATTR_SERVICE_VERSION,
serviceVersion: '2026.1.0',
});
expect(resource.attributes[ATTR_SERVICE_NAME]).toBe('config-service');
expect(resource.attributes[ATTR_SERVICE_INSTANCE_ID]).toBe('env-instance');
expect(resource.attributes[ATTR_SERVICE_VERSION]).toBe('2026.1.0');
expect(resource.attributes['deployment.environment']).toBe('production');
expect(resource.attributes['misskey.process.role']).toBe('env-role');
expect(resource.attributes['env.only']).toBe('value');
expect(resource.attributes['config.only']).toBe('value');
} finally {
if (previousServiceName == null) {
delete process.env['OTEL_SERVICE_NAME'];
} else {
process.env['OTEL_SERVICE_NAME'] = previousServiceName;
}
if (previousResourceAttributes == null) {
delete process.env['OTEL_RESOURCE_ATTRIBUTES'];
} else {
process.env['OTEL_RESOURCE_ATTRIBUTES'] = previousResourceAttributes;
}
}
});
});
describe('getMisskeyProcessRole', () => {
beforeEach(() => {
mocks.envOption.disableClustering = false;
mocks.envOption.onlyServer = false;
mocks.envOption.onlyQueue = false;
mocks.isPrimary = false;
});
test('labels non-clustered onlyServer as primary-server', () => {
mocks.envOption.disableClustering = true;
mocks.envOption.onlyServer = true;
expect(getMisskeyProcessRole()).toBe('primary-server');
});
test('labels clustered primary with onlyServer as fork-only', () => {
mocks.isPrimary = true;
mocks.envOption.onlyServer = true;
expect(getMisskeyProcessRole()).toBe('fork-only');
});
test('labels clustered worker running the HTTP server (onlyServer) as worker-server, not worker-queue', () => {
mocks.isPrimary = false;
mocks.envOption.onlyServer = true;
expect(getMisskeyProcessRole()).toBe('worker-server');
});
test('labels clustered worker without onlyServer as worker-queue', () => {
mocks.isPrimary = false;
mocks.envOption.onlyServer = false;
expect(getMisskeyProcessRole()).toBe('worker-queue');
});
});
@@ -4,11 +4,9 @@
*/
import { describe, expect, test, vi } from 'vitest';
import { normalizeSentryBackendConfig } from '@/config.js';
import { SentryTelemetryAdapter, buildSentryIntegrations, buildSentryNodeOptions, buildSentryOtlpInitOptions } from '@/core/telemetry/adapters/SentryTelemetryAdapter.js';
import { SentryTelemetryAdapter, buildSentryIntegrations, buildSentryNodeOptions } from '@/core/telemetry/adapters/SentryTelemetryAdapter.js';
type TestIntegration = Parameters<ReturnType<typeof buildSentryIntegrations>>[0][number];
type TestSpanProcessor = NonNullable<ReturnType<typeof buildSentryOtlpInitOptions>['openTelemetrySpanProcessors']>[number];
function testIntegration(name: string): TestIntegration {
return { name };
@@ -37,9 +35,7 @@ describe('SentryTelemetryAdapter', () => {
nodeProfilingIntegration: () => testIntegration('ProfilingIntegration'),
});
const result = integrations([
testIntegration('Http'),
]);
const result = integrations([testIntegration('Http')]);
expect(result.map((integration: TestIntegration) => integration.name)).toEqual(['Http', 'ProfilingIntegration']);
});
@@ -52,9 +48,7 @@ describe('SentryTelemetryAdapter', () => {
warn,
});
const result = integrations([
testIntegration('Http'),
]);
const result = integrations([testIntegration('Http')]);
expect(result.map((integration: TestIntegration) => integration.name)).toEqual(['Http']);
expect(warn).toHaveBeenCalledWith('Unknown Sentry integration configured in sentryForBackend.disabledIntegrations: Unknown');
@@ -79,104 +73,6 @@ describe('SentryTelemetryAdapter', () => {
expect(options.tracePropagationTargets).toEqual(['^https://internal\\.example/']);
});
test('builds Sentry options that export spans to both Sentry and OTLP', () => {
const existingProcessor = { name: 'existingProcessor' } as unknown as TestSpanProcessor;
const otlpProcessor = { name: 'otlpProcessor' };
const result = buildSentryOtlpInitOptions({
sentryConfig: {
enableNodeProfiling: false,
disabledIntegrations: ['Redis'],
options: {
openTelemetrySpanProcessors: [existingProcessor],
tracesSampleRate: 0.25,
},
},
otelConfig: { serviceVersion: '2026.1.0' },
otlpProcessor,
});
expect(result.tracesSampleRate).toBe(0.25);
expect(result.openTelemetrySpanProcessors).toEqual([existingProcessor, otlpProcessor]);
// OTel併存時もremoteへtrace headerを漏らさないデフォルトはSentry単体時と揃える。
expect(result.tracePropagationTargets).toEqual([]);
expect(result.integrations).toBeTypeOf('function');
if (typeof result.integrations !== 'function') throw new Error('Expected integrations to be a function');
expect(result.integrations([
testIntegration('Http'),
testIntegration('Redis'),
testIntegration('Postgres'),
]).map((integration: TestIntegration) => integration.name)).toEqual(['Http', 'Postgres']);
});
test('uses safe defaults when Sentry options are omitted', () => {
const otlpProcessor = { name: 'otlpProcessor' };
const result = buildSentryOtlpInitOptions({
sentryConfig: normalizeSentryBackendConfig({
enableNodeProfiling: false,
}),
otelConfig: { serviceVersion: '2026.1.0' },
otlpProcessor,
});
expect(result.tracePropagationTargets).toEqual([]);
expect(result.openTelemetrySpanProcessors).toEqual([otlpProcessor]);
});
test('does not disable Sentry trace propagation when explicitly enabled for OTel coexistence', () => {
const result = buildSentryOtlpInitOptions({
sentryConfig: {
enableNodeProfiling: false,
options: {},
},
otelConfig: {
serviceVersion: '2026.1.0',
propagateTraceToRemote: true,
},
otlpProcessor: { name: 'otlpProcessor' },
});
expect(result.tracePropagationTargets).toBeUndefined();
});
test('honors explicit tracePropagationTargets for OTel coexistence even without propagateTraceToRemote', () => {
const result = buildSentryOtlpInitOptions({
sentryConfig: {
enableNodeProfiling: false,
options: {
tracePropagationTargets: ['^https://internal\\.example/'],
},
},
otelConfig: { serviceVersion: '2026.1.0' },
otlpProcessor: { name: 'otlpProcessor' },
});
expect(result.tracePropagationTargets).toEqual(['^https://internal\\.example/']);
});
test('warns when OTel-only options are ignored in Sentry coexistence mode', () => {
const warn = vi.fn();
buildSentryOtlpInitOptions({
sentryConfig: {
enableNodeProfiling: false,
options: {},
},
otelConfig: {
serviceVersion: '2026.1.0',
sampleRate: 0.25,
resourceAttributes: {
'deployment.environment': 'production',
},
},
otlpProcessor: { name: 'otlpProcessor' },
warn,
});
expect(warn).toHaveBeenCalledWith(expect.stringContaining('otelForBackend.sampleRate is ignored'));
expect(warn).toHaveBeenCalledWith(expect.stringContaining('otelForBackend.resourceAttributes is ignored'));
});
});
describe('SentryTelemetryAdapter trace context', () => {
@@ -238,69 +134,3 @@ describe('SentryTelemetryAdapter.shutdown', () => {
vi.doUnmock('@sentry/profiling-node');
});
});
describe('SentryTelemetryAdapter.createWithOtlpExport', () => {
test('registers the OTel diag logger before creating the OTLP exporter', async () => {
const init = vi.fn();
const close = vi.fn();
const setLogger = vi.fn();
const nodeProfilingIntegration = vi.fn();
const BatchSpanProcessor = vi.fn(function (this: { exporter: unknown }, exporter: unknown) {
this.exporter = exporter;
});
const OTLPTraceExporter = vi.fn(function (this: { options: unknown }, options: unknown) {
this.options = options;
});
vi.doMock('@sentry/node', () => ({
init,
close,
}));
vi.doMock('@sentry/profiling-node', () => ({
nodeProfilingIntegration,
}));
vi.doMock('@opentelemetry/api', () => ({
context: { active: vi.fn() },
diag: { setLogger },
DiagLogLevel: { WARN: 50 },
propagation: { inject: vi.fn(), extract: vi.fn() },
ROOT_CONTEXT: {},
SpanStatusCode: { ERROR: 2 },
trace: { getTracer: vi.fn(), getSpanContext: vi.fn() },
}));
vi.doMock('@opentelemetry/sdk-trace-base', () => ({
BatchSpanProcessor,
}));
vi.doMock('@opentelemetry/exporter-trace-otlp-proto', () => ({
OTLPTraceExporter,
}));
await SentryTelemetryAdapter.createWithOtlpExport({
enableNodeProfiling: false,
options: {},
}, {
serviceVersion: '2026.1.0',
endpoint: 'http://collector:4318/v1/traces',
});
expect(setLogger).toHaveBeenCalledWith(expect.objectContaining({
error: expect.any(Function),
warn: expect.any(Function),
}), {
logLevel: 50,
suppressOverrideMessage: true,
});
expect(OTLPTraceExporter).toHaveBeenCalledWith({
url: 'http://collector:4318/v1/traces',
});
expect(init).toHaveBeenCalledWith(expect.objectContaining({
openTelemetrySpanProcessors: [expect.any(Object)],
}));
vi.doUnmock('@sentry/node');
vi.doUnmock('@sentry/profiling-node');
vi.doUnmock('@opentelemetry/api');
vi.doUnmock('@opentelemetry/sdk-trace-base');
vi.doUnmock('@opentelemetry/exporter-trace-otlp-proto');
});
});
@@ -1,188 +0,0 @@
/*
* SPDX-FileCopyrightText: syuilo and misskey-project
* SPDX-License-Identifier: AGPL-3.0-only
*/
import { describe, expect, test, vi } from 'vitest';
import { SpanKind, SpanStatusCode } from '@opentelemetry/api';
import { createHttpClientInstrumentation } from '@/core/telemetry/http-client-instrumentation.js';
function request() {
return {
method: 'POST',
protocol: 'https:',
path: '/inbox?token=secret',
host: 'remote.example',
getHeader: vi.fn((name: string) => name === 'host' ? 'user:password@remote.example:8443' : undefined),
};
}
describe('http-client-instrumentation', () => {
test('creates and completes a sanitized CLIENT span from diagnostics channels', () => {
const listeners = new Map<string, (message: unknown) => void>();
const span = {
end: vi.fn(),
recordException: vi.fn(),
setAttribute: vi.fn(),
setStatus: vi.fn(),
};
const tracer = { startSpan: vi.fn(() => span) } as any;
const unsubscribe = createHttpClientInstrumentation({
tracer,
spanKindClient: SpanKind.CLIENT,
spanStatusCodeError: SpanStatusCode.ERROR,
subscribe: (name, listener) => {
listeners.set(name, listener);
return () => listeners.delete(name);
},
});
const clientRequest = request();
listeners.get('http.client.request.created')!({ request: clientRequest });
listeners.get('http.client.response.finish')!({
request: clientRequest,
response: { statusCode: 201, httpVersion: '1.1' },
});
expect(tracer.startSpan).toHaveBeenCalledWith('POST', {
kind: SpanKind.CLIENT,
attributes: {
'http.request.method': 'POST',
'url.full': 'https://remote.example:8443/inbox',
'server.address': 'remote.example',
'server.port': 8443,
},
});
expect(span.setAttribute).toHaveBeenCalledWith('http.response.status_code', 201);
expect(span.setAttribute).toHaveBeenCalledWith('network.protocol.version', '1.1');
expect(span.end).toHaveBeenCalledTimes(1);
unsubscribe();
expect(listeners.size).toBe(0);
});
test('records a request error and ends the span once', () => {
const listeners = new Map<string, (message: unknown) => void>();
const span = { end: vi.fn(), recordException: vi.fn(), setAttribute: vi.fn(), setStatus: vi.fn() };
const error = Object.assign(new Error('connection refused'), { code: 'ECONNREFUSED' });
const clientRequest = request();
createHttpClientInstrumentation({
tracer: { startSpan: vi.fn(() => span) } as any,
spanKindClient: SpanKind.CLIENT,
spanStatusCodeError: SpanStatusCode.ERROR,
subscribe: (name, listener) => {
listeners.set(name, listener);
return () => listeners.delete(name);
},
});
listeners.get('http.client.request.created')!({ request: clientRequest });
listeners.get('http.client.request.error')!({ request: clientRequest, error });
listeners.get('http.client.response.finish')!({ request: clientRequest, response: { statusCode: 200 } });
expect(span.recordException).toHaveBeenCalledWith(error);
expect(span.setAttribute).toHaveBeenCalledWith('error.type', 'ECONNREFUSED');
expect(span.setStatus).toHaveBeenCalledWith({ code: SpanStatusCode.ERROR });
expect(span.end).toHaveBeenCalledTimes(1);
});
test('records the response status code as error.type for an error response', () => {
const listeners = new Map<string, (message: unknown) => void>();
const span = { end: vi.fn(), recordException: vi.fn(), setAttribute: vi.fn(), setStatus: vi.fn() };
const clientRequest = request();
createHttpClientInstrumentation({
tracer: { startSpan: vi.fn(() => span) } as any,
spanKindClient: SpanKind.CLIENT,
spanStatusCodeError: SpanStatusCode.ERROR,
subscribe: (name, listener) => {
listeners.set(name, listener);
return () => listeners.delete(name);
},
});
listeners.get('http.client.request.created')!({ request: clientRequest });
listeners.get('http.client.response.finish')!({ request: clientRequest, response: { statusCode: 502 } });
expect(span.setAttribute).toHaveBeenCalledWith('error.type', '502');
expect(span.setStatus).toHaveBeenCalledWith({ code: SpanStatusCode.ERROR });
});
test('uses safe request details when the request URL cannot be constructed', () => {
const listeners = new Map<string, (message: unknown) => void>();
const span = { end: vi.fn(), recordException: vi.fn(), setAttribute: vi.fn(), setStatus: vi.fn() };
const tracer = { startSpan: vi.fn(() => span) };
createHttpClientInstrumentation({
tracer: tracer as any,
spanKindClient: SpanKind.CLIENT,
spanStatusCodeError: SpanStatusCode.ERROR,
subscribe: (name, listener) => {
listeners.set(name, listener);
return () => listeners.delete(name);
},
});
const clientRequest = {
...request(),
host: '::1',
getHeader: vi.fn(() => undefined),
};
listeners.get('http.client.request.created')!({ request: clientRequest });
expect(tracer.startSpan).toHaveBeenCalledWith('POST', {
kind: SpanKind.CLIENT,
attributes: {
'http.request.method': 'POST',
'url.full': 'https://localhost/',
'server.address': 'localhost',
'server.port': 443,
},
});
});
test('does not propagate errors from diagnostics listeners', () => {
const listeners = new Map<string, (message: unknown) => void>();
const reportError = vi.fn();
const responseError = new Error('response instrumentation failed');
const requestError = new Error('request instrumentation failed');
const startError = new Error('span creation failed');
const span = {
end: vi.fn(),
recordException: vi.fn(() => { throw requestError; }),
setAttribute: vi.fn(() => { throw responseError; }),
setStatus: vi.fn(),
};
const tracer = {
startSpan: vi.fn()
.mockReturnValueOnce(span)
.mockReturnValueOnce(span)
.mockImplementationOnce(() => { throw startError; }),
};
createHttpClientInstrumentation({
tracer: tracer as any,
spanKindClient: SpanKind.CLIENT,
spanStatusCodeError: SpanStatusCode.ERROR,
subscribe: (name, listener) => {
listeners.set(name, listener);
return () => listeners.delete(name);
},
reportError,
});
const responseRequest = request();
const errorRequest = request();
listeners.get('http.client.request.created')!({ request: responseRequest });
expect(() => listeners.get('http.client.response.finish')!({
request: responseRequest,
response: { statusCode: 200 },
})).not.toThrow();
listeners.get('http.client.request.created')!({ request: errorRequest });
expect(() => listeners.get('http.client.request.error')!({
request: errorRequest,
error: new Error('connection failed'),
})).not.toThrow();
expect(() => listeners.get('http.client.request.created')!({ request: request() })).not.toThrow();
expect(reportError).toHaveBeenCalledWith(responseError);
expect(reportError).toHaveBeenCalledWith(requestError);
expect(reportError).toHaveBeenCalledWith(startError);
});
});
@@ -1,57 +0,0 @@
/*
* SPDX-FileCopyrightText: syuilo and misskey-project
* SPDX-License-Identifier: AGPL-3.0-only
*/
import { beforeEach, describe, expect, test, vi } from 'vitest';
import type * as Bull from 'bullmq';
import { instrumentQueue } from '@/core/telemetry/queue-instrumentation.js';
const mocks = vi.hoisted(() => ({
injectTraceContext: vi.fn((carrier: Record<string, string>) => {
carrier['traceparent'] = '00-0123456789abcdef0123456789abcdef-0123456789abcdef-01';
}),
}));
vi.mock('@/core/telemetry/telemetry-registry.js', () => ({
injectTraceContext: mocks.injectTraceContext,
}));
describe('queue-instrumentation', () => {
beforeEach(() => {
mocks.injectTraceContext.mockClear();
});
test('injects the active trace context for add()', () => {
const add = vi.fn();
const queue = instrumentQueue({ add, addBulk: vi.fn() } as unknown as Bull.Queue<{ noteId: string }>);
const data = { noteId: '9d6b9a65-46c9-4e1b-a640-9589693893c9' };
queue.add('endedPollNotification', data);
expect(mocks.injectTraceContext).toHaveBeenCalledTimes(1);
expect(data).toMatchObject({
__misskeyTraceContext: {
traceparent: '00-0123456789abcdef0123456789abcdef-0123456789abcdef-01',
},
});
expect(add).toHaveBeenCalledWith('endedPollNotification', data, undefined);
});
test('injects every job passed to addBulk()', () => {
const addBulk = vi.fn();
const queue = instrumentQueue({ add: vi.fn(), addBulk } as unknown as Bull.Queue<{ to: string }>);
const jobs = [
{ name: 'deliver', data: { to: 'https://remote.example/inbox' } },
{ name: 'deliver', data: { to: 'https://remote2.example/inbox' } },
];
queue.addBulk(jobs);
expect(mocks.injectTraceContext).toHaveBeenCalledTimes(2);
expect(jobs).toEqual(expect.arrayContaining([
expect.objectContaining({ data: expect.objectContaining({ __misskeyTraceContext: expect.any(Object) }) }),
]));
expect(addBulk).toHaveBeenCalledWith(jobs);
});
});
@@ -1,233 +0,0 @@
/*
* SPDX-FileCopyrightText: syuilo and misskey-project
* SPDX-License-Identifier: AGPL-3.0-only
*/
import { describe, expect, test, vi } from 'vitest';
import { SpanStatusCode } from '@opentelemetry/api';
import type { Context, SpanContext } from '@opentelemetry/api';
import { executeSpan, getQueueSpanContext, getQueueTraceContextMode, injectActiveTraceContext, injectQueueTraceContext, recordSpanError, startSpanWithQueueTraceContext } from '@/core/telemetry/queue-trace-context.js';
const rootContext = { kind: 'root' } as unknown as Context;
const extractedContext = { kind: 'extracted' } as unknown as Context;
const sourceSpanContext: SpanContext = {
traceId: '0123456789abcdef0123456789abcdef',
spanId: '0123456789abcdef',
traceFlags: 1,
isRemote: true,
};
function jobData() {
return {
name: 'deliver',
__misskeyTraceContext: {
traceparent: '00-0123456789abcdef0123456789abcdef-0123456789abcdef-01',
},
};
}
describe('queue-trace-context', () => {
test('stores only a non-empty carrier in the job data', () => {
const data = { noteId: '9d6b9a65-46c9-4e1b-a640-9589693893c9' };
injectQueueTraceContext(data, carrier => {
carrier['traceparent'] = '00-0123456789abcdef0123456789abcdef-0123456789abcdef-01';
});
expect(data).toMatchObject({
__misskeyTraceContext: {
traceparent: '00-0123456789abcdef0123456789abcdef-0123456789abcdef-01',
},
});
});
test('does not store an empty carrier when no active trace exists', () => {
const data = { noteId: '9d6b9a65-46c9-4e1b-a640-9589693893c9' };
injectQueueTraceContext(data, () => {});
expect(data).not.toHaveProperty('__misskeyTraceContext');
});
test('ignores non-object job data', () => {
const inject = vi.fn();
injectQueueTraceContext(null, inject);
injectQueueTraceContext('not a job object', inject);
expect(inject).not.toHaveBeenCalled();
});
test('injects the active context with the configured propagator', () => {
const activeContext = {} as Context;
const carrier = {};
const inject = vi.fn();
injectActiveTraceContext({
tracer: { startActiveSpan: vi.fn() } as any,
propagation: { inject, extract: vi.fn() } as any,
trace: { getSpanContext: vi.fn() },
getActiveContext: () => activeContext,
rootContext,
mode: 'link',
spanStatusCodeError: 2 as any,
}, carrier);
expect(inject).toHaveBeenCalledWith(activeContext, carrier);
});
test('starts a new root trace with a link by default', () => {
const extract = vi.fn(() => extractedContext);
const getSpanContext = vi.fn(() => sourceSpanContext);
const result = getQueueSpanContext(jobData(), {
rootContext,
propagation: { inject: vi.fn(), extract },
trace: { getSpanContext },
mode: 'link',
});
expect(extract).toHaveBeenCalledWith(rootContext, jobData().__misskeyTraceContext);
expect(result).toEqual({
options: {
root: true,
links: [{ context: sourceSpanContext }],
},
parentContext: rootContext,
});
expect(result?.parentContext).toBe(rootContext);
expect(getSpanContext).toHaveBeenCalledWith(extractedContext);
});
test('uses the extracted context as the parent when parent mode is selected', () => {
const result = getQueueSpanContext(jobData(), {
rootContext,
propagation: { inject: vi.fn(), extract: () => extractedContext },
trace: { getSpanContext: () => sourceSpanContext },
mode: 'parent',
});
expect(result).toEqual({
options: {},
parentContext: extractedContext,
});
expect(result?.parentContext).toBe(extractedContext);
});
test('ignores malformed or missing carriers', () => {
const extract = vi.fn(() => extractedContext);
const deps = {
rootContext,
propagation: { inject: vi.fn(), extract },
trace: { getSpanContext: () => sourceSpanContext },
mode: 'link' as const,
};
expect(getQueueSpanContext({}, deps)).toBeUndefined();
expect(getQueueSpanContext({ __misskeyTraceContext: { traceparent: 1 } }, deps)).toBeUndefined();
expect(extract).not.toHaveBeenCalled();
});
test('defaults to link mode and rejects invalid configuration', () => {
expect(getQueueTraceContextMode(undefined)).toBe('link');
expect(getQueueTraceContextMode('parent')).toBe('parent');
expect(() => getQueueTraceContextMode('children')).toThrow('otelForBackend.jobTraceContextMode');
});
test('executes synchronous work and ends the span once', () => {
const span = { end: vi.fn(), recordException: vi.fn(), setStatus: vi.fn() };
expect(executeSpan(span as any, () => 'ok', SpanStatusCode.ERROR)).toBe('ok');
expect(span.end).toHaveBeenCalledOnce();
expect(span.recordException).not.toHaveBeenCalled();
});
test('waits for resolved Promise work before ending the span', async () => {
const span = { end: vi.fn(), recordException: vi.fn(), setStatus: vi.fn() };
let resolveWork: ((value: string) => void) | undefined;
const work = new Promise<string>(resolve => {
resolveWork = resolve;
});
const result = executeSpan(span as any, () => work, SpanStatusCode.ERROR);
expect(span.end).not.toHaveBeenCalled();
if (resolveWork == null) throw new Error('work resolver was not initialized');
resolveWork('ok');
await expect(result).resolves.toBe('ok');
expect(span.end).toHaveBeenCalledOnce();
expect(span.recordException).not.toHaveBeenCalled();
});
test('records and propagates synchronous failures', () => {
const span = { end: vi.fn(), recordException: vi.fn(), setStatus: vi.fn() };
const error = new Error('boom');
expect(() => executeSpan(span as any, () => { throw error; }, SpanStatusCode.ERROR)).toThrow(error);
expect(span.recordException).toHaveBeenCalledWith(error);
expect(span.setStatus).toHaveBeenCalledWith({ code: SpanStatusCode.ERROR, message: error.message });
expect(span.end).toHaveBeenCalledOnce();
});
test('records and propagates rejected Promise work', async () => {
const span = { end: vi.fn(), recordException: vi.fn(), setStatus: vi.fn() };
const error = new Error('boom');
await expect(executeSpan(span as any, () => Promise.reject(error), SpanStatusCode.ERROR)).rejects.toBe(error);
expect(span.recordException).toHaveBeenCalledWith(error);
expect(span.setStatus).toHaveBeenCalledWith({ code: SpanStatusCode.ERROR, message: error.message });
expect(span.end).toHaveBeenCalledOnce();
});
test('normalizes non-Error failures before recording them', () => {
const span = { recordException: vi.fn(), setStatus: vi.fn() };
recordSpanError(span as any, 'boom', SpanStatusCode.ERROR);
expect(span.recordException).toHaveBeenCalledWith(expect.any(Error));
expect(span.recordException).toHaveBeenCalledWith(expect.objectContaining({ message: 'boom' }));
expect(span.setStatus).toHaveBeenCalledWith({ code: SpanStatusCode.ERROR, message: 'boom' });
});
test('starts a span with the extracted queue context', () => {
const span = { end: vi.fn(), recordException: vi.fn(), setStatus: vi.fn() };
const startActiveSpan = vi.fn((_name: string, _options: unknown, _context: unknown, fn: (spanArg: typeof span) => string) => fn(span));
const getSpanContext = vi.fn(() => sourceSpanContext);
const deps = {
tracer: { startActiveSpan } as any,
rootContext,
propagation: { inject: vi.fn(), extract: () => extractedContext } as any,
trace: { getSpanContext },
getActiveContext: () => rootContext,
mode: 'link' as const,
spanStatusCodeError: SpanStatusCode.ERROR,
};
expect(startSpanWithQueueTraceContext(deps, 'Queue: Deliver', jobData(), () => 'ok', () => 'fallback')).toBe('ok');
expect(startActiveSpan).toHaveBeenCalledWith('Queue: Deliver', {
root: true,
links: [{ context: sourceSpanContext }],
}, rootContext, expect.any(Function));
expect(getSpanContext).toHaveBeenCalledOnce();
expect(getSpanContext).toHaveBeenCalledWith(extractedContext);
expect(span.end).toHaveBeenCalledOnce();
});
test('uses the fallback without starting a span when queue context is missing', () => {
const startActiveSpan = vi.fn();
const fallback = vi.fn(() => 'fallback');
const deps = {
tracer: { startActiveSpan } as any,
rootContext,
propagation: { inject: vi.fn(), extract: vi.fn() } as any,
trace: { getSpanContext: vi.fn() },
getActiveContext: () => rootContext,
mode: 'link' as const,
spanStatusCodeError: SpanStatusCode.ERROR,
};
expect(startSpanWithQueueTraceContext(deps, 'Queue: Deliver', {}, () => 'ok', fallback)).toBe('fallback');
expect(fallback).toHaveBeenCalledOnce();
expect(startActiveSpan).not.toHaveBeenCalled();
});
});
@@ -4,13 +4,13 @@
*/
import { describe, expect, test, vi } from 'vitest';
import { runQueueJobWithTraceContext } from '@/queue/queue-job-runner.js';
import { runQueueJob } from '@/queue/queue-job-runner.js';
import { TelemetryService } from '@/core/telemetry/TelemetryService.js';
describe('runQueueJobWithTraceContext', () => {
describe('runQueueJob', () => {
test('returns the processor result without invoking the error handler', async () => {
let spanActive = false;
const startSpanWithTraceContext = vi.fn(<T>(_name: string, _jobData: object, fn: () => T): T => {
const startSpan = vi.fn(<T>(_name: string, fn: () => T): T => {
spanActive = true;
const result = fn();
if (result instanceof Promise) return result.finally(() => { spanActive = false; }) as T;
@@ -18,11 +18,11 @@ describe('runQueueJobWithTraceContext', () => {
return result;
});
const telemetryService = {
startSpanWithTraceContext,
startSpan,
} as unknown as TelemetryService;
const onError = vi.fn();
await expect(runQueueJobWithTraceContext(telemetryService, 'Queue: test', {}, () => 'ok', onError)).resolves.toBe('ok');
await expect(runQueueJob(telemetryService, 'Queue: test', () => 'ok', onError)).resolves.toBe('ok');
expect(onError).not.toHaveBeenCalled();
expect(spanActive).toBe(false);
@@ -30,7 +30,7 @@ describe('runQueueJobWithTraceContext', () => {
test('handles failures while the processor span is active and rethrows the original error', async () => {
let spanActive = false;
const startSpanWithTraceContext = vi.fn(<T>(_name: string, _jobData: object, fn: () => T): T => {
const startSpan = vi.fn(<T>(_name: string, fn: () => T): T => {
spanActive = true;
const result = fn();
if (result instanceof Promise) return result.finally(() => { spanActive = false; }) as T;
@@ -38,7 +38,7 @@ describe('runQueueJobWithTraceContext', () => {
return result;
});
const telemetryService = {
startSpanWithTraceContext,
startSpan,
} as unknown as TelemetryService;
const onError = vi.fn((error: Error) => {
expect(spanActive).toBe(true);
@@ -46,7 +46,7 @@ describe('runQueueJobWithTraceContext', () => {
});
const originalError = new Error('failed');
await expect(runQueueJobWithTraceContext(telemetryService, 'Queue: test', {}, async () => {
await expect(runQueueJob(telemetryService, 'Queue: test', async () => {
throw originalError;
}, onError)).rejects.toBe(originalError);
@@ -1,53 +0,0 @@
/*
* SPDX-FileCopyrightText: syuilo and misskey-project
* SPDX-License-Identifier: AGPL-3.0-only
*/
import { describe, expect, test, vi } from 'vitest';
import { shouldRegisterHttpServerInstrumentation, registerHttpServerInstrumentation } from '@/server/http-server-instrumentation.js';
const mocks = vi.hoisted(() => ({
plugin: vi.fn(),
instrumentation: vi.fn(),
}));
vi.mock('@fastify/otel', () => ({
FastifyOtelInstrumentation: class {
public plugin = mocks.plugin;
public constructor(options: unknown) {
mocks.instrumentation(options);
}
},
}));
describe('http-server-instrumentation', () => {
test('registers Fastify instrumentation when only OpenTelemetry is configured', async () => {
const plugin = vi.fn();
const fastify = { register: vi.fn().mockResolvedValue(undefined) };
mocks.plugin.mockReturnValue(plugin);
await registerHttpServerInstrumentation(fastify as any, { otelForBackend: {} } as any);
expect(mocks.instrumentation).toHaveBeenCalledTimes(1);
expect(fastify.register).toHaveBeenCalledWith(plugin);
const requestHook = mocks.instrumentation.mock.calls[0][0].requestHook;
const span = { updateName: vi.fn() };
requestHook(span, { method: 'POST', routeOptions: { url: '/notes/create' } });
expect(span.updateName).toHaveBeenCalledWith('POST /notes/create');
});
test('does not register duplicate request instrumentation with Sentry', async () => {
const fastify = { register: vi.fn() };
await registerHttpServerInstrumentation(fastify as any, { otelForBackend: {}, sentryForBackend: {} } as any);
expect(fastify.register).not.toHaveBeenCalled();
expect(shouldRegisterHttpServerInstrumentation({ otelForBackend: {}, sentryForBackend: {} } as any)).toBe(false);
});
test('does not register instrumentation without OpenTelemetry', () => {
expect(shouldRegisterHttpServerInstrumentation({} as any)).toBe(false);
});
});
@@ -1,184 +0,0 @@
/*
* SPDX-FileCopyrightText: syuilo and misskey-project
* SPDX-License-Identifier: AGPL-3.0-only
*/
import { createRequire } from 'node:module';
import { context, propagation, trace } from '@opentelemetry/api';
import { NodeTracerProvider } from '@opentelemetry/sdk-trace-node';
import { describe, expect, test, vi } from 'vitest';
import type { Span as SdkSpan, SpanProcessor } from '@opentelemetry/sdk-trace-base';
import { installDatabaseInstrumentation, installInstrumentation } from '@/core/telemetry/database-instrumentation.js';
const backendRequire = createRequire(import.meta.url);
describe('database-instrumentation', () => {
test('does not install PostgreSQL instrumentation when disabled', async () => {
const uninstall = await installDatabaseInstrumentation({} as any, {
capturePgSpans: false,
capturePgStatement: false,
capturePgConnectionSpans: false,
});
expect(uninstall).toBeTypeOf('function');
expect(() => uninstall()).not.toThrow();
});
test('registers pg instrumentation with the active provider', () => {
const provider = {};
const pg = { setTracerProvider: vi.fn(), enable: vi.fn(), disable: vi.fn() };
const config = vi.fn();
const PgInstrumentation = class {
public constructor(options: unknown) {
config(options);
return pg as any;
}
};
const uninstall = installInstrumentation(provider as any, {
PgInstrumentation: PgInstrumentation as any,
}, {
capturePgStatement: false,
capturePgConnectionSpans: false,
});
expect(pg.setTracerProvider).toHaveBeenCalledWith(provider);
expect(pg.enable).toHaveBeenCalledOnce();
expect(config).toHaveBeenCalledWith(expect.objectContaining({
enhancedDatabaseReporting: false,
requireParentSpan: true,
ignoreConnectSpans: true,
requestHook: expect.any(Function),
}));
const span = { setAttribute: vi.fn() };
(config.mock.calls[0][0] as { requestHook: (span: any) => void }).requestHook(span);
expect(span.setAttribute).toHaveBeenCalledWith('db.statement', '[REDACTED]');
expect(span.setAttribute).toHaveBeenCalledWith('db.query.text', '[REDACTED]');
uninstall();
expect(pg.disable).toHaveBeenCalledOnce();
});
test('redacts SQL attributes after pg instrumentation records the statement', async () => {
const statement = 'SELECT very_secret_literal';
const attributeWrites: Array<{ key: string; value: unknown }> = [];
let attributesAtStart: Record<string, unknown> | undefined;
let pgSpan: SdkSpan | undefined;
const spanProcessor: SpanProcessor = {
forceFlush: async () => {},
onStart: span => {
if (span.instrumentationScope.name !== '@opentelemetry/instrumentation-pg') return;
pgSpan = span;
attributesAtStart = { ...span.attributes };
const setAttribute = span.setAttribute.bind(span);
span.setAttribute = (key, value) => {
attributeWrites.push({ key, value });
return setAttribute(key, value);
};
},
onEnd: () => {},
shutdown: async () => {},
};
const provider = new NodeTracerProvider({ spanProcessors: [spanProcessor] });
provider.register();
let uninstall: (() => void) | undefined;
try {
uninstall = await installDatabaseInstrumentation(provider, {
capturePgSpans: true,
capturePgStatement: false,
capturePgConnectionSpans: false,
});
const { Client } = backendRequire('pg') as typeof import('pg');
const tracer = provider.getTracer('database-instrumentation-test');
tracer.startActiveSpan('parent', parentSpan => {
try {
// pg instrumentation runs before an unconnected client queues the query,
// so no database is required.
new Client().query(statement, () => {});
} finally {
parentSpan.end();
}
});
const recordedPgSpan = pgSpan;
expect(recordedPgSpan).toBeDefined();
if (recordedPgSpan == null) throw new Error('PostgreSQL span was not started.');
expect(attributesAtStart).not.toHaveProperty('db.statement');
expect(attributesAtStart).not.toHaveProperty('db.query.text');
const rawSqlIndex = attributeWrites.findIndex(({ key, value }) =>
(key === 'db.statement' || key === 'db.query.text') && value === statement);
expect(rawSqlIndex).toBeGreaterThanOrEqual(0);
const rawSqlKey = attributeWrites[rawSqlIndex].key;
expect(attributeWrites.findIndex(({ key, value }, index) =>
index > rawSqlIndex && key === rawSqlKey && value === '[REDACTED]')).toBeGreaterThan(rawSqlIndex);
expect(recordedPgSpan.attributes['db.statement']).toBe('[REDACTED]');
expect(recordedPgSpan.attributes['db.query.text']).toBe('[REDACTED]');
expect(Object.values(recordedPgSpan.attributes)).not.toContain(statement);
} finally {
pgSpan?.end();
uninstall?.();
await provider.shutdown();
trace.disable();
context.disable();
propagation.disable();
}
});
test('keeps SQL statement attributes when explicitly enabled', () => {
const pg = { setTracerProvider: vi.fn(), enable: vi.fn(), disable: vi.fn() };
const config = vi.fn();
installInstrumentation({} as any, {
PgInstrumentation: class {
public constructor(options: unknown) {
config(options);
return pg as any;
}
} as any,
}, {
capturePgStatement: true,
capturePgConnectionSpans: false,
});
expect(config.mock.calls[0][0]).not.toHaveProperty('requestHook');
});
test('enables connection spans when explicitly configured', () => {
const pg = { setTracerProvider: vi.fn(), enable: vi.fn(), disable: vi.fn() };
const config = vi.fn();
installInstrumentation({} as any, {
PgInstrumentation: class {
public constructor(options: unknown) {
config(options);
return pg as any;
}
} as any,
}, {
capturePgStatement: false,
capturePgConnectionSpans: true,
});
expect(config).toHaveBeenCalledWith(expect.objectContaining({
ignoreConnectSpans: false,
}));
});
test('cleans up the pg instrumentation when initialization fails', () => {
const pg = { setTracerProvider: vi.fn(), enable: vi.fn(), disable: vi.fn() };
pg.enable.mockImplementation(() => { throw new Error('failed'); });
expect(() => installInstrumentation({} as any, {
PgInstrumentation: class { public constructor() { return pg as any; } } as any,
})).toThrow('failed');
expect(pg.disable).toHaveBeenCalledOnce();
});
});
@@ -1,209 +0,0 @@
/*
* SPDX-FileCopyrightText: syuilo and misskey-project
* SPDX-License-Identifier: AGPL-3.0-only
*/
import { describe, expect, test, vi } from 'vitest';
import { SpanKind, SpanStatusCode } from '@opentelemetry/api';
import { createRedisInstrumentation } from '@/core/telemetry/redis-instrumentation.js';
describe('redis-instrumentation', () => {
test('creates and completes a span for an ioredis command', () => {
let subscribers: any;
const unsubscribe = vi.fn();
const span = { end: vi.fn(), recordException: vi.fn(), setStatus: vi.fn(), setAttribute: vi.fn() };
const tracer = { startSpan: vi.fn(() => span) };
const uninstall = createRedisInstrumentation({
tracingChannel: () => ({ subscribe: (value) => { subscribers = value; }, unsubscribe }),
tracer: tracer as any,
getActiveSpan: () => ({}) as any,
spanKindClient: SpanKind.CLIENT,
spanStatusCodeError: SpanStatusCode.ERROR,
}, { captureCommandSpans: true });
const command = { command: 'get', args: ['key'], database: 0, serverAddress: 'redis', serverPort: 6379 };
subscribers.start(command);
subscribers.asyncEnd(command);
expect(tracer.startSpan).toHaveBeenCalledWith('get', expect.objectContaining({
kind: SpanKind.CLIENT,
attributes: expect.objectContaining({
'db.system.name': 'redis',
'db.operation.name': 'get',
'server.address': 'redis',
'server.port': 6379,
}),
}));
expect(span.end).toHaveBeenCalledOnce();
uninstall();
expect(unsubscribe).toHaveBeenCalledWith(subscribers);
});
test('records rejected Redis commands as errors', () => {
let subscribers: any;
const span = { end: vi.fn(), recordException: vi.fn(), setStatus: vi.fn(), setAttribute: vi.fn() };
createRedisInstrumentation({
tracingChannel: () => ({ subscribe: (value) => { subscribers = value; }, unsubscribe: vi.fn() }),
tracer: { startSpan: vi.fn(() => span) } as any,
getActiveSpan: () => ({}) as any,
spanKindClient: SpanKind.CLIENT,
spanStatusCodeError: SpanStatusCode.ERROR,
}, { captureCommandSpans: true });
const command = { command: 'get', args: ['key'], database: 0, serverAddress: 'redis', serverPort: undefined };
const error = Object.assign(new Error('ERR Redis failed'), { code: 'ERR' });
subscribers.start(command);
Object.assign(command, { error });
subscribers.error(command);
subscribers.asyncEnd(command);
expect(span.recordException).toHaveBeenCalledWith(error);
expect(span.recordException).toHaveBeenCalledOnce();
expect(span.setStatus).toHaveBeenCalledWith({ code: SpanStatusCode.ERROR, message: 'ERR Redis failed' });
expect(span.setAttribute).toHaveBeenCalledWith('error.type', 'ERR');
expect(span.setAttribute).toHaveBeenCalledWith('db.response.status_code', 'ERR');
expect(span.end).toHaveBeenCalledOnce();
});
test('does not create a root span when no parent span is active', () => {
let subscribers: any;
const tracer = { startSpan: vi.fn() };
createRedisInstrumentation({
tracingChannel: () => ({ subscribe: (value) => { subscribers = value; }, unsubscribe: vi.fn() }),
tracer: tracer as any,
getActiveSpan: () => undefined,
spanKindClient: SpanKind.CLIENT,
spanStatusCodeError: SpanStatusCode.ERROR,
}, { captureCommandSpans: true });
subscribers.start({ command: 'get', args: ['key'], database: 0, serverAddress: 'redis', serverPort: 6379 });
expect(tracer.startSpan).not.toHaveBeenCalled();
});
test('creates a root span when explicitly enabled', () => {
let subscribers: any;
const span = { end: vi.fn(), recordException: vi.fn(), setStatus: vi.fn(), setAttribute: vi.fn() };
const tracer = { startSpan: vi.fn(() => span) };
createRedisInstrumentation({
tracingChannel: () => ({ subscribe: (value) => { subscribers = value; }, unsubscribe: vi.fn() }),
tracer: tracer as any,
getActiveSpan: () => undefined,
spanKindClient: SpanKind.CLIENT,
spanStatusCodeError: SpanStatusCode.ERROR,
}, { captureCommandSpans: true, requireParentSpan: false });
const command = { command: 'get', args: ['key'], database: 0, serverAddress: 'redis', serverPort: 6379 };
subscribers.start(command);
subscribers.asyncEnd(command);
expect(tracer.startSpan).toHaveBeenCalledOnce();
expect(span.end).toHaveBeenCalledOnce();
});
test('records connection spans only when explicitly enabled', () => {
const subscribers = new Map<string, any>();
const unsubscribe = vi.fn();
const span = { end: vi.fn(), recordException: vi.fn(), setStatus: vi.fn(), setAttribute: vi.fn() };
const tracer = { startSpan: vi.fn(() => span) };
const uninstall = createRedisInstrumentation({
tracingChannel: (name) => ({ subscribe: (value) => { subscribers.set(name, value); }, unsubscribe }),
tracer: tracer as any,
getActiveSpan: () => undefined,
spanKindClient: SpanKind.CLIENT,
spanStatusCodeError: SpanStatusCode.ERROR,
}, { captureConnectionSpans: true });
const connection = { serverAddress: 'redis', serverPort: 6379 };
subscribers.get('ioredis:connect').start(connection);
subscribers.get('ioredis:connect').asyncEnd(connection);
expect(tracer.startSpan).toHaveBeenCalledWith('connect', expect.objectContaining({
kind: SpanKind.CLIENT,
attributes: expect.objectContaining({
'db.operation.name': 'connect',
'server.address': 'redis',
'server.port': 6379,
}),
}));
expect(span.end).toHaveBeenCalledOnce();
uninstall();
expect(unsubscribe).toHaveBeenCalledWith(subscribers.get('ioredis:connect'));
});
test('does not subscribe to Redis command diagnostics unless explicitly enabled', () => {
const tracingChannel = vi.fn();
createRedisInstrumentation({
tracingChannel,
tracer: { startSpan: vi.fn() } as any,
getActiveSpan: () => undefined,
spanKindClient: SpanKind.CLIENT,
spanStatusCodeError: SpanStatusCode.ERROR,
});
expect(tracingChannel).not.toHaveBeenCalled();
});
test('omits the database attribute when diagnostics context has no database', () => {
let subscribers: any;
const span = { end: vi.fn(), recordException: vi.fn(), setStatus: vi.fn(), setAttribute: vi.fn() };
const tracer = { startSpan: vi.fn(() => span) };
createRedisInstrumentation({
tracingChannel: () => ({ subscribe: (value) => { subscribers = value; }, unsubscribe: vi.fn() }),
tracer: tracer as any,
getActiveSpan: () => ({}) as any,
spanKindClient: SpanKind.CLIENT,
spanStatusCodeError: SpanStatusCode.ERROR,
}, { captureCommandSpans: true });
subscribers.start({ command: 'get', args: ['key'], database: undefined, serverAddress: 'redis', serverPort: 6379 });
expect(tracer.startSpan).toHaveBeenCalledWith('get', expect.objectContaining({
attributes: expect.not.objectContaining({
'db.namespace': expect.anything(),
}),
}));
});
test('does not propagate errors from Redis tracing subscribers', () => {
const subscribers = new Map<string, any>();
const reportError = vi.fn();
const endError = new Error('span end failed');
const startError = new Error('span creation failed');
const span = {
end: vi.fn(() => { throw endError; }),
recordException: vi.fn(),
setStatus: vi.fn(),
setAttribute: vi.fn(),
};
const tracer = {
startSpan: vi.fn()
.mockReturnValueOnce(span)
.mockImplementationOnce(() => { throw startError; }),
};
createRedisInstrumentation({
tracingChannel: (name) => ({
subscribe: (value) => { subscribers.set(name, value); },
unsubscribe: vi.fn(),
}),
tracer: tracer as any,
getActiveSpan: () => ({}) as any,
spanKindClient: SpanKind.CLIENT,
spanStatusCodeError: SpanStatusCode.ERROR,
reportError,
}, {
captureCommandSpans: true,
captureConnectionSpans: true,
});
const command = { command: 'get', args: ['key'], database: 0, serverAddress: 'redis', serverPort: 6379 };
subscribers.get('ioredis:command').start(command);
expect(() => subscribers.get('ioredis:command').asyncEnd(command)).not.toThrow();
expect(() => subscribers.get('ioredis:connect').start({ serverAddress: 'redis', serverPort: 6379 })).not.toThrow();
expect(reportError).toHaveBeenCalledWith(endError);
expect(reportError).toHaveBeenCalledWith(startError);
});
});
+22 -110
View File
@@ -9,8 +9,6 @@ import type { Config } from '@/config.js';
const mocks = vi.hoisted(() => {
return {
sentryCreate: vi.fn(),
sentryCreateWithOtlpExport: vi.fn(),
otelCreate: vi.fn(),
setLogTraceContextProvider: vi.fn(),
};
});
@@ -22,13 +20,6 @@ vi.mock('@/logging/logging-runtime.js', () => ({
vi.mock('@/core/telemetry/adapters/SentryTelemetryAdapter.js', () => ({
SentryTelemetryAdapter: {
create: mocks.sentryCreate,
createWithOtlpExport: mocks.sentryCreateWithOtlpExport,
},
}));
vi.mock('@/core/telemetry/adapters/OpenTelemetryAdapter.js', () => ({
OpenTelemetryAdapter: {
create: mocks.otelCreate,
},
}));
@@ -43,44 +34,37 @@ describe('telemetry-registry', () => {
beforeEach(() => {
vi.resetModules();
mocks.sentryCreate.mockReset();
mocks.sentryCreateWithOtlpExport.mockReset();
mocks.otelCreate.mockReset();
mocks.setLogTraceContextProvider.mockReset();
mocks.sentryCreate.mockResolvedValue({ shutdown: vi.fn(), captureMessage: vi.fn(), startSpan: vi.fn() });
mocks.sentryCreateWithOtlpExport.mockResolvedValue({ shutdown: vi.fn(), captureMessage: vi.fn(), startSpan: vi.fn() });
mocks.otelCreate.mockResolvedValue({ shutdown: vi.fn(), captureMessage: vi.fn(), startSpan: vi.fn() });
});
test('uses OpenTelemetryAdapter when only otelForBackend is configured', async () => {
test('does not initialize an adapter when Sentry is not configured', async () => {
const { initTelemetry } = await import('@/core/telemetry/telemetry-registry.js');
const otelForBackend = { endpoint: 'http://collector:4318/v1/traces' };
await initTelemetry(config({ otelForBackend }));
await initTelemetry(config({}));
expect(mocks.otelCreate).toHaveBeenCalledWith({
...otelForBackend,
serviceVersion: '2026.1.0',
});
expect(mocks.sentryCreate).not.toHaveBeenCalled();
expect(mocks.sentryCreateWithOtlpExport).not.toHaveBeenCalled();
expect(mocks.setLogTraceContextProvider).not.toHaveBeenCalled();
});
test('registers the adapter trace context provider after telemetry initialization', async () => {
test('initializes Sentry and registers its trace context provider', async () => {
const { initTelemetry } = await import('@/core/telemetry/telemetry-registry.js');
const sentryForBackend = { options: {}, enableNodeProfiling: false };
const getActiveTraceContext = vi.fn(() => ({
traceId: '0123456789abcdef0123456789abcdef',
spanId: '0123456789abcdef',
traceFlags: 0,
}));
mocks.otelCreate.mockResolvedValue({
mocks.sentryCreate.mockResolvedValue({
shutdown: vi.fn(),
captureMessage: vi.fn(),
startSpan: vi.fn(),
getActiveTraceContext,
});
await initTelemetry(config({ otelForBackend: { endpoint: 'http://collector:4318/v1/traces' } }));
await initTelemetry(config({ sentryForBackend }));
expect(mocks.sentryCreate).toHaveBeenCalledWith(sentryForBackend);
expect(mocks.setLogTraceContextProvider).toHaveBeenCalledWith(expect.any(Function));
const provider = mocks.setLogTraceContextProvider.mock.calls[0][0] as () => unknown;
expect(provider()).toEqual({
@@ -91,21 +75,6 @@ describe('telemetry-registry', () => {
expect(getActiveTraceContext).toHaveBeenCalledOnce();
});
test('adds OTLP export to the Sentry provider when both Sentry and OTel are configured', async () => {
const { initTelemetry } = await import('@/core/telemetry/telemetry-registry.js');
const sentryForBackend = { options: {}, enableNodeProfiling: false };
const otelForBackend = { endpoint: 'http://collector:4318/v1/traces' };
await initTelemetry(config({ sentryForBackend, otelForBackend }));
expect(mocks.sentryCreateWithOtlpExport).toHaveBeenCalledWith(sentryForBackend, {
...otelForBackend,
serviceVersion: '2026.1.0',
});
expect(mocks.sentryCreate).not.toHaveBeenCalled();
expect(mocks.otelCreate).not.toHaveBeenCalled();
});
test('startSpan runs fn directly when no adapter is registered', async () => {
const { startSpan } = await import('@/core/telemetry/telemetry-registry.js');
@@ -114,89 +83,32 @@ describe('telemetry-registry', () => {
expect(fn).toHaveBeenCalledTimes(1);
});
test('startSpan delegates directly to the single registered adapter without extra wrapping', async () => {
test('startSpan delegates to the Sentry adapter', async () => {
const { initTelemetry, startSpan } = await import('@/core/telemetry/telemetry-registry.js');
const otelForBackend = { endpoint: 'http://collector:4318/v1/traces' };
const adapterStartSpan = vi.fn((_name: string, fn: () => string) => fn());
mocks.otelCreate.mockResolvedValue({ shutdown: vi.fn(), captureMessage: vi.fn(), startSpan: adapterStartSpan });
mocks.sentryCreate.mockResolvedValue({ shutdown: vi.fn(), captureMessage: vi.fn(), startSpan: adapterStartSpan });
await initTelemetry(config({ otelForBackend }));
await initTelemetry(config({ sentryForBackend: { options: {}, enableNodeProfiling: false } }));
const fn = vi.fn().mockReturnValue('result');
expect(startSpan('test', fn)).toBe('result');
expect(adapterStartSpan).toHaveBeenCalledWith('test', fn);
});
test('startSpanWithTraceContext executes an undefined-returning callback only once', async () => {
const { initTelemetry, startSpanWithTraceContext } = await import('@/core/telemetry/telemetry-registry.js');
const otelForBackend = { endpoint: 'http://collector:4318/v1/traces' };
const adapterStartSpan = vi.fn((_name: string, fn: () => undefined) => fn());
const adapterStartSpanWithTraceContext = vi.fn((_name: string, _jobData: object, fn: () => undefined) => fn());
mocks.otelCreate.mockResolvedValue({
shutdown: vi.fn(),
captureMessage: vi.fn(),
startSpan: adapterStartSpan,
startSpanWithTraceContext: adapterStartSpanWithTraceContext,
});
await initTelemetry(config({ otelForBackend }));
const jobData = { id: 'job' };
const fn = vi.fn(() => undefined);
expect(startSpanWithTraceContext('test', jobData, fn)).toBeUndefined();
expect(adapterStartSpanWithTraceContext).toHaveBeenCalledWith('test', jobData, fn);
expect(adapterStartSpan).not.toHaveBeenCalled();
expect(fn).toHaveBeenCalledTimes(1);
});
test('startSpan wraps work through multiple registered adapters in order for future adapter combinations', async () => {
const { initTelemetry, startSpan } = await import('@/core/telemetry/telemetry-registry.js');
const calls: string[] = [];
mocks.sentryCreate.mockResolvedValue({
shutdown: vi.fn(),
captureMessage: vi.fn(),
startSpan: vi.fn((_name: string, fn: () => string) => {
calls.push('sentry:start');
const result = fn();
calls.push('sentry:end');
return result;
}),
});
mocks.otelCreate.mockResolvedValue({
shutdown: vi.fn(),
captureMessage: vi.fn(),
startSpan: vi.fn((_name: string, fn: () => string) => {
calls.push('otel:start');
const result = fn();
calls.push('otel:end');
return result;
}),
});
await initTelemetry(config({ sentryForBackend: { options: {}, enableNodeProfiling: false } }));
await initTelemetry(config({ otelForBackend: { endpoint: 'http://collector:4318/v1/traces' } }));
const fn = vi.fn(() => {
calls.push('work');
return 'result';
});
expect(startSpan('test', fn)).toBe('result');
expect(calls).toEqual(['sentry:start', 'otel:start', 'work', 'otel:end', 'sentry:end']);
});
test('shutdownTelemetry waits for every adapter even when one shutdown rejects', async () => {
test('shutdownTelemetry waits for all registered adapters even when one rejects', async () => {
const { initTelemetry, shutdownTelemetry } = await import('@/core/telemetry/telemetry-registry.js');
const sentryShutdown = vi.fn().mockRejectedValue(new Error('sentry failed'));
const otelShutdown = vi.fn().mockResolvedValue(undefined);
mocks.sentryCreate.mockResolvedValue({ shutdown: sentryShutdown, captureMessage: vi.fn(), startSpan: vi.fn() });
mocks.otelCreate.mockResolvedValue({ shutdown: otelShutdown, captureMessage: vi.fn(), startSpan: vi.fn() });
const firstShutdown = vi.fn().mockRejectedValue(new Error('first failed'));
const secondShutdown = vi.fn().mockResolvedValue(undefined);
mocks.sentryCreate
.mockResolvedValueOnce({ shutdown: firstShutdown, captureMessage: vi.fn(), startSpan: vi.fn() })
.mockResolvedValueOnce({ shutdown: secondShutdown, captureMessage: vi.fn(), startSpan: vi.fn() });
await initTelemetry(config({ sentryForBackend: { options: {}, enableNodeProfiling: false } }));
await initTelemetry(config({ otelForBackend: { endpoint: 'http://collector:4318/v1/traces' } }));
const sentryForBackend = { options: {}, enableNodeProfiling: false };
await initTelemetry(config({ sentryForBackend }));
await initTelemetry(config({ sentryForBackend }));
await expect(shutdownTelemetry()).resolves.toBeUndefined();
expect(sentryShutdown).toHaveBeenCalledTimes(1);
expect(otelShutdown).toHaveBeenCalledTimes(1);
expect(firstShutdown).toHaveBeenCalledTimes(1);
expect(secondShutdown).toHaveBeenCalledTimes(1);
});
});
-205
View File
@@ -188,9 +188,6 @@ importers:
'@fastify/multipart':
specifier: 10.1.0
version: 10.1.0
'@fastify/otel':
specifier: 0.20.1
version: 0.20.1(@opentelemetry/api@1.9.1)
'@fastify/static':
specifier: 10.1.2
version: 10.1.2
@@ -221,33 +218,6 @@ importers:
'@nestjs/testing':
specifier: 11.1.28
version: 11.1.28(@nestjs/common@11.1.28(reflect-metadata@0.2.2)(rxjs@7.8.2))(@nestjs/core@11.1.28)(@nestjs/platform-express@11.1.28)
'@opentelemetry/api':
specifier: 1.9.1
version: 1.9.1
'@opentelemetry/core':
specifier: 2.9.0
version: 2.9.0(@opentelemetry/api@1.9.1)
'@opentelemetry/exporter-trace-otlp-proto':
specifier: 0.220.0
version: 0.220.0(@opentelemetry/api@1.9.1)
'@opentelemetry/instrumentation':
specifier: 0.220.0
version: 0.220.0(@opentelemetry/api@1.9.1)
'@opentelemetry/instrumentation-pg':
specifier: 0.72.0
version: 0.72.0(@opentelemetry/api@1.9.1)
'@opentelemetry/resources':
specifier: 2.9.0
version: 2.9.0(@opentelemetry/api@1.9.1)
'@opentelemetry/sdk-trace-base':
specifier: 2.9.0
version: 2.9.0(@opentelemetry/api@1.9.1)
'@opentelemetry/sdk-trace-node':
specifier: 2.9.0
version: 2.9.0(@opentelemetry/api@1.9.1)
'@opentelemetry/semantic-conventions':
specifier: 1.43.0
version: 1.43.0
'@oxc-project/runtime':
specifier: 0.139.0
version: 0.139.0
@@ -1941,11 +1911,6 @@ packages:
'@fastify/multipart@10.1.0':
resolution: {integrity: sha512-b2r6CovmLQvLFJ5HJDtxVigZjcO9TwkRGslhiNQxsGJwilm5B3eqvXPOVLudC44zXMXJITpQ/iEATIE7zoXqcQ==}
'@fastify/otel@0.20.1':
resolution: {integrity: sha512-UG9gQUbqfBWcJNeO0bxj03xkoEE8NGJkIedN/red11elHFeZuVgoHKAx0x+wO4kGedyyQv8NWC62AdpGEvyM6g==}
peerDependencies:
'@opentelemetry/api': ^1.9.0
'@fastify/proxy-addr@5.1.0':
resolution: {integrity: sha512-INS+6gh91cLUjB+PVHfu1UqcB76Sqtpyp7bnL+FYojhjygvOPA9ctiD/JDKsyD9Xgu4hUhCSJBPig/w7duNajw==}
@@ -2649,10 +2614,6 @@ packages:
'@open-draft/until@2.1.0':
resolution: {integrity: sha512-U69T3ItWHvLwGg5eJ0n3I62nWuE6ilHlmz7zM0npLBRvPRd7e6NYmg54vvRtP5mZG7kZqZCFVdsTWo7BPtBujg==}
'@opentelemetry/api-logs@0.219.0':
resolution: {integrity: sha512-FFx7YnaYJlIjqWW/AG/yAZ0L/NEY724PipXXXQLdtZPbLwBGbUMTGL1i/esI56TWfTUXxhLfpgrnWJCG8aUJyg==}
engines: {node: '>=8.0.0'}
'@opentelemetry/api-logs@0.220.0':
resolution: {integrity: sha512-CmVa4ImJ+ynfrPMNaAXHET6Bhb44SwzmfyVJFq9ni2jgXJR/l7C6gfVFddNmHP+ZOkP9cf4f9DBe68qVLTHc9w==}
engines: {node: '>=8.0.0'}
@@ -2661,84 +2622,30 @@ packages:
resolution: {integrity: sha512-gLyJlPHPZYdAk1JENA9LeHejZe1Ti77/pTeFm/nMXmQH/HFZlcS/O2XJB+L8fkbrNSqhdtlvjBVjxwUYanNH5Q==}
engines: {node: '>=8.0.0'}
'@opentelemetry/context-async-hooks@2.9.0':
resolution: {integrity: sha512-OQ0vzvbZBiUhjqLnUaoNfYmP8553Crr3aggB4y0ZUi815mZ7idpdJXQmoKdeBKJelYttoBlLSSHubmyw3wvX4w==}
engines: {node: ^18.19.0 || >=20.6.0}
peerDependencies:
'@opentelemetry/api': '>=1.0.0 <1.10.0'
'@opentelemetry/core@2.9.0':
resolution: {integrity: sha512-m2nckMT80NnmjTYSPjJQObBJ+8dgkoajEOUbznL8AHZ3T3yHRk2P7gI1PhEBc1+lOnrYE9UWrWHqJDsmqjmNbw==}
engines: {node: ^18.19.0 || >=20.6.0}
peerDependencies:
'@opentelemetry/api': '>=1.0.0 <1.10.0'
'@opentelemetry/exporter-trace-otlp-proto@0.220.0':
resolution: {integrity: sha512-voTAD8XgJxlK7zLkXh8EzMB09zrQr3tyY/BsnDTlDiQU/UdK58MZ63A3mUjdEDrxMjCVmBHU3WQJhRmQe+Dvzg==}
engines: {node: ^18.19.0 || >=20.6.0}
peerDependencies:
'@opentelemetry/api': ^1.3.0
'@opentelemetry/instrumentation-pg@0.72.0':
resolution: {integrity: sha512-p9xrFc/6R8t6Y293sTYLZ83LnzZo/qY0bBPA4xabdQt0Qjt8i1SlYFsIeGY2Jmf5WcESNUdjQB3NxWnt5Ox7zw==}
engines: {node: ^18.19.0 || >=20.6.0}
peerDependencies:
'@opentelemetry/api': ^1.3.0
'@opentelemetry/instrumentation@0.219.0':
resolution: {integrity: sha512-X5t7I8GyIO9rmGHwoedZLREpQqrF1WW2nxzNNym6HOKpFiE+rvqV3ngC0xcZVO2YwIGf3KKmRdWrYwdwz3H9RQ==}
engines: {node: ^18.19.0 || >=20.6.0}
peerDependencies:
'@opentelemetry/api': ^1.3.0
'@opentelemetry/instrumentation@0.220.0':
resolution: {integrity: sha512-xQx3E2WxP1mDvKzxLxX+CTCtNLa560YJZ3087qYHerl2YmiKpv7AH+dAy7vmx+eVrZ5BwhfWUAVoKOoxCNHcpw==}
engines: {node: ^18.19.0 || >=20.6.0}
peerDependencies:
'@opentelemetry/api': ^1.3.0
'@opentelemetry/otlp-exporter-base@0.220.0':
resolution: {integrity: sha512-CXYo8UD5Mn9YbgebO2EL4wejtA+gxLmLiu6HCk2KH2BR7XhFN6/6p1UlCb23DYCjeYkndevLHuejCCN1yx4+OQ==}
engines: {node: ^18.19.0 || >=20.6.0}
peerDependencies:
'@opentelemetry/api': ^1.3.0
'@opentelemetry/otlp-transformer@0.220.0':
resolution: {integrity: sha512-lXGrv7KXZ0gNH9SVNUaa6vv6phVYGvJxfXAlMbzbakiXru75f5MZl8Z7oqiMMQD77riVHJCFlQvbZs/VVN2/4A==}
engines: {node: ^18.19.0 || >=20.6.0}
peerDependencies:
'@opentelemetry/api': ^1.3.0
'@opentelemetry/resources@2.9.0':
resolution: {integrity: sha512-jyA5MBLQ+Dkl3+JsZkUoUvL7yHvU64kLsvpXKarWm6347Sl1t1bXFTFykUePNpT5WH5pm9a2Qtt03iIYQhZ1Fg==}
engines: {node: ^18.19.0 || >=20.6.0}
peerDependencies:
'@opentelemetry/api': '>=1.3.0 <1.10.0'
'@opentelemetry/sdk-logs@0.220.0':
resolution: {integrity: sha512-WywcTkQtv2iNmt+6y5Kcd4rzvx9bLVsBa2Nwcmg01IUaBTkTow3W4d9KE5vNBpEDtb9tp21WcRBY/lANRrApYA==}
engines: {node: ^18.19.0 || >=20.6.0}
peerDependencies:
'@opentelemetry/api': '>=1.4.0 <1.10.0'
'@opentelemetry/sdk-metrics@2.9.0':
resolution: {integrity: sha512-Xx8RGS4H5XEBl01WuCreMIpiah9cCXMbSkeuIePPdD2cUpq/vUzYmj8E/MK1OsbOc93FuAD4jfn2WOacKwLn7Q==}
engines: {node: ^18.19.0 || >=20.6.0}
peerDependencies:
'@opentelemetry/api': '>=1.9.0 <1.10.0'
'@opentelemetry/sdk-trace-base@2.9.0':
resolution: {integrity: sha512-cp9zmTl62R8PJrpvFcmc8N2JQU/xfa0S+61q511Nji+QxCfZ8Ifvg7H27G8cANe4crg4RTrWsVvanHiXjSp6ag==}
engines: {node: ^18.19.0 || >=20.6.0}
peerDependencies:
'@opentelemetry/api': '>=1.3.0 <1.10.0'
'@opentelemetry/sdk-trace-node@2.9.0':
resolution: {integrity: sha512-ec9a7ps37huy5itYk0MalaZdSLlM6AXWp/FhtEjgMpp5leEGojBDvAl/UWttQnkMZOvFHKzRESn8TD3yKTF5nQ==}
engines: {node: ^18.19.0 || >=20.6.0}
peerDependencies:
'@opentelemetry/api': '>=1.0.0 <1.10.0'
'@opentelemetry/sdk-trace@2.9.0':
resolution: {integrity: sha512-sGA19HvtrrSKYsseHphluH6j3p6Xa3fqc7c7y8f/7mYWejc1lyDFcpSdD1kYa50HCLUeEo4zA5bW0pniaPszuw==}
engines: {node: ^18.19.0 || >=20.6.0}
@@ -2749,12 +2656,6 @@ packages:
resolution: {integrity: sha512-eSYWTm620tTk45EKSedaUL8MFYI8hW164hIXsgIHyxu3VobUB3fFCu5t0hQby6OoWRPsG1KkKUG2M5UadiLiVg==}
engines: {node: '>=14'}
'@opentelemetry/sql-common@0.42.0':
resolution: {integrity: sha512-nwUwUU+8O8a4bnLqk6CodWeegGMEANgC94KTAhXcpGWLrW/2/hek/0ajNbjXnSOoNuCX+nteUPs46HFHhou9Xw==}
engines: {node: ^18.19.0 || >=20.6.0}
peerDependencies:
'@opentelemetry/api': ^1.1.0
'@oxc-parser/binding-android-arm-eabi@0.127.0':
resolution: {integrity: sha512-0LC7ye4hvqbIKxAzThzvswgHLFu2AURKzYLeSVvLdu2TBOYWQDmHnTqPLeA597BcUCxiLqLsS4CJ5uoI5WYWCQ==}
engines: {node: ^20.19.0 || >=22.12.0}
@@ -4036,12 +3937,6 @@ packages:
'@types/offscreencanvas@2019.7.3':
resolution: {integrity: sha512-ieXiYmgSRXUDeOntE1InxjWyvEelZGP63M+cGuquuRLuIKKT1osnkXjxev9B7d1nXSug5vpunx+gNlbVxMlC9A==}
'@types/pg-pool@2.0.7':
resolution: {integrity: sha512-U4CwmGVQcbEuqpyju8/ptOKg6gEC+Tqsvj2xS9o1g71bUh8twxnC6ZL5rZKCsGN0iyH0CwgUyc9VR5owNQF9Ng==}
'@types/pg@8.15.6':
resolution: {integrity: sha512-NoaMtzhxOrubeL/7UZuNTrejB4MPAJ0RpxZqXQf2qXuVlTPuG6Y8p4u9dKRaue4yjmC7ZhzVO2/Yyyn25znrPQ==}
'@types/pg@8.20.0':
resolution: {integrity: sha512-bEPFOaMAHTEP1EzpvHTbmwR8UsFyHSKsRisLIHVMXnpNefSbGA1bD6CVy+qKjGSqmZqNqBDV2azOBo8TgkcVow==}
@@ -10055,16 +9950,6 @@ snapshots:
fastify-plugin: 6.0.0
secure-json-parse: 4.1.0
'@fastify/otel@0.20.1(@opentelemetry/api@1.9.1)':
dependencies:
'@opentelemetry/api': 1.9.1
'@opentelemetry/core': 2.9.0(@opentelemetry/api@1.9.1)
'@opentelemetry/instrumentation': 0.219.0(@opentelemetry/api@1.9.1)
'@opentelemetry/semantic-conventions': 1.43.0
minimatch: 10.2.5
transitivePeerDependencies:
- supports-color
'@fastify/proxy-addr@5.1.0':
dependencies:
'@fastify/forwarded': 3.0.1
@@ -10673,55 +10558,17 @@ snapshots:
'@open-draft/until@2.1.0': {}
'@opentelemetry/api-logs@0.219.0':
dependencies:
'@opentelemetry/api': 1.9.1
'@opentelemetry/api-logs@0.220.0':
dependencies:
'@opentelemetry/api': 1.9.1
'@opentelemetry/api@1.9.1': {}
'@opentelemetry/context-async-hooks@2.9.0(@opentelemetry/api@1.9.1)':
dependencies:
'@opentelemetry/api': 1.9.1
'@opentelemetry/core@2.9.0(@opentelemetry/api@1.9.1)':
dependencies:
'@opentelemetry/api': 1.9.1
'@opentelemetry/semantic-conventions': 1.43.0
'@opentelemetry/exporter-trace-otlp-proto@0.220.0(@opentelemetry/api@1.9.1)':
dependencies:
'@opentelemetry/api': 1.9.1
'@opentelemetry/core': 2.9.0(@opentelemetry/api@1.9.1)
'@opentelemetry/otlp-exporter-base': 0.220.0(@opentelemetry/api@1.9.1)
'@opentelemetry/otlp-transformer': 0.220.0(@opentelemetry/api@1.9.1)
'@opentelemetry/resources': 2.9.0(@opentelemetry/api@1.9.1)
'@opentelemetry/sdk-trace': 2.9.0(@opentelemetry/api@1.9.1)
'@opentelemetry/instrumentation-pg@0.72.0(@opentelemetry/api@1.9.1)':
dependencies:
'@opentelemetry/api': 1.9.1
'@opentelemetry/core': 2.9.0(@opentelemetry/api@1.9.1)
'@opentelemetry/instrumentation': 0.220.0(@opentelemetry/api@1.9.1)
'@opentelemetry/semantic-conventions': 1.43.0
'@opentelemetry/sql-common': 0.42.0(@opentelemetry/api@1.9.1)
'@types/pg': 8.15.6
'@types/pg-pool': 2.0.7
transitivePeerDependencies:
- supports-color
'@opentelemetry/instrumentation@0.219.0(@opentelemetry/api@1.9.1)':
dependencies:
'@opentelemetry/api': 1.9.1
'@opentelemetry/api-logs': 0.219.0
import-in-the-middle: 3.0.1
require-in-the-middle: 8.0.1
transitivePeerDependencies:
- supports-color
'@opentelemetry/instrumentation@0.220.0(@opentelemetry/api@1.9.1)':
dependencies:
'@opentelemetry/api': 1.9.1
@@ -10731,42 +10578,12 @@ snapshots:
transitivePeerDependencies:
- supports-color
'@opentelemetry/otlp-exporter-base@0.220.0(@opentelemetry/api@1.9.1)':
dependencies:
'@opentelemetry/api': 1.9.1
'@opentelemetry/core': 2.9.0(@opentelemetry/api@1.9.1)
'@opentelemetry/otlp-transformer': 0.220.0(@opentelemetry/api@1.9.1)
'@opentelemetry/otlp-transformer@0.220.0(@opentelemetry/api@1.9.1)':
dependencies:
'@opentelemetry/api': 1.9.1
'@opentelemetry/api-logs': 0.220.0
'@opentelemetry/core': 2.9.0(@opentelemetry/api@1.9.1)
'@opentelemetry/resources': 2.9.0(@opentelemetry/api@1.9.1)
'@opentelemetry/sdk-logs': 0.220.0(@opentelemetry/api@1.9.1)
'@opentelemetry/sdk-metrics': 2.9.0(@opentelemetry/api@1.9.1)
'@opentelemetry/sdk-trace': 2.9.0(@opentelemetry/api@1.9.1)
'@opentelemetry/resources@2.9.0(@opentelemetry/api@1.9.1)':
dependencies:
'@opentelemetry/api': 1.9.1
'@opentelemetry/core': 2.9.0(@opentelemetry/api@1.9.1)
'@opentelemetry/semantic-conventions': 1.43.0
'@opentelemetry/sdk-logs@0.220.0(@opentelemetry/api@1.9.1)':
dependencies:
'@opentelemetry/api': 1.9.1
'@opentelemetry/api-logs': 0.220.0
'@opentelemetry/core': 2.9.0(@opentelemetry/api@1.9.1)
'@opentelemetry/resources': 2.9.0(@opentelemetry/api@1.9.1)
'@opentelemetry/semantic-conventions': 1.43.0
'@opentelemetry/sdk-metrics@2.9.0(@opentelemetry/api@1.9.1)':
dependencies:
'@opentelemetry/api': 1.9.1
'@opentelemetry/core': 2.9.0(@opentelemetry/api@1.9.1)
'@opentelemetry/resources': 2.9.0(@opentelemetry/api@1.9.1)
'@opentelemetry/sdk-trace-base@2.9.0(@opentelemetry/api@1.9.1)':
dependencies:
'@opentelemetry/api': 1.9.1
@@ -10775,13 +10592,6 @@ snapshots:
'@opentelemetry/sdk-trace': 2.9.0(@opentelemetry/api@1.9.1)
'@opentelemetry/semantic-conventions': 1.43.0
'@opentelemetry/sdk-trace-node@2.9.0(@opentelemetry/api@1.9.1)':
dependencies:
'@opentelemetry/api': 1.9.1
'@opentelemetry/context-async-hooks': 2.9.0(@opentelemetry/api@1.9.1)
'@opentelemetry/core': 2.9.0(@opentelemetry/api@1.9.1)
'@opentelemetry/sdk-trace-base': 2.9.0(@opentelemetry/api@1.9.1)
'@opentelemetry/sdk-trace@2.9.0(@opentelemetry/api@1.9.1)':
dependencies:
'@opentelemetry/api': 1.9.1
@@ -10791,11 +10601,6 @@ snapshots:
'@opentelemetry/semantic-conventions@1.43.0': {}
'@opentelemetry/sql-common@0.42.0(@opentelemetry/api@1.9.1)':
dependencies:
'@opentelemetry/api': 1.9.1
'@opentelemetry/core': 2.9.0(@opentelemetry/api@1.9.1)
'@oxc-parser/binding-android-arm-eabi@0.127.0':
optional: true
@@ -12027,16 +11832,6 @@ snapshots:
'@types/offscreencanvas@2019.7.3': {}
'@types/pg-pool@2.0.7':
dependencies:
'@types/pg': 8.20.0
'@types/pg@8.15.6':
dependencies:
'@types/node': 26.1.1
pg-protocol: 1.15.0
pg-types: 2.2.0
'@types/pg@8.20.0':
dependencies:
'@types/node': 26.1.1