Saltar al contenido principal

Go SDK

El SDK nativo de Go (client/client.go) abstrae las llamadas HTTP, maneja el enrutamiento avanzado (Prioridades, Consumer Groups) y provee un cliente WebSocket altamente resiliente con backoff exponencial automático.


Instalación

go get github.com/x-name15/tinymq/client

Tu go.mod no tendrá dependencias indirectas — el cliente de TinyMQ usa únicamente la librería estándar de Go.


HTTP Client: Publicación y Administración

El cliente HTTP estándar usa patrones modernos de Go, requiriendo un context.Context para cancelación segura y timeouts. Puedes usar PublishOptions para enrutamiento avanzado de mensajes.

package main

import (
"context"
"log"
"time"

"github.com/x-name15/tinymq/client"
)

func main() {
// Inicializa el cliente con un timeout seguro de 60s por defecto
mq := client.NewClient("http://127.0.0.1:7800", "optional_api_key")
ctx := context.Background()

// 1. Administración de Topics y Grupos (Opcional)
// retain=0 crea el topic sin expiración automática
mq.CreateTopic(ctx, "orders", "reject", 0)
// retain=24h auto-expira todos los mensajes después de 24 horas
mq.CreateTopic(ctx, "sensor.data", "drop-oldest", 24*time.Hour)
mq.CreateGroup(ctx, "orders", "billing-service")

// 2. Publicación estándar
payload := []byte(`{"event": "user_signup", "id": 99}`)
if err := mq.Publish(ctx, "users.new", payload, nil); err != nil {
log.Fatalf("publish failed: %v", err)
}

// 3. Publicación avanzada (Priority, Idempotency, Broadcast y Headers)
opts := &client.PublishOptions{
Priority: "high",
Broadcast: true, // Fan-out a todas las queues
Idempotency: "txn_987654321",
Headers: map[string]string{
"X-Source": "api-gateway",
},
}

if err := mq.Publish(ctx, "users.premium", payload, opts); err != nil {
log.Fatalf("advanced publish failed: %v", err)
}
}

Referencia de PublishOptions

CampoTipoDescripción
PrioritystringPrioridad del mensaje: "high", "normal" o "low"
BroadcastboolFan-out a todos los consumidores esperando simultáneamente (efímero)
IdempotencystringClave de idempotencia personalizada para evitar duplicados
TTLtime.DurationTime-To-Live — el mensaje expira si no se consume a tiempo
Delaytime.DurationRetraso de entrega — el mensaje permanece oculto hasta que pase la duración
Headersmap[string]stringHeaders personalizados X-MQ-* almacenados junto al mensaje

CreateTopic — Cambio de Firma (Breaking)

CreateTopic ahora se mapea directamente a POST /api/topics. El parámetro maxQueueSize int fue reemplazado por retain time.Duration:

// Antes (ya no compila)
mq.CreateTopic(ctx, "orders", "durable", 10000)

// Ahora
mq.CreateTopic(ctx, "orders", "reject", 0) // sin expiración automática
mq.CreateTopic(ctx, "sensor.data", "drop-oldest", 24*time.Hour) // retención de 24h

Pasar retain=0 omite el campo de retención en el cuerpo de la petición, dejando el topic sin TTL por defecto.


Polling HTTP de Alta Resiliencia (Consumer Groups)

Para consumidores HTTP estándar, Subscribe actúa como un fetcher de long-polling sincrónico. Soporta de forma nativa Consumer Groups para balanceo de carga distribuido y seguro entre múltiples workers.

package main

import (
"context"
"fmt"
"log"

"github.com/x-name15/tinymq/client"
)

func main() {
mq := client.NewClient("http://127.0.0.1:7800")
ctx := context.Background()

// Configura el Long-Polling con Consumer Groups
opts := &client.SubscriptionOptions{
Timeout: "30s",
Group: "billing-service", // Los mensajes se balancean entre los workers
}

log.Println("Worker iniciado, esperando orders...")

// Loop estándar de consumo
for {
msgs, err := mq.Subscribe(ctx, "orders", opts)
if err != nil {
log.Printf("Error de red o broker: %v (reintentando...)", err)
continue
}

for _, msg := range msgs {
fmt.Printf("Procesando order ID: %s | Payload: %s\n", msg.ID, string(msg.Payload))
// Implementa tu lógica de negocio / manejo de DLQ aquí
}
}
}
Integración con Re-queue & DLQ

Cuando tu lógica de procesamiento falla, llama a mq.Requeue(ctx, msg). Después de 3 fallos totales, el broker aísla automáticamente el mensaje en {topic}.dlq para su inspección posterior.


Cliente WebSocket en Tiempo Real (Con Auto-Reconexión)

Para latencia por debajo del milisegundo sin overhead HTTP, usa el cliente WebSocket nativo. Es ideal para conexiones de larga duración y alto throughput.

El WSClient cuenta con un Read-Loop de Hilo Único Thread-Safe y un Pipeline de Reconexión Automatizada. Si el servidor cae o la red se particiona, el SDK cicla automáticamente los sockets y se re-suscribe a todos los topics activos usando una estrategia de backoff exponencial (de 1s hasta 32s).

package main

import (
"fmt"
"log"

"github.com/x-name15/tinymq/client"
"github.com/x-name15/tinymq/internal/message"
)

func main() {
// NewWSClient marca automáticamente e inicia los loops de auto-reconexión y keepalive
ws, err := client.NewWSClient("127.0.0.1:7800")
if err != nil {
log.Fatalf("Conexión inicial fallida: %v", err)
}
defer ws.Close() // Libera recursos y goroutines de forma segura

// 1. Suscripción asíncrona (Thread-Safe)
err = ws.Subscribe("iot.sensors.*", func(msg message.Message) {
// Los handlers corren de forma sincrónica por conexión por defecto.
// Para concurrencia masiva, despacha a tu propio worker pool aquí.
fmt.Printf("Push instantáneo -> Topic: %s | Payload: %s\n", msg.Topic, string(msg.Payload))
})
if err != nil {
log.Fatalf("Suscripción fallida: %v", err)
}

// 2. Desinscripción dinámica (Opcional)
// ws.Unsubscribe("iot.sensors.*")

select {} // Bloquea indefinidamente mientras WS maneja el tráfico en segundo plano
}

Arquitectura del WSClient

ComponenteDescripción
Read-Loop Thread-SafeUna única goroutine maneja todas las lecturas del socket, eliminando race conditions
Pipeline de Auto-ReconexiónAl desconectarse, cicla el socket y re-suscribe a todos los topics activos
Backoff ExponencialEl retraso de reintento empieza en 1s y se duplica hasta 32s
KeepaliveResponde automáticamente a los frames Ping del servidor para mantener la conexión
defer ws.Close()Libera todas las goroutines y recursos de forma segura

Referencia de API

HTTP Client

  • NewClient(baseURL string, apiKey ...string) *Client
  • Publish(ctx context.Context, topic string, payload []byte, opts *PublishOptions) error
  • Subscribe(ctx context.Context, topic string, opts *SubscriptionOptions) ([]message.Message, error)
  • CreateTopic(ctx context.Context, name, policy string, retain time.Duration) error ⚠️ Breaking change — maxQueueSize int reemplazado por retain time.Duration
  • CreateGroup(ctx context.Context, topic, group string) error
  • Peek(ctx context.Context, topic string, limit int) ([]message.Message, error)

WSClient

  • NewWSClient(addr string) (*WSClient, error)
  • Subscribe(topic string, handler func(message.Message)) error
  • Unsubscribe(topic string)
  • Close()