1
0
mirror of https://github.com/misskey-dev/misskey.git synced 2026-07-29 02:34:38 +02:00

enhance: OTel設定時、JSONログにtrace_id, span_id, trace_flagがつくように (#17748)

* enhance: OTel設定時、JSONログにtrace_id, span_id, trace_flagがつくように

* fix
This commit is contained in:
おさむのひと
2026-07-20 19:00:54 +09:00
committed by GitHub
parent 36c442a7b4
commit bbd57c751c
17 changed files with 449 additions and 89 deletions

View File

@@ -11,6 +11,7 @@ import { installHttpClientInstrumentation } from '@/core/telemetry/http-client-i
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';
@@ -161,6 +162,15 @@ export class OpenTelemetryAdapter implements TelemetryAdapter {
});
}
/** 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));

View File

@@ -6,6 +6,7 @@
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';
@@ -175,6 +176,15 @@ export class SentryTelemetryAdapter implements TelemetryAdapter {
});
}
/** activeなSpanの識別子を、Logging基盤で扱える形式へ変換します。 */
public getActiveTraceContext(): LogTraceContext | undefined {
const activeSpan = this.Sentry.getActiveSpan();
if (activeSpan == null) return undefined;
const { traceId, spanId, traceFlags } = activeSpan.spanContext();
return { traceId, spanId, traceFlags };
}
public startSpan<T>(name: string, fn: () => T): T {
return this.Sentry.startSpan({ name }, fn);
}

View File

@@ -4,6 +4,7 @@
*/
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']>;
@@ -35,6 +36,9 @@ export interface TelemetryAdapter {
*/
captureMessage(message: string, opts: TelemetryCaptureMessageOptions): void;
/** 現在のactive Spanからログへ付加するTrace Contextを取得する。 */
getActiveTraceContext?(): LogTraceContext | undefined;
/**
* API endpointやqueue jobなど、呼び出し側の処理単位をspanで包む。
* fnの戻り値・例外はそのまま呼び出し側へ返し、Promiseの場合はsettleまでspanを閉じない。

View File

@@ -4,6 +4,7 @@
*/
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';
@@ -23,14 +24,21 @@ export async function initTelemetry(config: Config): Promise<void> {
};
// SentryとOTelを同時に使う場合はproviderを分けず、Sentry側へOTLP processorを追加する。
let adapter: TelemetryAdapter | undefined;
if (config.sentryForBackend && otelForBackend) {
adapters.push(await SentryTelemetryAdapter.createWithOtlpExport(config.sentryForBackend, otelForBackend));
adapter = await SentryTelemetryAdapter.createWithOtlpExport(config.sentryForBackend, otelForBackend);
} else if (config.sentryForBackend) {
// Sentry単体時は既存のSentry adapterだけを登録する。
adapters.push(await SentryTelemetryAdapter.create(config.sentryForBackend));
adapter = await SentryTelemetryAdapter.create(config.sentryForBackend);
} else if (otelForBackend) {
// OTel単体時だけMisskey自身でNodeTracerProviderを立てる。
adapters.push(await OpenTelemetryAdapter.create(otelForBackend));
adapter = await OpenTelemetryAdapter.create(otelForBackend);
}
if (adapter != null) {
adapters.push(adapter);
// Telemetryの初期化後に登録し、初期化前のBootstrapログは従来どおり出力する。
setLogTraceContextProvider(() => adapter.getActiveTraceContext?.());
}
}

View File

@@ -23,6 +23,9 @@ type JsonLogRecord = {
readonly processId: number;
readonly isPrimary: boolean;
readonly workerId: number | null;
readonly trace_id?: string;
readonly span_id?: string;
readonly trace_flags?: number;
};
const defaultDependencies: JsonConsoleBackendDependencies = {
@@ -45,6 +48,10 @@ function createJsonLogRecord(record: LogRecord): JsonLogRecord {
processId: record.processId,
isPrimary: record.isPrimary,
workerId: record.workerId,
// active Spanの情報は値が存在するログだけ、標準のsnake_case名で出力します。
...(record.traceId != null ? { trace_id: record.traceId } : {}),
...(record.spanId != null ? { span_id: record.spanId } : {}),
...(record.traceFlags != null ? { trace_flags: record.traceFlags } : {}),
};
}

View File

@@ -13,7 +13,7 @@ import {
type LogNormalizationProfile,
} from './LogNormalizer.js';
import type { LogBackend } from './LogBackend.js';
import type { LogLevel, LogLevelSetting, LogRecord, LogRecordInput } from './types.js';
import type { LogLevel, LogLevelSetting, LogRecord, LogRecordInput, LogTraceContextProvider } from './types.js';
/** ログを出力したプロセスを識別するための情報です。 */
export type LogProcessInfo = {
@@ -113,6 +113,7 @@ export class LogManager {
private backend: LogBackend;
private readonly dependencies: LogManagerDependencies;
private normalizationProfile: LogNormalizationProfile;
private traceContextProvider: LogTraceContextProvider | undefined;
private configuredLevel: LogLevelSetting | undefined;
private configuredDomains: readonly (readonly [string, LogLevelSetting])[];
private shutdownPromise: Promise<void> | undefined;
@@ -132,6 +133,7 @@ export class LogManager {
...dependencies,
};
this.normalizationProfile = options.normalizationProfile ?? 'standard';
this.traceContextProvider = undefined;
this.configuredLevel = undefined;
this.configuredDomains = [];
}
@@ -156,6 +158,11 @@ export class LogManager {
this.normalizationProfile = profile;
}
/** ログ出力時にactiveなTrace Contextを取得する処理を登録します。 */
public setTraceContextProvider(provider?: LogTraceContextProvider): void {
this.traceContextProvider = provider;
}
/** backendに残っているログをflushしてから終了処理を行います。 */
public shutdown(): Promise<void> {
if (this.shutdownPromise != null) return this.shutdownPromise;
@@ -220,6 +227,8 @@ export class LogManager {
const normalizedError = typeof error !== 'undefined'
? serializeLogError(error, { profile: this.normalizationProfile })
: undefined;
// 実際に出力するログだけ、TelemetryからactiveなTrace Contextを取得します。
const traceContext = this.traceContextProvider?.();
const record = {
...inputWithoutStructuredValues,
context,
@@ -228,6 +237,7 @@ export class LogManager {
processId: processInfo.processId,
isPrimary: processInfo.isPrimary,
workerId: processInfo.workerId,
...(traceContext ?? {}),
...(normalizedAttributes ? { attributes: normalizedAttributes } : {}),
...(normalizedError ? { error: normalizedError } : {}),
} as LogRecord;

View File

@@ -8,8 +8,8 @@ import { BootstrapConsoleBackend } from './BootstrapConsoleBackend.js';
import { JsonConsoleBackend } from './JsonConsoleBackend.js';
import { PrettyConsoleBackend } from './PrettyConsoleBackend.js';
import type { LogManagerConfiguration } from './LogManager.js';
import type { LogTraceContextProvider, LogFormat } from './types.js';
import type { LogBackend } from './LogBackend.js';
import type { LogFormat } from './types.js';
/**
* プロセス内のすべてのLoggerが共有するLogManagerです。
@@ -32,6 +32,11 @@ export function configureLogging(configuration?: LogManagerConfiguration & { rea
logManager.setBackend(backend);
}
/** Telemetry初期化後に、ログへTrace Contextを付加する取得処理を登録します。 */
export function setLogTraceContextProvider(provider?: LogTraceContextProvider): void {
logManager.setTraceContextProvider(provider);
}
/** プロセス終了前に現在のログ出力処理を保留分まで書き出して閉じます。 */
export function shutdownLogging(): Promise<void> {
return logManager.shutdown();

View File

@@ -26,6 +26,16 @@ export type LogAttributeValue =
/** 正規化後のログ属性です。 */
export type LogAttributes = Readonly<Record<string, LogAttributeValue>>;
/** activeなSpanとログを関連付けるためのTrace Contextです。 */
export type LogTraceContext = {
readonly traceId: string;
readonly spanId: string;
readonly traceFlags: number;
};
/** LogManagerが出力直前に呼び出すTrace Context取得処理です。 */
export type LogTraceContextProvider = () => LogTraceContext | undefined;
/** ロガーの呼び出し側が構造化ログとして指定する入力です。 */
export type LogWriteInput = {
readonly level: LogLevel;
@@ -81,6 +91,9 @@ export type LogRecord = Omit<LogRecordInput, 'attributes' | 'error'> & {
readonly processId: number;
readonly isPrimary: boolean;
readonly workerId: number | null;
readonly traceId?: string;
readonly spanId?: string;
readonly traceFlags?: number;
readonly attributes?: LogAttributes;
readonly error?: SerializedError;
};

View File

@@ -11,6 +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 { UserWebhookDeliverProcessorService } from './processors/UserWebhookDeliverProcessorService.js';
import { SystemWebhookDeliverProcessorService } from './processors/SystemWebhookDeliverProcessorService.js';
import { EndedPollNotificationProcessorService } from './processors/EndedPollNotificationProcessorService.js';
@@ -176,26 +177,30 @@ export class QueueProcessorService implements OnApplicationShutdown {
default: throw new Error(`unrecognized job type ${job.name} for system`);
}
};
const logger = this.logger.createSubLogger('system');
this.systemQueueWorker = new Bull.Worker(QUEUE.SYSTEM, (job) => {
return this.telemetryService.startSpanWithTraceContext('Queue: System: ' + job.name, job.data, () => processer(job));
return runQueueJobWithTraceContext(
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) });
this.telemetryService.captureMessage(`Queue: System: ${job.name}: ${err.name}: ${err.message}`, {
level: 'error',
extra: { job, err },
});
},
);
}, {
...baseWorkerOptions(this.config, QUEUE.SYSTEM),
autorun: false,
});
const logger = this.logger.createSubLogger('system');
this.systemQueueWorker
.on('active', (job) => logger.debug(`active id=${job.id}`))
.on('completed', (job, result) => logger.debug(`completed(${result}) id=${job.id}`))
.on('failed', (job, err: Error) => {
logger.error(`failed(${err.name}: ${err.message}) id=${job?.id ?? '?'}`, { job: renderJob(job), e: renderError(err) });
this.telemetryService.captureMessage(`Queue: System: ${job?.name ?? '?'}: ${err.name}: ${err.message}`, {
level: 'error',
extra: { job, err },
});
})
.on('error', (err: Error) => logger.error(`error ${err.name}: ${err.message}`, { e: renderError(err) }))
.on('stalled', (jobId) => logger.warn(`stalled id=${jobId}`));
}
@@ -227,26 +232,30 @@ export class QueueProcessorService implements OnApplicationShutdown {
default: throw new Error(`unrecognized job type ${job.name} for db`);
}
};
const logger = this.logger.createSubLogger('db');
this.dbQueueWorker = new Bull.Worker(QUEUE.DB, (job) => {
return this.telemetryService.startSpanWithTraceContext('Queue: DB: ' + job.name, job.data, () => processer(job));
return runQueueJobWithTraceContext(
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) });
this.telemetryService.captureMessage(`Queue: DB: ${job.name}: ${err.name}: ${err.message}`, {
level: 'error',
extra: { job, err },
});
},
);
}, {
...baseWorkerOptions(this.config, QUEUE.DB),
autorun: false,
});
const logger = this.logger.createSubLogger('db');
this.dbQueueWorker
.on('active', (job) => logger.debug(`active id=${job.id}`))
.on('completed', (job, result) => logger.debug(`completed(${result}) id=${job.id}`))
.on('failed', (job, err) => {
logger.error(`failed(${err.name}: ${err.message}) id=${job?.id ?? '?'}`, { job: renderJob(job), e: renderError(err) });
this.telemetryService.captureMessage(`Queue: DB: ${job?.name ?? '?'}: ${err.name}: ${err.message}`, {
level: 'error',
extra: { job, err },
});
})
.on('error', (err: Error) => logger.error(`error ${err.name}: ${err.message}`, { e: renderError(err) }))
.on('stalled', (jobId) => logger.warn(`stalled id=${jobId}`));
}
@@ -254,8 +263,22 @@ export class QueueProcessorService implements OnApplicationShutdown {
//#region deliver
{
const logger = this.logger.createSubLogger('deliver');
this.deliverQueueWorker = new Bull.Worker(QUEUE.DELIVER, (job) => {
return this.telemetryService.startSpanWithTraceContext('Queue: Deliver', job.data, () => this.deliverProcessorService.process(job));
return runQueueJobWithTraceContext(
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) });
this.telemetryService.captureMessage(`Queue: Deliver: ${err.name}: ${err.message}`, {
level: 'error',
extra: { job, err },
});
},
);
}, {
...baseWorkerOptions(this.config, QUEUE.DELIVER),
autorun: false,
@@ -269,18 +292,9 @@ export class QueueProcessorService implements OnApplicationShutdown {
},
});
const logger = this.logger.createSubLogger('deliver');
this.deliverQueueWorker
.on('active', (job) => logger.debug(`active ${getJobInfo(job, true)} to=${job.data.to}`))
.on('completed', (job, result) => logger.debug(`completed(${result}) ${getJobInfo(job, true)} to=${job.data.to}`))
.on('failed', (job, err) => {
logger.error(`failed(${err.name}: ${err.message}) ${getJobInfo(job)} to=${job ? job.data.to : '-'}`);
this.telemetryService.captureMessage(`Queue: Deliver: ${err.name}: ${err.message}`, {
level: 'error',
extra: { job, err },
});
})
.on('error', (err: Error) => logger.error(`error ${err.name}: ${err.message}`, { e: renderError(err) }))
.on('stalled', (jobId) => logger.warn(`stalled id=${jobId}`));
}
@@ -288,8 +302,23 @@ export class QueueProcessorService implements OnApplicationShutdown {
//#region inbox
{
const logger = this.logger.createSubLogger('inbox');
this.inboxQueueWorker = new Bull.Worker(QUEUE.INBOX, (job) => {
return this.telemetryService.startSpanWithTraceContext('Queue: Inbox', job.data, () => this.inboxProcessorService.process(job));
return runQueueJobWithTraceContext(
this.telemetryService,
'Queue: Inbox',
job.data,
() => this.inboxProcessorService.process(job),
err => {
const activityId = job.data.activity ? job.data.activity.id : 'none';
logger.error(`failed(${err.name}: ${err.message}) ${getJobInfo(job)} activity=${activityId}`, { job: renderJob(job), e: renderError(err) });
this.telemetryService.captureMessage(`Queue: Inbox: ${err.name}: ${err.message}`, {
level: 'error',
extra: { job, err },
});
},
);
}, {
...baseWorkerOptions(this.config, QUEUE.INBOX),
autorun: false,
@@ -303,18 +332,9 @@ export class QueueProcessorService implements OnApplicationShutdown {
},
});
const logger = this.logger.createSubLogger('inbox');
this.inboxQueueWorker
.on('active', (job) => logger.debug(`active ${getJobInfo(job, true)}`))
.on('completed', (job, result) => logger.debug(`completed(${result}) ${getJobInfo(job, true)}`))
.on('failed', (job, err) => {
logger.error(`failed(${err.name}: ${err.message}) ${getJobInfo(job)} activity=${job ? (job.data.activity ? job.data.activity.id : 'none') : '-'}`, { job: renderJob(job), e: renderError(err) });
this.telemetryService.captureMessage(`Queue: Inbox: ${err.name}: ${err.message}`, {
level: 'error',
extra: { job, err },
});
})
.on('error', (err: Error) => logger.error(`error ${err.name}: ${err.message}`, { e: renderError(err) }))
.on('stalled', (jobId) => logger.warn(`stalled id=${jobId}`));
}
@@ -322,8 +342,22 @@ export class QueueProcessorService implements OnApplicationShutdown {
//#region user-webhook deliver
{
const logger = this.logger.createSubLogger('user-webhook');
this.userWebhookDeliverQueueWorker = new Bull.Worker(QUEUE.USER_WEBHOOK_DELIVER, (job) => {
return this.telemetryService.startSpanWithTraceContext('Queue: UserWebhookDeliver', job.data, () => this.userWebhookDeliverProcessorService.process(job));
return runQueueJobWithTraceContext(
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) });
this.telemetryService.captureMessage(`Queue: UserWebhookDeliver: ${err.name}: ${err.message}`, {
level: 'error',
extra: { job, err },
});
},
);
}, {
...baseWorkerOptions(this.config, QUEUE.USER_WEBHOOK_DELIVER),
autorun: false,
@@ -337,18 +371,9 @@ export class QueueProcessorService implements OnApplicationShutdown {
},
});
const logger = this.logger.createSubLogger('user-webhook');
this.userWebhookDeliverQueueWorker
.on('active', (job) => logger.debug(`active ${getJobInfo(job, true)} to=${job.data.to}`))
.on('completed', (job, result) => logger.debug(`completed(${result}) ${getJobInfo(job, true)} to=${job.data.to}`))
.on('failed', (job, err) => {
logger.error(`failed(${err.name}: ${err.message}) ${getJobInfo(job)} to=${job ? job.data.to : '-'}`);
this.telemetryService.captureMessage(`Queue: UserWebhookDeliver: ${err.name}: ${err.message}`, {
level: 'error',
extra: { job, err },
});
})
.on('error', (err: Error) => logger.error(`error ${err.name}: ${err.message}`, { e: renderError(err) }))
.on('stalled', (jobId) => logger.warn(`stalled id=${jobId}`));
}
@@ -356,8 +381,22 @@ export class QueueProcessorService implements OnApplicationShutdown {
//#region system-webhook deliver
{
const logger = this.logger.createSubLogger('system-webhook');
this.systemWebhookDeliverQueueWorker = new Bull.Worker(QUEUE.SYSTEM_WEBHOOK_DELIVER, (job) => {
return this.telemetryService.startSpanWithTraceContext('Queue: SystemWebhookDeliver', job.data, () => this.systemWebhookDeliverProcessorService.process(job));
return runQueueJobWithTraceContext(
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) });
this.telemetryService.captureMessage(`Queue: SystemWebhookDeliver: ${err.name}: ${err.message}`, {
level: 'error',
extra: { job, err },
});
},
);
}, {
...baseWorkerOptions(this.config, QUEUE.SYSTEM_WEBHOOK_DELIVER),
autorun: false,
@@ -371,18 +410,9 @@ export class QueueProcessorService implements OnApplicationShutdown {
},
});
const logger = this.logger.createSubLogger('system-webhook');
this.systemWebhookDeliverQueueWorker
.on('active', (job) => logger.debug(`active ${getJobInfo(job, true)} to=${job.data.to}`))
.on('completed', (job, result) => logger.debug(`completed(${result}) ${getJobInfo(job, true)} to=${job.data.to}`))
.on('failed', (job, err) => {
logger.error(`failed(${err.name}: ${err.message}) ${getJobInfo(job)} to=${job ? job.data.to : '-'}`);
this.telemetryService.captureMessage(`Queue: SystemWebhookDeliver: ${err.name}: ${err.message}`, {
level: 'error',
extra: { job, err },
});
})
.on('error', (err: Error) => logger.error(`error ${err.name}: ${err.message}`, { e: renderError(err) }))
.on('stalled', (jobId) => logger.warn(`stalled id=${jobId}`));
}
@@ -399,9 +429,22 @@ export class QueueProcessorService implements OnApplicationShutdown {
default: throw new Error(`unrecognized job type ${job.name} for relationship`);
}
};
const logger = this.logger.createSubLogger('relationship');
this.relationshipQueueWorker = new Bull.Worker(QUEUE.RELATIONSHIP, (job) => {
return this.telemetryService.startSpanWithTraceContext('Queue: Relationship: ' + job.name, job.data, () => processer(job));
return runQueueJobWithTraceContext(
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) });
this.telemetryService.captureMessage(`Queue: Relationship: ${job.name}: ${err.name}: ${err.message}`, {
level: 'error',
extra: { job, err },
});
},
);
}, {
...baseWorkerOptions(this.config, QUEUE.RELATIONSHIP),
autorun: false,
@@ -412,18 +455,9 @@ export class QueueProcessorService implements OnApplicationShutdown {
},
});
const logger = this.logger.createSubLogger('relationship');
this.relationshipQueueWorker
.on('active', (job) => logger.debug(`active id=${job.id}`))
.on('completed', (job, result) => logger.debug(`completed(${result}) id=${job.id}`))
.on('failed', (job, err) => {
logger.error(`failed(${err.name}: ${err.message}) id=${job?.id ?? '?'}`, { job: renderJob(job), e: renderError(err) });
this.telemetryService.captureMessage(`Queue: Relationship: ${job?.name ?? '?'}: ${err.name}: ${err.message}`, {
level: 'error',
extra: { job, err },
});
})
.on('error', (err: Error) => logger.error(`error ${err.name}: ${err.message}`, { e: renderError(err) }))
.on('stalled', (jobId) => logger.warn(`stalled id=${jobId}`));
}
@@ -438,27 +472,31 @@ export class QueueProcessorService implements OnApplicationShutdown {
default: throw new Error(`unrecognized job type ${job.name} for objectStorage`);
}
};
const logger = this.logger.createSubLogger('objectStorage');
this.objectStorageQueueWorker = new Bull.Worker(QUEUE.OBJECT_STORAGE, (job) => {
return this.telemetryService.startSpanWithTraceContext('Queue: ObjectStorage: ' + job.name, job.data, () => processer(job));
return runQueueJobWithTraceContext(
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) });
this.telemetryService.captureMessage(`Queue: ObjectStorage: ${job.name}: ${err.name}: ${err.message}`, {
level: 'error',
extra: { job, err },
});
},
);
}, {
...baseWorkerOptions(this.config, QUEUE.OBJECT_STORAGE),
autorun: false,
concurrency: 16,
});
const logger = this.logger.createSubLogger('objectStorage');
this.objectStorageQueueWorker
.on('active', (job) => logger.debug(`active id=${job.id}`))
.on('completed', (job, result) => logger.debug(`completed(${result}) id=${job.id}`))
.on('failed', (job, err) => {
logger.error(`failed(${err.name}: ${err.message}) id=${job?.id ?? '?'}`, { job: renderJob(job), e: renderError(err) });
this.telemetryService.captureMessage(`Queue: ObjectStorage: ${job?.name ?? '?'}: ${err.name}: ${err.message}`, {
level: 'error',
extra: { job, err },
});
})
.on('error', (err: Error) => logger.error(`error ${err.name}: ${err.message}`, { e: renderError(err) }))
.on('stalled', (jobId) => logger.warn(`stalled id=${jobId}`));
}
@@ -466,8 +504,22 @@ export class QueueProcessorService implements OnApplicationShutdown {
//#region ended poll notification
{
const logger = this.logger.createSubLogger('ended-poll-notification');
this.endedPollNotificationQueueWorker = new Bull.Worker(QUEUE.ENDED_POLL_NOTIFICATION, (job) => {
return this.telemetryService.startSpanWithTraceContext('Queue: EndedPollNotification', job.data, () => this.endedPollNotificationProcessorService.process(job));
return runQueueJobWithTraceContext(
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) });
this.telemetryService.captureMessage(`Queue: EndedPollNotification: ${err.name}: ${err.message}`, {
level: 'error',
extra: { job, err },
});
},
);
}, {
...baseWorkerOptions(this.config, QUEUE.ENDED_POLL_NOTIFICATION),
autorun: false,
@@ -477,8 +529,22 @@ export class QueueProcessorService implements OnApplicationShutdown {
//#region post scheduled note
{
this.postScheduledNoteQueueWorker = new Bull.Worker(QUEUE.POST_SCHEDULED_NOTE, async (job) => {
return this.telemetryService.startSpanWithTraceContext('Queue: PostScheduledNote', job.data, () => this.postScheduledNoteProcessorService.process(job));
const logger = this.logger.createSubLogger('post-scheduled-note');
this.postScheduledNoteQueueWorker = new Bull.Worker(QUEUE.POST_SCHEDULED_NOTE, (job) => {
return runQueueJobWithTraceContext(
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) });
this.telemetryService.captureMessage(`Queue: PostScheduledNote: ${err.name}: ${err.message}`, {
level: 'error',
extra: { job, err },
});
},
);
}, {
...baseWorkerOptions(this.config, QUEUE.POST_SCHEDULED_NOTE),
autorun: false,

View File

@@ -0,0 +1,32 @@
/*
* SPDX-FileCopyrightText: syuilo and misskey-project
* SPDX-License-Identifier: AGPL-3.0-only
*/
import type { TelemetryService } from '@/core/telemetry/TelemetryService.js';
type QueueTelemetryService = Pick<TelemetryService, 'startSpanWithTraceContext'>;
/** QueueのprocessorをTrace Context付きで実行し、失敗処理をSpan内で行います。 */
export function runQueueJobWithTraceContext<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> => {
try {
return await processJob();
} catch (error) {
// 失敗イベントを待たず、processor Spanがactiveな間にログと通知を行います。
const normalizedError = error instanceof Error ? error : new Error(String(error));
try {
onError(normalizedError);
} catch {
// 失敗ログの処理が例外を投げても、Queueへは元のエラーを返します。
}
throw error;
}
});
}