basalt / queue-kafka/src / KafkaQueuePluginOptions
Interface: KafkaQueuePluginOptions
Defined in: queue-kafka/src/index.ts:268
Everything queuePlugin accepts, minus driver (this plugin IS the driver choice), plus every Kafka driver option.
Extends
Omit<QueuePluginOptions,"driver">.KafkaDriverOptions
Properties
brokers
> brokers: string[]
Defined in: queue-kafka/src/index.ts:60
Inherited from
client?
> optional client?: KafkaClient
Defined in: queue-kafka/src/index.ts:69
Injectable client — defaults to kafkajs. Tests pass a fake.
Inherited from
clientId?
> optional clientId?: string
Defined in: queue-kafka/src/index.ts:61
Inherited from
deadSuffix?
> optional deadSuffix?: string
Defined in: queue-kafka/src/index.ts:67
Suffix for the dead-letter topic. Default '.dead'.
Inherited from
groupId?
> optional groupId?: string
Defined in: queue-kafka/src/index.ts:63
Consumer group used by workers. Default 'basalt-queue'.
Inherited from
jobs?
> optional jobs?: JobDefinition<unknown>[]
Defined in: queue/src/index.ts:61
Jobs known to this process (producer and/or worker).
Inherited from
onError?
> optional onError?: (error, info) => void
Defined in: queue-kafka/src/index.ts:77
Called on infrastructure faults the driver cannot recover in line — a worker's connect/subscribe/run failing at boot (otherwise the app reports healthy with ZERO workers and the rejection is process-fatal), or a retry/dead-letter re-publish failing inside the consume callback. Same pattern as the rabbitmq/sqs drivers. Default: console.error with context.
Parameters
error
unknown
info
queue?
string
source
"consumer" | "producer"
Returns
void
Inherited from
onUnsupported?
> optional onUnsupported?: UnsupportedPolicy
Defined in: queue/src/index.ts:86
What to do when a job uses an option the driver can't honor (e.g. a delayed job on a driver without delayed delivery). Default 'warn' — set 'throw' in production for a hard guarantee, 'ignore' for the old behavior.
Inherited from
QueuePluginOptions.onUnsupported
removeOnComplete?
> optional removeOnComplete?: JobRetention
Defined in: queue/src/index.ts:93
Default retention for completed jobs, for drivers whose backend keeps them (BullMQ/Redis today). true removes on finish, a number keeps that many, { age: '7d', count: 500 } caps both. Default: keep the last 1000. A job can override via defineJob.
Inherited from
QueuePluginOptions.removeOnComplete
removeOnFail?
> optional removeOnFail?: JobRetention
Defined in: queue/src/index.ts:98
Default retention for failed jobs. Default false (keep all, for inspection and retries) — set e.g. { age: '14d' } so failures don't grow unbounded.
Inherited from
QueuePluginOptions.removeOnFail
retrySuffix?
> optional retrySuffix?: string
Defined in: queue-kafka/src/index.ts:65
Suffix for the retry topic. Default '.retry'.
Inherited from
KafkaDriverOptions.retrySuffix
signingKey?
> optional signingKey?: QueueSigningKey | readonly QueueSigningKey[]
Defined in: queue/src/index.ts:104
HMAC key(s) that sign every job envelope; the worker rejects unsigned or tampered jobs. Without it, anyone who can write to the broker can enqueue jobs and choose their tenant/user. See QueueManagerOptions.signingKey.
Inherited from
validateTenantId?
> optional validateTenantId?: (id) => boolean
Defined in: queue/src/index.ts:109
Tenant-id grammar accepted from a job's context. Default: tenancy's default grammar. Pass the same function you gave tenancyPlugin.
Parameters
id
string
Returns
boolean
Inherited from
QueuePluginOptions.validateTenantId
workers?
> optional workers?: object[]
Defined in: queue/src/index.ts:80
Queues to start workers for in this process at boot.
concurrency?
> optional concurrency?: number
queue
> queue: string