From 3d111f9008a91cfb335f6dc7f647c1d68684a5e4 Mon Sep 17 00:00:00 2001 From: Alexander Date: Thu, 7 May 2026 16:50:40 +0200 Subject: [PATCH] Refactor logging --- cmd/server/main.go | 68 ++++++++---- config.example.yaml | 8 ++ flake.nix | 2 +- go.mod | 24 ++++- go.sum | 79 +++++++++++--- internal/config/config.go | 20 ++++ internal/logging/interceptor.go | 137 ++++++++++++++++++++++++ internal/logging/logger.go | 21 ++++ internal/metrics/metrics.go | 115 ++++++++++++++++++++ internal/provider/musicbrainz/client.go | 36 +++++++ internal/server/server.go | 86 +++++++++++---- internal/service/metadata.go | 66 ++++++++++++ tests/e2e/metadata_test.go | 13 ++- tests/e2e/noop_repo_test.go | 66 ++++++++++++ 14 files changed, 683 insertions(+), 58 deletions(-) create mode 100644 internal/logging/interceptor.go create mode 100644 internal/logging/logger.go create mode 100644 internal/metrics/metrics.go create mode 100644 tests/e2e/noop_repo_test.go diff --git a/cmd/server/main.go b/cmd/server/main.go index e5b1dc1..fa8d2ea 100644 --- a/cmd/server/main.go +++ b/cmd/server/main.go @@ -4,18 +4,21 @@ import ( "context" "flag" "fmt" - "log" "net" "os" "os/signal" "syscall" "time" + "github.com/grpc-ecosystem/go-grpc-middleware/v2/interceptors/logging" "github.com/jackc/pgx/v5/pgxpool" + "github.com/rs/zerolog/log" "google.golang.org/grpc" "google.golang.org/grpc/reflection" "github.com/metadata-agregator/internal/config" + applog "github.com/metadata-agregator/internal/logging" + "github.com/metadata-agregator/internal/metrics" "github.com/metadata-agregator/internal/provider/musicbrainz" "github.com/metadata-agregator/internal/repository/postgres" "github.com/metadata-agregator/internal/server" @@ -29,30 +32,54 @@ func main() { cfg, err := config.Load(*configPath) if err != nil { - log.Fatalf("failed to load config: %v", err) + fmt.Fprintf(os.Stderr, "failed to load config: %v\n", err) + os.Exit(1) } - ctx := context.Background() + applog.Init(cfg.Logging) + + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() services, cleanup := buildServices(ctx, cfg) defer cleanup() + srvMetrics := metrics.NewServerMetrics() + addr := fmt.Sprintf(":%d", cfg.Server.Port) lis, err := net.Listen("tcp", addr) if err != nil { - log.Fatalf("failed to listen: %v", err) + log.Fatal().Err(err).Str("addr", addr).Msg("failed to listen") } - grpcServer := grpc.NewServer() + grpcServer := grpc.NewServer( + grpc.ChainUnaryInterceptor( + applog.RequestIDInterceptor(), + applog.PeerInfoInterceptor(), + srvMetrics.UnaryServerInterceptor(), + logging.UnaryServerInterceptor(applog.InterceptorLogger(), applog.LoggingOpts()...), + ), + grpc.ChainStreamInterceptor( + applog.StreamRequestIDInterceptor(), + srvMetrics.StreamServerInterceptor(), + logging.StreamServerInterceptor(applog.InterceptorLogger(), applog.LoggingOpts()...), + ), + ) metadatav1.RegisterMetadataServiceServer(grpcServer, server.NewMetadataServer(services)) reflection.Register(grpcServer) - go gracefulShutdown(grpcServer) + srvMetrics.InitializeMetrics(grpcServer) - log.Printf("gRPC server listening on %s", addr) + if cfg.Metrics.Enabled { + go srvMetrics.StartHTTPServer(ctx, cfg.Metrics.Port) + } + + go gracefulShutdown(grpcServer, cancel) + + log.Info().Str("addr", addr).Msg("gRPC server listening") if err := grpcServer.Serve(lis); err != nil { - log.Fatalf("failed to serve: %v", err) + log.Fatal().Err(err).Msg("failed to serve") } } @@ -66,7 +93,7 @@ func buildServices(ctx context.Context, cfg *config.Config) (map[metadatav1.Prov } if dbURL == "" { - log.Println("no database configured, running in provider-only mode") + log.Warn().Msg("no database configured, running in provider-only mode") services[metadatav1.Provider_PROVIDER_MUSICBRAINZ] = service.NewMetadataService( &noopArtistRepo{}, &noopAlbumRepo{}, @@ -78,7 +105,7 @@ func buildServices(ctx context.Context, cfg *config.Config) (map[metadatav1.Prov pool, err := connectDB(ctx, dbURL) if err != nil { - log.Fatalf("database connection failed: %v", err) + log.Fatal().Err(err).Msg("database connection failed") } artistRepo := postgres.NewArtistRepository(pool) @@ -92,7 +119,7 @@ func buildServices(ctx context.Context, cfg *config.Config) (map[metadatav1.Prov mb, ) - log.Println("database connected, caching enabled") + log.Info().Msg("database connected, caching enabled") return services, func() { pool.Close() } } @@ -100,15 +127,17 @@ func connectDB(ctx context.Context, dbURL string) (*pgxpool.Pool, error) { ctx, cancel := context.WithTimeout(ctx, 10*time.Second) defer cancel() - config, err := pgxpool.ParseConfig(dbURL) + pgxCfg, err := pgxpool.ParseConfig(dbURL) if err != nil { return nil, err } - config.MaxConns = 10 - config.MinConns = 2 + pgxCfg.MaxConns = 10 + pgxCfg.MinConns = 2 - pool, err := pgxpool.NewWithConfig(ctx, config) + log.Debug().Msg("connecting to database") + + pool, err := pgxpool.NewWithConfig(ctx, pgxCfg) if err != nil { return nil, err } @@ -121,10 +150,11 @@ func connectDB(ctx context.Context, dbURL string) (*pgxpool.Pool, error) { return pool, nil } -func gracefulShutdown(server *grpc.Server) { +func gracefulShutdown(srv *grpc.Server, cancel context.CancelFunc) { sigCh := make(chan os.Signal, 1) signal.Notify(sigCh, syscall.SIGINT, syscall.SIGTERM) - <-sigCh - log.Println("shutting down...") - server.GracefulStop() + sig := <-sigCh + log.Info().Str("signal", sig.String()).Msg("shutting down") + cancel() + srv.GracefulStop() } diff --git a/config.example.yaml b/config.example.yaml index 7861a28..ac26e83 100644 --- a/config.example.yaml +++ b/config.example.yaml @@ -8,3 +8,11 @@ database: password: metadata name: metadata sslmode: disable + +logging: + level: info + format: json + +metrics: + enabled: true + port: 9090 diff --git a/flake.nix b/flake.nix index bfff0da..461c578 100644 --- a/flake.nix +++ b/flake.nix @@ -41,7 +41,7 @@ pname = "metadata-server"; version = "0.1.0"; src = ./.; - vendorHash = null; + vendorHash = "sha256-qNNFHRKqGoo/GveSl9+Bm+DkTKDoSQeGtG8p1edmCJk="; subPackages = [ "cmd/server" ]; meta = { diff --git a/go.mod b/go.mod index 4f918fc..de91156 100644 --- a/go.mod +++ b/go.mod @@ -3,22 +3,36 @@ module github.com/metadata-agregator go 1.25.0 require ( + github.com/google/uuid v1.6.0 + github.com/grpc-ecosystem/go-grpc-middleware/providers/prometheus v1.1.0 + github.com/grpc-ecosystem/go-grpc-middleware/v2 v2.3.3 github.com/jackc/pgx/v5 v5.9.2 + github.com/prometheus/client_golang v1.23.2 + github.com/rs/zerolog v1.35.1 golang.org/x/time v0.15.0 - google.golang.org/grpc v1.68.0 - google.golang.org/protobuf v1.35.2 + google.golang.org/grpc v1.74.2 + google.golang.org/protobuf v1.36.8 gopkg.in/yaml.v3 v3.0.1 ) require ( + github.com/beorn7/perks v1.0.1 // indirect + github.com/cespare/xxhash/v2 v2.3.0 // indirect github.com/jackc/pgpassfile v1.0.0 // indirect github.com/jackc/pgservicefile v0.0.0-20240606120523-5a60cdf6a761 // indirect github.com/jackc/puddle/v2 v2.2.2 // indirect github.com/kr/text v0.2.0 // indirect + github.com/mattn/go-colorable v0.1.14 // indirect + github.com/mattn/go-isatty v0.0.20 // indirect + github.com/munnerz/goautoneg v0.0.0-20191010083416-a7dc8b61c822 // indirect + github.com/prometheus/client_model v0.6.2 // indirect + github.com/prometheus/common v0.66.1 // indirect + github.com/prometheus/procfs v0.16.1 // indirect github.com/rogpeppe/go-internal v1.14.1 // indirect - golang.org/x/net v0.29.0 // indirect + go.yaml.in/yaml/v2 v2.4.2 // indirect + golang.org/x/net v0.43.0 // indirect golang.org/x/sync v0.17.0 // indirect - golang.org/x/sys v0.26.0 // indirect + golang.org/x/sys v0.35.0 // indirect golang.org/x/text v0.29.0 // indirect - google.golang.org/genproto/googleapis/rpc v0.0.0-20240903143218-8af14fe29dc1 // indirect + google.golang.org/genproto/googleapis/rpc v0.0.0-20250528174236-200df99c418a // indirect ) diff --git a/go.sum b/go.sum index 256cad0..d85a063 100644 --- a/go.sum +++ b/go.sum @@ -1,11 +1,25 @@ +github.com/beorn7/perks v1.0.1 h1:VlbKKnNfV8bJzeqoa4cOKqO6bYr3WgKZxO8Z16+hsOM= +github.com/beorn7/perks v1.0.1/go.mod h1:G2ZrVWU2WbWT9wwq4/hrbKbnv/1ERSJQ0ibhJ6rlkpw= +github.com/cespare/xxhash/v2 v2.3.0 h1:UL815xU9SqsFlibzuggzjXhog7bL6oX9BbNZnL2UFvs= +github.com/cespare/xxhash/v2 v2.3.0/go.mod h1:VGX0DQ3Q6kWi7AoAeZDth3/j3BFtOZR5XLFGgcrjCOs= github.com/creack/pty v1.1.9/go.mod h1:oKZEueFk5CKHvIhNR5MUki03XCEU+Q6VDXinZuGJ33E= github.com/davecgh/go-spew v1.1.0/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= github.com/davecgh/go-spew v1.1.1 h1:vj9j/u1bqnvCEfJOwUhtlOARqs3+rkHYY13jYWTU97c= github.com/davecgh/go-spew v1.1.1/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= +github.com/go-logr/logr v1.4.3 h1:CjnDlHq8ikf6E492q6eKboGOC0T8CDaOvkHCIg8idEI= +github.com/go-logr/logr v1.4.3/go.mod h1:9T104GzyrTigFIr8wt5mBrctHMim0Nb2HLGrmQ40KvY= +github.com/go-logr/stdr v1.2.2 h1:hSWxHoqTgW2S2qGc0LTAI563KZ5YKYRhT3MFKZMbjag= +github.com/go-logr/stdr v1.2.2/go.mod h1:mMo/vtBO5dYbehREoey6XUKy/eSumjCCveDpRre4VKE= github.com/golang/protobuf v1.5.4 h1:i7eJL8qZTpSEXOPTxNKhASYpMn+8e5Q6AdndVa1dWek= github.com/golang/protobuf v1.5.4/go.mod h1:lnTiLA8Wa4RWRcIUkrtSVa5nRhsEGBg48fD6rSs7xps= -github.com/google/go-cmp v0.6.0 h1:ofyhxvXcZhMsU5ulbFiLKl/XBFqE1GSq7atu8tAmTRI= -github.com/google/go-cmp v0.6.0/go.mod h1:17dUlkBOakJ0+DkrSSNjCkIjxS6bF9zb3elmeNGIjoY= +github.com/google/go-cmp v0.7.0 h1:wk8382ETsv4JYUZwIsn6YpYiWiBsYLSJiTsyBybVuN8= +github.com/google/go-cmp v0.7.0/go.mod h1:pXiqmnSA92OHEEa9HXL2W4E7lf9JzCmGVUdgjX3N/iU= +github.com/google/uuid v1.6.0 h1:NIvaJDMOsjHA8n1jAhLSgzrAzy1Hgr+hNrb57e+94F0= +github.com/google/uuid v1.6.0/go.mod h1:TIyPZe4MgqvfeYDBFedMoGGpEw/LqOeaOT+nhxU+yHo= +github.com/grpc-ecosystem/go-grpc-middleware/providers/prometheus v1.1.0 h1:QGLs/O40yoNK9vmy4rhUGBVyMf1lISBGtXRpsu/Qu/o= +github.com/grpc-ecosystem/go-grpc-middleware/providers/prometheus v1.1.0/go.mod h1:hM2alZsMUni80N33RBe6J0e423LB+odMj7d3EMP9l20= +github.com/grpc-ecosystem/go-grpc-middleware/v2 v2.3.3 h1:B+8ClL/kCQkRiU82d9xajRPKYMrB7E0MbtzWVi1K4ns= +github.com/grpc-ecosystem/go-grpc-middleware/v2 v2.3.3/go.mod h1:NbCUVmiS4foBGBHOYlCT25+YmGpJ32dZPi75pGEUpj4= github.com/jackc/pgpassfile v1.0.0 h1:/6Hmqy13Ss2zCq62VdNG8tM1wchn8zjSGOBJ6icpsIM= github.com/jackc/pgpassfile v1.0.0/go.mod h1:CEx0iS5ambNFdcRtxPj5JhEz+xB6uRky5eyVu/W2HEg= github.com/jackc/pgservicefile v0.0.0-20240606120523-5a60cdf6a761 h1:iCEnooe7UlwOQYpKFhBabPMi4aNAfoODPEFNiAnClxo= @@ -14,35 +28,72 @@ github.com/jackc/pgx/v5 v5.9.2 h1:3ZhOzMWnR4yJ+RW1XImIPsD1aNSz4T4fyP7zlQb56hw= github.com/jackc/pgx/v5 v5.9.2/go.mod h1:mal1tBGAFfLHvZzaYh77YS/eC6IX9OWbRV1QIIM0Jn4= github.com/jackc/puddle/v2 v2.2.2 h1:PR8nw+E/1w0GLuRFSmiioY6UooMp6KJv0/61nB7icHo= github.com/jackc/puddle/v2 v2.2.2/go.mod h1:vriiEXHvEE654aYKXXjOvZM39qJ0q+azkZFrfEOc3H4= -github.com/kr/pretty v0.3.0 h1:WgNl7dwNpEZ6jJ9k1snq4pZsg7DOEN8hP9Xw0Tsjwk0= -github.com/kr/pretty v0.3.0/go.mod h1:640gp4NfQd8pI5XOwp5fnNeVWj67G7CFk/SaSQn7NBk= +github.com/klauspost/compress v1.18.0 h1:c/Cqfb0r+Yi+JtIEq73FWXVkRonBlf0CRNYc8Zttxdo= +github.com/klauspost/compress v1.18.0/go.mod h1:2Pp+KzxcywXVXMr50+X0Q/Lsb43OQHYWRCY2AiWywWQ= +github.com/kr/pretty v0.3.1 h1:flRD4NNwYAUpkphVc1HcthR4KEIFJ65n8Mw5qdRn3LE= +github.com/kr/pretty v0.3.1/go.mod h1:hoEshYVHaxMs3cyo3Yncou5ZscifuDolrwPKZanG3xk= github.com/kr/text v0.2.0 h1:5Nx0Ya0ZqY2ygV366QzturHI13Jq95ApcVaJBhpS+AY= github.com/kr/text v0.2.0/go.mod h1:eLer722TekiGuMkidMxC/pM04lWEeraHUUmBw8l2grE= +github.com/kylelemons/godebug v1.1.0 h1:RPNrshWIDI6G2gRW9EHilWtl7Z6Sb1BR0xunSBf0SNc= +github.com/kylelemons/godebug v1.1.0/go.mod h1:9/0rRGxNHcop5bhtWyNeEfOS8JIWk580+fNqagV/RAw= +github.com/mattn/go-colorable v0.1.14 h1:9A9LHSqF/7dyVVX6g0U9cwm9pG3kP9gSzcuIPHPsaIE= +github.com/mattn/go-colorable v0.1.14/go.mod h1:6LmQG8QLFO4G5z1gPvYEzlUgJ2wF+stgPZH1UqBm1s8= +github.com/mattn/go-isatty v0.0.20 h1:xfD0iDuEKnDkl03q4limB+vH+GxLEtL/jb4xVJSWWEY= +github.com/mattn/go-isatty v0.0.20/go.mod h1:W+V8PltTTMOvKvAeJH7IuucS94S2C6jfK/D7dTCTo3Y= +github.com/munnerz/goautoneg v0.0.0-20191010083416-a7dc8b61c822 h1:C3w9PqII01/Oq1c1nUAm88MOHcQC9l5mIlSMApZMrHA= +github.com/munnerz/goautoneg v0.0.0-20191010083416-a7dc8b61c822/go.mod h1:+n7T8mK8HuQTcFwEeznm/DIxMOiR9yIdICNftLE1DvQ= github.com/pmezard/go-difflib v1.0.0 h1:4DBwDE0NGyQoBHbLQYPwSUPoCMWR5BEzIk/f1lZbAQM= github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4= +github.com/prometheus/client_golang v1.23.2 h1:Je96obch5RDVy3FDMndoUsjAhG5Edi49h0RJWRi/o0o= +github.com/prometheus/client_golang v1.23.2/go.mod h1:Tb1a6LWHB3/SPIzCoaDXI4I8UHKeFTEQ1YCr+0Gyqmg= +github.com/prometheus/client_model v0.6.2 h1:oBsgwpGs7iVziMvrGhE53c/GrLUsZdHnqNwqPLxwZyk= +github.com/prometheus/client_model v0.6.2/go.mod h1:y3m2F6Gdpfy6Ut/GBsUqTWZqCUvMVzSfMLjcu6wAwpE= +github.com/prometheus/common v0.66.1 h1:h5E0h5/Y8niHc5DlaLlWLArTQI7tMrsfQjHV+d9ZoGs= +github.com/prometheus/common v0.66.1/go.mod h1:gcaUsgf3KfRSwHY4dIMXLPV0K/Wg1oZ8+SbZk/HH/dA= +github.com/prometheus/procfs v0.16.1 h1:hZ15bTNuirocR6u0JZ6BAHHmwS1p8B4P6MRqxtzMyRg= +github.com/prometheus/procfs v0.16.1/go.mod h1:teAbpZRB1iIAJYREa1LsoWUXykVXA1KlTmWl8x/U+Is= github.com/rogpeppe/go-internal v1.14.1 h1:UQB4HGPB6osV0SQTLymcB4TgvyWu6ZyliaW0tI/otEQ= github.com/rogpeppe/go-internal v1.14.1/go.mod h1:MaRKkUm5W0goXpeCfT7UZI6fk/L7L7so1lCWt35ZSgc= +github.com/rs/zerolog v1.35.1 h1:m7xQeoiLIiV0BCEY4Hs+j2NG4Gp2o2KPKmhnnLiazKI= +github.com/rs/zerolog v1.35.1/go.mod h1:EjML9kdfa/RMA7h/6z6pYmq1ykOuA8/mjWaEvGI+jcw= github.com/stretchr/objx v0.1.0/go.mod h1:HFkY916IF+rwdDfMAkV7OtwuqBVzrE8GR6GFx+wExME= github.com/stretchr/testify v1.3.0/go.mod h1:M5WIy9Dh21IEIfnGCwXGc5bZfKNJtfHm1UVUgZn+9EI= github.com/stretchr/testify v1.7.0/go.mod h1:6Fq8oRcR53rry900zMqJjRRixrwX3KX962/h/Wwjteg= github.com/stretchr/testify v1.11.1 h1:7s2iGBzp5EwR7/aIZr8ao5+dra3wiQyKjjFuvgVKu7U= github.com/stretchr/testify v1.11.1/go.mod h1:wZwfW3scLgRK+23gO65QZefKpKQRnfz6sD981Nm4B6U= -golang.org/x/net v0.29.0 h1:5ORfpBpCs4HzDYoodCDBbwHzdR5UrLBZ3sOnUJmFoHo= -golang.org/x/net v0.29.0/go.mod h1:gLkgy8jTGERgjzMic6DS9+SP0ajcu6Xu3Orq/SpETg0= +go.opentelemetry.io/auto/sdk v1.1.0 h1:cH53jehLUN6UFLY71z+NDOiNJqDdPRaXzTel0sJySYA= +go.opentelemetry.io/auto/sdk v1.1.0/go.mod h1:3wSPjt5PWp2RhlCcmmOial7AvC4DQqZb7a7wCow3W8A= +go.opentelemetry.io/otel v1.36.0 h1:UumtzIklRBY6cI/lllNZlALOF5nNIzJVb16APdvgTXg= +go.opentelemetry.io/otel v1.36.0/go.mod h1:/TcFMXYjyRNh8khOAO9ybYkqaDBb/70aVwkNML4pP8E= +go.opentelemetry.io/otel/metric v1.36.0 h1:MoWPKVhQvJ+eeXWHFBOPoBOi20jh6Iq2CcCREuTYufE= +go.opentelemetry.io/otel/metric v1.36.0/go.mod h1:zC7Ks+yeyJt4xig9DEw9kuUFe5C3zLbVjV2PzT6qzbs= +go.opentelemetry.io/otel/sdk v1.36.0 h1:b6SYIuLRs88ztox4EyrvRti80uXIFy+Sqzoh9kFULbs= +go.opentelemetry.io/otel/sdk v1.36.0/go.mod h1:+lC+mTgD+MUWfjJubi2vvXWcVxyr9rmlshZni72pXeY= +go.opentelemetry.io/otel/sdk/metric v1.36.0 h1:r0ntwwGosWGaa0CrSt8cuNuTcccMXERFwHX4dThiPis= +go.opentelemetry.io/otel/sdk/metric v1.36.0/go.mod h1:qTNOhFDfKRwX0yXOqJYegL5WRaW376QbB7P4Pb0qva4= +go.opentelemetry.io/otel/trace v1.36.0 h1:ahxWNuqZjpdiFAyrIoQ4GIiAIhxAunQR6MUoKrsNd4w= +go.opentelemetry.io/otel/trace v1.36.0/go.mod h1:gQ+OnDZzrybY4k4seLzPAWNwVBBVlF2szhehOBB/tGA= +go.uber.org/goleak v1.3.0 h1:2K3zAYmnTNqV73imy9J1T3WC+gmCePx2hEGkimedGto= +go.uber.org/goleak v1.3.0/go.mod h1:CoHD4mav9JJNrW/WLlf7HGZPjdw8EucARQHekz1X6bE= +go.yaml.in/yaml/v2 v2.4.2 h1:DzmwEr2rDGHl7lsFgAHxmNz/1NlQ7xLIrlN2h5d1eGI= +go.yaml.in/yaml/v2 v2.4.2/go.mod h1:081UH+NErpNdqlCXm3TtEran0rJZGxAYx9hb/ELlsPU= +golang.org/x/net v0.43.0 h1:lat02VYK2j4aLzMzecihNvTlJNQUq316m2Mr9rnM6YE= +golang.org/x/net v0.43.0/go.mod h1:vhO1fvI4dGsIjh73sWfUVjj3N7CA9WkKJNQm2svM6Jg= golang.org/x/sync v0.17.0 h1:l60nONMj9l5drqw6jlhIELNv9I0A4OFgRsG9k2oT9Ug= golang.org/x/sync v0.17.0/go.mod h1:9KTHXmSnoGruLpwFjVSX0lNNA75CykiMECbovNTZqGI= -golang.org/x/sys v0.26.0 h1:KHjCJyddX0LoSTb3J+vWpupP9p0oznkqVk/IfjymZbo= -golang.org/x/sys v0.26.0/go.mod h1:/VUhepiaJMQUp4+oa/7Zr1D23ma6VTLIYjOOTFZPUcA= +golang.org/x/sys v0.6.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= +golang.org/x/sys v0.35.0 h1:vz1N37gP5bs89s7He8XuIYXpyY0+QlsKmzipCbUtyxI= +golang.org/x/sys v0.35.0/go.mod h1:BJP2sWEmIv4KK5OTEluFJCKSidICx8ciO85XgH3Ak8k= golang.org/x/text v0.29.0 h1:1neNs90w9YzJ9BocxfsQNHKuAT4pkghyXc4nhZ6sJvk= golang.org/x/text v0.29.0/go.mod h1:7MhJOA9CD2qZyOKYazxdYMF85OwPdEr9jTtBpO7ydH4= golang.org/x/time v0.15.0 h1:bbrp8t3bGUeFOx08pvsMYRTCVSMk89u4tKbNOZbp88U= golang.org/x/time v0.15.0/go.mod h1:Y4YMaQmXwGQZoFaVFk4YpCt4FLQMYKZe9oeV/f4MSno= -google.golang.org/genproto/googleapis/rpc v0.0.0-20240903143218-8af14fe29dc1 h1:pPJltXNxVzT4pK9yD8vR9X75DaWYYmLGMsEvBfFQZzQ= -google.golang.org/genproto/googleapis/rpc v0.0.0-20240903143218-8af14fe29dc1/go.mod h1:UqMtugtsSgubUsoxbuAoiCXvqvErP7Gf0so0mK9tHxU= -google.golang.org/grpc v1.68.0 h1:aHQeeJbo8zAkAa3pRzrVjZlbz6uSfeOXlJNQM0RAbz0= -google.golang.org/grpc v1.68.0/go.mod h1:fmSPC5AsjSBCK54MyHRx48kpOti1/jRfOlwEWywNjWA= -google.golang.org/protobuf v1.35.2 h1:8Ar7bF+apOIoThw1EdZl0p1oWvMqTHmpA2fRTyZO8io= -google.golang.org/protobuf v1.35.2/go.mod h1:9fA7Ob0pmnwhb644+1+CVWFRbNajQ6iRojtC/QF5bRE= +google.golang.org/genproto/googleapis/rpc v0.0.0-20250528174236-200df99c418a h1:v2PbRU4K3llS09c7zodFpNePeamkAwG3mPrAery9VeE= +google.golang.org/genproto/googleapis/rpc v0.0.0-20250528174236-200df99c418a/go.mod h1:qQ0YXyHHx3XkvlzUtpXDkS29lDSafHMZBAZDc03LQ3A= +google.golang.org/grpc v1.74.2 h1:WoosgB65DlWVC9FqI82dGsZhWFNBSLjQ84bjROOpMu4= +google.golang.org/grpc v1.74.2/go.mod h1:CtQ+BGjaAIXHs/5YS3i473GqwBBa1zGQNevxdeBEXrM= +google.golang.org/protobuf v1.36.8 h1:xHScyCOEuuwZEc6UtSOvPbAT4zRh0xcNRYekJwfqyMc= +google.golang.org/protobuf v1.36.8/go.mod h1:fuxRtAxBytpl4zzqUh6/eyUujkJdNiuEkXntxiD/uRU= gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0= gopkg.in/check.v1 v1.0.0-20201130134442-10cb98267c6c h1:Hei/4ADfdWqJk1ZMxUNpqntNwaWcugrBjAiHlqqRiVk= gopkg.in/check.v1 v1.0.0-20201130134442-10cb98267c6c/go.mod h1:JHkPIbrfpd72SG/EVd6muEfDQjcINNoR0C8j2r3qZ4Q= diff --git a/internal/config/config.go b/internal/config/config.go index bfc9a8e..ed61923 100644 --- a/internal/config/config.go +++ b/internal/config/config.go @@ -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 == "" { diff --git a/internal/logging/interceptor.go b/internal/logging/interceptor.go new file mode 100644 index 0000000..764ea37 --- /dev/null +++ b/internal/logging/interceptor.go @@ -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 +} diff --git a/internal/logging/logger.go b/internal/logging/logger.go new file mode 100644 index 0000000..b87b8c2 --- /dev/null +++ b/internal/logging/logger.go @@ -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() +} diff --git a/internal/metrics/metrics.go b/internal/metrics/metrics.go new file mode 100644 index 0000000..3c93be6 --- /dev/null +++ b/internal/metrics/metrics.go @@ -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") + } +} diff --git a/internal/provider/musicbrainz/client.go b/internal/provider/musicbrainz/client.go index e928343..54c563f 100644 --- a/internal/provider/musicbrainz/client.go +++ b/internal/provider/musicbrainz/client.go @@ -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)) diff --git a/internal/server/server.go b/internal/server/server.go index e6b230b..3c023eb 100644 --- a/internal/server/server.go +++ b/internal/server/server.go @@ -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) } diff --git a/internal/service/metadata.go b/internal/service/metadata.go index 497270f..47ecbcd 100644 --- a/internal/service/metadata.go +++ b/internal/service/metadata.go @@ -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) } diff --git a/tests/e2e/metadata_test.go b/tests/e2e/metadata_test.go index 8c6db8e..84780f0 100644 --- a/tests/e2e/metadata_test.go +++ b/tests/e2e/metadata_test.go @@ -9,7 +9,9 @@ import ( "google.golang.org/grpc" "google.golang.org/grpc/credentials/insecure" + "github.com/metadata-agregator/internal/provider/musicbrainz" "github.com/metadata-agregator/internal/server" + "github.com/metadata-agregator/internal/service" metadatav1 "github.com/metadata-agregator/pkg/gen/metadata/v1" ) @@ -32,8 +34,17 @@ func startTestServer(t *testing.T) *testServer { t.Fatalf("failed to listen: %v", err) } + mb := musicbrainz.New() + noopArtist := &noopArtistRepo{} + noopAlbum := &noopAlbumRepo{} + noopTrack := &noopTrackRepo{} + svc := service.NewMetadataService(noopArtist, noopAlbum, noopTrack, mb) + services := map[metadatav1.Provider]*service.MetadataService{ + metadatav1.Provider_PROVIDER_MUSICBRAINZ: svc, + } + grpcServer := grpc.NewServer() - metadatav1.RegisterMetadataServiceServer(grpcServer, server.NewMetadataServer()) + metadatav1.RegisterMetadataServiceServer(grpcServer, server.NewMetadataServer(services)) go func() { if err := grpcServer.Serve(lis); err != nil { diff --git a/tests/e2e/noop_repo_test.go b/tests/e2e/noop_repo_test.go new file mode 100644 index 0000000..e2e80b0 --- /dev/null +++ b/tests/e2e/noop_repo_test.go @@ -0,0 +1,66 @@ +package e2e + +import ( + "context" + + "github.com/metadata-agregator/internal/domain" + "github.com/metadata-agregator/internal/repository" +) + +type noopArtistRepo struct{} + +func (r *noopArtistRepo) GetByID(ctx context.Context, id string) (*domain.Artist, error) { + return nil, repository.ErrNotFound +} + +func (r *noopArtistRepo) GetByExternalID(ctx context.Context, source, sourceID string) (*domain.Artist, error) { + return nil, repository.ErrNotFound +} + +func (r *noopArtistRepo) Search(ctx context.Context, query string, limit, offset int) (*domain.SearchResult[domain.Artist], error) { + return &domain.SearchResult[domain.Artist]{}, nil +} + +func (r *noopArtistRepo) Save(ctx context.Context, artist *domain.Artist) error { + return nil +} + +type noopAlbumRepo struct{} + +func (r *noopAlbumRepo) GetByID(ctx context.Context, id string) (*domain.Album, error) { + return nil, repository.ErrNotFound +} + +func (r *noopAlbumRepo) GetByExternalID(ctx context.Context, source, sourceID string) (*domain.Album, error) { + return nil, repository.ErrNotFound +} + +func (r *noopAlbumRepo) GetByArtistID(ctx context.Context, artistID string, limit, offset int) (*domain.SearchResult[domain.Album], error) { + return &domain.SearchResult[domain.Album]{}, nil +} + +func (r *noopAlbumRepo) Save(ctx context.Context, album *domain.Album) error { + return nil +} + +type noopTrackRepo struct{} + +func (r *noopTrackRepo) GetByID(ctx context.Context, id string) (*domain.Track, error) { + return nil, repository.ErrNotFound +} + +func (r *noopTrackRepo) GetByExternalID(ctx context.Context, source, sourceID string) (*domain.Track, error) { + return nil, repository.ErrNotFound +} + +func (r *noopTrackRepo) GetByISRC(ctx context.Context, isrc string) (*domain.Track, error) { + return nil, repository.ErrNotFound +} + +func (r *noopTrackRepo) GetByAlbumID(ctx context.Context, albumID string) ([]domain.Track, error) { + return nil, nil +} + +func (r *noopTrackRepo) Save(ctx context.Context, track *domain.Track) error { + return nil +}