Skip to content

FutureQ

A high-performance, distributed delayed-message queue broker written in Go.

FutureQ lets producers publish messages with a relative delay and guarantees reliable dispatch to consumers when the delay expires. It combines durable embedded storage, Raft-based replication, and bidirectional gRPC streaming into a single, easy-to-operate binary.

:octicons-arrow-right-24: Get started :octicons-mark-github-24: GitHub repository :octicons-file-code-24: Protocol definition

Features

  • ⏱ Delayed messaging β€” enqueue a message with a delay_ms; it becomes visible to consumers only after the delay expires. Optional per-message ttl_ms discards messages nobody consumed.
  • πŸ’Ύ Durable storage β€” disk-backed by Pebble (CockroachDB's LSM store) with a time-optimized 24-byte key schema, a pure in-memory mode, and an alternative bbolt engine.
  • πŸ”€ High availability β€” multi-node replication via Dragonboat (multi-group Raft): an event shard for messages plus a metadata group for cluster topology.
  • βž• Dynamic membership β€” nodes join and leave a running cluster over gRPC (JoinCluster / LeaveCluster) with catch-up before voting. No static bootstrap list needed after the first node.
  • πŸ‘₯ Consumer groups & topics β€” topic-based routing with fan-out across groups and round-robin dispatch within a group. "Universal" consumers (no group) each receive every message.
  • βœ… At-least-once delivery β€” in-flight tracking with automatic re-dispatch of unacknowledged messages; batched deletes (a single Raft proposal per batch) amortize LSM tombstone costs.
  • βš™οΈ Tunable durability β€” publish acknowledgements at three levels: quorum commit, leader-local persistence (fast path), or fire-and-forget β€” with a broker-side minimum enforced by policy.
  • πŸ“ˆ Observability β€” Prometheus metrics endpoint and structured logging (zap) out of the box.

Technology stack

Concern Choice
Language Go 1.26
Storage Pebble (LSM tree) Β· bbolt (alternative) Β· in-memory
Consensus Dragonboat v4 (multi-group Raft)
Transport gRPC (bidirectional streaming)
Metrics Prometheus
Protocol futureq-io/protocol (Protobuf, v0.2.0)

How it works

                β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”   PublishStream   β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”
                β”‚  Producer    β”‚ ─────────────────▢│                          β”‚
                β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜                   β”‚        FutureQ Node       β”‚
                                                   β”‚                          β”‚
                β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”   Subscribe       β”‚  β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”  β”‚
                β”‚  Consumer    β”‚ ◀──────────────── β”‚  β”‚ gRPC API (8443)    β”‚  β”‚
                β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜                   β”‚  β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”¬β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜  β”‚
                                                   β”‚            β–Ό             β”‚
                β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”   Raft (50005)    β”‚  β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”  β”‚
                β”‚  Other Nodes β”‚ ◀───────────────▢ β”‚  β”‚ Dispatcher / Hub   β”‚  β”‚
                β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜                   β”‚  β”‚ Deleter Β· Janitor  β”‚  β”‚
                                                   β”‚  β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”¬β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜  β”‚
                β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”   Prometheus      β”‚            β–Ό             β”‚
                β”‚  Metrics     β”‚ ◀── (9090) ──────│  β”‚ Pebble + Raft log    β”‚  β”‚
                β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜                   β”‚  β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜  β”‚
                                                   β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜
  • Writes go through Raft as one log entry per batch (or straight to Pebble in standalone mode) and are stored under time-bucketed keys for efficient expiry scans.
  • Reads are push-based: the dispatcher scans for due messages and routes them to connected consumers through a hub using a round-robin strategy per consumer group.
  • Acks are batched by a deleter and committed as a single Raft proposal, keeping write amplification low.

Delivery is at-least-once β€” consumers should be idempotent.