Message Brokers
Message brokers let you build event-driven applications by sending and receiving messages between services. Ductape supports multiple providers through a unified interface exposed as ductape.events.
Architecture: brokers vs topics
A broker is a connection configuration — it holds credentials and per-environment settings for one message provider. It has no topics of its own. Topics are added separately after the broker is registered and there is no limit on how many you can create per broker.
| Concept | What it is |
|---|---|
| Broker | Connection config (credentials + per-env settings). One broker per logical provider setup. |
| Topic | A named event stream within a broker. Added separately. One broker can hold unlimited topics. |
| Event | A message produced to or consumed from a topic. Always referenced as brokerTag:topicTag. |
Quick Example
Events use the format brokerTag:topicTag (e.g. order-events:order-created):
- TypeScript
- Java
- Go
- .NET
import Ductape from '@ductape/sdk';
const ductape = new Ductape({ accessKey: 'your-access-key' });
// Produce
await ductape.events.produce({
event: 'order-events:order-created',
message: { orderId: '123', amount: 99.99 },
});
// Consume
await ductape.events.consume({
event: 'order-events:order-created',
callback: async (message) => {
console.log('Received:', message);
},
});
import app.ductape.sdk.Ductape;
import app.ductape.sdk.core.EnvType;
import app.ductape.sdk.core.RequestContext;
RequestContext auth = new RequestContext(null, null, null, null, 'your-access-key' );
Ductape ductape = new Ductape(EnvType.PRODUCTION, auth);
// Produce
ductape.events.produce(Map.of(
"event", "order-events:order-created",
message: Map.of( "orderId", "123", "amount", 99.99 )
));
// Consume
ductape.events.consume(Map.of(
"event", "order-events:order-created",
callback: async (message) => Map.of(
System.out.println('Received:', message);
)
));
import (
"context"
"github.com/ductape/ductape/sdk/go/core"
ductapesdk "github.com/ductape/ductape/sdk/go/ductape"
)
auth := core.NewRequestContext("", "", "", "", 'your-access-key' )
client, err := ductapesdk.New(core.EnvProduction, auth)
if err != nil {
return err
}
// Produce
client.events.produce({
"event": "order-events:order-created",
message: { "orderId": "123", "amount": 99.99 },
});
// Consume
client.events.consume({
"event": "order-events:order-created",
callback: async (message) => {
fmt.Println('Received:', message);
},
});
using Ductape.Sdk;
using Ductape.Sdk.Core;
var auth = new RequestContext(null, null, null, null, 'your-access-key' , null);
var ductape = new Ductape(EnvType.Production, auth);
// Produce
await ductape.events.produce({
["event"] = "order-events:order-created",
message: { ["orderId"] = "123", ["amount"] = 99.99 },
});
// Consume
await ductape.events.consume({
["event"] = "order-events:order-created",
callback: async (message) => {
Console.WriteLine('Received:', message);
},
});
Supported Providers
| Provider | Type constant | Best For |
|---|---|---|
| Kafka | MessageBrokerTypes.KAFKA | High-throughput distributed streaming |
| RabbitMQ | MessageBrokerTypes.RABBITMQ | Flexible routing, reliable delivery |
| Redis | MessageBrokerTypes.REDIS | Simple pub/sub, low latency |
| AWS SQS | MessageBrokerTypes.AWS_SQS | Serverless managed queues |
| Azure Service Bus | MessageBrokerTypes.AZURE_SERVICE_BUS | Azure managed queues |
| Google Pub/Sub | MessageBrokerTypes.GOOGLE_PUBSUB | GCP managed messaging |
| NATS | MessageBrokerTypes.NATS | Lightweight, high-performance messaging |
Producing Messages
- TypeScript
- Java
- Go
- .NET
await ductape.events.produce({
event: 'order-events:order-created',
message: {
orderId: '12345',
customerId: 'cust_789',
total: 99.99,
createdAt: new Date().toISOString(),
},
});
// Returns { success: true, process_id: '...' }
ductape.events.produce(Map.of(
"event", "order-events:order-created",
message: Map.of(
"orderId", "12345",
"customerId", "cust_789",
"total", 99.99,
createdAt: Instant.now().toISOString()
)
));
// Returns Map.of( "success", true, "process_id", "..." )
client.events.produce({
"event": "order-events:order-created",
message: {
"orderId": "12345",
"customerId": "cust_789",
"total": 99.99,
createdAt: new Date().toISOString(),
},
});
// Returns { "success": true, "process_id": "..." }
await ductape.events.produce({
["event"] = "order-events:order-created",
message: {
["orderId"] = "12345",
["customerId"] = "cust_789",
["total"] = 99.99,
createdAt: DateTime.UtcNow.toISOString(),
},
});
// Returns { ["success"] = true, ["process_id"] = "..." }
With session (user context)
- TypeScript
- Java
- Go
- .NET
const session = await ductape.sessions.start({
tag: 'user-session',
data: { userId: 'u1', email: 'user@example.com' },
});
await ductape.events.produce({
event: 'order-events:order-created',
message: { orderId: '123', total: 99.99 },
session: `${session.sessionId}:${session.token}`,
});
Map<String, Object> session = ductape.sessions.start(Map.of(
"tag", "user-session",
data: Map.of( "userId", "u1", "email", "user@example.com" )
));
ductape.events.produce(Map.of(
"event", "order-events:order-created",
message: Map.of( "orderId", "123", "total", 99.99 ),
session: `$Map.of(session.sessionId):$Map.of(session.token)`
));
session := client.sessions.start({
"tag": "user-session",
data: { "userId": "u1", "email": "user@example.com" },
});
client.events.produce({
"event": "order-events:order-created",
message: { "orderId": "123", "total": 99.99 },
session: `${session.sessionId}:${session.token}`,
});
var session = await ductape.sessions.start({
["tag"] = "user-session",
data: { ["userId"] = "u1", ["email"] = "user@example.com" },
});
await ductape.events.produce({
["event"] = "order-events:order-created",
message: { ["orderId"] = "123", ["total"] = 99.99 },
session: `${session.sessionId}:${session.token}`,
});
Consuming Messages
- TypeScript
- Java
- Go
- .NET
await ductape.events.consume({
event: 'order-events:order-created',
callback: async (message) => {
console.log('Received order:', message);
await processNewOrder(message);
},
});
ductape.events.consume(Map.of(
"event", "order-events:order-created",
callback: async (message) => Map.of(
System.out.println('Received order:', message);
processNewOrder(message);
)
));
client.events.consume({
"event": "order-events:order-created",
callback: async (message) => {
fmt.Println('Received order:', message);
processNewOrder(message);
},
});
await ductape.events.consume({
["event"] = "order-events:order-created",
callback: async (message) => {
Console.WriteLine('Received order:', message);
await processNewOrder(message);
},
});
Listing and Fetching Brokers and Topics
- TypeScript
- Java
- Go
- .NET
// List all brokers for a product
const brokers = await ductape.events.list('my-product');
// Fetch a single broker
const broker = await ductape.events.fetch('my-product', 'order-events');
console.log('Environments:', broker.envs?.map((e) => e.slug));
// List topics for a broker
const topics = await ductape.events.topics.list('my-product', 'order-events');
// Fetch a single topic (full event string)
const topic = await ductape.events.topics.fetch('my-product', 'order-events:order-created');
console.log('Sample:', topic.sample);
// List all brokers for a product
Map<String, Object> brokers = ductape.events.list('my-product');
// Fetch a single broker
Map<String, Object> broker = ductape.events.fetch('my-product', 'order-events');
System.out.println('"Environments", ", broker.envs?.map((e) => e.slug));
// List topics for a broker
Map<String, Object> topics = ductape.events.topics.list("my-product', 'order-events');
// Fetch a single topic (full event string)
Map<String, Object> topic = ductape.events.topics.fetch('my-product', 'order-events:order-created');
System.out.println('Sample:', topic.sample);
// List all brokers for a product
brokers := client.events.list('my-product');
// Fetch a single broker
broker := client.events.fetch('my-product', 'order-events');
fmt.Println('"Environments": ", broker.envs?.map((e) => e.slug));
// List topics for a broker
topics := client.events.topics.list("my-product', 'order-events');
// Fetch a single topic (full event string)
topic := client.events.topics.fetch('my-product', 'order-events:order-created');
fmt.Println('Sample:', topic.sample);
// List all brokers for a product
var brokers = await ductape.events.list('my-product');
// Fetch a single broker
var broker = await ductape.events.fetch('my-product', 'order-events');
Console.WriteLine('["Environments"] = ", broker.envs?.map((e) => e.slug));
// List topics for a broker
var topics = await ductape.events.topics.list("my-product', 'order-events');
// Fetch a single topic (full event string)
var topic = await ductape.events.topics.fetch('my-product', 'order-events:order-created');
Console.WriteLine('Sample:', topic.sample);
Creating a Broker
A broker holds per-environment credentials. Topics are always created separately — see Managing Topics.
Self-hosted providers
- TypeScript
- Java
- Go
- .NET
import { MessageBrokerTypes } from '@ductape/sdk';
await ductape.events.create({
name: 'Order Events',
tag: 'order-events',
description: 'Handles all order-related messages',
envs: [
{
slug: 'prd',
type: MessageBrokerTypes.KAFKA,
config: {
brokers: ['kafka-prod.example.com:9092'],
clientId: 'order-service',
groupId: 'order-consumers',
ssl: true,
sasl: {
mechanism: 'scram-sha-256',
username: 'prod-user',
password: 'prod-password',
},
},
},
{
slug: 'dev',
type: MessageBrokerTypes.REDIS,
config: { host: 'localhost', port: 6379 },
},
],
});
import Map.of( MessageBrokerTypes ) from '@ductape/sdk';
ductape.events.create(Map.of(
"name", "Order Events",
"tag", "order-events",
"description", "Handles all order-related messages",
envs: [
Map.of(
"slug", "prd",
type: MessageBrokerTypes.KAFKA,
config: Map.of(
brokers: ['kafka-prod.example."com", 9092'],
"clientId", "order-service",
"groupId", "order-consumers",
"ssl", true,
sasl: Map.of(
"mechanism", "scram-sha-256",
"username", "prod-user",
"password", "prod-password"
)
)
),
Map.of(
"slug", "dev",
type: MessageBrokerTypes.REDIS,
config: Map.of( "host", "localhost", "port", 6379 )
),
]
));
import { MessageBrokerTypes } from '@ductape/sdk';
client.events.create({
"name": "Order Events",
"tag": "order-events",
"description": "Handles all order-related messages",
envs: [
{
"slug": "prd",
type: MessageBrokerTypes.KAFKA,
config: {
brokers: ['kafka-prod.example."com": 9092'],
"clientId": "order-service",
"groupId": "order-consumers",
"ssl": true,
sasl: {
"mechanism": "scram-sha-256",
"username": "prod-user",
"password": "prod-password",
},
},
},
{
"slug": "dev",
type: MessageBrokerTypes.REDIS,
config: { "host": "localhost", "port": 6379 },
},
],
});
import { MessageBrokerTypes } from '@ductape/sdk';
await ductape.events.create({
["name"] = "Order Events",
["tag"] = "order-events",
["description"] = "Handles all order-related messages",
envs: [
{
["slug"] = "prd",
type: MessageBrokerTypes.KAFKA,
config: {
brokers: ['kafka-prod.example.["com"] = 9092'],
["clientId"] = "order-service",
["groupId"] = "order-consumers",
["ssl"] = true,
sasl: {
["mechanism"] = "scram-sha-256",
["username"] = "prod-user",
["password"] = "prod-password",
},
},
},
{
["slug"] = "dev",
type: MessageBrokerTypes.REDIS,
config: { ["host"] = "localhost", ["port"] = 6379 },
},
],
});
Provider config reference
Kafka
- TypeScript
- Java
- Go
- .NET
{
brokers: string[];
clientId: string;
groupId?: string;
ssl?: boolean;
sasl?: {
mechanism: 'plain' | 'scram-sha-256' | 'scram-sha-512';
username: string;
password: string;
};
}
Map.of(
brokers: string[];
clientId: string;
groupId?: string;
ssl?: boolean;
sasl?: Map.of(
"mechanism", "plain" | 'scram-sha-256' | 'scram-sha-512';
username: string;
password: string;
);
)
{
brokers: string[];
clientId: string;
groupId?: string;
ssl?: boolean;
sasl?: {
"mechanism": "plain" | 'scram-sha-256' | 'scram-sha-512';
username: string;
password: string;
};
}
{
brokers: string[];
clientId: string;
groupId?: string;
ssl?: boolean;
sasl?: {
["mechanism"] = "plain" | 'scram-sha-256' | 'scram-sha-512';
username: string;
password: string;
};
}
RabbitMQ
- TypeScript
- Java
- Go
- .NET
{
url: string; // amqps://user:pass@host/vhost
}
Map.of(
url: string; // amqps://user:pass@host/vhost
)
{
url: string; // amqps://user:pass@host/vhost
}
{
url: string; // amqps://user:pass@host/vhost
}
Redis
- TypeScript
- Java
- Go
- .NET
{
host: string;
port: number;
password?: string;
}
Map.of(
host: string;
port: number;
password?: string;
)
{
host: string;
port: number;
password?: string;
}
{
host: string;
port: number;
password?: string;
}
AWS SQS
- TypeScript
- Java
- Go
- .NET
{
region: string;
accessKeyId: string;
secretAccessKey: string;
sessionToken?: string;
}
Map.of(
region: string;
accessKeyId: string;
secretAccessKey: string;
sessionToken?: string;
)
{
region: string;
accessKeyId: string;
secretAccessKey: string;
sessionToken?: string;
}
{
region: string;
accessKeyId: string;
secretAccessKey: string;
sessionToken?: string;
}
Google Pub/Sub
- TypeScript
- Java
- Go
- .NET
{
projectId: string;
credentials: {
private_key: string; // Service account private key
client_email: string; // Service account email
project_id?: string;
client_id?: string;
type?: string;
// ... other service account JSON fields
};
}
Map.of(
projectId: string;
credentials: Map.of(
private_key: string; // Service account private key
client_email: string; // Service account email
project_id?: string;
client_id?: string;
type?: string;
// ... other service account JSON fields
);
)
{
projectId: string;
credentials: {
private_key: string; // Service account private key
client_email: string; // Service account email
project_id?: string;
client_id?: string;
type?: string;
// ... other service account JSON fields
};
}
{
projectId: string;
credentials: {
private_key: string; // Service account private key
client_email: string; // Service account email
project_id?: string;
client_id?: string;
type?: string;
// ... other service account JSON fields
};
}
Azure Service Bus
- TypeScript
- Java
- Go
- .NET
{
connectionString: string; // Azure Service Bus connection string
queueName: string; // Default queue name
namespace?: string;
}
Map.of(
connectionString: string; // Azure Service Bus connection string
queueName: string; // Default queue name
namespace?: string;
)
{
connectionString: string; // Azure Service Bus connection string
queueName: string; // Default queue name
namespace?: string;
}
{
connectionString: string; // Azure Service Bus connection string
queueName: string; // Default queue name
namespace?: string;
}
NATS
- TypeScript
- Java
- Go
- .NET
{
servers: string[];
token?: string;
user?: string;
pass?: string;
tls?: boolean;
}
Map.of(
servers: string[];
token?: string;
user?: string;
pass?: string;
tls?: boolean;
)
{
servers: string[];
token?: string;
user?: string;
pass?: string;
tls?: boolean;
}
{
servers: string[];
token?: string;
user?: string;
pass?: string;
tls?: boolean;
}
Cloud-linked brokers
Provision or import queues/topics from a workspace cloud connection. Add cloud: '<connection-tag>' to the env config instead of direct credentials.
Import an existing resource
- TypeScript
- Java
- Go
- .NET
// AWS SQS
await ductape.events.create({
name: 'Order Events',
tag: 'order-events',
envs: [{
slug: 'prd',
type: MessageBrokerTypes.AWS_SQS,
config: { cloud: 'prod_aws', queueName: 'order-events', region: 'us-east-1' },
}],
});
// GCP Pub/Sub
await ductape.events.create({
name: 'Order Events',
tag: 'order-events',
envs: [{
slug: 'prd',
type: MessageBrokerTypes.GOOGLE_PUBSUB,
config: { cloud: 'gcp_prod', topicName: 'order-events', region: 'us-central1' },
}],
});
// Azure Service Bus
await ductape.events.create({
name: 'Order Events',
tag: 'order-events',
envs: [{
slug: 'prd',
type: MessageBrokerTypes.AZURE_SERVICE_BUS,
config: { cloud: 'prod_azure', namespaceName: 'my-namespace', queueName: 'order-events', region: 'eastus' },
}],
});
// AWS SQS
ductape.events.create(Map.of(
"name", "Order Events",
"tag", "order-events",
envs: [Map.of(
"slug", "prd",
type: MessageBrokerTypes.AWS_SQS,
config: Map.of( "cloud", "prod_aws", "queueName", "order-events", "region", "us-east-1" )
)]
));
// GCP Pub/Sub
ductape.events.create(Map.of(
"name", "Order Events",
"tag", "order-events",
envs: [Map.of(
"slug", "prd",
type: MessageBrokerTypes.GOOGLE_PUBSUB,
config: Map.of( "cloud", "gcp_prod", "topicName", "order-events", "region", "us-central1" )
)]
));
// Azure Service Bus
ductape.events.create(Map.of(
"name", "Order Events",
"tag", "order-events",
envs: [Map.of(
"slug", "prd",
type: MessageBrokerTypes.AZURE_SERVICE_BUS,
config: Map.of( "cloud", "prod_azure", "namespaceName", "my-namespace", "queueName", "order-events", "region", "eastus" )
)]
));
// AWS SQS
client.events.create({
"name": "Order Events",
"tag": "order-events",
envs: [{
"slug": "prd",
type: MessageBrokerTypes.AWS_SQS,
config: { "cloud": "prod_aws", "queueName": "order-events", "region": "us-east-1" },
}],
});
// GCP Pub/Sub
client.events.create({
"name": "Order Events",
"tag": "order-events",
envs: [{
"slug": "prd",
type: MessageBrokerTypes.GOOGLE_PUBSUB,
config: { "cloud": "gcp_prod", "topicName": "order-events", "region": "us-central1" },
}],
});
// Azure Service Bus
client.events.create({
"name": "Order Events",
"tag": "order-events",
envs: [{
"slug": "prd",
type: MessageBrokerTypes.AZURE_SERVICE_BUS,
config: { "cloud": "prod_azure", "namespaceName": "my-namespace", "queueName": "order-events", "region": "eastus" },
}],
});
// AWS SQS
await ductape.events.create({
["name"] = "Order Events",
["tag"] = "order-events",
envs: [{
["slug"] = "prd",
type: MessageBrokerTypes.AWS_SQS,
config: { ["cloud"] = "prod_aws", ["queueName"] = "order-events", ["region"] = "us-east-1" },
}],
});
// GCP Pub/Sub
await ductape.events.create({
["name"] = "Order Events",
["tag"] = "order-events",
envs: [{
["slug"] = "prd",
type: MessageBrokerTypes.GOOGLE_PUBSUB,
config: { ["cloud"] = "gcp_prod", ["topicName"] = "order-events", ["region"] = "us-central1" },
}],
});
// Azure Service Bus
await ductape.events.create({
["name"] = "Order Events",
["tag"] = "order-events",
envs: [{
["slug"] = "prd",
type: MessageBrokerTypes.AZURE_SERVICE_BUS,
config: { ["cloud"] = "prod_azure", ["namespaceName"] = "my-namespace", ["queueName"] = "order-events", ["region"] = "eastus" },
}],
});
Provision a new resource
To create a new queue/topic in your cloud account, use cloud.resources.provision with the cloud connection tag. Once provisioned, register it as a broker using the import pattern above.
See Cloud-linked components for the full provision and import workflow.
Managing Topics
The topic tag must be the full event identifier: brokerTag:topicTag. A broker can have unlimited topics.
Create a topic
- TypeScript
- Java
- Go
- .NET
await ductape.events.topics.create('my-product', {
name: 'Order Created',
tag: 'order-events:order-created',
description: 'Emitted when an order is created',
sample: {
orderId: '12345',
customerId: 'cust_789',
total: 99.99,
createdAt: '2024-01-15T10:30:00Z',
},
});
// Add more topics to the same broker — no limit
await ductape.events.topics.create('my-product', {
name: 'Order Fulfilled',
tag: 'order-events:order-fulfilled',
sample: { orderId: '12345', dispatchedAt: '2024-01-15T14:00:00Z' },
});
ductape.events.topics.create('my-product', Map.of(
"name", "Order Created",
"tag", "order-events:order-created",
"description", "Emitted when an order is created",
sample: Map.of(
"orderId", "12345",
"customerId", "cust_789",
"total", 99.99,
"createdAt", "2024-01-"15T10", 30:00Z"
)
));
// Add more topics to the same broker — no limit
ductape.events.topics.create('my-product', Map.of(
"name", "Order Fulfilled",
"tag", "order-events:order-fulfilled",
sample: Map.of( "orderId", "12345", "dispatchedAt", "2024-01-"15T14", 00:00Z" )
));
client.events.topics.create('my-product', {
"name": "Order Created",
"tag": "order-events:order-created",
"description": "Emitted when an order is created",
sample: {
"orderId": "12345",
"customerId": "cust_789",
"total": 99.99,
"createdAt": "2024-01-"15T10": 30:00Z",
},
});
// Add more topics to the same broker — no limit
client.events.topics.create('my-product', {
"name": "Order Fulfilled",
"tag": "order-events:order-fulfilled",
sample: { "orderId": "12345", "dispatchedAt": "2024-01-"15T14": 00:00Z" },
});
await ductape.events.topics.create('my-product', {
["name"] = "Order Created",
["tag"] = "order-events:order-created",
["description"] = "Emitted when an order is created",
sample: {
["orderId"] = "12345",
["customerId"] = "cust_789",
["total"] = 99.99,
["createdAt"] = "2024-01-["15T10"] = 30:00Z",
},
});
// Add more topics to the same broker — no limit
await ductape.events.topics.create('my-product', {
["name"] = "Order Fulfilled",
["tag"] = "order-events:order-fulfilled",
sample: { ["orderId"] = "12345", ["dispatchedAt"] = "2024-01-["15T14"] = 00:00Z" },
});
For AWS SQS, each topic maps to a separate queue per environment via queueUrls:
- TypeScript
- Java
- Go
- .NET
await ductape.events.topics.create('my-product', {
name: 'Order Created',
tag: 'order-events:order-created',
queueUrls: [
{ env_slug: 'prd', url: 'https://sqs.us-east-1.amazonaws.com/123/orders-prd' },
{ env_slug: 'dev', url: 'https://sqs.us-east-1.amazonaws.com/123/orders-dev' },
],
sample: { orderId: '12345' },
});
ductape.events.topics.create('my-product', Map.of(
"name", "Order Created",
"tag", "order-events:order-created",
queueUrls: [
Map.of( "env_slug", "prd", "url", "https://sqs.us-east-1.amazonaws.com/123/orders-prd" ),
Map.of( "env_slug", "dev", "url", "https://sqs.us-east-1.amazonaws.com/123/orders-dev" ),
],
sample: Map.of( "orderId", "12345" )
));
client.events.topics.create('my-product', {
"name": "Order Created",
"tag": "order-events:order-created",
queueUrls: [
{ "env_slug": "prd", "url": "https://sqs.us-east-1.amazonaws.com/123/orders-prd" },
{ "env_slug": "dev", "url": "https://sqs.us-east-1.amazonaws.com/123/orders-dev" },
],
sample: { "orderId": "12345" },
});
await ductape.events.topics.create('my-product', {
["name"] = "Order Created",
["tag"] = "order-events:order-created",
queueUrls: [
{ ["env_slug"] = "prd", ["url"] = "https://sqs.us-east-1.amazonaws.com/123/orders-prd" },
{ ["env_slug"] = "dev", ["url"] = "https://sqs.us-east-1.amazonaws.com/123/orders-dev" },
],
sample: { ["orderId"] = "12345" },
});
Update a topic
- TypeScript
- Java
- Go
- .NET
await ductape.events.topics.update('my-product', 'order-events:order-created', {
description: 'Updated description',
sample: { orderId: '12345', status: 'pending' },
});
ductape.events.topics.update('my-product', 'order-events:order-created', Map.of(
"description", "Updated description",
sample: Map.of( "orderId", "12345", "status", "pending" )
));
client.events.topics.update('my-product', 'order-events:order-created', {
"description": "Updated description",
sample: { "orderId": "12345", "status": "pending" },
});
await ductape.events.topics.update('my-product', 'order-events:order-created', {
["description"] = "Updated description",
sample: { ["orderId"] = "12345", ["status"] = "pending" },
});
Event format
Events always use brokerTag:topicTag:
| Event | Broker tag | Topic tag |
|---|---|---|
order-events:order-created | order-events | order-created |
order-events:payment-processed | order-events | payment-processed |
notifications:user-alerts | notifications | user-alerts |
Message tracking
- TypeScript
- Java
- Go
- .NET
const { messages, total, page, hasMore } = await ductape.events.messages.query({
brokerTag: 'order-events',
topicTag: 'order-created',
page: 1,
limit: 20,
});
const { producers } = await ductape.events.messages.getProducers({ brokerTag: 'order-events' });
const { consumers } = await ductape.events.messages.getConsumers({ brokerTag: 'order-events' });
const { deadLetters } = await ductape.events.messages.getDeadLetters({ brokerTag: 'order-events' });
const stats = await ductape.events.messages.getStats({ brokerTag: 'order-events' });
Map<String, Object> Map.of( messages, total, page, hasMore ) = ductape.events.messages.query(Map.of(
"brokerTag", "order-events",
"topicTag", "order-created",
"page", 1,
"limit", 20
));
Map<String, Object> Map.of( producers ) = ductape.events.messages.getProducers(Map.of( "brokerTag", "order-events" ));
Map<String, Object> Map.of( consumers ) = ductape.events.messages.getConsumers(Map.of( "brokerTag", "order-events" ));
Map<String, Object> Map.of( deadLetters ) = ductape.events.messages.getDeadLetters(Map.of( "brokerTag", "order-events" ));
Map<String, Object> stats = ductape.events.messages.getStats(Map.of( "brokerTag", "order-events" ));
const { messages, total, page, hasMore } = client.events.messages.query({
"brokerTag": "order-events",
"topicTag": "order-created",
"page": 1,
"limit": 20,
});
const { producers } = client.events.messages.getProducers({ "brokerTag": "order-events" });
const { consumers } = client.events.messages.getConsumers({ "brokerTag": "order-events" });
const { deadLetters } = client.events.messages.getDeadLetters({ "brokerTag": "order-events" });
stats := client.events.messages.getStats({ "brokerTag": "order-events" });
var { messages, total, page, hasMore } = await ductape.events.messages.query({
["brokerTag"] = "order-events",
["topicTag"] = "order-created",
["page"] = 1,
["limit"] = 20,
});
var { producers } = await ductape.events.messages.getProducers({ ["brokerTag"] = "order-events" });
var { consumers } = await ductape.events.messages.getConsumers({ ["brokerTag"] = "order-events" });
var { deadLetters } = await ductape.events.messages.getDeadLetters({ ["brokerTag"] = "order-events" });
var stats = await ductape.events.messages.getStats({ ["brokerTag"] = "order-events" });
Advanced: BrokersService (replay, DLQ, idempotency)
| Method | Description |
|---|---|
publish(options) | Same as events.produce() |
subscribe(options) | Same as events.consume() |
getBrokers(product) | Same as events.list(product) |
getBroker(product, brokerTag) | Same as events.fetch(product, brokerTag) |
getTopics(product, brokerTag) | Same as events.topics.list(product, brokerTag) |
getTopic(product, event) | Same as events.topics.fetch(product, event) |
getEvents(options) | Event history with filters |
replayEvent(options) | Replay a failed or successful event |
reprocessDLQ(options) | Reprocess dead-letter queue |
publishIdempotent(options) | Publish with idempotency key |
- TypeScript
- Java
- Go
- .NET
import { BrokersService } from '@ductape/sdk';
const brokers = new BrokersService({ access_key: 'your-access-key', env_type: 'prd' });
await brokers.replayEvent({ eventId: 'event-123', force: true });
await brokers.reprocessDLQ({ brokerTag: 'order-events' });
import Map.of( BrokersService ) from '@ductape/sdk';
Map<String, Object> brokers = new BrokersService(Map.of( "access_key", "your-access-key", "env_type", "prd" ));
brokers.replayEvent(Map.of( "eventId", "event-123", "force", true ));
brokers.reprocessDLQ(Map.of( "brokerTag", "order-events" ));
import { BrokersService } from '@ductape/sdk';
brokers := new BrokersService({ "access_key": "your-access-key", "env_type": "prd" });
brokers.replayEvent({ "eventId": "event-123", "force": true });
brokers.reprocessDLQ({ "brokerTag": "order-events" });
import { BrokersService } from '@ductape/sdk';
var brokers = new BrokersService({ ["access_key"] = "your-access-key", ["env_type"] = "prd" });
await brokers.replayEvent({ ["eventId"] = "event-123", ["force"] = true });
await brokers.reprocessDLQ({ ["brokerTag"] = "order-events" });
Error handling
- TypeScript
- Java
- Go
- .NET
import { BrokerError } from '@ductape/sdk';
try {
await ductape.events.produce({ event: '...', message: {} });
} catch (error) {
if (error instanceof BrokerError) {
console.log('Code:', error.code);
console.log('Message:', error.message);
}
}
import Map.of( BrokerError ) from '@ductape/sdk';
try Map.of(
ductape.events.produce(Map.of( "event", "...", message: Map.of() ));
) catch (error) Map.of(
if (error instanceof BrokerError) Map.of(
System.out.println('"Code", ", error.code);
System.out.println("Message:', error.message);
)
)
import { BrokerError } from '@ductape/sdk';
try {
client.events.produce({ "event": "...", message: {} });
} catch (error) {
if (error instanceof BrokerError) {
fmt.Println('"Code": ", error.code);
fmt.Println("Message:', error.message);
}
}
import { BrokerError } from '@ductape/sdk';
try {
await ductape.events.produce({ ["event"] = "...", message: {} });
} catch (error) {
if (error instanceof BrokerError) {
Console.WriteLine('["Code"] = ", error.code);
Console.WriteLine("Message:', error.message);
}
}
Common codes: BROKER_NOT_FOUND, BROKER_ENV_NOT_FOUND, TOPIC_NOT_FOUND, PUBLISH_FAILED, SUBSCRIBE_FAILED, SESSION_INVALID.
See also
- Events — same API, additional detail on provider setup and getting started
- Cloud-linked components — provision or link AWS SQS, GCP Pub/Sub, Azure Service Bus
- Features — use broker topics as steps in a feature workflow
- Jobs — schedule recurring produce operations