Single-node · In-memory · gRPC

Heads up: I ended up switching to NATS + PostgreSQL. If you need consumer groups, bare NATS won't help — you need JetStream, which requires persistence. tinybroker fills the gap when you want in-memory consumer groups with no persistence overhead. The code works and is free to use; just know why you're reaching for it.

Node.js

Dependencies

npm install @grpc/grpc-js @grpc/proto-loader

No code-generation step required — proto-loader loads proto files at runtime.

Connect

const grpc = require('@grpc/grpc-js');
const protoLoader = require('@grpc/proto-loader');
const path = require('path');

const PROTO_DIR = path.join(__dirname, 'proto');

const packageDef = protoLoader.loadSync(
  path.join(PROTO_DIR, 'tinybroker/v1/service.proto'),
  {
    keepCase: true,
    longs: String,
    defaults: true,
    oneofs: true,
    includeDirs: [PROTO_DIR],
  }
);

const { tinybroker: { v1: { Broker } } } =
  grpc.loadPackageDefinition(packageDef);

const client = new Broker(
  'tinybroker:50051',
  grpc.credentials.createInsecure()
);

Publish

client.publish(
  { topic: 'events.user.signup', payload: Buffer.from('{"user_id":"abc"}') },
  (err, response) => {
    if (err) console.error(err);
  }
);

// Or with async/await via util.promisify:
const { promisify } = require('util');
const publish = promisify(client.publish.bind(client));
await publish({ topic: 'events.user.signup', payload: Buffer.from('{}') });

Subscribe (fan-out)

const stream = client.subscribe();

// Open subscription
stream.write({
  open: {
    topic_patterns: ['events.user.*', 'events.order.*'],
    mode: 1, // DELIVERY_MODE_MULTI
  },
});

let subscriptionId;

stream.on('data', (event) => {
  if (event.subscription_id) {
    subscriptionId = event.subscription_id;
    console.log('subscription confirmed:', subscriptionId);
  } else if (event.message) {
    const { topic, payload, system_id } = event.message;
    console.log('message:', topic, payload.toString());
  } else if (event.patterns_updated) {
    console.log('patterns now:', event.patterns_updated.topic_patterns);
  }
});

stream.on('error', (err) => {
  console.error('stream error:', err);
  // implement reconnect
});

stream.on('end', () => {
  console.log('stream ended');
});

Subscribe (consumer group)

const stream = client.subscribe();

stream.write({
  open: {
    subscription_id: 'job-workers',
    topic_patterns: ['jobs.*'],
    mode: 2, // DELIVERY_MODE_BALANCED
    consumer_group: 'job-workers',
  },
});

stream.on('data', async (event) => {
  if (event.message) {
    const { topic, payload, system_id } = event.message;
    await processJob(payload);
    // Acknowledge so the next message is delivered
    await promisify(client.ack.bind(client))({
      subscription_id: 'job-workers',
      system_id,
    });
  }
});

Dynamic pattern management

// Add patterns on a live subscription
stream.write({
  add_patterns: { topic_patterns: ['alerts.*'] },
});

// Remove patterns
stream.write({
  remove_patterns: { topic_patterns: ['events.user.*'] },
});

TypeScript

Install types:

npm install --save-dev @types/node grpc-tools

Generate typed stubs with protoc + the grpc TypeScript plugin, or use ts-proto:

npm install ts-proto
protoc \
  --plugin=./node_modules/.bin/protoc-gen-ts_proto \
  --ts_proto_out=gen \
  --ts_proto_opt=outputServices=grpc-js \
  -I proto \
  tinybroker/v1/service.proto tinybroker/v1/client.proto \
  tinybroker/v1/server.proto tinybroker/v1/shared.proto

Reconnect pattern

function connectWithReconnect() {
  const stream = client.subscribe();
  stream.write({ open: { topic_patterns: ['events.*'], mode: 1 } });

  stream.on('data', handleEvent);
  stream.on('error', (err) => {
    console.error('reconnecting in 2s:', err.message);
    setTimeout(connectWithReconnect, 2000);
  });
  stream.on('end', connectWithReconnect);
}