2024-09-17 01:14:12 -07:00
|
|
|
import { GlobalConfig } from '@n8n/config';
|
2024-09-18 02:16:17 -07:00
|
|
|
import type { Application } from 'express';
|
2024-09-17 01:14:12 -07:00
|
|
|
import express from 'express';
|
2024-09-17 01:38:47 -07:00
|
|
|
import { InstanceSettings } from 'n8n-core';
|
2024-09-17 01:14:12 -07:00
|
|
|
import { ensureError } from 'n8n-workflow';
|
|
|
|
import { strict as assert } from 'node:assert';
|
|
|
|
import http from 'node:http';
|
|
|
|
import type { Server } from 'node:http';
|
|
|
|
import { Service } from 'typedi';
|
|
|
|
|
|
|
|
import { CredentialsOverwrites } from '@/credentials-overwrites';
|
|
|
|
import * as Db from '@/db';
|
|
|
|
import { CredentialsOverwritesAlreadySetError } from '@/errors/credentials-overwrites-already-set.error';
|
|
|
|
import { NonJsonBodyError } from '@/errors/non-json-body.error';
|
|
|
|
import { PortTakenError } from '@/errors/port-taken.error';
|
|
|
|
import { ServiceUnavailableError } from '@/errors/response-errors/service-unavailable.error';
|
|
|
|
import { ExternalHooks } from '@/external-hooks';
|
|
|
|
import type { ICredentialsOverwrite } from '@/interfaces';
|
2024-10-01 03:16:09 -07:00
|
|
|
import { Logger } from '@/logging/logger.service';
|
2024-09-18 02:16:17 -07:00
|
|
|
import { PrometheusMetricsService } from '@/metrics/prometheus-metrics.service';
|
2024-09-17 01:14:12 -07:00
|
|
|
import { rawBodyReader, bodyParser } from '@/middlewares';
|
|
|
|
import * as ResponseHelper from '@/response-helper';
|
|
|
|
import { ScalingService } from '@/scaling/scaling.service';
|
|
|
|
|
2024-09-18 02:16:17 -07:00
|
|
|
export type WorkerServerEndpointsConfig = {
|
|
|
|
/** Whether the `/healthz` endpoint is enabled. */
|
|
|
|
health: boolean;
|
|
|
|
|
|
|
|
/** Whether the [credentials overwrites endpoint](https://docs.n8n.io/embed/configuration/#credential-overwrites) is enabled. */
|
|
|
|
overwrites: boolean;
|
|
|
|
|
|
|
|
/** Whether the `/metrics` endpoint is enabled. */
|
|
|
|
metrics: boolean;
|
|
|
|
};
|
|
|
|
|
2024-09-17 01:14:12 -07:00
|
|
|
/**
|
|
|
|
* Responsible for handling HTTP requests sent to a worker.
|
|
|
|
*/
|
|
|
|
@Service()
|
|
|
|
export class WorkerServer {
|
|
|
|
private readonly port: number;
|
|
|
|
|
|
|
|
private readonly server: Server;
|
|
|
|
|
2024-09-18 02:16:17 -07:00
|
|
|
private readonly app: Application;
|
|
|
|
|
|
|
|
private endpointsConfig: WorkerServerEndpointsConfig;
|
|
|
|
|
2024-09-17 01:14:12 -07:00
|
|
|
private overwritesLoaded = false;
|
|
|
|
|
|
|
|
constructor(
|
|
|
|
private readonly globalConfig: GlobalConfig,
|
|
|
|
private readonly logger: Logger,
|
|
|
|
private readonly scalingService: ScalingService,
|
|
|
|
private readonly credentialsOverwrites: CredentialsOverwrites,
|
|
|
|
private readonly externalHooks: ExternalHooks,
|
2024-09-17 01:38:47 -07:00
|
|
|
private readonly instanceSettings: InstanceSettings,
|
2024-09-18 02:16:17 -07:00
|
|
|
private readonly prometheusMetricsService: PrometheusMetricsService,
|
2024-09-17 01:14:12 -07:00
|
|
|
) {
|
2024-09-17 01:38:47 -07:00
|
|
|
assert(this.instanceSettings.instanceType === 'worker');
|
2024-09-17 01:14:12 -07:00
|
|
|
|
2024-09-18 02:16:17 -07:00
|
|
|
this.app = express();
|
2024-09-17 01:14:12 -07:00
|
|
|
|
2024-09-18 02:16:17 -07:00
|
|
|
this.app.disable('x-powered-by');
|
2024-09-17 01:14:12 -07:00
|
|
|
|
2024-09-18 02:16:17 -07:00
|
|
|
this.server = http.createServer(this.app);
|
2024-09-17 01:14:12 -07:00
|
|
|
|
|
|
|
this.port = this.globalConfig.queue.health.port;
|
|
|
|
|
|
|
|
this.server.on('error', (error: NodeJS.ErrnoException) => {
|
|
|
|
if (error.code === 'EADDRINUSE') throw new PortTakenError(this.port);
|
|
|
|
});
|
2024-09-18 02:16:17 -07:00
|
|
|
}
|
2024-09-17 01:14:12 -07:00
|
|
|
|
2024-09-18 02:16:17 -07:00
|
|
|
async init(endpointsConfig: WorkerServerEndpointsConfig) {
|
|
|
|
assert(Object.values(endpointsConfig).some((e) => e));
|
2024-09-17 01:14:12 -07:00
|
|
|
|
2024-09-18 02:16:17 -07:00
|
|
|
this.endpointsConfig = endpointsConfig;
|
|
|
|
|
|
|
|
await this.mountEndpoints();
|
2024-09-17 01:14:12 -07:00
|
|
|
|
|
|
|
await new Promise<void>((resolve) => this.server.listen(this.port, resolve));
|
|
|
|
|
|
|
|
await this.externalHooks.run('worker.ready');
|
|
|
|
|
|
|
|
this.logger.info(`\nn8n worker server listening on port ${this.port}`);
|
|
|
|
}
|
|
|
|
|
2024-09-18 02:16:17 -07:00
|
|
|
private async mountEndpoints() {
|
|
|
|
if (this.endpointsConfig.health) {
|
|
|
|
this.app.get('/healthz', async (req, res) => await this.healthcheck(req, res));
|
|
|
|
}
|
|
|
|
|
|
|
|
if (this.endpointsConfig.overwrites) {
|
|
|
|
const { endpoint } = this.globalConfig.credentials.overwrite;
|
|
|
|
|
|
|
|
this.app.post(`/${endpoint}`, rawBodyReader, bodyParser, (req, res) =>
|
|
|
|
this.handleOverwrites(req, res),
|
|
|
|
);
|
|
|
|
}
|
|
|
|
|
|
|
|
if (this.endpointsConfig.metrics) {
|
|
|
|
await this.prometheusMetricsService.init(this.app);
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
2024-09-17 01:14:12 -07:00
|
|
|
private async healthcheck(_req: express.Request, res: express.Response) {
|
|
|
|
this.logger.debug('[WorkerServer] Health check started');
|
|
|
|
|
|
|
|
try {
|
|
|
|
await Db.getConnection().query('SELECT 1');
|
|
|
|
} catch (value) {
|
|
|
|
this.logger.error('[WorkerServer] No database connection', ensureError(value));
|
|
|
|
|
|
|
|
return ResponseHelper.sendErrorResponse(
|
|
|
|
res,
|
|
|
|
new ServiceUnavailableError('No database connection'),
|
|
|
|
);
|
|
|
|
}
|
|
|
|
|
|
|
|
try {
|
|
|
|
await this.scalingService.pingQueue();
|
|
|
|
} catch (value) {
|
|
|
|
this.logger.error('[WorkerServer] No Redis connection', ensureError(value));
|
|
|
|
|
|
|
|
return ResponseHelper.sendErrorResponse(
|
|
|
|
res,
|
|
|
|
new ServiceUnavailableError('No Redis connection'),
|
|
|
|
);
|
|
|
|
}
|
|
|
|
|
|
|
|
this.logger.debug('[WorkerServer] Health check succeeded');
|
|
|
|
|
|
|
|
ResponseHelper.sendSuccessResponse(res, { status: 'ok' }, true, 200);
|
|
|
|
}
|
|
|
|
|
|
|
|
private handleOverwrites(
|
|
|
|
req: express.Request<{}, {}, ICredentialsOverwrite>,
|
|
|
|
res: express.Response,
|
|
|
|
) {
|
|
|
|
if (this.overwritesLoaded) {
|
|
|
|
ResponseHelper.sendErrorResponse(res, new CredentialsOverwritesAlreadySetError());
|
|
|
|
return;
|
|
|
|
}
|
|
|
|
|
|
|
|
if (req.contentType !== 'application/json') {
|
|
|
|
ResponseHelper.sendErrorResponse(res, new NonJsonBodyError());
|
|
|
|
return;
|
|
|
|
}
|
|
|
|
|
|
|
|
this.credentialsOverwrites.setData(req.body);
|
|
|
|
|
|
|
|
this.overwritesLoaded = true;
|
|
|
|
|
|
|
|
ResponseHelper.sendSuccessResponse(res, { success: true }, true, 200);
|
|
|
|
}
|
|
|
|
}
|