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
| Campo | Tipo | Descripción |
|---|---|---|
Priority | string | Prioridad del mensaje: "high", "normal" o "low" |
Broadcast | bool | Fan-out a todos los consumidores esperando simultáneamente (efímero) |
Idempotency | string | Clave de idempotencia personalizada para evitar duplicados |
TTL | time.Duration | Time-To-Live — el mensaje expira si no se consume a tiempo |
Delay | time.Duration | Retraso de entrega — el mensaje permanece oculto hasta que pase la duración |
Headers | map[string]string | Headers 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í
}
}
}
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
| Componente | Descripción |
|---|---|
| Read-Loop Thread-Safe | Una única goroutine maneja todas las lecturas del socket, eliminando race conditions |
| Pipeline de Auto-Reconexión | Al desconectarse, cicla el socket y re-suscribe a todos los topics activos |
| Backoff Exponencial | El retraso de reintento empieza en 1s y se duplica hasta 32s |
| Keepalive | Responde 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) *ClientPublish(ctx context.Context, topic string, payload []byte, opts *PublishOptions) errorSubscribe(ctx context.Context, topic string, opts *SubscriptionOptions) ([]message.Message, error)CreateTopic(ctx context.Context, name, policy string, retain time.Duration) error⚠️ Breaking change —maxQueueSize intreemplazado porretain time.DurationCreateGroup(ctx context.Context, topic, group string) errorPeek(ctx context.Context, topic string, limit int) ([]message.Message, error)
WSClient
NewWSClient(addr string) (*WSClient, error)Subscribe(topic string, handler func(message.Message)) errorUnsubscribe(topic string)Close()