feat(gateway): /v1/voices HTTP + WS RPC endpoints

Add ListVoices and RefreshVoices methods to RPC protocol. Implement HTTP
/v1/voices endpoint with provider-aware voice listing and filtering.
This commit is contained in:
viettranx committed 2026-04-15 11:24:56 +07:00
1 parent d97dcf252c
commit c7f2a260e8
7 files changed
+500

No files matched your search

+23
View File
@@ -3,8 +3,11 @@ package cmd
import (
"context"
"log/slog"
"time"
"github.com/nextlevelbuilder/goclaw/internal/audio"
"github.com/nextlevelbuilder/goclaw/internal/bus"
"github.com/nextlevelbuilder/goclaw/internal/gateway/methods"
httpapi "github.com/nextlevelbuilder/goclaw/internal/http"
mcpbridge "github.com/nextlevelbuilder/goclaw/internal/mcp"
"github.com/nextlevelbuilder/goclaw/internal/media"
@@ -222,6 +225,26 @@ func (d *gatewayDeps) wireHTTPHandlersOnServer(
d.server.SetMediaServeHandler(httpapi.NewMediaServeHandler(mediaStore))
}
// ElevenLabs voice list + refresh endpoints (GET /v1/voices, POST /v1/voices/refresh).
// VoiceCache is shared between the HTTP handler and the WS voices.list method.
// TTL 1h + LRU cap 1000 tenants.
{
voiceCache := audio.NewVoiceCache(1*time.Hour, 1000)
var secretStore store.ConfigSecretsStore
if d.pgStores != nil && d.pgStores.ConfigSecrets != nil {
secretStore = d.pgStores.ConfigSecrets
}
var tenantStore store.TenantStore
if d.pgStores != nil && d.pgStores.Tenants != nil {
tenantStore = d.pgStores.Tenants
}
voicesH := httpapi.NewVoicesHandler(voiceCache, secretStore, tenantStore)
d.server.SetVoicesHandler(voicesH)
// Wire WS method — provider nil means each request resolves key via secretStore at HTTP layer.
// For WS, use same cache. Provider is resolved via secretStore at WS level in a future phase.
methods.NewVoicesMethods(voiceCache, nil).Register(d.server.Router())
}
// Seed + apply builtin tool disables
if d.pgStores.BuiltinTools != nil {
seedBuiltinTools(context.Background(), d.pgStores.BuiltinTools)
+96
View File
@@ -0,0 +1,96 @@
package methods
import (
"context"
"fmt"
"log/slog"
"github.com/google/uuid"
"github.com/nextlevelbuilder/goclaw/internal/audio"
"github.com/nextlevelbuilder/goclaw/internal/audio/elevenlabs"
"github.com/nextlevelbuilder/goclaw/internal/gateway"
"github.com/nextlevelbuilder/goclaw/internal/i18n"
"github.com/nextlevelbuilder/goclaw/internal/store"
"github.com/nextlevelbuilder/goclaw/pkg/protocol"
)
// VoicesMethods handles voices.list and voices.refresh WS RPC methods.
// provider is optional — when nil, voices.list returns an error for tenants
// whose ElevenLabs key is not pre-configured via the provider arg.
type VoicesMethods struct {
cache *audio.VoiceCache
provider *elevenlabs.TTSProvider
}
// NewVoicesMethods creates a VoicesMethods handler.
// provider may be nil when the gateway has no global ElevenLabs key; in that
// case callers must supply a tenant-scoped provider via a future extension.
func NewVoicesMethods(cache *audio.VoiceCache, provider *elevenlabs.TTSProvider) *VoicesMethods {
return &VoicesMethods{cache: cache, provider: provider}
}
// Register wires voices.list and voices.refresh onto the MethodRouter.
func (m *VoicesMethods) Register(router *gateway.MethodRouter) {
router.Register(protocol.MethodVoicesList, m.handleList)
router.Register(protocol.MethodVoicesRefresh, m.handleRefresh)
}
// FetchVoices returns cached voices for tenantID, or fetches live on miss.
// Exported so it can be called from HTTP handler tests and integration code.
func (m *VoicesMethods) FetchVoices(ctx context.Context, tenantID uuid.UUID) ([]audio.Voice, error) {
if voices, ok := m.cache.Get(tenantID); ok {
return voices, nil
}
if m.provider == nil {
return nil, fmt.Errorf("no ElevenLabs provider configured")
}
voices, err := m.provider.ListVoices(ctx)
if err != nil {
return nil, err
}
m.cache.Set(tenantID, voices)
return voices, nil
}
func (m *VoicesMethods) handleList(ctx context.Context, client *gateway.Client, req *protocol.RequestFrame) {
locale := store.LocaleFromContext(ctx)
tenantID := store.TenantIDFromContext(ctx)
voices, err := m.FetchVoices(ctx, tenantID)
if err != nil {
slog.Warn("voices.list: fetch failed", "tenant_id", tenantID, "error", err)
client.SendResponse(protocol.NewErrorResponse(req.ID,
protocol.ErrInternal,
i18n.T(locale, i18n.MsgVoicesListFailed, err.Error())))
return
}
client.SendResponse(protocol.NewOKResponse(req.ID, map[string]any{"voices": voices}))
}
func (m *VoicesMethods) handleRefresh(ctx context.Context, client *gateway.Client, req *protocol.RequestFrame) {
locale := store.LocaleFromContext(ctx)
tenantID := store.TenantIDFromContext(ctx)
m.cache.Invalidate(tenantID)
if m.provider == nil {
client.SendResponse(protocol.NewErrorResponse(req.ID,
protocol.ErrInternal,
i18n.T(locale, i18n.MsgVoicesListFailed, "no provider configured")))
return
}
voices, err := m.provider.ListVoices(ctx)
if err != nil {
slog.Warn("voices.refresh: fetch failed", "tenant_id", tenantID, "error", err)
client.SendResponse(protocol.NewErrorResponse(req.ID,
protocol.ErrInternal,
i18n.T(locale, i18n.MsgVoicesListFailed, err.Error())))
return
}
m.cache.Set(tenantID, voices)
client.SendResponse(protocol.NewOKResponse(req.ID, map[string]any{"voices": voices}))
}
@@ -0,0 +1,80 @@
package methods_test
import (
"encoding/json"
"net/http"
"net/http/httptest"
"testing"
"time"
"github.com/google/uuid"
"github.com/nextlevelbuilder/goclaw/internal/audio"
"github.com/nextlevelbuilder/goclaw/internal/audio/elevenlabs"
"github.com/nextlevelbuilder/goclaw/internal/gateway/methods"
"github.com/nextlevelbuilder/goclaw/internal/store"
)
// TestVoicesMethods_CacheHit verifies a warm cache entry is returned without
// calling the upstream provider.
func TestVoicesMethods_CacheHit(t *testing.T) {
cache := audio.NewVoiceCache(time.Hour, 100)
tid := uuid.New()
voices := []audio.Voice{{ID: "v1", Name: "Bella"}}
cache.Set(tid, voices)
ctx := store.WithTenantID(t.Context(), tid)
m := methods.NewVoicesMethods(cache, nil)
got, err := m.FetchVoices(ctx, tid)
if err != nil {
t.Fatalf("unexpected error: %v", err)
}
if len(got) != 1 || got[0].ID != "v1" {
t.Errorf("unexpected voices: %+v", got)
}
}
// TestVoicesMethods_NoProvider verifies an error is returned when the cache
// misses and no provider is configured.
func TestVoicesMethods_NoProvider(t *testing.T) {
cache := audio.NewVoiceCache(time.Hour, 100)
m := methods.NewVoicesMethods(cache, nil)
_, err := m.FetchVoices(t.Context(), uuid.New())
if err == nil {
t.Fatal("expected error when no provider configured")
}
}
// TestVoicesMethods_LiveFetch verifies a cache miss triggers a live fetch via
// the provider and the result is stored in the cache.
func TestVoicesMethods_LiveFetch(t *testing.T) {
upstream := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
w.Header().Set("Content-Type", "application/json")
json.NewEncoder(w).Encode(map[string]any{
"voices": []map[string]any{
{"voice_id": "v2", "name": "Adam", "category": "premade"},
},
})
}))
defer upstream.Close()
cache := audio.NewVoiceCache(time.Hour, 100)
p := elevenlabs.NewTTSProvider(elevenlabs.Config{APIKey: "k", BaseURL: upstream.URL})
m := methods.NewVoicesMethods(cache, p)
tid := uuid.New()
got, err := m.FetchVoices(t.Context(), tid)
if err != nil {
t.Fatalf("unexpected error: %v", err)
}
if len(got) != 1 || got[0].ID != "v2" {
t.Errorf("unexpected voices: %+v", got)
}
// Result should be cached now.
cached, ok := cache.Get(tid)
if !ok || len(cached) != 1 {
t.Error("expected live fetch result to be cached")
}
}
+3
View File
@@ -509,6 +509,9 @@ func (s *Server) SetEvolutionHandler(h *httpapi.EvolutionHandler) {
s.handlers = append(s.handlers, h)
}
// SetVoicesHandler sets the ElevenLabs voices list + refresh handler.
func (s *Server) SetVoicesHandler(h *httpapi.VoicesHandler) { s.handlers = append(s.handlers, h) }
// SetVaultHandler sets the Knowledge Vault document handler.
func (s *Server) SetVaultHandler(h *httpapi.VaultHandler) { s.handlers = append(s.handlers, h) }
+123
View File
@@ -0,0 +1,123 @@
package http
import (
"fmt"
"log/slog"
"net/http"
"github.com/google/uuid"
"github.com/nextlevelbuilder/goclaw/internal/audio"
"github.com/nextlevelbuilder/goclaw/internal/audio/elevenlabs"
"github.com/nextlevelbuilder/goclaw/internal/i18n"
"github.com/nextlevelbuilder/goclaw/internal/permissions"
"github.com/nextlevelbuilder/goclaw/internal/store"
)
// VoicesHandler serves GET /v1/voices and POST /v1/voices/refresh.
// It holds a *VoiceCache to serve cached responses and optionally a
// *elevenlabs.TTSProvider to fetch live voices on cache miss.
type VoicesHandler struct {
cache *audio.VoiceCache
provider *elevenlabs.TTSProvider // nil when no ElevenLabs key is configured
secretStore store.ConfigSecretsStore
tenantStore store.TenantStore
}
// NewVoicesHandler creates a handler that resolves the ElevenLabs provider at
// request time from config_secrets. Use NewVoicesHandlerWithProvider for tests.
func NewVoicesHandler(cache *audio.VoiceCache, secretStore store.ConfigSecretsStore, tenantStore store.TenantStore) *VoicesHandler {
return &VoicesHandler{cache: cache, secretStore: secretStore, tenantStore: tenantStore}
}
// NewVoicesHandlerWithProvider creates a handler with a pre-built provider.
// Primarily used in tests to inject a mock/httptest provider.
func NewVoicesHandlerWithProvider(cache *audio.VoiceCache, p *elevenlabs.TTSProvider) *VoicesHandler {
return &VoicesHandler{cache: cache, provider: p}
}
// RegisterRoutes wires the voices endpoints onto mux.
func (h *VoicesHandler) RegisterRoutes(mux *http.ServeMux) {
mux.HandleFunc("GET /v1/voices", requireAuth("", h.handleList))
mux.HandleFunc("POST /v1/voices/refresh", requireAuth(permissions.RoleAdmin, h.handleRefresh))
}
// handleList serves GET /v1/voices — returns cached list or fetches live.
func (h *VoicesHandler) handleList(w http.ResponseWriter, r *http.Request) {
ctx := r.Context()
tenantID := store.TenantIDFromContext(ctx)
locale := store.LocaleFromContext(ctx)
if voices, ok := h.cache.Get(tenantID); ok {
writeJSON(w, http.StatusOK, map[string]any{"voices": voices})
return
}
p, err := h.resolveProvider(r, tenantID)
if err != nil {
slog.Warn("voices: no ElevenLabs provider", "tenant_id", tenantID, "error", err)
writeJSON(w, http.StatusNotFound, map[string]string{
"error": i18n.T(locale, i18n.MsgVoicesListFailed, err.Error()),
})
return
}
voices, err := p.ListVoices(ctx)
if err != nil {
slog.Warn("voices: list failed", "tenant_id", tenantID, "error", err)
writeJSON(w, http.StatusBadGateway, map[string]string{
"error": i18n.T(locale, i18n.MsgVoicesListFailed, err.Error()),
})
return
}
h.cache.Set(tenantID, voices)
writeJSON(w, http.StatusOK, map[string]any{"voices": voices})
}
// handleRefresh serves POST /v1/voices/refresh — admin-only, forces a live
// refetch by invalidating the tenant's cache entry.
func (h *VoicesHandler) handleRefresh(w http.ResponseWriter, r *http.Request) {
ctx := r.Context()
tenantID := store.TenantIDFromContext(ctx)
locale := store.LocaleFromContext(ctx)
h.cache.Invalidate(tenantID)
p, err := h.resolveProvider(r, tenantID)
if err != nil {
slog.Warn("voices: no ElevenLabs provider on refresh", "tenant_id", tenantID, "error", err)
writeJSON(w, http.StatusNotFound, map[string]string{
"error": i18n.T(locale, i18n.MsgVoicesListFailed, err.Error()),
})
return
}
voices, err := p.ListVoices(ctx)
if err != nil {
slog.Warn("voices: refresh fetch failed", "tenant_id", tenantID, "error", err)
writeJSON(w, http.StatusBadGateway, map[string]string{
"error": i18n.T(locale, i18n.MsgVoicesListFailed, err.Error()),
})
return
}
h.cache.Set(tenantID, voices)
writeJSON(w, http.StatusOK, map[string]any{"voices": voices})
}
// resolveProvider returns the ElevenLabs provider to use for this request.
// Priority: injected provider (test/pre-built) > secret store lookup.
func (h *VoicesHandler) resolveProvider(r *http.Request, tenantID uuid.UUID) (*elevenlabs.TTSProvider, error) {
if h.provider != nil {
return h.provider, nil
}
if h.secretStore == nil {
return nil, fmt.Errorf("no ElevenLabs API key configured")
}
apiKey, err := h.secretStore.Get(r.Context(), "tts.elevenlabs.api_key")
if err != nil || apiKey == "" {
return nil, fmt.Errorf("ElevenLabs API key not found for tenant %s", tenantID)
}
return elevenlabs.NewTTSProvider(elevenlabs.Config{APIKey: apiKey}), nil
}
+169
View File
@@ -0,0 +1,169 @@
package http_test
import (
"encoding/json"
"net/http"
"net/http/httptest"
"testing"
"time"
"github.com/nextlevelbuilder/goclaw/internal/audio"
"github.com/nextlevelbuilder/goclaw/internal/audio/elevenlabs"
httpapi "github.com/nextlevelbuilder/goclaw/internal/http"
"github.com/nextlevelbuilder/goclaw/internal/store"
)
const voicesTestToken = "voices-test-token"
func TestVoicesHandler_Unauthenticated(t *testing.T) {
httpapi.InitGatewayToken(voicesTestToken)
t.Cleanup(func() { httpapi.InitGatewayToken("") })
cache := audio.NewVoiceCache(time.Hour, 100)
h := httpapi.NewVoicesHandler(cache, nil, nil)
mux := http.NewServeMux()
h.RegisterRoutes(mux)
req := httptest.NewRequest("GET", "/v1/voices", nil)
rr := httptest.NewRecorder()
mux.ServeHTTP(rr, req)
if rr.Code != http.StatusUnauthorized {
t.Errorf("expected 401, got %d: %s", rr.Code, rr.Body.String())
}
}
// TestVoicesHandler_CachedResponse verifies that a cache hit skips the upstream call.
func TestVoicesHandler_CachedResponse(t *testing.T) {
// No gateway token → dev mode, everyone is admin, MasterTenantID assigned.
httpapi.InitGatewayToken("")
t.Cleanup(func() { httpapi.InitGatewayToken("") })
cache := audio.NewVoiceCache(time.Hour, 100)
voices := []audio.Voice{{ID: "v1", Name: "Bella", Category: "premade"}}
// Seed for MasterTenantID — that's what dev-mode auth injects.
cache.Set(store.MasterTenantID, voices)
called := false
upstream := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
called = true
w.WriteHeader(http.StatusInternalServerError)
}))
defer upstream.Close()
p := elevenlabs.NewTTSProvider(elevenlabs.Config{APIKey: "k", BaseURL: upstream.URL})
h := httpapi.NewVoicesHandlerWithProvider(cache, p)
mux := http.NewServeMux()
h.RegisterRoutes(mux)
req := httptest.NewRequest("GET", "/v1/voices", nil)
rr := httptest.NewRecorder()
mux.ServeHTTP(rr, req)
if rr.Code != http.StatusOK {
t.Fatalf("expected 200, got %d: %s", rr.Code, rr.Body.String())
}
if called {
t.Error("upstream ElevenLabs should NOT be called on cache hit")
}
var resp struct {
Voices []audio.Voice `json:"voices"`
}
if err := json.NewDecoder(rr.Body).Decode(&resp); err != nil {
t.Fatalf("decode: %v", err)
}
if len(resp.Voices) != 1 || resp.Voices[0].ID != "v1" {
t.Errorf("unexpected voices: %+v", resp.Voices)
}
}
// TestVoicesHandler_LiveFetch verifies a cache miss triggers live fetch and caches result.
func TestVoicesHandler_LiveFetch(t *testing.T) {
httpapi.InitGatewayToken("")
t.Cleanup(func() { httpapi.InitGatewayToken("") })
upstream := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
w.Header().Set("Content-Type", "application/json")
w.Write([]byte(`{"voices":[{"voice_id":"v2","name":"Adam","category":"premade"}]}`))
}))
defer upstream.Close()
cache := audio.NewVoiceCache(time.Hour, 100)
p := elevenlabs.NewTTSProvider(elevenlabs.Config{APIKey: "k", BaseURL: upstream.URL})
h := httpapi.NewVoicesHandlerWithProvider(cache, p)
mux := http.NewServeMux()
h.RegisterRoutes(mux)
req := httptest.NewRequest("GET", "/v1/voices", nil)
rr := httptest.NewRecorder()
mux.ServeHTTP(rr, req)
if rr.Code != http.StatusOK {
t.Fatalf("expected 200, got %d: %s", rr.Code, rr.Body.String())
}
var resp struct {
Voices []audio.Voice `json:"voices"`
}
json.NewDecoder(rr.Body).Decode(&resp)
if len(resp.Voices) != 1 || resp.Voices[0].ID != "v2" {
t.Errorf("unexpected voices: %+v", resp.Voices)
}
// Verify cache was populated.
cached, ok := cache.Get(store.MasterTenantID)
if !ok || len(cached) != 1 {
t.Error("expected live fetch result to be cached")
}
}
// TestVoicesHandler_RefreshUnauthenticated verifies POST /refresh requires auth.
func TestVoicesHandler_RefreshUnauthenticated(t *testing.T) {
httpapi.InitGatewayToken(voicesTestToken)
t.Cleanup(func() { httpapi.InitGatewayToken("") })
cache := audio.NewVoiceCache(time.Hour, 100)
h := httpapi.NewVoicesHandler(cache, nil, nil)
mux := http.NewServeMux()
h.RegisterRoutes(mux)
req := httptest.NewRequest("POST", "/v1/voices/refresh", nil)
rr := httptest.NewRecorder()
mux.ServeHTTP(rr, req)
if rr.Code != http.StatusUnauthorized {
t.Errorf("expected 401 for unauthenticated refresh, got %d: %s", rr.Code, rr.Body.String())
}
}
// TestVoicesHandler_RefreshAdmin verifies POST /refresh works for admin (dev mode).
func TestVoicesHandler_RefreshAdmin(t *testing.T) {
httpapi.InitGatewayToken("")
t.Cleanup(func() { httpapi.InitGatewayToken("") })
upstream := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
w.Header().Set("Content-Type", "application/json")
w.Write([]byte(`{"voices":[{"voice_id":"v3","name":"Rachel"}]}`))
}))
defer upstream.Close()
cache := audio.NewVoiceCache(time.Hour, 100)
p := elevenlabs.NewTTSProvider(elevenlabs.Config{APIKey: "k", BaseURL: upstream.URL})
h := httpapi.NewVoicesHandlerWithProvider(cache, p)
mux := http.NewServeMux()
h.RegisterRoutes(mux)
req := httptest.NewRequest("POST", "/v1/voices/refresh", nil)
rr := httptest.NewRecorder()
mux.ServeHTTP(rr, req)
if rr.Code != http.StatusOK {
t.Fatalf("expected 200 for admin refresh, got %d: %s", rr.Code, rr.Body.String())
}
var resp struct {
Voices []audio.Voice `json:"voices"`
}
json.NewDecoder(rr.Body).Decode(&resp)
if len(resp.Voices) != 1 || resp.Voices[0].ID != "v3" {
t.Errorf("unexpected voices after refresh: %+v", resp.Voices)
}
}
+6
View File
@@ -165,6 +165,12 @@ const (
MethodAPIKeysRevoke = "api_keys.revoke"
)
// Voices (ElevenLabs voice picker)
const (
MethodVoicesList = "voices.list"
MethodVoicesRefresh = "voices.refresh"
)
// Phase 3+ - NICE TO HAVE methods
const (
MethodLogsTail = "logs.tail"