Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
406 changes: 405 additions & 1 deletion README.md

Large diffs are not rendered by default.

283 changes: 198 additions & 85 deletions go/futureq/client.go

Large diffs are not rendered by default.

330 changes: 189 additions & 141 deletions go/futureq/consumer.go

Large diffs are not rendered by default.

84 changes: 45 additions & 39 deletions go/futureq/doc.go
Original file line number Diff line number Diff line change
@@ -1,62 +1,68 @@
// Package futureq provides a production-ready Go client SDK for the FutureQ
// scheduled message queue.
// Package futureq provides a production-ready Go client SDK for the
// FutureQ scheduled message queue.
//
// # Overview
//
// FutureQ is a distributed, time-bucket-based scheduled queue backed by Pebble
// (an LSM key-value store) and optionally replicated via the Dragonboat Raft
// library. This SDK abstracts the underlying gRPC bi-directional streaming
// protocol into two high-level, idiomatic Go clients:
// FutureQ is a distributed, time-bucket-based scheduled queue backed by
// Pebble and replicated via Dragonboat Raft. This SDK abstracts the
// gRPC bi-directional streaming protocol into two high-level clients:
//
// - [Producer] — schedules messages to be delivered at a specific time.
// - [Consumer] — subscribes to the queue and receives messages when they
// become due, acknowledging each one to prevent redelivery.
// - [Producer] — publishes batches of messages with optional delays,
// TTLs and secondary indexes.
// - [Consumer] — subscribes to a (topic, group) pair and invokes a
// handler for every delivered message, ACKing or NACKing each one.
//
// # Topology discovery
//
// The SDK tracks the cluster's Raft leader and routes streams to it.
// Two discovery mechanisms are available:
//
// 1. Polling — the SDK periodically calls GetClusterInfo on one of
// the seed addresses passed to [New]. This is the default.
// 2. Dragonboat discovery — when [WithDragonboatDiscovery] is set the
// SDK embeds a non-voting Dragonboat replica of the metadata Raft
// group inside the client process. Topology changes flow in over
// Raft as they commit, with no polling. The embedded replica stores
// all of its state in an in-memory VFS — nothing is ever written
// to disk.
//
// Dragonboat discovery is experimental.
//
// # Connecting
//
// Create a [Client] with [New] (or [NewWithConn] to supply your own
// [google.golang.org/grpc.ClientConn]):
// Create a [Client] with [New]:
//
// client, err := futureq.New("localhost:8443", futureq.WithInsecure())
// if err != nil {
// log.Fatal(err)
// }
// client, err := futureq.New(
// []string{"node1.internal:9000", "node2.internal:9000"},
// futureq.WithInsecure(),
// )
// if err != nil { log.Fatal(err) }
// defer client.Close()
//
// # Producing messages
//
// Obtain a [Producer] from the client and call [Producer.Publish]:
// # Producing
//
// producer, err := client.NewProducer(ctx)
// if err != nil {
// log.Fatal(err)
// }
// if err != nil { log.Fatal(err) }
// defer producer.Close()
//
// err = producer.Publish(ctx, futureq.Message{
// Topic: "notifications",
// Payload: []byte(`{"user": 42}`),
// ExecuteAt: time.Now().Add(5 * time.Minute),
// })
//
// # Consuming messages
// err = producer.PublishBatch(ctx, []futureq.Message{
// {Topic: "email", Payload: []byte("…"), Delay: 5 * time.Minute},
// }, futureq.AckQuorum)
//
// Obtain a [Consumer] from the client and call [Consumer.Subscribe]:
// # Consuming
//
// consumer, err := client.NewConsumer(ctx)
// if err != nil {
// log.Fatal(err)
// }
// consumer, err := client.NewConsumer(ctx, "email", "workers")
// if err != nil { log.Fatal(err) }
// defer consumer.Close()
//
// err = consumer.Subscribe(ctx, func(msg futureq.Delivery) error {
// fmt.Printf("received: %s\n", msg.Payload)
// return nil // returning nil ACKs the message
// err = consumer.Subscribe(ctx, func(d futureq.Delivery) error {
// process(d.Payload)
// return nil // nil ACKs the message
// })
//
// # Error handling
//
// All public methods return typed errors. Sentinel errors defined in this
// package (e.g. [ErrNotLeader], [ErrStreamClosed]) can be inspected with
// [errors.Is].
// All public methods return typed errors. Sentinel errors defined in
// this package (e.g. [ErrNotLeader], [ErrNoLeader], [ErrStreamClosed])
// can be inspected with [errors.Is].
package futureq
41 changes: 27 additions & 14 deletions go/futureq/errors.go
Original file line number Diff line number Diff line change
Expand Up @@ -8,34 +8,47 @@ import (
// Sentinel errors returned by the SDK.
// Use [errors.Is] to test for them:
//
// if errors.Is(err, futureq.ErrNotLeader) { … }
// if errors.Is(err, futureq.ErrNoLeader) { … }
var (
// ErrNotLeader is returned by [Producer.Publish] when the connected node is
// not the current Raft cluster leader and therefore cannot accept writes.
// The caller should retry against the leader node.
// ErrNotLeader is returned by [Producer.PublishBatch] when the
// connected node is not the current Raft cluster leader and therefore
// cannot accept writes. The SDK usually recovers from this
// automatically by re-resolving the leader through the topology
// tracker; callers may also retry manually.
ErrNotLeader = errors.New("futureq: node is not the cluster leader")

// ErrNoLeader is returned when the SDK does not currently know of any
// live leader to route a request to (e.g. immediately after startup
// before the first topology refresh, or during a rolling restart).
// Retrying after a short backoff usually succeeds.
ErrNoLeader = errors.New("futureq: no known cluster leader")

// ErrStreamClosed is returned when the underlying gRPC bi-directional
// stream has been closed by the server or the network. The [Producer] or
// [Consumer] should be discarded and a new one created.
// stream has been closed by the server or the network. The [Producer]
// or [Consumer] should be discarded and a new one created.
ErrStreamClosed = errors.New("futureq: stream closed")

// ErrPublishFailed is returned by [Producer.Publish] when the server
// acknowledged the message but reported an application-level error.
// The wrapped error message contains the server's error string.
// ErrPublishFailed is returned by [Producer.PublishBatch] when the
// server acknowledged the batch but reported an application-level
// error. Use [errors.As] to recover the structured [PublishError].
ErrPublishFailed = errors.New("futureq: publish failed")

// ErrHandlerPanic is returned by [Consumer.Subscribe] when the message
// handler panicked. The wrapped value contains the recovered panic value.
// handler panicked. The wrapped value contains the recovered panic value.
ErrHandlerPanic = errors.New("futureq: handler panicked")

// ErrClosed is returned when a method is called on a [Producer] or
// [Consumer] that has already been closed.
// ErrClosed is returned when a method is called on a [Client],
// [Producer] or [Consumer] that has already been closed.
ErrClosed = errors.New("futureq: client is closed")

// ErrTopologyUnavailable is returned when the SDK cannot obtain the
// cluster topology from any of the configured addresses (all seeds
// unreachable and no cached leader).
ErrTopologyUnavailable = errors.New("futureq: cluster topology unavailable")
)

// PublishError is the structured error type returned when a single Publish
// call is acknowledged by the server with success=false.
// PublishError is the structured error type returned when a batch is
// acknowledged by the server with success=false.
//
// It wraps [ErrPublishFailed] and additionally carries the server-supplied
// error message.
Expand Down
123 changes: 38 additions & 85 deletions go/futureq/example_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -9,131 +9,84 @@ import (
"github.com/futureq-io/sdk/go/futureq"
)

// ExampleClient_NewProducer demonstrates how to create a producer and
// schedule a single message.
func ExampleClient_NewProducer() {
// Example demonstrates the simplest produce-and-consume flow.
func Example() {
client, err := futureq.New(
"futureq.internal:8443",
futureq.WithTLS(nil),
[]string{"localhost:9000"},
futureq.WithInsecure(),
)
if err != nil {
log.Fatal(err)
}
defer client.Close()

ctx := context.Background()
producer, err := client.NewProducer(ctx, futureq.WithPublishTimeout(5*time.Second))

producer, err := client.NewProducer(ctx)
if err != nil {
log.Fatal(err)
}
defer producer.Close()

err = producer.Publish(ctx, futureq.Message{
Topic: "email-notifications",
Payload: []byte(`{"to":"alice@example.com","subject":"Welcome!"}`),
ExecuteAt: time.Now().Add(10 * time.Minute),
Topic: "email",
Payload: []byte(`{"to":"user@example.com"}`),
Delay: 5 * time.Minute,
})
if err != nil {
log.Printf("publish error: %v", err)
return
}

fmt.Println("message scheduled")
// Output: message scheduled
}

// ExampleProducer_PublishBatch shows how to schedule multiple messages
// in a single call.
// func ExampleProducer_PublishBatch() {
// client, err := futureq.New("futureq.internal:8443", futureq.WithTLS(nil))
// if err != nil {
// log.Fatal(err)
// }
// defer client.Close()

// ctx := context.Background()
// producer, err := client.NewProducer(ctx)
// if err != nil {
// log.Fatal(err)
// }
// defer producer.Close()

// now := time.Now()
// messages := []futureq.Message{
// {Topic: "reminders", Payload: []byte("reminder-1"), ExecuteAt: now.Add(1 * time.Minute)},
// {Topic: "reminders", Payload: []byte("reminder-2"), ExecuteAt: now.Add(2 * time.Minute)},
// {Topic: "reminders", Payload: []byte("reminder-3"), ExecuteAt: now.Add(3 * time.Minute)},
// }

// result, err := producer.PublishBatch(ctx, messages)
// if err != nil {
// log.Fatalf("transport error: %v", err)
// }

// for i, e := range result.Errors {
// if e != nil {
// log.Printf("message %d failed: %v", i, e)
// }
// }

// fmt.Printf("failed: %d/%d\n", len(result.FailedIndices()), len(messages))
// }

// ExampleClient_NewConsumer demonstrates how to subscribe to the queue
// and process messages with automatic ACK/NACK.
func ExampleClient_NewConsumer() {
client, err := futureq.New("futureq.internal:8443", futureq.WithTLS(nil))
if err != nil {
log.Fatal(err)
}
defer client.Close()

ctx, cancel := context.WithCancel(context.Background())
defer cancel()

consumer, err := client.NewConsumer(ctx,
futureq.WithConcurrency(4),
futureq.WithAckTimeout(3*time.Second),
)
consumer, err := client.NewConsumer(ctx, "email", "senders")
if err != nil {
log.Fatal(err)
}
defer consumer.Close()

err = consumer.Subscribe(ctx, func(d futureq.Delivery) error {
fmt.Printf("received on topic %q: %s\n", d.Topic, d.Payload)
// Return nil to ACK; return an error to NACK and trigger redelivery.
_ = consumer.Subscribe(ctx, func(d futureq.Delivery) error {
fmt.Printf("received on %s: %s\n", d.Topic, d.Payload)
return nil
})
}

// ExampleClient_dragonboatDiscovery shows how to enable the
// experimental Dragonboat-based topology discovery.
func ExampleClient_dragonboatDiscovery() {
client, err := futureq.New(
[]string{"node1.internal:9000", "node2.internal:9000"},
futureq.WithInsecure(),
// Push-based topology updates; no disk I/O.
futureq.WithDragonboatDiscovery(nil),
)
if err != nil {
log.Printf("consumer error: %v", err)
log.Fatal(err)
}
defer client.Close()
}

// ExampleProducer_PublishWithRetry demonstrates the built-in retry helper.
func ExampleProducer_PublishWithRetry() {
client, err := futureq.New("futureq.internal:8443", futureq.WithTLS(nil))
// ExampleProducer_PublishBatch demonstrates atomic batch publishing.
func ExampleProducer_PublishBatch() {
client, err := futureq.New([]string{"localhost:9000"}, futureq.WithInsecure())
if err != nil {
log.Fatal(err)
}
defer client.Close()

ctx := context.Background()
producer, err := client.NewProducer(ctx)
producer, err := client.NewProducer(context.Background())
if err != nil {
log.Fatal(err)
}
defer producer.Close()

policy := futureq.DefaultRetryPolicy()
policy.MaxAttempts = 5

err = producer.PublishWithRetry(ctx, futureq.Message{
Topic: "orders",
Payload: []byte(`{"order_id": 9001}`),
ExecuteAt: time.Now().Add(30 * time.Second),
}, policy)
err = producer.PublishBatch(context.Background(), []futureq.Message{
{Topic: "events", Payload: []byte("a"), Delay: time.Minute},
{Topic: "events", Payload: []byte("b"), Delay: 2 * time.Minute, TTL: time.Hour},
{Topic: "events", Payload: []byte("c"), Indexes: []futureq.Index{
futureq.StringIndex("user:42"),
futureq.Int64Index(1001),
}},
}, futureq.AckQuorum)
if err != nil {
log.Printf("all retries exhausted: %v", err)
log.Fatal(err)
}
}
Loading