Refactor logging
This commit is contained in:
@@ -10,12 +10,24 @@ import (
|
||||
type Config struct {
|
||||
Server ServerConfig `yaml:"server"`
|
||||
Database DatabaseConfig `yaml:"database"`
|
||||
Logging LoggingConfig `yaml:"logging"`
|
||||
Metrics MetricsConfig `yaml:"metrics"`
|
||||
}
|
||||
|
||||
type ServerConfig struct {
|
||||
Port int `yaml:"port"`
|
||||
}
|
||||
|
||||
type LoggingConfig struct {
|
||||
Level string `yaml:"level"`
|
||||
Format string `yaml:"format"` // "json" or "console"
|
||||
}
|
||||
|
||||
type MetricsConfig struct {
|
||||
Enabled bool `yaml:"enabled"`
|
||||
Port int `yaml:"port"`
|
||||
}
|
||||
|
||||
type DatabaseConfig struct {
|
||||
Host string `yaml:"host"`
|
||||
Port int `yaml:"port"`
|
||||
@@ -35,6 +47,14 @@ func Load(path string) (*Config, error) {
|
||||
Port: 5432,
|
||||
SSLMode: "disable",
|
||||
},
|
||||
Logging: LoggingConfig{
|
||||
Level: "info",
|
||||
Format: "json",
|
||||
},
|
||||
Metrics: MetricsConfig{
|
||||
Enabled: true,
|
||||
Port: 9090,
|
||||
},
|
||||
}
|
||||
|
||||
if path == "" {
|
||||
|
||||
@@ -0,0 +1,137 @@
|
||||
package logging
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"time"
|
||||
|
||||
"github.com/google/uuid"
|
||||
"github.com/grpc-ecosystem/go-grpc-middleware/v2/interceptors/logging"
|
||||
"github.com/rs/zerolog"
|
||||
"github.com/rs/zerolog/log"
|
||||
"google.golang.org/grpc"
|
||||
"google.golang.org/grpc/codes"
|
||||
"google.golang.org/grpc/metadata"
|
||||
"google.golang.org/grpc/status"
|
||||
)
|
||||
|
||||
func InterceptorLogger() logging.Logger {
|
||||
return logging.LoggerFunc(func(ctx context.Context, lvl logging.Level, msg string, fields ...any) {
|
||||
l := zerolog.Ctx(ctx).With().Fields(fields).Logger()
|
||||
|
||||
switch lvl {
|
||||
case logging.LevelDebug:
|
||||
l.Debug().Msg(msg)
|
||||
case logging.LevelInfo:
|
||||
l.Info().Msg(msg)
|
||||
case logging.LevelWarn:
|
||||
l.Warn().Msg(msg)
|
||||
case logging.LevelError:
|
||||
l.Error().Msg(msg)
|
||||
default:
|
||||
l.Info().Msg(msg)
|
||||
}
|
||||
})
|
||||
}
|
||||
|
||||
func LoggingOpts() []logging.Option {
|
||||
return []logging.Option{
|
||||
logging.WithLogOnEvents(logging.StartCall, logging.FinishCall),
|
||||
logging.WithLevels(levelFunc),
|
||||
logging.WithDurationField(func(d time.Duration) logging.Fields {
|
||||
return logging.Fields{"grpc.duration_ms", d.Milliseconds()}
|
||||
}),
|
||||
}
|
||||
}
|
||||
|
||||
func levelFunc(code codes.Code) logging.Level {
|
||||
switch code {
|
||||
case codes.OK, codes.NotFound, codes.Canceled:
|
||||
return logging.LevelInfo
|
||||
case codes.InvalidArgument, codes.AlreadyExists, codes.Unauthenticated:
|
||||
return logging.LevelWarn
|
||||
default:
|
||||
return logging.LevelError
|
||||
}
|
||||
}
|
||||
|
||||
func RequestIDInterceptor() grpc.UnaryServerInterceptor {
|
||||
return func(ctx context.Context, req any, info *grpc.UnaryServerInfo, handler grpc.UnaryHandler) (any, error) {
|
||||
requestID := uuid.New().String()
|
||||
|
||||
logger := log.Logger.With().
|
||||
Str("request_id", requestID).
|
||||
Str("grpc.method", info.FullMethod).
|
||||
Logger()
|
||||
|
||||
ctx = logger.WithContext(ctx)
|
||||
|
||||
resp, err := handler(ctx, req)
|
||||
if err != nil {
|
||||
st, _ := status.FromError(err)
|
||||
zerolog.Ctx(ctx).Debug().
|
||||
Str("grpc.code", st.Code().String()).
|
||||
Str("grpc.error", st.Message()).
|
||||
Msg("request failed")
|
||||
}
|
||||
|
||||
return resp, err
|
||||
}
|
||||
}
|
||||
|
||||
func PeerInfoInterceptor() grpc.UnaryServerInterceptor {
|
||||
return func(ctx context.Context, req any, info *grpc.UnaryServerInfo, handler grpc.UnaryHandler) (any, error) {
|
||||
md, ok := metadata.FromIncomingContext(ctx)
|
||||
if ok {
|
||||
evt := zerolog.Ctx(ctx).With()
|
||||
if ua := md.Get("user-agent"); len(ua) > 0 {
|
||||
evt = evt.Str("peer.user_agent", ua[0])
|
||||
}
|
||||
if auth := md.Get("authorization"); len(auth) > 0 {
|
||||
evt = evt.Bool("peer.authenticated", true)
|
||||
}
|
||||
logger := evt.Logger()
|
||||
ctx = logger.WithContext(ctx)
|
||||
}
|
||||
|
||||
return handler(ctx, req)
|
||||
}
|
||||
}
|
||||
|
||||
func StreamRequestIDInterceptor() grpc.StreamServerInterceptor {
|
||||
return func(srv any, ss grpc.ServerStream, info *grpc.StreamServerInfo, handler grpc.StreamHandler) error {
|
||||
requestID := uuid.New().String()
|
||||
|
||||
logger := log.Logger.With().
|
||||
Str("request_id", requestID).
|
||||
Str("grpc.method", info.FullMethod).
|
||||
Logger()
|
||||
|
||||
ctx := logger.WithContext(ss.Context())
|
||||
wrapped := &wrappedStream{ServerStream: ss, ctx: ctx}
|
||||
|
||||
return handler(srv, wrapped)
|
||||
}
|
||||
}
|
||||
|
||||
type wrappedStream struct {
|
||||
grpc.ServerStream
|
||||
ctx context.Context
|
||||
}
|
||||
|
||||
func (w *wrappedStream) Context() context.Context {
|
||||
return w.ctx
|
||||
}
|
||||
|
||||
func (w *wrappedStream) SendMsg(m any) error {
|
||||
zerolog.Ctx(w.ctx).Trace().Str("direction", "send").Msg(fmt.Sprintf("stream message: %T", m))
|
||||
return w.ServerStream.SendMsg(m)
|
||||
}
|
||||
|
||||
func (w *wrappedStream) RecvMsg(m any) error {
|
||||
err := w.ServerStream.RecvMsg(m)
|
||||
if err == nil {
|
||||
zerolog.Ctx(w.ctx).Trace().Str("direction", "recv").Msg(fmt.Sprintf("stream message: %T", m))
|
||||
}
|
||||
return err
|
||||
}
|
||||
@@ -0,0 +1,21 @@
|
||||
package logging
|
||||
|
||||
import (
|
||||
"os"
|
||||
|
||||
"github.com/rs/zerolog"
|
||||
"github.com/rs/zerolog/log"
|
||||
|
||||
"github.com/metadata-agregator/internal/config"
|
||||
)
|
||||
|
||||
func Init(cfg config.LoggingConfig) {
|
||||
level, err := zerolog.ParseLevel(cfg.Level)
|
||||
if err != nil {
|
||||
level = zerolog.InfoLevel
|
||||
}
|
||||
zerolog.SetGlobalLevel(level)
|
||||
|
||||
log.Logger = zerolog.New(zerolog.ConsoleWriter{Out: os.Stderr}).
|
||||
With().Timestamp().Logger()
|
||||
}
|
||||
@@ -0,0 +1,115 @@
|
||||
package metrics
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"net/http"
|
||||
"time"
|
||||
|
||||
grpcprom "github.com/grpc-ecosystem/go-grpc-middleware/providers/prometheus"
|
||||
"github.com/prometheus/client_golang/prometheus"
|
||||
"github.com/prometheus/client_golang/prometheus/promauto"
|
||||
"github.com/prometheus/client_golang/prometheus/promhttp"
|
||||
"github.com/rs/zerolog/log"
|
||||
"google.golang.org/grpc"
|
||||
)
|
||||
|
||||
var (
|
||||
CacheHits = promauto.NewCounterVec(prometheus.CounterOpts{
|
||||
Name: "metadata_cache_hits_total",
|
||||
Help: "Number of cache hits by entity type",
|
||||
}, []string{"entity"})
|
||||
|
||||
CacheMisses = promauto.NewCounterVec(prometheus.CounterOpts{
|
||||
Name: "metadata_cache_misses_total",
|
||||
Help: "Number of cache misses by entity type",
|
||||
}, []string{"entity"})
|
||||
|
||||
ProviderRequests = promauto.NewCounterVec(prometheus.CounterOpts{
|
||||
Name: "metadata_provider_requests_total",
|
||||
Help: "Number of requests to external providers",
|
||||
}, []string{"provider", "operation", "status"})
|
||||
|
||||
ProviderLatency = promauto.NewHistogramVec(prometheus.HistogramOpts{
|
||||
Name: "metadata_provider_request_duration_seconds",
|
||||
Help: "Latency of external provider requests",
|
||||
Buckets: []float64{0.01, 0.05, 0.1, 0.25, 0.5, 1, 2.5, 5, 10},
|
||||
}, []string{"provider", "operation"})
|
||||
|
||||
DBQueryLatency = promauto.NewHistogramVec(prometheus.HistogramOpts{
|
||||
Name: "metadata_db_query_duration_seconds",
|
||||
Help: "Latency of database queries",
|
||||
Buckets: prometheus.ExponentialBuckets(0.001, 2, 10),
|
||||
}, []string{"operation"})
|
||||
)
|
||||
|
||||
type ServerMetrics struct {
|
||||
grpcMetrics *grpcprom.ServerMetrics
|
||||
registry *prometheus.Registry
|
||||
}
|
||||
|
||||
func NewServerMetrics() *ServerMetrics {
|
||||
reg := prometheus.NewRegistry()
|
||||
reg.MustRegister(prometheus.NewGoCollector())
|
||||
reg.MustRegister(prometheus.NewProcessCollector(prometheus.ProcessCollectorOpts{}))
|
||||
|
||||
srvMetrics := grpcprom.NewServerMetrics(
|
||||
grpcprom.WithServerHandlingTimeHistogram(
|
||||
grpcprom.WithHistogramBuckets([]float64{
|
||||
0.001, 0.01, 0.05, 0.1, 0.3, 0.6, 1, 3, 6, 10, 30,
|
||||
}),
|
||||
),
|
||||
)
|
||||
reg.MustRegister(srvMetrics)
|
||||
reg.MustRegister(CacheHits)
|
||||
reg.MustRegister(CacheMisses)
|
||||
reg.MustRegister(ProviderRequests)
|
||||
reg.MustRegister(ProviderLatency)
|
||||
reg.MustRegister(DBQueryLatency)
|
||||
|
||||
return &ServerMetrics{
|
||||
grpcMetrics: srvMetrics,
|
||||
registry: reg,
|
||||
}
|
||||
}
|
||||
|
||||
func (m *ServerMetrics) UnaryServerInterceptor() grpc.UnaryServerInterceptor {
|
||||
return m.grpcMetrics.UnaryServerInterceptor()
|
||||
}
|
||||
|
||||
func (m *ServerMetrics) StreamServerInterceptor() grpc.StreamServerInterceptor {
|
||||
return m.grpcMetrics.StreamServerInterceptor()
|
||||
}
|
||||
|
||||
func (m *ServerMetrics) InitializeMetrics(srv *grpc.Server) {
|
||||
m.grpcMetrics.InitializeMetrics(srv)
|
||||
}
|
||||
|
||||
func (m *ServerMetrics) StartHTTPServer(ctx context.Context, port int) {
|
||||
mux := http.NewServeMux()
|
||||
mux.Handle("/metrics", promhttp.HandlerFor(m.registry, promhttp.HandlerOpts{
|
||||
EnableOpenMetrics: true,
|
||||
}))
|
||||
mux.HandleFunc("/health", func(w http.ResponseWriter, r *http.Request) {
|
||||
w.WriteHeader(http.StatusOK)
|
||||
_, _ = w.Write([]byte("OK"))
|
||||
})
|
||||
|
||||
srv := &http.Server{
|
||||
Addr: fmt.Sprintf(":%d", port),
|
||||
Handler: mux,
|
||||
ReadHeaderTimeout: 5 * time.Second,
|
||||
}
|
||||
|
||||
go func() {
|
||||
<-ctx.Done()
|
||||
shutdownCtx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
|
||||
defer cancel()
|
||||
_ = srv.Shutdown(shutdownCtx)
|
||||
}()
|
||||
|
||||
log.Info().Int("port", port).Msg("metrics HTTP server starting")
|
||||
if err := srv.ListenAndServe(); err != nil && err != http.ErrServerClosed {
|
||||
log.Error().Err(err).Msg("metrics HTTP server failed")
|
||||
}
|
||||
}
|
||||
@@ -9,7 +9,10 @@ import (
|
||||
"net/url"
|
||||
"time"
|
||||
|
||||
"github.com/rs/zerolog"
|
||||
"golang.org/x/time/rate"
|
||||
|
||||
"github.com/metadata-agregator/internal/metrics"
|
||||
)
|
||||
|
||||
const (
|
||||
@@ -32,7 +35,11 @@ func newClient() *client {
|
||||
}
|
||||
|
||||
func (c *client) get(ctx context.Context, endpoint string, params url.Values) ([]byte, error) {
|
||||
log := zerolog.Ctx(ctx)
|
||||
|
||||
log.Trace().Str("endpoint", endpoint).Msg("waiting for rate limiter")
|
||||
if err := c.limiter.Wait(ctx); err != nil {
|
||||
log.Debug().Err(err).Msg("rate limiter interrupted")
|
||||
return nil, fmt.Errorf("rate limiter: %w", err)
|
||||
}
|
||||
|
||||
@@ -51,29 +58,48 @@ func (c *client) get(ctx context.Context, endpoint string, params url.Values) ([
|
||||
req.Header.Set("User-Agent", userAgent)
|
||||
req.Header.Set("Accept", "application/json")
|
||||
|
||||
start := time.Now()
|
||||
log.Trace().Str("url", reqURL).Msg("sending HTTP request")
|
||||
|
||||
resp, err := c.http.Do(req)
|
||||
if err != nil {
|
||||
log.Debug().Err(err).Str("endpoint", endpoint).Dur("duration", time.Since(start)).Msg("HTTP request failed")
|
||||
metrics.ProviderRequests.WithLabelValues("musicbrainz", endpoint, "error").Inc()
|
||||
return nil, fmt.Errorf("do request: %w", err)
|
||||
}
|
||||
defer resp.Body.Close()
|
||||
|
||||
duration := time.Since(start)
|
||||
metrics.ProviderLatency.WithLabelValues("musicbrainz", endpoint).Observe(duration.Seconds())
|
||||
|
||||
if resp.StatusCode == http.StatusNotFound {
|
||||
log.Debug().Str("endpoint", endpoint).Int("status", resp.StatusCode).Dur("duration", duration).Msg("HTTP 404")
|
||||
metrics.ProviderRequests.WithLabelValues("musicbrainz", endpoint, "not_found").Inc()
|
||||
return nil, ErrNotFound
|
||||
}
|
||||
|
||||
if resp.StatusCode == http.StatusServiceUnavailable {
|
||||
log.Warn().Str("endpoint", endpoint).Dur("duration", duration).Msg("HTTP 503 rate limited by provider")
|
||||
metrics.ProviderRequests.WithLabelValues("musicbrainz", endpoint, "rate_limited").Inc()
|
||||
return nil, ErrRateLimited
|
||||
}
|
||||
|
||||
if resp.StatusCode != http.StatusOK {
|
||||
body, _ := io.ReadAll(resp.Body)
|
||||
log.Warn().Str("endpoint", endpoint).Int("status", resp.StatusCode).Str("body", string(body)).Dur("duration", duration).Msg("unexpected HTTP status")
|
||||
metrics.ProviderRequests.WithLabelValues("musicbrainz", endpoint, fmt.Sprintf("%d", resp.StatusCode)).Inc()
|
||||
return nil, fmt.Errorf("unexpected status %d: %s", resp.StatusCode, string(body))
|
||||
}
|
||||
|
||||
log.Trace().Str("endpoint", endpoint).Int("status", resp.StatusCode).Dur("duration", duration).Msg("HTTP request succeeded")
|
||||
metrics.ProviderRequests.WithLabelValues("musicbrainz", endpoint, "ok").Inc()
|
||||
|
||||
return io.ReadAll(resp.Body)
|
||||
}
|
||||
|
||||
func (c *client) lookup(ctx context.Context, entity, id string, inc []string) ([]byte, error) {
|
||||
zerolog.Ctx(ctx).Debug().Str("entity", entity).Str("id", id).Strs("includes", inc).Msg("provider lookup")
|
||||
|
||||
params := url.Values{}
|
||||
if len(inc) > 0 {
|
||||
incStr := ""
|
||||
@@ -90,6 +116,14 @@ func (c *client) lookup(ctx context.Context, entity, id string, inc []string) ([
|
||||
}
|
||||
|
||||
func (c *client) browse(ctx context.Context, entity, linkedEntity, linkedID string, limit, offset int, inc []string) ([]byte, error) {
|
||||
zerolog.Ctx(ctx).Debug().
|
||||
Str("entity", entity).
|
||||
Str("linked_entity", linkedEntity).
|
||||
Str("linked_id", linkedID).
|
||||
Int("limit", limit).
|
||||
Int("offset", offset).
|
||||
Msg("provider browse")
|
||||
|
||||
params := url.Values{}
|
||||
params.Set(linkedEntity, linkedID)
|
||||
params.Set("limit", fmt.Sprintf("%d", limit))
|
||||
@@ -110,6 +144,8 @@ func (c *client) browse(ctx context.Context, entity, linkedEntity, linkedID stri
|
||||
}
|
||||
|
||||
func (c *client) search(ctx context.Context, entity, query string, limit, offset int) ([]byte, error) {
|
||||
zerolog.Ctx(ctx).Debug().Str("entity", entity).Str("query", query).Int("limit", limit).Int("offset", offset).Msg("provider search")
|
||||
|
||||
params := url.Values{}
|
||||
params.Set("query", query)
|
||||
params.Set("limit", fmt.Sprintf("%d", limit))
|
||||
|
||||
+68
-18
@@ -4,6 +4,7 @@ import (
|
||||
"context"
|
||||
"errors"
|
||||
|
||||
"github.com/rs/zerolog"
|
||||
"google.golang.org/grpc/codes"
|
||||
"google.golang.org/grpc/status"
|
||||
|
||||
@@ -22,13 +23,16 @@ func NewMetadataServer(services map[metadatav1.Provider]*service.MetadataService
|
||||
return &MetadataServer{services: services}
|
||||
}
|
||||
|
||||
func (s *MetadataServer) getService(p metadatav1.Provider) (*service.MetadataService, error) {
|
||||
func (s *MetadataServer) getService(ctx context.Context, p metadatav1.Provider) (*service.MetadataService, error) {
|
||||
if p == metadatav1.Provider_PROVIDER_UNSPECIFIED {
|
||||
p = metadatav1.Provider_PROVIDER_MUSICBRAINZ
|
||||
}
|
||||
|
||||
zerolog.Ctx(ctx).Debug().Str("provider", p.String()).Msg("resolved provider")
|
||||
|
||||
svc, ok := s.services[p]
|
||||
if !ok {
|
||||
zerolog.Ctx(ctx).Warn().Str("provider", p.String()).Msg("unknown provider requested")
|
||||
return nil, status.Errorf(codes.InvalidArgument, "unknown provider: %v", p)
|
||||
}
|
||||
|
||||
@@ -36,7 +40,9 @@ func (s *MetadataServer) getService(p metadatav1.Provider) (*service.MetadataSer
|
||||
}
|
||||
|
||||
func (s *MetadataServer) GetArtist(ctx context.Context, req *metadatav1.GetArtistRequest) (*metadatav1.Artist, error) {
|
||||
svc, err := s.getService(req.Provider)
|
||||
log := zerolog.Ctx(ctx)
|
||||
|
||||
svc, err := s.getService(ctx, req.Provider)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
@@ -45,22 +51,28 @@ func (s *MetadataServer) GetArtist(ctx context.Context, req *metadatav1.GetArtis
|
||||
switch v := req.Identifier.(type) {
|
||||
case *metadatav1.GetArtistRequest_Id:
|
||||
id = v.Id
|
||||
log.Debug().Str("lookup", "id").Str("artist_id", id).Msg("getting artist")
|
||||
case *metadatav1.GetArtistRequest_External:
|
||||
id = v.External.SourceId
|
||||
log.Debug().Str("lookup", "external").Str("source_id", id).Str("source", v.External.Source).Msg("getting artist")
|
||||
default:
|
||||
log.Warn().Msg("get artist called without identifier")
|
||||
return nil, status.Error(codes.InvalidArgument, "identifier required")
|
||||
}
|
||||
|
||||
artist, err := svc.GetArtist(ctx, id)
|
||||
if err != nil {
|
||||
return nil, toGRPCError(err)
|
||||
return nil, toGRPCError(ctx, err)
|
||||
}
|
||||
|
||||
log.Trace().Str("artist_id", artist.ID).Str("name", artist.Name).Msg("artist found")
|
||||
return toProtoArtist(artist), nil
|
||||
}
|
||||
|
||||
func (s *MetadataServer) SearchArtists(ctx context.Context, req *metadatav1.SearchArtistsRequest) (*metadatav1.SearchArtistsResponse, error) {
|
||||
svc, err := s.getService(req.Provider)
|
||||
log := zerolog.Ctx(ctx)
|
||||
|
||||
svc, err := s.getService(ctx, req.Provider)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
@@ -70,9 +82,11 @@ func (s *MetadataServer) SearchArtists(ctx context.Context, req *metadatav1.Sear
|
||||
limit = 25
|
||||
}
|
||||
|
||||
log.Debug().Str("query", req.Query).Int("limit", limit).Int("offset", int(req.Offset)).Msg("searching artists")
|
||||
|
||||
result, err := svc.SearchArtists(ctx, req.Query, limit, int(req.Offset))
|
||||
if err != nil {
|
||||
return nil, toGRPCError(err)
|
||||
return nil, toGRPCError(ctx, err)
|
||||
}
|
||||
|
||||
resp := &metadatav1.SearchArtistsResponse{
|
||||
@@ -83,11 +97,14 @@ func (s *MetadataServer) SearchArtists(ctx context.Context, req *metadatav1.Sear
|
||||
resp.Artists = append(resp.Artists, toProtoArtist(&a))
|
||||
}
|
||||
|
||||
log.Trace().Int("total", result.Total).Int("returned", len(resp.Artists)).Msg("artist search complete")
|
||||
return resp, nil
|
||||
}
|
||||
|
||||
func (s *MetadataServer) SearchAlbums(ctx context.Context, req *metadatav1.SearchAlbumsRequest) (*metadatav1.SearchAlbumsResponse, error) {
|
||||
svc, err := s.getService(req.Provider)
|
||||
log := zerolog.Ctx(ctx)
|
||||
|
||||
svc, err := s.getService(ctx, req.Provider)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
@@ -97,9 +114,11 @@ func (s *MetadataServer) SearchAlbums(ctx context.Context, req *metadatav1.Searc
|
||||
limit = 25
|
||||
}
|
||||
|
||||
log.Debug().Str("query", req.Query).Str("artist", req.Artist).Int("limit", limit).Int("offset", int(req.Offset)).Msg("searching albums")
|
||||
|
||||
result, err := svc.SearchAlbums(ctx, req.Query, req.Artist, limit, int(req.Offset))
|
||||
if err != nil {
|
||||
return nil, toGRPCError(err)
|
||||
return nil, toGRPCError(ctx, err)
|
||||
}
|
||||
|
||||
resp := &metadatav1.SearchAlbumsResponse{
|
||||
@@ -110,11 +129,14 @@ func (s *MetadataServer) SearchAlbums(ctx context.Context, req *metadatav1.Searc
|
||||
resp.Albums = append(resp.Albums, toProtoAlbum(&a))
|
||||
}
|
||||
|
||||
log.Trace().Int("total", result.Total).Int("returned", len(resp.Albums)).Msg("album search complete")
|
||||
return resp, nil
|
||||
}
|
||||
|
||||
func (s *MetadataServer) GetAlbum(ctx context.Context, req *metadatav1.GetAlbumRequest) (*metadatav1.Album, error) {
|
||||
svc, err := s.getService(req.Provider)
|
||||
log := zerolog.Ctx(ctx)
|
||||
|
||||
svc, err := s.getService(ctx, req.Provider)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
@@ -123,22 +145,28 @@ func (s *MetadataServer) GetAlbum(ctx context.Context, req *metadatav1.GetAlbumR
|
||||
switch v := req.Identifier.(type) {
|
||||
case *metadatav1.GetAlbumRequest_Id:
|
||||
id = v.Id
|
||||
log.Debug().Str("lookup", "id").Str("album_id", id).Msg("getting album")
|
||||
case *metadatav1.GetAlbumRequest_External:
|
||||
id = v.External.SourceId
|
||||
log.Debug().Str("lookup", "external").Str("source_id", id).Str("source", v.External.Source).Msg("getting album")
|
||||
default:
|
||||
log.Warn().Msg("get album called without identifier")
|
||||
return nil, status.Error(codes.InvalidArgument, "identifier required")
|
||||
}
|
||||
|
||||
album, err := svc.GetAlbum(ctx, id)
|
||||
if err != nil {
|
||||
return nil, toGRPCError(err)
|
||||
return nil, toGRPCError(ctx, err)
|
||||
}
|
||||
|
||||
log.Trace().Str("album_id", album.ID).Str("title", album.Title).Msg("album found")
|
||||
return toProtoAlbum(album), nil
|
||||
}
|
||||
|
||||
func (s *MetadataServer) GetArtistAlbums(ctx context.Context, req *metadatav1.GetArtistAlbumsRequest) (*metadatav1.GetArtistAlbumsResponse, error) {
|
||||
svc, err := s.getService(req.Provider)
|
||||
log := zerolog.Ctx(ctx)
|
||||
|
||||
svc, err := s.getService(ctx, req.Provider)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
@@ -148,9 +176,11 @@ func (s *MetadataServer) GetArtistAlbums(ctx context.Context, req *metadatav1.Ge
|
||||
limit = 25
|
||||
}
|
||||
|
||||
log.Debug().Str("artist_id", req.ArtistId).Int("limit", limit).Int("offset", int(req.Offset)).Msg("getting artist albums")
|
||||
|
||||
result, err := svc.GetArtistAlbums(ctx, req.ArtistId, limit, int(req.Offset))
|
||||
if err != nil {
|
||||
return nil, toGRPCError(err)
|
||||
return nil, toGRPCError(ctx, err)
|
||||
}
|
||||
|
||||
resp := &metadatav1.GetArtistAlbumsResponse{
|
||||
@@ -161,11 +191,14 @@ func (s *MetadataServer) GetArtistAlbums(ctx context.Context, req *metadatav1.Ge
|
||||
resp.Albums = append(resp.Albums, toProtoAlbum(&a))
|
||||
}
|
||||
|
||||
log.Trace().Int("total", result.Total).Int("returned", len(resp.Albums)).Msg("artist albums retrieved")
|
||||
return resp, nil
|
||||
}
|
||||
|
||||
func (s *MetadataServer) GetTrack(ctx context.Context, req *metadatav1.GetTrackRequest) (*metadatav1.Track, error) {
|
||||
svc, err := s.getService(req.Provider)
|
||||
log := zerolog.Ctx(ctx)
|
||||
|
||||
svc, err := s.getService(ctx, req.Provider)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
@@ -174,42 +207,51 @@ func (s *MetadataServer) GetTrack(ctx context.Context, req *metadatav1.GetTrackR
|
||||
|
||||
switch v := req.Identifier.(type) {
|
||||
case *metadatav1.GetTrackRequest_Id:
|
||||
log.Debug().Str("lookup", "id").Str("track_id", v.Id).Msg("getting track")
|
||||
t, err := svc.GetTrack(ctx, v.Id)
|
||||
if err != nil {
|
||||
return nil, toGRPCError(err)
|
||||
return nil, toGRPCError(ctx, err)
|
||||
}
|
||||
track = toProtoTrack(t)
|
||||
|
||||
case *metadatav1.GetTrackRequest_External:
|
||||
log.Debug().Str("lookup", "external").Str("source_id", v.External.SourceId).Str("source", v.External.Source).Msg("getting track")
|
||||
t, err := svc.GetTrack(ctx, v.External.SourceId)
|
||||
if err != nil {
|
||||
return nil, toGRPCError(err)
|
||||
return nil, toGRPCError(ctx, err)
|
||||
}
|
||||
track = toProtoTrack(t)
|
||||
|
||||
case *metadatav1.GetTrackRequest_Isrc:
|
||||
log.Debug().Str("lookup", "isrc").Str("isrc", v.Isrc).Msg("getting track")
|
||||
t, err := svc.GetTrackByISRC(ctx, v.Isrc)
|
||||
if err != nil {
|
||||
return nil, toGRPCError(err)
|
||||
return nil, toGRPCError(ctx, err)
|
||||
}
|
||||
track = toProtoTrack(t)
|
||||
|
||||
default:
|
||||
log.Warn().Msg("get track called without identifier")
|
||||
return nil, status.Error(codes.InvalidArgument, "identifier required")
|
||||
}
|
||||
|
||||
log.Trace().Str("track_id", track.Id).Str("title", track.Title).Msg("track found")
|
||||
return track, nil
|
||||
}
|
||||
|
||||
func (s *MetadataServer) GetAlbumTracks(ctx context.Context, req *metadatav1.GetAlbumTracksRequest) (*metadatav1.GetAlbumTracksResponse, error) {
|
||||
svc, err := s.getService(req.Provider)
|
||||
log := zerolog.Ctx(ctx)
|
||||
|
||||
svc, err := s.getService(ctx, req.Provider)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
log.Debug().Str("album_id", req.AlbumId).Msg("getting album tracks")
|
||||
|
||||
tracks, err := svc.GetAlbumTracks(ctx, req.AlbumId)
|
||||
if err != nil {
|
||||
return nil, toGRPCError(err)
|
||||
return nil, toGRPCError(ctx, err)
|
||||
}
|
||||
|
||||
resp := &metadatav1.GetAlbumTracksResponse{}
|
||||
@@ -217,29 +259,37 @@ func (s *MetadataServer) GetAlbumTracks(ctx context.Context, req *metadatav1.Get
|
||||
resp.Tracks = append(resp.Tracks, toProtoTrack(&t))
|
||||
}
|
||||
|
||||
log.Trace().Int("track_count", len(resp.Tracks)).Msg("album tracks retrieved")
|
||||
return resp, nil
|
||||
}
|
||||
|
||||
func (s *MetadataServer) SyncArtist(ctx context.Context, req *metadatav1.SyncArtistRequest) (*metadatav1.SyncArtistResponse, error) {
|
||||
zerolog.Ctx(ctx).Warn().Msg("sync artist called but not implemented")
|
||||
return nil, status.Error(codes.Unimplemented, "sync not yet implemented")
|
||||
}
|
||||
|
||||
func toGRPCError(err error) error {
|
||||
func toGRPCError(ctx context.Context, err error) error {
|
||||
if err == nil {
|
||||
return nil
|
||||
}
|
||||
|
||||
log := zerolog.Ctx(ctx)
|
||||
|
||||
if errors.Is(err, repository.ErrNotFound) {
|
||||
log.Debug().Msg("entity not found")
|
||||
return status.Error(codes.NotFound, "not found")
|
||||
}
|
||||
|
||||
if errors.Is(err, musicbrainz.ErrNotFound) {
|
||||
log.Debug().Msg("provider entity not found")
|
||||
return status.Error(codes.NotFound, "not found")
|
||||
}
|
||||
|
||||
if errors.Is(err, musicbrainz.ErrRateLimited) {
|
||||
log.Warn().Msg("provider rate limited")
|
||||
return status.Error(codes.ResourceExhausted, "rate limited")
|
||||
}
|
||||
|
||||
log.Error().Err(err).Msg("internal error")
|
||||
return status.Errorf(codes.Internal, "internal error: %v", err)
|
||||
}
|
||||
|
||||
@@ -4,7 +4,10 @@ import (
|
||||
"context"
|
||||
"errors"
|
||||
|
||||
"github.com/rs/zerolog"
|
||||
|
||||
"github.com/metadata-agregator/internal/domain"
|
||||
"github.com/metadata-agregator/internal/metrics"
|
||||
"github.com/metadata-agregator/internal/provider"
|
||||
"github.com/metadata-agregator/internal/repository"
|
||||
)
|
||||
@@ -31,116 +34,179 @@ func NewMetadataService(
|
||||
}
|
||||
|
||||
func (s *MetadataService) GetArtist(ctx context.Context, id string) (*domain.Artist, error) {
|
||||
log := zerolog.Ctx(ctx)
|
||||
|
||||
artist, err := s.artists.GetByExternalID(ctx, s.provider.Name(), id)
|
||||
if err == nil {
|
||||
log.Debug().Str("artist_id", id).Msg("artist cache hit")
|
||||
metrics.CacheHits.WithLabelValues("artist").Inc()
|
||||
return artist, nil
|
||||
}
|
||||
if !errors.Is(err, repository.ErrNotFound) {
|
||||
log.Error().Err(err).Str("artist_id", id).Msg("artist cache lookup failed")
|
||||
return nil, err
|
||||
}
|
||||
|
||||
metrics.CacheMisses.WithLabelValues("artist").Inc()
|
||||
log.Debug().Str("artist_id", id).Str("provider", s.provider.Name()).Msg("artist cache miss, fetching from provider")
|
||||
|
||||
artist, err = s.provider.GetArtist(ctx, id)
|
||||
if err != nil {
|
||||
log.Debug().Err(err).Str("artist_id", id).Msg("provider fetch failed")
|
||||
return nil, err
|
||||
}
|
||||
|
||||
if saveErr := s.artists.Save(ctx, artist); saveErr != nil {
|
||||
log.Warn().Err(saveErr).Str("artist_id", id).Msg("failed to cache artist")
|
||||
return artist, nil
|
||||
}
|
||||
|
||||
log.Trace().Str("artist_id", id).Msg("artist cached successfully")
|
||||
return artist, nil
|
||||
}
|
||||
|
||||
func (s *MetadataService) SearchArtists(ctx context.Context, query string, limit, offset int) (*domain.SearchResult[domain.Artist], error) {
|
||||
log := zerolog.Ctx(ctx)
|
||||
|
||||
result, err := s.artists.Search(ctx, query, limit, offset)
|
||||
if err == nil && len(result.Items) > 0 {
|
||||
log.Debug().Int("count", len(result.Items)).Msg("artist search served from cache")
|
||||
metrics.CacheHits.WithLabelValues("artist_search").Inc()
|
||||
return result, nil
|
||||
}
|
||||
|
||||
metrics.CacheMisses.WithLabelValues("artist_search").Inc()
|
||||
log.Debug().Str("query", query).Str("provider", s.provider.Name()).Msg("artist search cache miss, querying provider")
|
||||
return s.provider.SearchArtists(ctx, query, limit, offset)
|
||||
}
|
||||
|
||||
func (s *MetadataService) SearchAlbums(ctx context.Context, query string, artist string, limit, offset int) (*domain.SearchResult[domain.Album], error) {
|
||||
zerolog.Ctx(ctx).Debug().Str("query", query).Str("artist", artist).Str("provider", s.provider.Name()).Msg("searching albums via provider")
|
||||
return s.provider.SearchAlbums(ctx, query, artist, limit, offset)
|
||||
}
|
||||
|
||||
func (s *MetadataService) GetAlbum(ctx context.Context, id string) (*domain.Album, error) {
|
||||
log := zerolog.Ctx(ctx)
|
||||
|
||||
album, err := s.albums.GetByExternalID(ctx, s.provider.Name(), id)
|
||||
if err == nil {
|
||||
log.Debug().Str("album_id", id).Msg("album cache hit")
|
||||
metrics.CacheHits.WithLabelValues("album").Inc()
|
||||
return album, nil
|
||||
}
|
||||
if !errors.Is(err, repository.ErrNotFound) {
|
||||
log.Error().Err(err).Str("album_id", id).Msg("album cache lookup failed")
|
||||
return nil, err
|
||||
}
|
||||
|
||||
metrics.CacheMisses.WithLabelValues("album").Inc()
|
||||
log.Debug().Str("album_id", id).Str("provider", s.provider.Name()).Msg("album cache miss, fetching from provider")
|
||||
|
||||
album, err = s.provider.GetAlbum(ctx, id)
|
||||
if err != nil {
|
||||
log.Debug().Err(err).Str("album_id", id).Msg("provider fetch failed")
|
||||
return nil, err
|
||||
}
|
||||
|
||||
if saveErr := s.albums.Save(ctx, album); saveErr != nil {
|
||||
log.Warn().Err(saveErr).Str("album_id", id).Msg("failed to cache album")
|
||||
return album, nil
|
||||
}
|
||||
|
||||
log.Trace().Str("album_id", id).Msg("album cached successfully")
|
||||
return album, nil
|
||||
}
|
||||
|
||||
func (s *MetadataService) GetArtistAlbums(ctx context.Context, artistID string, limit, offset int) (*domain.SearchResult[domain.Album], error) {
|
||||
log := zerolog.Ctx(ctx)
|
||||
|
||||
result, err := s.albums.GetByArtistID(ctx, artistID, limit, offset)
|
||||
if err == nil && len(result.Items) > 0 {
|
||||
log.Debug().Str("artist_id", artistID).Int("count", len(result.Items)).Msg("artist albums served from cache")
|
||||
metrics.CacheHits.WithLabelValues("artist_albums").Inc()
|
||||
return result, nil
|
||||
}
|
||||
|
||||
metrics.CacheMisses.WithLabelValues("artist_albums").Inc()
|
||||
log.Debug().Str("artist_id", artistID).Str("provider", s.provider.Name()).Msg("artist albums cache miss, querying provider")
|
||||
return s.provider.GetArtistAlbums(ctx, artistID, limit, offset)
|
||||
}
|
||||
|
||||
func (s *MetadataService) GetTrack(ctx context.Context, id string) (*domain.Track, error) {
|
||||
log := zerolog.Ctx(ctx)
|
||||
|
||||
track, err := s.tracks.GetByExternalID(ctx, s.provider.Name(), id)
|
||||
if err == nil {
|
||||
log.Debug().Str("track_id", id).Msg("track cache hit")
|
||||
metrics.CacheHits.WithLabelValues("track").Inc()
|
||||
return track, nil
|
||||
}
|
||||
if !errors.Is(err, repository.ErrNotFound) {
|
||||
log.Error().Err(err).Str("track_id", id).Msg("track cache lookup failed")
|
||||
return nil, err
|
||||
}
|
||||
|
||||
metrics.CacheMisses.WithLabelValues("track").Inc()
|
||||
log.Debug().Str("track_id", id).Str("provider", s.provider.Name()).Msg("track cache miss, fetching from provider")
|
||||
|
||||
track, err = s.provider.GetTrack(ctx, id)
|
||||
if err != nil {
|
||||
log.Debug().Err(err).Str("track_id", id).Msg("provider fetch failed")
|
||||
return nil, err
|
||||
}
|
||||
|
||||
if saveErr := s.tracks.Save(ctx, track); saveErr != nil {
|
||||
log.Warn().Err(saveErr).Str("track_id", id).Msg("failed to cache track")
|
||||
return track, nil
|
||||
}
|
||||
|
||||
log.Trace().Str("track_id", id).Msg("track cached successfully")
|
||||
return track, nil
|
||||
}
|
||||
|
||||
func (s *MetadataService) GetTrackByISRC(ctx context.Context, isrc string) (*domain.Track, error) {
|
||||
log := zerolog.Ctx(ctx)
|
||||
|
||||
track, err := s.tracks.GetByISRC(ctx, isrc)
|
||||
if err == nil {
|
||||
log.Debug().Str("isrc", isrc).Msg("track ISRC cache hit")
|
||||
metrics.CacheHits.WithLabelValues("track_isrc").Inc()
|
||||
return track, nil
|
||||
}
|
||||
if !errors.Is(err, repository.ErrNotFound) {
|
||||
log.Error().Err(err).Str("isrc", isrc).Msg("track ISRC cache lookup failed")
|
||||
return nil, err
|
||||
}
|
||||
|
||||
metrics.CacheMisses.WithLabelValues("track_isrc").Inc()
|
||||
log.Debug().Str("isrc", isrc).Str("provider", s.provider.Name()).Msg("track ISRC cache miss, fetching from provider")
|
||||
|
||||
track, err = s.provider.GetTrackByISRC(ctx, isrc)
|
||||
if err != nil {
|
||||
log.Debug().Err(err).Str("isrc", isrc).Msg("provider ISRC fetch failed")
|
||||
return nil, err
|
||||
}
|
||||
|
||||
if saveErr := s.tracks.Save(ctx, track); saveErr != nil {
|
||||
log.Warn().Err(saveErr).Str("isrc", isrc).Msg("failed to cache track")
|
||||
return track, nil
|
||||
}
|
||||
|
||||
log.Trace().Str("isrc", isrc).Msg("track cached successfully")
|
||||
return track, nil
|
||||
}
|
||||
|
||||
func (s *MetadataService) GetAlbumTracks(ctx context.Context, albumID string) ([]domain.Track, error) {
|
||||
log := zerolog.Ctx(ctx)
|
||||
|
||||
tracks, err := s.tracks.GetByAlbumID(ctx, albumID)
|
||||
if err == nil && len(tracks) > 0 {
|
||||
log.Debug().Str("album_id", albumID).Int("count", len(tracks)).Msg("album tracks served from cache")
|
||||
metrics.CacheHits.WithLabelValues("album_tracks").Inc()
|
||||
return tracks, nil
|
||||
}
|
||||
|
||||
metrics.CacheMisses.WithLabelValues("album_tracks").Inc()
|
||||
log.Debug().Str("album_id", albumID).Str("provider", s.provider.Name()).Msg("album tracks cache miss, querying provider")
|
||||
return s.provider.GetAlbumTracks(ctx, albumID)
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user