Skip to main content

Queue Worker SDK

@eclosion-tech/syntropy-queue provides a shared way to publish jobs, run workers, and monitor task lifecycle events.

Use this package when you want a central worker fleet and multiple services producing jobs.

Install

npm install @eclosion-tech/syntropy-queue

Publish Jobs

Infra-level publishers can write directly to SQS:

import {
SqsQueueBackend,
createQueueJobEnvelope,
} from "@eclosion-tech/syntropy-queue";

const backend = new SqsQueueBackend({
region: process.env.AWS_REGION!,
queueUrl: process.env.SQS_QUEUE_URL!,
});

const job = createQueueJobEnvelope({
tenant: "syntropy",
task: "workflow.runRuntime",
env: "prod",
payload: {
organizationId: "org_uuid",
accountId: "project_uuid",
workflowId: "wf_uuid",
},
});

await backend.publish(job);

Application producers that should not hold AWS credentials can enqueue via API key: the key needs queue_jobs:write (or project:admin).

import { createSyntropyApiClient } from "@eclosion-tech/syntropy-node";

const api = createSyntropyApiClient({
baseUrl: "https://syntropy.chat/api",
apiKey: process.env.SYNTROPY_API_KEY,
});

await api.queue.enqueue(
{ organizationId: "org_uuid", projectId: "project_uuid" },
{
task: "workflow.runRuntime",
payload: { limit: 100 },
idempotencyKey: "syntropy:workflow.runRuntime:manual",
}
);

To monitor lifecycle events from the shared queue API, use queue_jobs:read and either poll:

const { jobs } = await api.queue.list(scope, { status: "running", limit: 50 });
const { events } = await api.queue.listEvents(scope, {
since: new Date(Date.now() - 30_000).toISOString(),
});

or consume the SSE stream endpoint:

GET /api/organizations/{organizationId}/projects/{projectId}/queue/jobs/stream
Authorization: Bearer syn_sk_xxx

Run Workers

import {
QueueWorkerRuntime,
SqsQueueBackend,
} from "@eclosion-tech/syntropy-queue";

const runtime = new QueueWorkerRuntime({
backend: new SqsQueueBackend({
region: process.env.AWS_REGION!,
queueUrl: process.env.SQS_QUEUE_URL!,
}),
});

runtime.register("workflow.runRuntime", async (job, context) => {
context.logger.info("Running task", {
task: job.task,
jobId: job.jobId,
tenant: job.tenant,
});

// process payload
});

await runtime.start(new AbortController().signal);

Monitor Task Lifecycle

You can instrument worker jobs without custom wrapper code by adding observers.

import {
QueueWorkerRuntime,
SqsQueueBackend,
createQueueTelemetryObserver,
} from "@eclosion-tech/syntropy-queue";
import { Syntropy } from "@eclosion-tech/syntropy-node";

const runtime = new QueueWorkerRuntime({
backend: new SqsQueueBackend({
region: process.env.AWS_REGION!,
queueUrl: process.env.SQS_QUEUE_URL!,
}),
observers: [
createQueueTelemetryObserver({
client: {
captureEvent: (name, payload) => Syntropy.captureEvent(name, payload),
captureError: (error, context) => Syntropy.captureError(error, context),
},
eventPrefix: "worker",
includePayload: false,
}),
],
});

Emitted lifecycle events:

  • job.received
  • job.unhandled
  • task.started
  • task.succeeded
  • task.failed

Reliability Notes

  • If no handler is registered for job.task, the runtime discards the job.
  • If a handler throws, the runtime returns retry to the backend.
  • Observer failures are isolated and do not fail task execution.