npm package discovery and stats viewer.

Discover Tips

  • General search

    [free text search, go nuts!]

  • Package details

    pkg:[package-name]

  • User packages

    @[username]

Sponsor

Optimize Toolset

I’ve always been into building performant and accessible sites, but lately I’ve been taking it extremely seriously. So much so that I’ve been building a tool to help me optimize and monitor the sites that I build to make sure that I’m making an attempt to offer the best experience to those who visit them. If you’re into performant, accessible and SEO friendly sites, you might like it too! You can check it out at Optimize Toolset.

About

Hi, 👋, I’m Ryan Hefner  and I built this site for me, and you! The goal of this site was to provide an easy way for me to check the stats on my npm packages, both for prioritizing issues and updates, and to give me a little kick in the pants to keep up on stuff.

As I was building it, I realized that I was actually using the tool to build the tool, and figured I might as well put this out there and hopefully others will find it to be a fast and useful way to search and browse npm packages as I have.

If you’re interested in other things I’m working on, follow me on Twitter or check out the open source projects I’ve been publishing on GitHub.

I am also working on a Twitter bot for this site to tweet the most popular, newest, random packages from npm. Please follow that account now and it will start sending out packages soon–ish.

Open Software & Tools

This site wouldn’t be possible without the immense generosity and tireless efforts from the people who make contributions to the world and share their work via open source initiatives. Thank you 🙏

© 2026 – Pkg Stats / Ryan Hefner

@cogitator-ai/worker

v0.9.3

Published

Distributed job queue for Cogitator agent execution

Readme

@cogitator-ai/worker

Distributed job queue for Cogitator agent execution. Built on BullMQ for reliable, scalable background processing.

Installation

pnpm add @cogitator-ai/worker @cogitator-ai/core ioredis

BullMQ and @cogitator-ai/swarms come as dependencies; ioredis is a peer dependency. Website docs: Worker Queues, Distributed Swarms.

Features

  • BullMQ-Based - Reliable job processing with Redis
  • Job Types - Agents, workflow graphs, swarms and distributed swarm turns
  • Worker Runtime - Run jobs with your own Cogitator (provider keys, memory) and tool implementations
  • Distributed Swarms - DistributedSwarmWorker executes agent turns for @cogitator-ai/swarms
  • Auto-Retry - Exponential backoff for failed jobs
  • Priority Queue - Process important jobs first
  • Delayed Jobs - Schedule jobs for later execution
  • Prometheus Metrics - Built-in HPA support
  • Redis Cluster - Production-ready scalability
  • Graceful Shutdown - Wait for active jobs before stopping

Quick Start

Producer: Add Jobs

import { JobQueue } from '@cogitator-ai/worker';

const queue = new JobQueue({
  redis: { host: 'localhost', port: 6379 },
});

const agentConfig = {
  name: 'Assistant',
  instructions: 'You are a helpful assistant.',
  model: 'openai/gpt-6.1-sol',
  provider: 'openai' as const,
  tools: [],
};

const job = await queue.addAgentJob(agentConfig, 'Hello, world!', {
  threadId: 'user-123',
  priority: 1,
});

console.log(`Job added: ${job.id}`);

Consumer: Process Jobs

import { Cogitator } from '@cogitator-ai/core';
import { WorkerPool } from '@cogitator-ai/worker';

const pool = new WorkerPool({
  redis: { host: 'localhost', port: 6379 },
  concurrency: 5,
  workerCount: 2,
  cogitator: new Cogitator({
    llm: { providers: { openai: { apiKey: process.env.OPENAI_API_KEY! } } },
  }),
  tools: [searchTool], // implementations for tools referenced by serialized agents
});

await pool.start();

Serialized agents reference tools by name. The worker resolves them from its tools; a job whose agent needs a tool the worker does not provide fails with Tools not registered on this worker: <names>.


Job Queue

The JobQueue class manages job creation and status tracking.

Creating a Queue

import { JobQueue } from '@cogitator-ai/worker';

const queue = new JobQueue({
  name: 'my-queue',
  redis: {
    host: 'localhost',
    port: 6379,
    password: 'secret',
  },
  defaultJobOptions: {
    attempts: 3,
    backoff: { type: 'exponential', delay: 1000 },
    removeOnComplete: 100,
    removeOnFail: 500,
  },
});

Queue Configuration

interface QueueConfig {
  name?: string; // Default: 'cogitator-jobs'
  redis: {
    url?: string; // 'redis://user:password@host:port/db', 'rediss://' for TLS
    host?: string; // Default: 'localhost'
    port?: number; // Default: 6379
    username?: string;
    password?: string;
    db?: number;
    tls?: boolean; // implied by a rediss:// url
    cluster?: {
      nodes: { host: string; port: number }[];
    };
  };
  defaultJobOptions?: {
    attempts?: number; // Default: 3
    backoff?: {
      type: 'exponential' | 'fixed';
      delay: number; // Delay in ms
    };
    removeOnComplete?: boolean | number; // Default: 100
    removeOnFail?: boolean | number; // Default: 500
  };
}

Adding Jobs

Agent Jobs:

The simplest way is to serialize an agent you already have. serializeAgent writes it in the agent wire format of @cogitator-ai/core (toAgentWire), shared by agent, workflow and swarm jobs and by distributed swarm turns. It keeps the whole configuration (model and provider, instructions, sampling, stop sequences, reasoning effort, response format, iteration limit, timeout), the tool schemas and every agent it can hand over to, and turns a Zod response schema into JSON Schema, so the worker validates the structured answer the same way an in-process run does. A config with a key the worker does not know is refused rather than run without that setting:

import { Agent } from '@cogitator-ai/core';
import { serializeAgent } from '@cogitator-ai/worker';
import { z } from 'zod';

const Verdict = z.object({ verdict: z.enum(['run', 'hold']), reason: z.string() });
const chief = new Agent({
  name: 'chief',
  model: 'openrouter/openai/gpt-6-luna',
  instructions: 'Decide whether the pitch runs.',
  reasoning: { effort: 'medium' },
  responseFormat: { type: 'json_schema', schema: Verdict },
});

const job = await queue.addAgentJob(serializeAgent(chief), 'A pitch about ...');
// later: (await queue.getJob(job.id))?.returnvalue.structured -> { verdict, reason }

Or write the serialized form by hand:

const agentConfig: SerializedAgent = {
  name: 'Researcher',
  instructions: 'Research and summarize topics.',
  model: 'openai/gpt-6.1-sol', // or 'gpt-6.1-sol' - the provider is prepended when missing
  provider: 'openai',
  temperature: 0.7,
  maxTokens: 2048,
  maxIterations: 5,
  reasoning: { effort: 'low' }, // optional
  responseFormat: { type: 'json_schema', schema: { type: 'object', properties: {} } }, // optional, JSON Schema
  tools: [
    {
      name: 'search',
      description: 'Search the web',
      parameters: { type: 'object', properties: { query: { type: 'string' } } },
    },
  ],
};

const job = await queue.addAgentJob(agentConfig, 'Research quantum computing', {
  threadId: 'thread-123', // default: a new id per job
  userId: 'user-456', // the run's userId: owns the thread, reaches tools as context.userId
  priority: 1, // Lower = higher priority
  delay: 5000, // Delay 5 seconds
  metadata: { source: 'api' },
});

The worker routes model exactly like the same agent in-process: a prefix that names a built-in provider, a backend in the worker Cogitator's llm.backends or a registered plugin picks that provider, so 'openrouter/deepseek/deepseek-v4-pro' runs on an openrouter backend whatever provider says. provider (optional, any provider the worker routes to, custom backends and plugins included) is prepended only to a model whose prefix names none, such as 'meta-llama/llama-4-scout' with provider: 'openrouter'. Without provider such a model runs on the worker's llm.defaultProvider, and a provider the worker cannot route to fails the job. serializeAgent sends an agent with an explicit provider as <provider>/<model>, so it takes the same route on the worker as in-process.

Workflow Jobs:

Workflow jobs run a DAG over a shared state object initialised from the job input. Nodes run as soon as their predecessors settle; independent branches run concurrently.

const workflowConfig: SerializedWorkflow = {
  id: 'triage',
  name: 'Ticket triage',
  nodes: [
    {
      id: 'classify',
      type: 'agent',
      config: {
        agentConfig: classifierAgent, // SerializedAgent; a json_schema agent writes its structured answer
        prompt: 'Classify this ticket as BUG or QUESTION: {{ticket}}',
        outputKey: 'category',
      },
    },
    { id: 'normalize', type: 'transform', config: { transform: 'trim', inputKey: 'category' } },
    {
      id: 'is-bug',
      type: 'condition',
      config: { key: 'normalize', operator: 'contains', value: 'BUG' },
    },
    {
      id: 'summary',
      type: 'transform',
      config: { transform: 'template', template: 'Bug report: {{ticket}}' },
    },
  ],
  edges: [
    { from: 'classify', to: 'normalize' },
    { from: 'normalize', to: 'is-bug' },
    { from: 'is-bug', to: 'summary', condition: 'true' },
  ],
};

await queue.addWorkflowJob(workflowConfig, { ticket: 'App crashes on login' });

| Node type | Config | | ----------- | ----------------------------------------------------------------------------------------------------------------------------------------- | | agent | { agentConfig, prompt?, outputKey? } — prompt supports {{path}} placeholders; default prompt is the state as JSON | | transform | { transform, inputKey?, outputKey?, template? } — uppercase, lowercase, trim, json-parse, json-stringify, template | | condition | { key, operator, value? } — equals, not-equals, contains, exists, gt, lt; outgoing edges use condition: 'true' \| 'false' | | parallel | {} — fan-out marker; successors run concurrently |

Node outputs are stored in the state under outputKey (default: node id). Nodes reachable only through untaken condition branches are skipped ({ skipped: true } in nodeResults). Graphs are validated before execution (unknown nodes, invalid configs, cycles).

Swarm Jobs:

const swarmConfig: SerializedSwarm = {
  topology: 'voting',
  agents: [researcherConfig, writerConfig, editorConfig],
  coordinator: coordinatorConfig, // decides when no consensus is reached
  maxRounds: 3,
  consensusThreshold: 0.8,
};

await queue.addSwarmJob(swarmConfig, 'Write an article about AI', {
  priority: 1,
  metadata: { project: 'blog' },
});

| Topology | Swarm strategy | Notes | | --------------- | -------------- | ------------------------------------------------------------------------------------------------------------------------------------------------------------------------------ | | sequential | pipeline | One stage per agent, in order | | hierarchical | hierarchical | coordinator is required and becomes the supervisor | | collaborative | pipeline | Every agent contributes in each of maxRounds rounds (default 3), seeing all contributions so far, coordinator combines them, otherwise the last contribution is the answer | | debate | debate | maxRounds rounds, coordinator moderates | | voting | consensus | consensusThreshold, maxRounds, coordinator breaks ties |

Queue Methods

const job = await queue.getJob('job-id');

const state: JobState = await queue.getJobState('job-id');
// 'waiting' | 'prioritized' | 'delayed' | 'active' | 'completed' | 'failed'
// | 'waiting-children' | 'unknown'

await queue.resumeAgentJob(agentConfig, pausedResult, { decisions }); // see Job Results

const metrics = await queue.getMetrics();
// { waiting, active, completed, failed, delayed, depth, workerCount }

await queue.pause();
await queue.resume();

await queue.clean(60 * 60 * 1000, 1000, 'completed');
await queue.clean(24 * 60 * 60 * 1000, 100, 'failed');

const bullQueue = queue.getQueue();

await queue.close();

A job added with a priority waits in BullMQ's prioritized state instead of waiting. getMetrics() counts it in waiting and depth all the same, so cogitator_queue_depth covers every job that still has to run. waiting-children is the parent of a BullMQ flow (added through getQueue()) that waits for its children and is not counted, as its children are.


Worker Pool

The WorkerPool processes jobs with configurable concurrency.

Creating a Worker Pool

import { WorkerPool } from '@cogitator-ai/worker';

const pool = new WorkerPool(
  {
    redis: { host: 'localhost', port: 6379 },
    workerCount: 2,
    concurrency: 5,
    lockDuration: 30000,
    stalledInterval: 30000,
  },
  {
    onJobStarted: (jobId, type) => {
      console.log(`Job ${jobId} (${type}) started`);
    },
    onJobCompleted: (jobId, result) => {
      console.log(`Job ${jobId} completed:`, result);
    },
    onJobFailed: (jobId, error) => {
      console.error(`Job ${jobId} failed:`, error);
    },
    onWorkerError: (error) => {
      console.error('Worker error:', error);
    },
  }
);

await pool.start();

Worker Configuration

interface WorkerConfig extends QueueConfig {
  workerCount?: number; // Default: 1
  concurrency?: number; // Default: 5
  lockDuration?: number; // Default: 30000ms
  stalledInterval?: number; // Default: 30000ms
  cogitator?: Cogitator; // Default: new Cogitator()
  tools?: Tool[]; // Tool implementations, resolved by name
}

| Option | Default | Description | | ----------------- | ------- | ------------------------------------------ | | workerCount | 1 | Number of worker instances | | concurrency | 5 | Concurrent jobs per worker | | lockDuration | 30000 | Lock timeout before job considered stalled | | stalledInterval | 30000 | Interval to check for stalled jobs |

Worker Events

interface WorkerPoolEvents {
  onJobStarted?: (jobId: string, type: 'agent' | 'workflow' | 'swarm' | 'swarm-agent') => void;
  onJobCompleted?: (jobId: string, result: JobResult) => void;
  onJobFailed?: (jobId: string, error: Error) => void;
  onWorkerError?: (error: Error) => void;
}

Pool Methods

await pool.start();

pool.isPoolRunning();

pool.getWorkerCount();

const metrics = await pool.getMetrics(await queue.getMetrics());

// Job duration histogram and per-type counters
pool.metrics.format(await pool.getMetrics(await queue.getMetrics()));

// Cancel a job this pool is running: its agent runs stop and the job fails without retries
pool.cancelJob(jobId, 'No longer needed');

// Graceful shutdown (waits up to 30s for active jobs, then aborts them and force-closes;
// BullMQ hands the aborted jobs to another worker once their locks expire)
await pool.stop(30000);

// Force shutdown (aborts the jobs still running)
await pool.forceStop();

Every job gets the abort signal BullMQ gives it, down to each agent run, workflow node and swarm turn of the job. The processors take it too: processAgentJob(payload, runtime, { signal }), and the same third argument for processWorkflowJob, processSwarmJob and executeSwarmAgentJob (processSwarmAgentJob reads signal from its options).


Job Processors

Built-in processors handle each job type.

Using Processors Directly

Processors take the job payload and an optional runtime ({ cogitator, tools }):

import { processAgentJob, processWorkflowJob, processSwarmJob } from '@cogitator-ai/worker';

const runtime = { cogitator, tools: [searchTool] };

const agentResult = await processAgentJob(
  { type: 'agent', jobId: 'job-1', agentConfig: myAgentConfig, input: 'Hello!', threadId: 't-1' },
  runtime
);

const workflowResult = await processWorkflowJob(
  { type: 'workflow', jobId: 'job-2', runId: 'run-1', workflowConfig, input: { ticket: '...' } },
  runtime
);

const swarmResult = await processSwarmJob(
  { type: 'swarm', jobId: 'job-3', swarmConfig: mySwarmConfig, input: 'Solve this problem' },
  runtime
);

processSwarmAgentJob(payload, { publisher, isFinalAttempt, ...runtime }) executes one distributed swarm turn and publishes the result (tagged with the job id) to payload.stateKeys.results; executeSwarmAgentJob returns the result without publishing.


Distributed Swarm Workers

Swarms created with distributed.enabled (see @cogitator-ai/swarms) dispatch every agent turn to a Redis queue. DistributedSwarmWorker consumes those turns:

import { Cogitator } from '@cogitator-ai/core';
import { DistributedSwarmWorker } from '@cogitator-ai/worker';

const worker = new DistributedSwarmWorker(
  {
    redis: { host: 'localhost', port: 6379 },
    keyPrefix: 'swarm', // must match the swarm's distributed.redis.keyPrefix
    queue: 'swarm-agent-jobs', // must match distributed.queue
    concurrency: 4,
    cogitator: new Cogitator({ llm: { defaultModel: 'ollama/llama3.2' } }),
    tools: [searchTool],
  },
  {
    onJobCompleted: (job) => console.log('done', job.agentName),
    onJobFailed: (job, error) => console.error(job.agentName, error.message),
    onJobCancelled: (job) => console.log('the swarm gave up on', job.jobId),
    onError: (error) => console.error(error),
  }
);

await worker.start();
process.on('SIGTERM', () => void worker.stop()); // waits for in-flight turns

Failed turns are reported back to the swarm as errors, so the swarm's own errorHandling (retry, failover, skip) applies. A turn the swarm gave up on (its run aborted or timed out, the swarm closed, another attempt answered) is skipped, or aborted while it runs (checked every cancelCheckInterval ms, default 1000), and reported to onJobCancelled instead of being published. Each turn runs on the model its agent would use in-process, routed by the worker's cogitator, so give the worker the same llm.backends, plugins and provider keys as the process that runs the swarm.


Job Results

Each job type returns a specific result structure.

Agent Job Result

interface AgentJobResult {
  type: 'agent';
  output: string;
  structured?: unknown; // the validated answer of a json_schema agent, when it fits
  reasoning?: string; // the reasoning summary, when the agent asked for one
  usage: {
    inputTokens: number;
    outputTokens: number;
    totalTokens: number;
    cost: number; // USD, as the run reported it
    reasoningTokens?: number;
    cachedInputTokens?: number;
  };
  toolCalls: {
    name: string;
    input: unknown;
    output: unknown;
  }[];
  /** 'paused' when tool calls wait for approval: output is then not the answer */
  status?: 'completed' | 'paused';
  pendingApprovals?: ToolApprovalRequest[]; // the calls a paused run waits on
  checkpoint?: RunCheckpoint; // what a paused run continues from, see resumeAgentJob
  /** @deprecated use usage */
  tokenUsage?: {
    prompt: number;
    completion: number;
    total: number;
  };
}

A run that paused for tool approvals completes its job with status: 'paused'. Continue it with the decisions as a new job, on any worker:

if (result.status === 'paused') {
  await queue.resumeAgentJob(agentConfig, result, {
    decisions: { [result.pendingApprovals![0].toolCallId]: { approved: true } },
  });
}

Workflow and swarm jobs cannot wait for a decision: when one of their agents pauses, the job fails without retries with an AgentRunPausedError.

Workflow Job Result

interface WorkflowJobResult {
  type: 'workflow';
  output: Record<string, unknown>;
  nodeResults: Record<string, unknown>;
  duration: number;
}

Swarm Job Result

interface SwarmJobResult {
  type: 'swarm';
  output: string;
  rounds: number;
  agentOutputs: {
    agent: string;
    output: string;
  }[];
}

Prometheus Metrics

Built-in metrics for monitoring and Kubernetes HPA.

Exposing Metrics

import { JobQueue, WorkerPool } from '@cogitator-ai/worker';
import express from 'express';

const queue = new JobQueue({ redis: { host: 'localhost', port: 6379 } });
const pool = new WorkerPool({ redis: { host: 'localhost', port: 6379 } });
await pool.start();

const app = express();

app.get('/metrics', async (req, res) => {
  const queueMetrics = await queue.getMetrics(); // workerCount = workers connected to the queue
  res.type('text/plain').send(pool.metrics.format(queueMetrics));
});

app.listen(9090);

Available Metrics

| Metric | Type | Description | | -------------------------------- | --------- | ----------------------------------------------------------- | | cogitator_queue_depth | gauge | Total waiting + delayed jobs | | cogitator_queue_waiting | gauge | Jobs ready to run, prioritized jobs included | | cogitator_queue_active | gauge | Jobs currently being processed | | cogitator_queue_completed | gauge | Completed jobs kept in Redis (capped by removeOnComplete) | | cogitator_queue_failed | gauge | Failed jobs kept in Redis (capped by removeOnFail) | | cogitator_queue_delayed | gauge | Scheduled/delayed jobs | | cogitator_workers_total | gauge | Workers connected to the queue | | cogitator_job_duration_seconds | histogram | Job processing time | | cogitator_jobs_by_type_total | counter | Jobs by type | | cogitator_jobs_failed_total | counter | Jobs that failed their last attempt, by type |

cogitator_queue_completed and cogitator_queue_failed can go down as BullMQ trims old jobs, so alert on cogitator_jobs_failed_total instead (for example increase(cogitator_jobs_failed_total[5m]) > 5). The per-type counters appear after the first job of that kind.

Duration Histogram

import { DurationHistogram } from '@cogitator-ai/worker';

const histogram = new DurationHistogram('my_duration_seconds', 'Custom duration tracking');

histogram.observe(0.5);
histogram.observe(1.2);
histogram.observe(0.8);

console.log(histogram.format({ queue: 'main' }));

histogram.reset();

Metrics Collector

WorkerPool keeps one in pool.metrics: it calls recordJob(type, durationMs) for each completed job and recordFailure(type) for each job that failed its last attempt (retried attempts are not counted). Use your own collector when you process jobs outside the pool.

import { MetricsCollector } from '@cogitator-ai/worker';

const collector = new MetricsCollector();

collector.recordJob('agent', 1500);
collector.recordJob('workflow', 3200);
collector.recordFailure('agent');

const output = collector.format(queueMetrics, { queue: 'main' });

Kubernetes HPA Example

apiVersion: autoscaling/v2
kind: HorizontalPodAutoscaler
metadata:
  name: cogitator-workers
spec:
  scaleTargetRef:
    apiVersion: apps/v1
    kind: Deployment
    name: cogitator-workers
  minReplicas: 1
  maxReplicas: 10
  metrics:
    - type: External
      external:
        metric:
          name: cogitator_queue_depth
        target:
          type: AverageValue
          averageValue: 10

Redis Configuration

Single Node

const queue = new JobQueue({
  redis: {
    host: 'localhost',
    port: 6379,
    password: 'secret',
  },
});

Redis Cluster

Queues and workers connect to all cluster nodes (keys use the {cogitator} hash tag so they live in one slot):

const queue = new JobQueue({
  redis: {
    cluster: {
      nodes: [
        { host: 'redis-1', port: 6379 },
        { host: 'redis-2', port: 6379 },
        { host: 'redis-3', port: 6379 },
      ],
    },
    password: 'secret',
  },
});

Serialized Types

Jobs use serialized configurations that can be stored in Redis.

SerializedAgent

SerializedAgent is AgentWireConfig from @cogitator-ai/types:

interface SerializedAgent {
  id?: string;
  name: string;
  description?: string;
  instructions: string;
  model: string; // routed like in-process: a known provider prefix picks the provider
  provider?: LLMBackendProvider; // used when model names no provider the worker routes to
  temperature?: number;
  topP?: number;
  maxTokens?: number;
  stopSequences?: string[];
  responseFormat?:
    { type: 'text' } | { type: 'json' } | { type: 'json_schema'; schema: JSONSchema };
  reasoning?: ReasoningConfig;
  maxIterations?: number;
  onIterationLimit?: 'answer' | 'stop';
  timeout?: number;
  tools: ToolSchema[]; // resolved by name against the worker's tools
  handoffs?: { agent: string; toolName?: string; description?: string }[];
  handoffAgents?: Record<string, Omit<SerializedAgent, 'handoffAgents'>>; // handoff targets by name
}

Unknown keys are refused. Agent job results (AgentJobResult) carry output, structured, structuredError, reasoning, usage (tokens, cost, duration, reasoning and cache tokens), toolCalls with their outputs, and truncated, blocked and iterationLimitReached.

SerializedWorkflow

interface SerializedWorkflow {
  id: string;
  name: string;
  nodes: SerializedWorkflowNode[];
  edges: SerializedWorkflowEdge[];
}

interface SerializedWorkflowNode {
  id: string;
  type: 'agent' | 'transform' | 'condition' | 'parallel';
  config: Record<string, unknown>;
}

interface SerializedWorkflowEdge {
  from: string;
  to: string;
  condition?: string; // 'true' | 'false' for edges leaving condition nodes
}

Node configs are typed as AgentNodeConfig, TransformNodeConfig and ConditionNodeConfig.

SerializedSwarm

interface SerializedSwarm {
  topology: 'sequential' | 'hierarchical' | 'collaborative' | 'debate' | 'voting';
  agents: SerializedAgent[];
  coordinator?: SerializedAgent;
  maxRounds?: number;
  consensusThreshold?: number;
}

Examples

Complete Producer/Consumer

Producer (producer.ts):

import { JobQueue } from '@cogitator-ai/worker';

const queue = new JobQueue({
  redis: { host: 'localhost', port: 6379 },
});

async function main() {
  const agentConfig = {
    name: 'Summarizer',
    instructions: 'Summarize the given text concisely.',
    model: 'openai/gpt-6.1-sol',
    provider: 'openai' as const,
    tools: [],
  };

  const texts = [
    'Long article about technology...',
    'Research paper on climate change...',
    'News story about economics...',
  ];

  for (const text of texts) {
    const job = await queue.addAgentJob(agentConfig, text, {
      priority: 1,
    });
    console.log(`Queued job: ${job.id}`);
  }

  await queue.close();
}

main();

Consumer (consumer.ts):

import { WorkerPool } from '@cogitator-ai/worker';

const pool = new WorkerPool(
  {
    redis: { host: 'localhost', port: 6379 },
    concurrency: 5,
  },
  {
    onJobStarted: (id, type) => console.log(`Starting ${type} job: ${id}`),
    onJobCompleted: (id, result) => console.log(`Completed: ${id}`, result),
    onJobFailed: (id, error) => console.error(`Failed: ${id}`, error),
  }
);

async function main() {
  await pool.start();
  console.log('Worker pool started');

  process.on('SIGTERM', async () => {
    console.log('Shutting down...');
    await pool.stop(30000);
    process.exit(0);
  });
}

main();

Job Status Monitoring

import { JobQueue } from '@cogitator-ai/worker';

const queue = new JobQueue({
  redis: { host: 'localhost', port: 6379 },
});

async function monitorJob(jobId: string) {
  let lastState = '';

  while (true) {
    const state = await queue.getJobState(jobId);

    if (state !== lastState) {
      console.log(`Job ${jobId}: ${state}`);
      lastState = state;
    }

    if (state === 'completed' || state === 'failed') {
      const job = await queue.getJob(jobId);
      if (job) {
        console.log('Result:', await job.returnvalue);
      }
      break;
    }

    await new Promise((r) => setTimeout(r, 1000));
  }
}

Priority Processing

await queue.addAgentJob(config, 'Low priority', { priority: 10 });
await queue.addAgentJob(config, 'Medium priority', { priority: 5 });
await queue.addAgentJob(config, 'High priority', { priority: 1 });
await queue.addAgentJob(config, 'Critical', { priority: 0 });

Delayed Jobs

await queue.addAgentJob(config, 'Run in 5 seconds', { delay: 5000 });
await queue.addAgentJob(config, 'Run in 1 minute', { delay: 60000 });
await queue.addAgentJob(config, 'Run in 1 hour', { delay: 3600000 });

Type Reference

import type {
  // Serialized configs
  SerializedAgent,
  SerializedWorkflow,
  SerializedWorkflowNode,
  SerializedWorkflowEdge,
  AgentNodeConfig,
  TransformNodeConfig,
  ConditionNodeConfig,
  SerializedSwarm,

  // Job payloads
  JobPayload,
  AgentJobPayload,
  WorkflowJobPayload,
  SwarmJobPayload,
  SwarmAgentJobPayload,

  // Job results
  JobResult,
  AgentJobResult,
  WorkflowJobResult,
  SwarmJobResult,
  SwarmAgentJobResult,

  // Configuration
  QueueConfig,
  WorkerConfig,
  WorkerRuntime,
  QueueMetrics,
  JobState,
  DistributedSwarmWorkerConfig,
  DistributedSwarmWorkerEvents,
} from '@cogitator-ai/worker';

License

MIT