refactor(self-hosting): forward every setting, drop dead config (#3047)

This commit is contained in:
Hampus
2026-09-30 00:58:43 +02:00
committed by GitHub
parent 39f9beda5a
commit 2b8a743dc5
96 changed files with 1743 additions and 2391 deletions
+2 -6
View File
@@ -302,10 +302,8 @@ export class KVClient implements IKVProvider {
}
private createClusterClient(clusterConfig: ResolvedKVClientConfig): Cluster {
const {nodes, redisOptions} = resolveKVClusterConnection(clusterConfig.url, clusterConfig.clusterNodes);
const natMap = clusterConfig.clusterNatMap;
const hasNatMap = Object.keys(natMap).length > 0;
return new Cluster(nodes, {
const {node, redisOptions} = resolveKVClusterConnection(clusterConfig.url);
return new Cluster([node], {
clusterRetryStrategy: createRetryStrategy(),
redisOptions: {
...redisOptions,
@@ -315,7 +313,6 @@ export class KVClient implements IKVProvider {
protocol: 2,
},
scaleReads: 'master',
...(hasNatMap ? {natMap} : {}),
});
}
@@ -591,7 +588,6 @@ export class KVClient implements IKVProvider {
return new KVSubscription({
url: this.url,
mode: this.config.mode,
clusterNodes: this.config.clusterNodes,
timeoutMs: this.timeoutMs,
logger: this.logger,
});
@@ -9,16 +9,9 @@ export interface IKVLogger {
export type KVClientMode = 'standalone' | 'cluster';
export interface KVClusterNode {
host: string;
port: number;
}
export interface KVClientConfig {
url: string;
mode?: KVClientMode;
clusterNodes?: Array<KVClusterNode>;
clusterNatMap?: Record<string, KVClusterNode>;
timeoutMs?: number;
logger?: IKVLogger;
}
@@ -26,8 +19,6 @@ export interface KVClientConfig {
export interface ResolvedKVClientConfig {
url: string;
mode: KVClientMode;
clusterNodes: Array<KVClusterNode>;
clusterNatMap: Record<string, KVClusterNode>;
timeoutMs: number;
logger: IKVLogger;
}
@@ -42,8 +33,6 @@ export function resolveKVClientConfig(config: KVClientConfig | string): Resolved
return {
url: normalizeUrl(options.url),
mode: options.mode ?? 'standalone',
clusterNodes: options.clusterNodes ?? [],
clusterNatMap: options.clusterNatMap ?? {},
timeoutMs: options.timeoutMs ?? DEFAULT_KV_TIMEOUT_MS,
logger: options.logger ?? noopLogger,
};
@@ -1,13 +1,12 @@
import {domainToASCII} from 'node:url';
import type {KVClusterNode} from '@pkgs/kv_client/src/KVClientConfig';
import type {RedisOptions} from 'ioredis';
interface KVClusterConnection {
nodes: Array<KVClusterNode>;
node: {host: string; port: number};
redisOptions: RedisOptions;
}
export function resolveKVClusterConnection(url: string, nodes: ReadonlyArray<KVClusterNode>): KVClusterConnection {
export function resolveKVClusterConnection(url: string): KVClusterConnection {
const normalizedUrl = url.trim();
for (let index = 0; index < normalizedUrl.length; index++) {
const code = normalizedUrl.charCodeAt(index);
@@ -46,17 +45,8 @@ export function resolveKVClusterConnection(url: string, nodes: ReadonlyArray<KVC
redisOptions.tls = {};
}
const host = resolveClusterHost(authority, parsed);
const resolvedNodes = nodes.length > 0 ? [...nodes] : [{host, port: Number(parsed.port || '6379')}];
for (const node of resolvedNodes) {
if (node.host.trim().length === 0) {
throw new Error('KV cluster node must include a host');
}
if (!Number.isInteger(node.port) || node.port < 1 || node.port > 65535) {
throw new Error('KV cluster node port must be an integer between 1 and 65535');
}
}
return {
nodes: resolvedNodes,
node: {host, port: Number(parsed.port || '6379')},
redisOptions,
};
}
@@ -1,14 +1,13 @@
// SPDX-License-Identifier: AGPL-3.0-or-later
import type {IKVSubscription} from '@pkgs/kv_client/src/IKVProvider';
import type {IKVLogger, KVClientMode, KVClusterNode} from '@pkgs/kv_client/src/KVClientConfig';
import type {IKVLogger, KVClientMode} from '@pkgs/kv_client/src/KVClientConfig';
import {resolveKVClusterConnection} from '@pkgs/kv_client/src/KVClusterConnection';
import Redis, {type RedisOptions} from 'ioredis';
interface KVSubscriptionConfig {
url: string;
mode?: KVClientMode;
clusterNodes?: Array<KVClusterNode>;
timeoutMs: number;
logger: IKVLogger;
}
@@ -21,7 +20,6 @@ interface KVSubscriptionConnect {
export class KVSubscription implements IKVSubscription {
private readonly url: string;
private readonly mode: KVClientMode;
private readonly clusterNodes: Array<KVClusterNode>;
private readonly timeoutMs: number;
private readonly logger: IKVLogger;
private readonly desiredChannels = new Set<string>();
@@ -34,7 +32,6 @@ export class KVSubscription implements IKVSubscription {
constructor(config: KVSubscriptionConfig) {
this.url = config.url;
this.mode = config.mode ?? 'standalone';
this.clusterNodes = config.clusterNodes ?? [];
this.timeoutMs = config.timeoutMs;
this.logger = config.logger;
}
@@ -75,9 +72,9 @@ export class KVSubscription implements IKVSubscription {
protocol: 2,
retryStrategy: createRetryStrategy(),
};
const connection = this.mode === 'cluster' ? resolveKVClusterConnection(this.url, this.clusterNodes) : null;
const connection = this.mode === 'cluster' ? resolveKVClusterConnection(this.url) : null;
const client = connection
? new Redis({...connection.redisOptions, ...connection.nodes[0], db: 0, ...options})
? new Redis({...connection.redisOptions, ...connection.node, db: 0, ...options})
: new Redis(this.url, options);
client.on('message', (channel: string, message: string) => {
if (this.client !== client || this.closing !== null) {
-3
View File
@@ -2,10 +2,7 @@
"name": "@pkgs/mime_utils",
"version": "0.0.0",
"type": "module",
"main": "./src/index.ts",
"types": "./src/index.ts",
"exports": {
".": "./src/index.ts",
"./src/*": "./src/*",
"./*": "./*"
},
-7
View File
@@ -1,7 +0,0 @@
// SPDX-License-Identifier: AGPL-3.0-or-later
export {
getContentTypeFromFilename,
isSupportedMediaContentType,
normalizeContentType,
} from '@pkgs/mime_utils/src/ContentTypeUtils';
+1 -31
View File
@@ -4,10 +4,6 @@
"private": true,
"type": "module",
"exports": {
"./WorkerFactory": {
"import": "./src/runtime/WorkerFactory.ts",
"types": "./src/runtime/WorkerFactory.ts"
},
"./WorkerContext": {
"import": "./src/context/WorkerContext.ts",
"types": "./src/context/WorkerContext.ts"
@@ -24,39 +20,13 @@
"import": "./src/contracts/WorkerTypes.ts",
"types": "./src/contracts/WorkerTypes.ts"
},
"./IQueueProvider": {
"import": "./src/providers/IQueueProvider.ts",
"types": "./src/providers/IQueueProvider.ts"
},
"./HttpWorkerQueue": {
"import": "./src/providers/HttpWorkerQueue.ts",
"types": "./src/providers/HttpWorkerQueue.ts"
},
"./QueueProviderFactory": {
"import": "./src/providers/QueueProviderFactory.ts",
"types": "./src/providers/QueueProviderFactory.ts"
},
"./WorkerRunner": {
"import": "./src/runtime/WorkerRunner.ts",
"types": "./src/runtime/WorkerRunner.ts"
},
"./WorkerService": {
"import": "./src/services/WorkerService.ts",
"types": "./src/services/WorkerService.ts"
},
"./WorkerTaskRegistry": {
"import": "./src/runtime/WorkerTaskRegistry.ts",
"types": "./src/runtime/WorkerTaskRegistry.ts"
},
"./*": "./*"
},
"scripts": {
"typecheck": "tsc --noEmit"
},
"dependencies": {
"@fluxer/constants": "workspace:*",
"@fluxer/logger": "workspace:*",
"itty-time": "catalog:"
"@fluxer/logger": "workspace:*"
},
"devDependencies": {
"@types/node": "catalog:",
@@ -2,56 +2,6 @@
export type WorkerJobPayload = Record<string, unknown>;
export interface WorkerRuntimeConfig {
workerId?: string | undefined;
concurrency?: number | undefined;
taskTypes?: Array<string> | undefined;
}
export interface WorkerQueueConfig {
queueBaseUrl: string;
requestTimeoutMs?: number | undefined;
}
export interface WorkerConfig extends WorkerRuntimeConfig, WorkerQueueConfig {}
export interface TracingInterface {
withSpan<T>(
options: {
name: string;
attributes?: Record<string, unknown>;
},
fn: () => Promise<T>,
): Promise<T>;
addSpanEvent(name: string, attributes?: Record<string, unknown>): void;
setSpanAttributes(attributes: Record<string, unknown>): void;
}
export interface QueueJob {
id: string;
task_type: string;
payload: WorkerJobPayload;
priority: number;
run_at: string;
created_at: string;
attempts: number;
max_attempts: number;
error?: string | null;
deduplication_id?: string | null;
}
export interface LeasedQueueJob {
receipt: string;
visibility_deadline: string;
job: QueueJob;
}
export interface EnqueueOptions {
runAt?: Date | undefined;
maxAttempts?: number | undefined;
priority?: number | undefined;
}
export interface WorkerJobOptions {
queueName?: string | undefined;
runAt?: Date | undefined;
@@ -1,234 +0,0 @@
// SPDX-License-Identifier: AGPL-3.0-or-later
import {DEFAULT_HTTP_WORKER_TIMEOUT_MS} from '@fluxer/constants/src/Timeouts';
import type {
EnqueueOptions,
LeasedQueueJob,
TracingInterface,
WorkerJobPayload,
} from '@pkgs/worker/src/contracts/WorkerTypes';
import type {IQueueProvider} from '@pkgs/worker/src/providers/IQueueProvider';
export class HttpWorkerQueue implements IQueueProvider {
private readonly baseUrl: string;
private readonly timeoutMs: number;
private readonly tracing: TracingInterface | undefined;
constructor(options: {
baseUrl: string;
timeoutMs?: number | undefined;
tracing?: TracingInterface | undefined;
}) {
this.baseUrl = options.baseUrl;
this.timeoutMs = options.timeoutMs ?? DEFAULT_HTTP_WORKER_TIMEOUT_MS;
this.tracing = options.tracing;
}
private async withResponse<T>(
input: string | URL,
init: RequestInit,
consume: (response: Response) => Promise<T>,
): Promise<T> {
const controller = new AbortController();
const timeoutId = setTimeout(() => controller.abort(), this.timeoutMs);
try {
const response = await fetch(input, {
...init,
signal: controller.signal,
});
return await consume(response);
} finally {
clearTimeout(timeoutId);
controller.abort();
}
}
private async requireSuccess(response: Response, action: string): Promise<void> {
if (response.ok) {
return;
}
const text = await response.text();
throw new Error(`Failed to ${action}: ${response.status} ${text}`);
}
private async withOptionalSpan<T>(
options: {
name: string;
attributes?: Record<string, unknown>;
},
fn: () => Promise<T>,
): Promise<T> {
if (this.tracing) {
return this.tracing.withSpan(options, fn);
}
return fn();
}
private addSpanEvent(name: string, attributes?: Record<string, unknown>): void {
if (this.tracing) {
this.tracing.addSpanEvent(name, attributes);
}
}
private setSpanAttributes(attributes: Record<string, unknown>): void {
if (this.tracing) {
this.tracing.setSpanAttributes(attributes);
}
}
async enqueue(taskType: string, payload: WorkerJobPayload, options?: EnqueueOptions): Promise<string> {
return await this.withOptionalSpan(
{
name: 'queue.enqueue',
attributes: {
'queue.task_type': taskType,
'queue.priority': options?.priority ?? 0,
'queue.max_attempts': options?.maxAttempts ?? 5,
'queue.scheduled': options?.runAt !== undefined,
'net.peer.name': new URL(this.baseUrl).hostname,
},
},
async () => {
const body = {
task_type: taskType,
payload,
priority: options?.priority ?? 0,
run_at: options?.runAt?.toISOString(),
max_attempts: options?.maxAttempts ?? 5,
};
return this.withResponse(
`${this.baseUrl}/enqueue`,
{
method: 'POST',
headers: {'Content-Type': 'application/json'},
body: JSON.stringify(body),
},
async (response) => {
await this.requireSuccess(response, 'enqueue job');
const jobIdResult = (await response.json()) as {
job_id: string;
};
this.setSpanAttributes({'queue.job_id': jobIdResult.job_id});
return jobIdResult.job_id;
},
);
},
);
}
async dequeue(taskTypes: Array<string>, limit = 1): Promise<Array<LeasedQueueJob>> {
return await this.withOptionalSpan(
{
name: 'queue.dequeue',
attributes: {
'queue.task_types': taskTypes.join(','),
'queue.limit': limit,
'queue.service': 'fluxer-queue',
},
},
async () => {
this.addSpanEvent('dequeue.start');
const url = new URL(`${this.baseUrl}/dequeue`);
url.searchParams.set('task_types', taskTypes.join(','));
url.searchParams.set('limit', limit.toString());
url.searchParams.set('wait_time_ms', '0');
return this.withResponse(url, {method: 'GET'}, async (response) => {
await this.requireSuccess(response, 'dequeue job');
this.addSpanEvent('dequeue.parse_response');
const jobs = (await response.json()) as Array<LeasedQueueJob>;
const jobCount = jobs?.length ?? 0;
this.setSpanAttributes({
'queue.jobs_returned': jobCount,
'queue.empty': jobCount === 0,
});
this.addSpanEvent('dequeue.complete');
return jobs ?? [];
});
},
);
}
async upsertCron(id: string, taskType: string, payload: WorkerJobPayload, cronExpression: string): Promise<void> {
await this.withResponse(
`${this.baseUrl}/cron`,
{
method: 'POST',
headers: {'Content-Type': 'application/json'},
body: JSON.stringify({id, task_type: taskType, payload, cron_expression: cronExpression}),
},
(response) => this.requireSuccess(response, 'upsert cron job'),
);
}
async complete(receipt: string): Promise<void> {
return await this.withOptionalSpan(
{
name: 'queue.complete',
attributes: {
'queue.receipt': receipt,
},
},
async () =>
this.withResponse(
`${this.baseUrl}/ack`,
{
method: 'POST',
headers: {'Content-Type': 'application/json'},
body: JSON.stringify({receipt}),
},
(response) => this.requireSuccess(response, 'complete job'),
),
);
}
async fail(receipt: string, error: string): Promise<void> {
return await this.withOptionalSpan(
{
name: 'queue.fail',
attributes: {
'queue.receipt': receipt,
'queue.error_message': error,
},
},
async () =>
this.withResponse(
`${this.baseUrl}/nack`,
{
method: 'POST',
headers: {'Content-Type': 'application/json'},
body: JSON.stringify({receipt, error}),
},
(response) => this.requireSuccess(response, 'fail job'),
),
);
}
async cancelJob(jobId: string): Promise<boolean> {
return this.withResponse(`${this.baseUrl}/job/${jobId}`, {method: 'DELETE'}, async (response) => {
if (!response.ok) {
const text = await response.text();
if (response.status === 404) {
return false;
}
throw new Error(`Failed to cancel job: ${response.status} ${text}`);
}
const result = (await response.json()) as {
success: boolean;
};
return result.success ?? true;
});
}
async retryDeadLetterJob(jobId: string): Promise<boolean> {
return this.withResponse(`${this.baseUrl}/retry/${jobId}`, {method: 'POST'}, async (response) => {
if (!response.ok) {
const text = await response.text();
if (response.status === 404) {
return false;
}
throw new Error(`Failed to retry job: ${response.status} ${text}`);
}
return true;
});
}
}
@@ -1,13 +0,0 @@
// SPDX-License-Identifier: AGPL-3.0-or-later
import type {EnqueueOptions, LeasedQueueJob, WorkerJobPayload} from '@pkgs/worker/src/contracts/WorkerTypes';
export interface IQueueProvider {
enqueue(taskType: string, payload: WorkerJobPayload, options?: EnqueueOptions): Promise<string>;
dequeue(taskTypes: Array<string>, limit?: number): Promise<Array<LeasedQueueJob>>;
upsertCron(id: string, taskType: string, payload: WorkerJobPayload, cronExpression: string): Promise<void>;
complete(receipt: string): Promise<void>;
fail(receipt: string, error: string): Promise<void>;
cancelJob(jobId: string): Promise<boolean>;
retryDeadLetterJob(jobId: string): Promise<boolean>;
}
@@ -1,26 +0,0 @@
// SPDX-License-Identifier: AGPL-3.0-or-later
import type {TracingInterface} from '@pkgs/worker/src/contracts/WorkerTypes';
import {HttpWorkerQueue} from '@pkgs/worker/src/providers/HttpWorkerQueue';
import type {IQueueProvider} from '@pkgs/worker/src/providers/IQueueProvider';
export interface QueueProviderFactoryOptions {
queueProvider?: IQueueProvider | undefined;
queueBaseUrl?: string | undefined;
timeoutMs?: number | undefined;
tracing?: TracingInterface | undefined;
}
export function createQueueProvider(options: QueueProviderFactoryOptions): IQueueProvider {
if (options.queueProvider) {
return options.queueProvider;
}
if (!options.queueBaseUrl) {
throw new Error('Queue provider requires either queueProvider or queueBaseUrl');
}
return new HttpWorkerQueue({
baseUrl: options.queueBaseUrl,
timeoutMs: options.timeoutMs,
tracing: options.tracing,
});
}
@@ -1,175 +0,0 @@
// SPDX-License-Identifier: AGPL-3.0-or-later
import type {LoggerInterface} from '@fluxer/logger/src/LoggerInterface';
import {setWorkerDependencies} from '@pkgs/worker/src/context/WorkerContext';
import type {IWorkerService} from '@pkgs/worker/src/contracts/IWorkerService';
import type {WorkerTaskHandler} from '@pkgs/worker/src/contracts/WorkerTask';
import type {
LeasedQueueJob,
TracingInterface,
WorkerConfig,
WorkerQueueConfig,
WorkerRuntimeConfig,
} from '@pkgs/worker/src/contracts/WorkerTypes';
import type {IQueueProvider} from '@pkgs/worker/src/providers/IQueueProvider';
import {WorkerRunner} from '@pkgs/worker/src/runtime/WorkerRunner';
import {WorkerTaskRegistry} from '@pkgs/worker/src/runtime/WorkerTaskRegistry';
import {WorkerService} from '@pkgs/worker/src/services/WorkerService';
export interface CreateWorkerOptions {
queue: WorkerQueueOptions;
runtime?: WorkerRuntimeConfig | undefined;
logger: LoggerInterface;
dependencies?: unknown;
taskRegistry?: WorkerTaskRegistry | undefined;
tracing?: TracingInterface | undefined;
}
export interface CreateWorkerLegacyOptions {
config: WorkerConfig;
queueProvider?: IQueueProvider | undefined;
logger: LoggerInterface;
dependencies?: unknown;
taskRegistry?: WorkerTaskRegistry | undefined;
tracing?: TracingInterface | undefined;
}
export interface WorkerQueueOptions {
queueProvider?: IQueueProvider | undefined;
queueBaseUrl?: string | undefined;
requestTimeoutMs?: number | undefined;
}
interface ResolvedWorkerFactoryOptions {
queue: WorkerQueueOptions;
runtime: WorkerRuntimeConfig;
logger: LoggerInterface;
dependencies?: unknown;
taskRegistry?: WorkerTaskRegistry | undefined;
tracing?: TracingInterface | undefined;
}
export interface WorkerResult {
start: () => Promise<void>;
shutdown: () => Promise<void>;
processTask: (job: LeasedQueueJob) => Promise<void>;
getRunner: () => WorkerRunner;
getWorkerService: () => IWorkerService;
registerTask: <TPayload = Record<string, unknown>>(name: string, handler: WorkerTaskHandler<TPayload>) => void;
registerTasks: (tasks: Record<string, WorkerTaskHandler>) => void;
}
type WorkerFactoryOptions = CreateWorkerOptions | CreateWorkerLegacyOptions;
function isLegacyCreateWorkerOptions(options: WorkerFactoryOptions): options is CreateWorkerLegacyOptions {
return 'config' in options;
}
function resolveLegacyQueueOptions(config: WorkerQueueConfig, queueProvider?: IQueueProvider): WorkerQueueOptions {
return {
queueProvider,
queueBaseUrl: config.queueBaseUrl,
requestTimeoutMs: config.requestTimeoutMs,
};
}
function resolveWorkerFactoryOptions(options: WorkerFactoryOptions): ResolvedWorkerFactoryOptions {
if (isLegacyCreateWorkerOptions(options)) {
return {
queue: resolveLegacyQueueOptions(options.config, options.queueProvider),
runtime: {
workerId: options.config.workerId,
taskTypes: options.config.taskTypes,
concurrency: options.config.concurrency,
},
logger: options.logger,
dependencies: options.dependencies,
taskRegistry: options.taskRegistry,
tracing: options.tracing,
};
}
return {
queue: options.queue,
runtime: options.runtime ?? {},
logger: options.logger,
dependencies: options.dependencies,
taskRegistry: options.taskRegistry,
tracing: options.tracing,
};
}
function assertTaskRegistryMutable(runner: WorkerRunner | null): void {
if (runner?.isRunning()) {
throw new Error('Cannot register tasks after worker start. Register tasks before starting the worker.');
}
}
export function createWorker(options: WorkerFactoryOptions): WorkerResult {
const resolvedOptions = resolveWorkerFactoryOptions(options);
const {queue, runtime, logger, dependencies, taskRegistry: providedRegistry, tracing} = resolvedOptions;
if (dependencies !== undefined) {
setWorkerDependencies(dependencies);
}
const taskRegistry = providedRegistry ?? new WorkerTaskRegistry();
let runner: WorkerRunner | null = null;
let workerService: WorkerService | null = null;
function ensureRunner(): WorkerRunner {
if (!runner) {
runner = new WorkerRunner({
tasks: taskRegistry.getTasks(),
queueBaseUrl: queue.queueBaseUrl,
queueProvider: queue.queueProvider,
logger,
workerId: runtime.workerId,
taskTypes: runtime.taskTypes,
concurrency: runtime.concurrency,
tracing,
requestTimeoutMs: queue.requestTimeoutMs,
});
}
return runner;
}
function ensureWorkerService(): WorkerService {
if (!workerService) {
workerService = new WorkerService({
queueBaseUrl: queue.queueBaseUrl,
queueProvider: queue.queueProvider,
logger,
tracing,
timeoutMs: queue.requestTimeoutMs,
});
}
return workerService;
}
return {
async start() {
const r = ensureRunner();
await r.start();
},
async shutdown() {
if (runner) {
await runner.stop();
}
},
async processTask(job: LeasedQueueJob) {
const r = ensureRunner();
await r.processJob(job);
},
getRunner() {
return ensureRunner();
},
getWorkerService() {
return ensureWorkerService();
},
registerTask<TPayload = Record<string, unknown>>(name: string, handler: WorkerTaskHandler<TPayload>) {
assertTaskRegistryMutable(runner);
taskRegistry.register(name, handler);
runner = null;
},
registerTasks(tasks: Record<string, WorkerTaskHandler>) {
assertTaskRegistryMutable(runner);
taskRegistry.registerAll(tasks);
runner = null;
},
};
}
@@ -1,187 +0,0 @@
// SPDX-License-Identifier: AGPL-3.0-or-later
import {randomUUID} from 'node:crypto';
import type {LoggerInterface} from '@fluxer/logger/src/LoggerInterface';
import type {WorkerTaskHandler} from '@pkgs/worker/src/contracts/WorkerTask';
import type {LeasedQueueJob, TracingInterface} from '@pkgs/worker/src/contracts/WorkerTypes';
import type {IQueueProvider} from '@pkgs/worker/src/providers/IQueueProvider';
import {createQueueProvider} from '@pkgs/worker/src/providers/QueueProviderFactory';
import {WorkerService} from '@pkgs/worker/src/services/WorkerService';
import {ms} from 'itty-time';
export interface WorkerRunnerOptions {
tasks: Record<string, WorkerTaskHandler>;
queueBaseUrl?: string | undefined;
queueProvider?: IQueueProvider | undefined;
logger: LoggerInterface;
workerId?: string | undefined;
taskTypes?: Array<string> | undefined;
concurrency?: number | undefined;
tracing?: TracingInterface | undefined;
requestTimeoutMs?: number | undefined;
}
export class WorkerRunner {
private readonly tasks: Record<string, WorkerTaskHandler>;
private readonly workerId: string;
private readonly taskTypes: Array<string>;
private readonly concurrency: number;
private readonly queue: IQueueProvider;
private readonly workerService: WorkerService;
private readonly logger: LoggerInterface;
private readonly tracing: TracingInterface | undefined;
private running = false;
private abortController: AbortController | null = null;
private workerLoopPromises: Array<Promise<void>> = [];
constructor(options: WorkerRunnerOptions) {
this.tasks = options.tasks;
this.workerId = options.workerId ?? `worker-${randomUUID()}`;
this.taskTypes = options.taskTypes ?? Object.keys(options.tasks);
this.concurrency = options.concurrency ?? 1;
this.logger = options.logger;
this.tracing = options.tracing;
this.queue = createQueueProvider({
queueProvider: options.queueProvider,
queueBaseUrl: options.queueBaseUrl,
timeoutMs: options.requestTimeoutMs,
tracing: options.tracing,
});
this.workerService = new WorkerService({
queueProvider: this.queue,
logger: options.logger,
});
}
async start(): Promise<void> {
if (this.running) {
this.logger.warn({workerId: this.workerId}, 'Worker already running');
return;
}
this.running = true;
this.abortController = new AbortController();
this.logger.info(
{workerId: this.workerId, taskTypes: this.taskTypes, concurrency: this.concurrency},
'Worker starting',
);
this.workerLoopPromises = Array.from({length: this.concurrency}, (_, i) =>
this.workerLoop(i, this.abortController!.signal),
);
Promise.all(this.workerLoopPromises).catch((error) => {
this.logger.error({workerId: this.workerId, error}, 'Worker loop failed unexpectedly');
});
}
async stop(): Promise<void> {
if (!this.running) {
return;
}
this.running = false;
this.abortController?.abort();
const stopTimeout = new Promise<void>((resolve) => setTimeout(resolve, ms('2 seconds')));
await Promise.race([Promise.all(this.workerLoopPromises), stopTimeout]);
this.workerLoopPromises = [];
this.logger.info({workerId: this.workerId}, 'Worker stopped');
}
async processJob(leasedJob: LeasedQueueJob): Promise<void> {
await this.executeJob(leasedJob);
}
getWorkerService(): WorkerService {
return this.workerService;
}
getQueue(): IQueueProvider {
return this.queue;
}
isRunning(): boolean {
return this.running;
}
private async workerLoop(workerIndex: number, signal: AbortSignal): Promise<void> {
this.logger.info({workerId: this.workerId, workerIndex}, 'Worker loop started');
while (!signal.aborted) {
try {
const leasedJobs = await this.queue.dequeue(this.taskTypes, 1);
if (!leasedJobs || leasedJobs.length === 0) {
await this.sleep(100);
continue;
}
const leasedJob = leasedJobs[0]!;
const job = leasedJob.job;
this.logger.info(
{
workerId: this.workerId,
workerIndex,
jobId: job.id,
taskType: job.task_type,
attempts: job.attempts,
receipt: leasedJob.receipt,
},
'Processing job',
);
await this.executeJob(leasedJob);
this.logger.info({workerId: this.workerId, workerIndex, jobId: job.id}, 'Job completed successfully');
} catch (error) {
this.logger.error({workerId: this.workerId, workerIndex, error}, 'Worker loop error');
await this.sleep(ms('1 second'));
}
}
this.logger.info({workerId: this.workerId, workerIndex}, 'Worker loop stopped');
}
private async executeJob(leasedJob: LeasedQueueJob): Promise<void> {
const execute = async () => {
const task = this.tasks[leasedJob.job.task_type];
if (!task) {
throw new Error(`Unknown task: ${leasedJob.job.task_type}`);
}
this.tracing?.addSpanEvent('job.execution.start');
try {
await task(leasedJob.job.payload, {
logger: this.logger.child({jobId: leasedJob.job.id}),
jobId: 0n,
addJob: this.workerService.addJob.bind(this.workerService),
reportProgress: async () => {},
shouldCancel: async () => false,
setContextLink: async () => {},
});
this.tracing?.addSpanEvent('job.execution.success');
this.tracing?.setSpanAttributes({'job.status': 'success'});
await this.queue.complete(leasedJob.receipt);
} catch (error) {
this.logger.error({jobId: leasedJob.job.id, error}, 'Job failed');
this.tracing?.setSpanAttributes({
'job.status': 'failed',
'job.error': error instanceof Error ? error.message : String(error),
});
this.tracing?.addSpanEvent('job.execution.failed', {
error: error instanceof Error ? error.message : String(error),
});
await this.queue.fail(leasedJob.receipt, String(error));
}
};
if (this.tracing) {
await this.tracing.withSpan(
{
name: 'worker.process_job',
attributes: {
'worker.id': this.workerId,
'job.id': leasedJob.job.id,
'job.task_type': leasedJob.job.task_type,
'job.attempts': leasedJob.job.attempts,
},
},
execute,
);
} else {
await execute();
}
}
private async sleep(ms: number): Promise<void> {
return new Promise((resolve) => setTimeout(resolve, ms));
}
}
@@ -1,43 +0,0 @@
// SPDX-License-Identifier: AGPL-3.0-or-later
import type {WorkerTaskHandler} from '@pkgs/worker/src/contracts/WorkerTask';
export class WorkerTaskRegistry {
private readonly tasks: Map<string, WorkerTaskHandler> = new Map();
register<TPayload = Record<string, unknown>>(name: string, handler: WorkerTaskHandler<TPayload>): this {
this.tasks.set(name, handler as WorkerTaskHandler);
return this;
}
registerAll(tasks: Record<string, WorkerTaskHandler>): this {
for (const [name, handler] of Object.entries(tasks)) {
this.tasks.set(name, handler);
}
return this;
}
get(name: string): WorkerTaskHandler | undefined {
return this.tasks.get(name);
}
has(name: string): boolean {
return this.tasks.has(name);
}
getTaskNames(): Array<string> {
return Array.from(this.tasks.keys());
}
getTasks(): Record<string, WorkerTaskHandler> {
return Object.fromEntries(this.tasks);
}
get size(): number {
return this.tasks.size;
}
}
export function createTaskRegistry(): WorkerTaskRegistry {
return new WorkerTaskRegistry();
}
@@ -1,83 +0,0 @@
// SPDX-License-Identifier: AGPL-3.0-or-later
import type {LoggerInterface} from '@fluxer/logger/src/LoggerInterface';
import type {IWorkerService} from '@pkgs/worker/src/contracts/IWorkerService';
import type {TracingInterface, WorkerJobOptions, WorkerJobPayload} from '@pkgs/worker/src/contracts/WorkerTypes';
import type {IQueueProvider} from '@pkgs/worker/src/providers/IQueueProvider';
import {createQueueProvider} from '@pkgs/worker/src/providers/QueueProviderFactory';
export interface WorkerServiceOptions {
queueBaseUrl?: string | undefined;
queueProvider?: IQueueProvider | undefined;
logger: LoggerInterface;
tracing?: TracingInterface | undefined;
timeoutMs?: number | undefined;
}
export class WorkerService implements IWorkerService {
private readonly queue: IQueueProvider;
private readonly logger: LoggerInterface;
constructor(options: WorkerServiceOptions) {
this.queue = createQueueProvider({
queueProvider: options.queueProvider,
queueBaseUrl: options.queueBaseUrl,
timeoutMs: options.timeoutMs,
tracing: options.tracing,
});
this.logger = options.logger;
}
async addJob<TPayload extends WorkerJobPayload = WorkerJobPayload>(
taskType: string,
payload: TPayload,
options?: WorkerJobOptions,
): Promise<bigint> {
try {
await this.queue.enqueue(taskType, payload, {
runAt: options?.runAt,
maxAttempts: options?.maxAttempts,
priority: options?.priority,
});
this.logger.debug({taskType, payload}, 'Job queued successfully');
} catch (error) {
this.logger.error({error, taskType, payload}, 'Failed to queue job');
throw error;
}
return 0n;
}
async cancelJob(jobId: bigint): Promise<boolean> {
try {
const cancelled = await this.queue.cancelJob(jobId.toString());
if (cancelled) {
this.logger.info({jobId: jobId.toString()}, 'Job cancelled successfully');
} else {
this.logger.debug({jobId: jobId.toString()}, 'Job not found (may have already been processed)');
}
return cancelled;
} catch (error) {
this.logger.error({error, jobId: jobId.toString()}, 'Failed to cancel job');
throw error;
}
}
async retryDeadLetterJob(jobId: bigint): Promise<boolean> {
try {
const retried = await this.queue.retryDeadLetterJob(jobId.toString());
if (retried) {
this.logger.info({jobId: jobId.toString()}, 'Dead letter job retried successfully');
} else {
this.logger.debug({jobId: jobId.toString()}, 'Job not found in dead letter queue');
}
return retried;
} catch (error) {
this.logger.error({error, jobId: jobId.toString()}, 'Failed to retry dead letter job');
throw error;
}
}
getQueue(): IQueueProvider {
return this.queue;
}
}