orbit-engine
v1.1.0
Published
Production-grade Node.js concurrency engine with worker thread pools, bulkhead isolation, backpressure, circuit breakers, and event loop protection.
Maintainers
Readme
Orbit Engine
A production-grade Node.js concurrency engine built on worker_threads. Orbit provides isolated worker pools, backpressure, circuit breakers, admission control, and event loop protection — all in a single, zero-external-dependency package.
Features
- Worker Thread Pool — fixed-size pools with round-robin or least-busy scheduling
- Bulkhead Isolation — named pools (
cpu,io,critical) with independent queues and workers - Backpressure — bounded task queues with configurable policies (reject, drop-oldest, drop-newest)
- Circuit Breaker — automatic failure detection with configurable thresholds and recovery
- Admission Control — 5-point pipeline: runtime state, circuit breaker, queue capacity, event loop lag, per-task concurrency
- Event Loop Protection —
perf_hooks-based lag monitoring with configurable thresholds - Work Stealing — idle pools can steal from overloaded neighbors (opt-in per pool)
- Crash Recovery — supervisor pattern with restart budgets and automatic worker respawning
- Dead Letter Queue — failed/rejected tasks are captured for inspection
- Runtime State Machine — INITIALIZING → RUNNING → DEGRADED/OVERLOADED → STOPPING → STOPPED
- NestJS Integration — optional
@Parallel()decorator
Installation
npm install orbit-engineRequires Node.js >= 18.
Contributors: See docs/README.md for internal architecture, task lifecycle, and contribution guides.
Quick Start
import { Scheduler } from "orbit-engine";
// 1. Initialize with pool configuration
const scheduler = Scheduler.init({
workers: 4,
maxQueue: 200,
eventLoopLagThresholdMs: 100,
pools: {
cpu: { workers: 2, allowStealing: true },
io: { workers: 2, allowStealing: true },
critical: { workers: 1, allowStealing: false },
},
});
// 2. Register named tasks (handlers run inside worker threads)
scheduler.registerTask("processImage", processImage, { pool: "cpu" });
scheduler.registerTask("fetchData", fetchData, { pool: "io", maxConcurrent: 10 });
scheduler.registerTask("healthCheck", healthCheck, { pool: "critical" });
// 3. Run tasks
const result = await scheduler.run("processImage", { path: "/img/photo.jpg" });Configuration
interface SchedulerConfig {
workers: number; // fallback pool size
maxQueue: number; // per-pool queue capacity
eventLoopLagThresholdMs?: number; // lag threshold for admission rejection
circuitBreaker?: {
failureThreshold: number; // consecutive failures before opening
resetTimeoutMs: number; // time before half-open retry
};
backpressure?: {
policy: "reject" | "drop_oldest" | "drop_newest";
};
pools?: Record<string, {
workers: number;
allowStealing?: boolean; // default: true
}>;
maxPools?: number; // guard against pool explosion (default: 10)
observability?: {
hooks?: ObservabilityHooks; // lifecycle event handlers
enableDetailedMetrics?: boolean; // emit pool state change events
correlationIdHeaderName?: string; // extract correlation ID from payload key
};
}Observability
Orbit exposes optional hooks for integrating with production logging, metrics, and tracing systems (Pino, Winston, OpenTelemetry, Prometheus, etc.). Observability is fully optional — existing code works without it.
Hook overview
| Hook | When it fires |
|------|---------------|
| onTaskAdmitted | Task passes admission control and is queued |
| onTaskRejected | Task rejected by admission or enqueue failure |
| onTaskStarted | Task dispatched to a worker thread |
| onTaskCompleted | Task finished successfully |
| onTaskFailed | Task failed with an error |
| onTaskTimeout | Task exceeded timeoutMs |
| onPoolQueueFull | Pool queue reached capacity |
| onPoolStateChange | Queue or worker stats changed (requires enableDetailedMetrics) |
| onWorkerStart | Worker thread spawned |
| onWorkerCrash | Worker crashed |
| onWorkerRestart | Worker respawned after crash |
| onSchedulerStateChange | Runtime state transition |
| onCircuitBreakerStateChange | Circuit breaker opened/closed/half-open |
| onEventLoopLagWarning | Event loop lag exceeded threshold |
| onBackpressureTriggered | Backpressure policy activated |
| onOverloadDetected | Scheduler entered OVERLOADED state |
| onTaskDLQd | Task saved to dead letter queue |
Each event includes timestamp, taskId, taskName, and optional correlationId (extracted from payload).
Pino logging
import pino from "pino";
import { Scheduler } from "orbit-engine";
const logger = pino();
const scheduler = Scheduler.init({
workers: 4,
maxQueue: 200,
observability: {
hooks: {
onTaskCompleted: (event) => {
logger.info({ ...event }, `Task ${event.taskName} completed`);
},
onTaskFailed: (event) => {
logger.error({ ...event }, `Task ${event.taskName} failed`);
},
onCircuitBreakerStateChange: (event) => {
logger.warn({ ...event }, `Circuit breaker ${event.state}`);
},
},
},
});Prometheus metrics
import { Counter, Histogram } from "prom-client";
import { Scheduler } from "orbit-engine";
const taskCounter = new Counter({
name: "orbit_tasks_total",
help: "Total tasks by status",
labelNames: ["task_name", "status", "pool_name"],
});
const taskDuration = new Histogram({
name: "orbit_task_duration_ms",
help: "Task duration in milliseconds",
labelNames: ["task_name", "pool_name"],
});
const scheduler = Scheduler.init({
workers: 4,
maxQueue: 200,
observability: {
hooks: {
onTaskCompleted: (event) => {
taskCounter.inc({
task_name: event.taskName,
status: "success",
pool_name: event.poolName ?? "default",
});
taskDuration.observe(
{ task_name: event.taskName, pool_name: event.poolName ?? "default" },
event.durationMs,
);
},
onTaskFailed: (event) => {
taskCounter.inc({
task_name: event.taskName,
status: "failed",
pool_name: event.poolName ?? "default",
});
},
},
},
});OpenTelemetry tracing
import { trace } from "@opentelemetry/api";
import { Scheduler } from "orbit-engine";
const tracer = trace.getTracer("orbit-engine");
const activeSpans = new Map<string, { end: () => void }>();
const scheduler = Scheduler.init({
workers: 4,
maxQueue: 200,
observability: {
hooks: {
onTaskStarted: (event) => {
const span = tracer.startSpan(`task.${event.taskName}`);
span.setAttributes({
"task.id": event.taskId,
"pool.name": event.poolName ?? "default",
});
activeSpans.set(event.taskId, span);
},
onTaskCompleted: (event) => {
const span = activeSpans.get(event.taskId);
if (span) {
span.setAttribute("task.duration_ms", event.durationMs);
span.end();
activeSpans.delete(event.taskId);
}
},
onTaskFailed: (event) => {
const span = activeSpans.get(event.taskId);
if (span) {
span.recordException(new Error(event.error));
span.end();
activeSpans.delete(event.taskId);
}
},
},
},
});Correlation IDs
Pass correlationId on the task payload, or configure a custom header-style key:
Scheduler.init({
workers: 4,
maxQueue: 200,
observability: {
correlationIdHeaderName: "x-request-id",
hooks: { /* ... */ },
},
});
await scheduler.run("fetchData", {
url: "/api/users",
"x-request-id": req.headers["x-request-id"],
});Best practices
- Don't throw in hooks — hook errors are logged but should not affect task execution.
- Keep hooks fast — synchronous hooks are preferred; async hooks run fire-and-forget.
- Use structured events — each hook receives a rich event object; format in your integration layer.
- Enable detailed metrics sparingly —
enableDetailedMetricsemitsonPoolStateChangeon every queue mutation.
See examples/observability.ts for full integration examples.
API
Scheduler.init(config): Scheduler
Initialize the singleton scheduler. Call once at application startup.
scheduler.registerTask(name, handler, options?)
Register a named task. Options:
| Option | Type | Description |
|--------|------|-------------|
| pool | string | Target pool name (e.g. "cpu", "io") |
| maxConcurrent | number | Max simultaneous executions of this task |
scheduler.run<P, R>(name, payload, options?): Promise<R>
Submit a task for execution. Options:
| Option | Type | Description |
|--------|------|-------------|
| timeoutMs | number | Task timeout in milliseconds |
scheduler.getHealth()
Returns per-pool queue sizes, worker counts, busy workers, restart counts, circuit breaker state, and event loop lag.
scheduler.getState(): RuntimeState
Returns the current runtime state: RUNNING, DEGRADED, OVERLOADED, STOPPING, or STOPPED.
scheduler.shutdown(): Promise<void>
Graceful shutdown — waits for in-flight tasks, then terminates all workers.
Architecture
Producer (HTTP / job / decorator)
│
▼
┌──────────────────────────────┐
│ Scheduler │ ← thin orchestrator
│ ┌────────────────────────┐ │
│ │ Admission Controller │ │ ← 5-point pipeline
│ └────────────────────────┘ │
│ ┌────────────────────────┐ │
│ │ Runtime Controller │ │ ← state machine
│ └────────────────────────┘ │
│ ┌────────────────────────┐ │
│ │ Pool Manager │ │
│ │ ┌──────┐ ┌──────┐ │ │
│ │ │ cpu │ │ io │ … │ │ ← bulkheads
│ │ │Queue │ │Queue │ │ │
│ │ │Pool │ │Pool │ │ │
│ │ └──────┘ └──────┘ │ │
│ └────────────────────────┘ │
└──────────────────────────────┘
│
▼
Worker Threads (static task registry, no eval)Project Structure
orbit/
├── src/
│ ├── index.ts # package entry (re-exports public API)
│ ├── public/ # stable user-facing API
│ │ ├── index.ts # errors, types, logger exports
│ │ ├── types.ts
│ │ ├── errors.ts
│ │ └── logger.ts
│ ├── scheduler/ # orchestration
│ │ ├── scheduler.ts
│ │ ├── admission-controller.ts
│ │ ├── runtime-controller.ts
│ │ └── dead-letter-store.ts
│ ├── pool/ # bulkheads & worker pools
│ │ ├── pool-manager.ts
│ │ ├── worker-pool.ts
│ │ ├── task-queue.ts
│ │ └── scheduling-strategy.ts
│ ├── config/ # validation & defaults
│ │ ├── schema.ts
│ │ └── input-validation.ts
│ ├── monitoring/ # health & telemetry
│ │ ├── event-loop-monitor.ts
│ │ └── observability.ts
│ ├── runtime/ # worker thread entry (compiled to dist/runtime/)
│ │ ├── worker.ts
│ │ ├── task-registry.ts
│ │ └── builtin-tasks.ts
│ └── integrations/
│ └── nestjs/
│ └── parallel.decorator.ts
├── examples/
│ └── observability.ts # Pino, Prometheus, OTel examples
├── docs/ # contributor internals (why/what/how/lifecycle)
│ ├── README.md # doc index
│ ├── 01-why-orbit-exists.md
│ ├── 03-task-lifecycle.md
│ └── …
├── __tests__/
│ ├── fixtures/ # shared test task handlers
│ ├── unit/
│ ├── integration/
│ └── stress/
├── benchmarks/
│ └── run-all.js
├── server.js
├── tsconfig.json
├── vitest.config.ts
└── package.jsonError Handling
All Orbit errors inherit from OrbitError with a semantic code and context for fine-grained error handling:
import {
OrbitError,
TaskNotRegisteredError,
CircuitBreakerOpenError,
QueueFullError,
} from "orbit-engine";
try {
await scheduler.run("process", data);
} catch (err) {
if (err instanceof CircuitBreakerOpenError) {
return cachedResult;
} else if (err instanceof QueueFullError) {
return queue.push(data);
} else if (err instanceof OrbitError) {
console.error(`Orbit error [${err.code}]:`, err.message, err.context);
}
}Logging
Provide a logger to monitor task execution:
import { Scheduler, ConsoleLogger } from "orbit-engine";
const logger = new ConsoleLogger();
const scheduler = Scheduler.init({ workers: 4, maxQueue: 200 }, logger);Configuration Validation
Configuration is validated on initialization:
import { ConfigValidationError } from "orbit-engine";
try {
Scheduler.init({ workers: -1, maxQueue: 100 });
} catch (err) {
if (err instanceof ConfigValidationError) {
console.error("Invalid config:", err.errors);
}
}Input Validation
Payloads are validated for size (max 10MB) and serializability:
// OK
await scheduler.run("task", { data: "hello" });
// Error: Payload too large (>10MB)
await scheduler.run("task", { data: "x".repeat(11 * 1024 * 1024) });Health & Monitoring
Get real-time health snapshots for dashboards and adaptive control:
const health = scheduler.getHealth();
const totalQueued = Object.values(health.pools).reduce(
(sum, p) => sum + p.queue.size,
0,
);
if (health.circuit.open) {
console.warn("Circuit breaker is open");
}
if (health.eventLoopLagMs !== null && health.eventLoopLagMs > 100) {
console.warn("Event loop is lagged");
}Troubleshooting
"Task rejected by admission controller"
The system rejected your task. Common reasons:
- Queue is full — too many pending tasks
- Circuit breaker is open — too many recent failures
- Event loop is lagged — main thread is overloaded
- Task concurrency limit — max concurrent executions reached
Solution: Reduce load, increase pool size, or add backoff/retry logic.
"Task timed out"
The task exceeded timeoutMs. Increase the timeout or optimize the task.
Memory usage growing
Verify tasks complete successfully and queues drain via scheduler.getHealth().
Testing
Coverage target is 60%+:
npm run test:coverage# Run all tests
npm test
# Run by category
npm run test:unit
npm run test:integration
npm run test:stress
# Watch mode
npm run test:watchBenchmarking
npm run benchRuns a full production-readiness suite: load benchmarks, event loop health, memory soak, and chaos testing. Results are saved to benchmarks/report.json.
Design Patterns
| Pattern | Implementation |
|---------|---------------|
| Master–Worker | Scheduler → WorkerPool |
| Producer–Consumer | run() → TaskQueue → WorkerPool |
| Bulkhead | PoolManager with isolated pools |
| Circuit Breaker | Configurable failure threshold + auto-reset |
| Supervisor | Worker restart with per-pool budgets |
| Command | Explicit TaskCommand objects |
| Strategy | Pluggable SchedulingStrategy (round-robin, least-busy) |
| Observer | EventEmitter-based worker lifecycle events |
| Static Registry | Workers load handlers by name (no eval) |
License
MIT — see LICENSE.
