From 235d44fabf5a2cb35b6220b0a8cf2c819ba03c93 Mon Sep 17 00:00:00 2001 From: Marcus Pasell <3690498+rickyrombo@users.noreply.github.com> Date: Tue, 4 Nov 2025 23:05:01 -0800 Subject: [PATCH 1/2] Add server w/ own health check to solana indexer --- main.go | 4 + solana/indexer/server.go | 168 +++++++++++++++++++++++++++++++++++++++ 2 files changed, 172 insertions(+) create mode 100644 solana/indexer/server.go diff --git a/main.go b/main.go index b8560362..624c8d31 100644 --- a/main.go +++ b/main.go @@ -68,10 +68,14 @@ func main() { solanaIndexer := solana_indexer.New(config.Cfg) defer solanaIndexer.Close() + healthServer := solana_indexer.NewServer(config.Cfg) + // Capture termination signals for graceful shutdown of the indexer ctx, stop := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM, syscall.SIGINT, os.Interrupt) defer stop() + go healthServer.Start(ctx) + if err := solanaIndexer.Start(ctx); err != nil { if !errors.Is(err, context.Canceled) { panic(err) diff --git a/solana/indexer/server.go b/solana/indexer/server.go new file mode 100644 index 00000000..ef4f35d3 --- /dev/null +++ b/solana/indexer/server.go @@ -0,0 +1,168 @@ +package indexer + +import ( + "context" + "encoding/json" + "fmt" + "net" + "net/http" + "time" + + "api.audius.co/api/dbv1" + "api.audius.co/config" + "api.audius.co/logging" + "github.com/gagliardetto/solana-go/rpc" + "github.com/gofiber/fiber/v2" + "github.com/jackc/pgx/v5" + "go.uber.org/zap" +) + +type Server struct { + app *fiber.App + pool *dbv1.DBPools + logger *zap.Logger +} + +const MAX_SLOT_DIFF = 100 +const MAX_RETRY_QUEUE = 10 + +func NewServer(config config.Config) *Server { + logger := logging.NewZapLogger(config).Named("Server") + + // Create DBPools from read replicas + var connectionStrings []string + if len(config.ReadDbReplicas) > 0 && config.ReadDbReplicas[0] != "" { + // Use read replicas if configured + connectionStrings = config.ReadDbReplicas + } else { + // Fall back to single read database + connectionStrings = []string{config.ReadDbUrl} + } + + pool, err := dbv1.NewDBPools(connectionStrings, logger, config.Env, config.ZapLevel) + if err != nil { + logger.Fatal("read db connect failed", zap.Error(err)) + } + + solanaRpc := rpc.New(config.SolanaConfig.RpcProviders[0]) + + app := fiber.New(fiber.Config{ + JSONEncoder: json.Marshal, + JSONDecoder: json.Unmarshal, + UnescapePath: true, + }) + + app.Get("/solana/health", func(c *fiber.Ctx) error { + + chainSlot, err := solanaRpc.GetSlot(c.Context(), rpc.CommitmentConfirmed) + if err != nil { + return fmt.Errorf("failed to get chain slot: %w", err) + } + + sql := ` + WITH retry_queue_by_indexer AS ( + SELECT + indexer, + COUNT(*) AS retry_queue_count + FROM sol_retry_queue + GROUP BY indexer + ) SELECT DISTINCT ON (name) + name, + to_slot AS indexed_slot, + @chain_slot - to_slot AS slot_diff, + COALESCE(retry_queue_count, 0) AS retry_queue_count, + updated_at + FROM sol_slot_checkpoints + LEFT JOIN retry_queue_by_indexer ON indexer = name + ORDER BY name, from_slot DESC + ; + ` + + rows, err := pool.Query(c.Context(), sql, pgx.NamedArgs{ + "chain_slot": chainSlot, + }) + if err != nil { + return fmt.Errorf("failed to query indexer health: %w", err) + } + + type indexerHealthRow struct { + Name string `json:"name"` + SlotDiff uint64 `json:"slot_diff"` + IndexedSlot uint64 `json:"indexed_slot"` + RetryQueueCount int `json:"retry_queue_count"` + UpdatedAt *time.Time `db:"updated_at"` + } + healths, err := pgx.CollectRows(rows, pgx.RowToStructByName[indexerHealthRow]) + if err != nil { + return fmt.Errorf("failed to collect indexer health rows: %w", err) + } + + type solanaHealth struct { + ChainSlot uint64 `json:"chain_slot"` + Errors []string `json:"errors,omitempty"` + Indexers []indexerHealthRow + } + health := solanaHealth{ + ChainSlot: chainSlot, + Errors: make([]string, 0), + Indexers: healths, + } + + for _, h := range health.Indexers { + if h.RetryQueueCount > MAX_RETRY_QUEUE { + c.Status(fiber.StatusInternalServerError) + health.Errors = append(health.Errors, fmt.Sprintf("indexer %s has high retry queue count: %d", h.Name, h.RetryQueueCount)) + } + + if h.SlotDiff > MAX_SLOT_DIFF { + c.Status(fiber.StatusInternalServerError) + health.Errors = append(health.Errors, fmt.Sprintf("indexer %s has high slot diff: %d", h.Name, h.SlotDiff)) + } + } + + return c.JSON(fiber.Map{ + "data": health, + }) + }) + + return &Server{ + app: app, + pool: pool, + logger: logger, + } +} + +func (s *Server) Start(ctx context.Context) { + flushTicker := time.NewTicker(time.Second * 15) + defer flushTicker.Stop() + + go func() { + for range flushTicker.C { + s.logger.Sync() + } + }() + + go func() { + <-ctx.Done() + s.logger.Info("received shutdown signal, stopping server") + s.Shutdown(context.Background()) + }() + + // Bind to both ipv4 and ipv6 + listener, err := net.Listen("tcp", "[::]:1324") + if err != nil { + s.logger.Fatal("Failed to create listener", zap.Error(err)) + } + + if err := s.app.Listener(listener); err != nil && err != http.ErrServerClosed { + s.logger.Fatal("Failed to start server", zap.Error(err)) + } +} + +func (s *Server) Shutdown(ctx context.Context) { + if err := s.app.Shutdown(); err != nil { + s.logger.Error("failed to shutdown app", zap.Error(err)) + } + s.pool.Close() + s.logger.Sync() +} From b33b293f2206decc8aa3813a586c755186a91a7b Mon Sep 17 00:00:00 2001 From: Marcus Pasell <3690498+rickyrombo@users.noreply.github.com> Date: Wed, 5 Nov 2025 15:16:08 -0800 Subject: [PATCH 2/2] better health check --- main.go | 2 +- solana/indexer/server.go | 140 +++++------------------ solana/indexer/solana_indexer.go | 184 ++++++++++++++++++------------- 3 files changed, 135 insertions(+), 191 deletions(-) diff --git a/main.go b/main.go index 624c8d31..7c35268e 100644 --- a/main.go +++ b/main.go @@ -68,7 +68,7 @@ func main() { solanaIndexer := solana_indexer.New(config.Cfg) defer solanaIndexer.Close() - healthServer := solana_indexer.NewServer(config.Cfg) + healthServer := solana_indexer.NewServer(solanaIndexer) // Capture termination signals for graceful shutdown of the indexer ctx, stop := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM, syscall.SIGINT, os.Interrupt) diff --git a/solana/indexer/server.go b/solana/indexer/server.go index ef4f35d3..78588905 100644 --- a/solana/indexer/server.go +++ b/solana/indexer/server.go @@ -3,49 +3,21 @@ package indexer import ( "context" "encoding/json" - "fmt" "net" "net/http" - "time" - "api.audius.co/api/dbv1" - "api.audius.co/config" - "api.audius.co/logging" - "github.com/gagliardetto/solana-go/rpc" "github.com/gofiber/fiber/v2" - "github.com/jackc/pgx/v5" + "github.com/mcuadros/go-defaults" "go.uber.org/zap" ) type Server struct { - app *fiber.App - pool *dbv1.DBPools + *fiber.App logger *zap.Logger } -const MAX_SLOT_DIFF = 100 -const MAX_RETRY_QUEUE = 10 - -func NewServer(config config.Config) *Server { - logger := logging.NewZapLogger(config).Named("Server") - - // Create DBPools from read replicas - var connectionStrings []string - if len(config.ReadDbReplicas) > 0 && config.ReadDbReplicas[0] != "" { - // Use read replicas if configured - connectionStrings = config.ReadDbReplicas - } else { - // Fall back to single read database - connectionStrings = []string{config.ReadDbUrl} - } - - pool, err := dbv1.NewDBPools(connectionStrings, logger, config.Env, config.ZapLevel) - if err != nil { - logger.Fatal("read db connect failed", zap.Error(err)) - } - - solanaRpc := rpc.New(config.SolanaConfig.RpcProviders[0]) - +func NewServer(indexer *SolanaIndexer) *Server { + logger := indexer.logger.Named("Server") app := fiber.New(fiber.Config{ JSONEncoder: json.Marshal, JSONDecoder: json.Unmarshal, @@ -53,99 +25,47 @@ func NewServer(config config.Config) *Server { }) app.Get("/solana/health", func(c *fiber.Ctx) error { - - chainSlot, err := solanaRpc.GetSlot(c.Context(), rpc.CommitmentConfirmed) - if err != nil { - return fmt.Errorf("failed to get chain slot: %w", err) + type queryParams struct { + MaxSlotDiff int64 `query:"max_slot_diff" default:"100"` + MaxRetryQueue int `query:"max_retry_queue" default:"10"` } - - sql := ` - WITH retry_queue_by_indexer AS ( - SELECT - indexer, - COUNT(*) AS retry_queue_count - FROM sol_retry_queue - GROUP BY indexer - ) SELECT DISTINCT ON (name) - name, - to_slot AS indexed_slot, - @chain_slot - to_slot AS slot_diff, - COALESCE(retry_queue_count, 0) AS retry_queue_count, - updated_at - FROM sol_slot_checkpoints - LEFT JOIN retry_queue_by_indexer ON indexer = name - ORDER BY name, from_slot DESC - ; - ` - - rows, err := pool.Query(c.Context(), sql, pgx.NamedArgs{ - "chain_slot": chainSlot, - }) + var q queryParams + err := c.QueryParser(&q) if err != nil { - return fmt.Errorf("failed to query indexer health: %w", err) + return c.Status(fiber.StatusBadRequest).JSON(fiber.Map{ + "error": err.Error(), + }) } - type indexerHealthRow struct { - Name string `json:"name"` - SlotDiff uint64 `json:"slot_diff"` - IndexedSlot uint64 `json:"indexed_slot"` - RetryQueueCount int `json:"retry_queue_count"` - UpdatedAt *time.Time `db:"updated_at"` - } - healths, err := pgx.CollectRows(rows, pgx.RowToStructByName[indexerHealthRow]) - if err != nil { - return fmt.Errorf("failed to collect indexer health rows: %w", err) - } + defaults.SetDefaults(&q) - type solanaHealth struct { - ChainSlot uint64 `json:"chain_slot"` - Errors []string `json:"errors,omitempty"` - Indexers []indexerHealthRow - } - health := solanaHealth{ - ChainSlot: chainSlot, - Errors: make([]string, 0), - Indexers: healths, + health, err := indexer.GetHealth(c.Context(), uint64(q.MaxSlotDiff), q.MaxRetryQueue) + if err != nil { + return c.Status(fiber.StatusInternalServerError).JSON(fiber.Map{ + "error": err.Error(), + }) } - - for _, h := range health.Indexers { - if h.RetryQueueCount > MAX_RETRY_QUEUE { - c.Status(fiber.StatusInternalServerError) - health.Errors = append(health.Errors, fmt.Sprintf("indexer %s has high retry queue count: %d", h.Name, h.RetryQueueCount)) - } - - if h.SlotDiff > MAX_SLOT_DIFF { - c.Status(fiber.StatusInternalServerError) - health.Errors = append(health.Errors, fmt.Sprintf("indexer %s has high slot diff: %d", h.Name, h.SlotDiff)) - } + if len(health.Errors) > 0 { + c.Status(fiber.StatusInternalServerError) } - return c.JSON(fiber.Map{ "data": health, }) }) - return &Server{ - app: app, - pool: pool, + App: app, logger: logger, } } func (s *Server) Start(ctx context.Context) { - flushTicker := time.NewTicker(time.Second * 15) - defer flushTicker.Stop() - - go func() { - for range flushTicker.C { - s.logger.Sync() - } - }() - go func() { <-ctx.Done() s.logger.Info("received shutdown signal, stopping server") - s.Shutdown(context.Background()) + if err := s.App.Shutdown(); err != nil { + s.logger.Error("failed to shutdown app", zap.Error(err)) + } + s.logger.Sync() }() // Bind to both ipv4 and ipv6 @@ -154,15 +74,7 @@ func (s *Server) Start(ctx context.Context) { s.logger.Fatal("Failed to create listener", zap.Error(err)) } - if err := s.app.Listener(listener); err != nil && err != http.ErrServerClosed { + if err := s.App.Listener(listener); err != nil && err != http.ErrServerClosed { s.logger.Fatal("Failed to start server", zap.Error(err)) } } - -func (s *Server) Shutdown(ctx context.Context) { - if err := s.app.Shutdown(); err != nil { - s.logger.Error("failed to shutdown app", zap.Error(err)) - } - s.pool.Close() - s.logger.Sync() -} diff --git a/solana/indexer/solana_indexer.go b/solana/indexer/solana_indexer.go index bb8ae0b7..7f09b157 100644 --- a/solana/indexer/solana_indexer.go +++ b/solana/indexer/solana_indexer.go @@ -17,26 +17,26 @@ import ( "api.audius.co/solana/indexer/token" "github.com/gagliardetto/solana-go" "github.com/gagliardetto/solana-go/rpc" + "github.com/jackc/pgx/v5" "github.com/jackc/pgx/v5/pgxpool" "github.com/maypok86/otter" + pb "github.com/rpcpool/yellowstone-grpc/examples/golang/proto" "go.uber.org/zap" ) +type Indexer interface { + Start(ctx context.Context) + HandleUpdate(ctx context.Context, updateMessage *pb.SubscribeUpdate) error +} + type SolanaIndexer struct { - rpcClient common.RpcClient - grpcClient common.GrpcClient + rpcClient common.RpcClient config config.Config pool database.DbPool workerCount int32 - dammV2Indexer *damm_v2.Indexer - tokenIndexer *token.Indexer - programIndexer *program.Indexer - dbcIndexer *dbc.Indexer - lockerIndexer *locker.Indexer - - checkpointId string + indexers map[string]Indexer logger *zap.Logger } @@ -93,18 +93,20 @@ func New(config config.Config) *SolanaIndexer { grpcConfig, rpcClient, pool, logger, ) + indexers := make(map[string]Indexer) + indexers[damm_v2.NAME] = dammV2Indexer + indexers[token.NAME] = tokenIndexer + indexers[program.NAME] = programIndexer + indexers[dbc.NAME] = dbcIndexer + indexers[locker.NAME] = lockerIndexer + s := &SolanaIndexer{ rpcClient: rpcClient, logger: logger, config: config, pool: pool, workerCount: workerCount, - - dammV2Indexer: dammV2Indexer, - tokenIndexer: tokenIndexer, - programIndexer: programIndexer, - dbcIndexer: dbcIndexer, - lockerIndexer: lockerIndexer, + indexers: indexers, } return s @@ -133,11 +135,9 @@ func (s *SolanaIndexer) Start(ctx context.Context) error { balanceHistoryJob.ScheduleEvery(balanceHistoryCtx, 1*time.Hour) go balanceHistoryJob.Run(balanceHistoryCtx) - go s.tokenIndexer.Start(ctx) - go s.dammV2Indexer.Start(ctx) - go s.programIndexer.Start(ctx) - go s.dbcIndexer.Start(ctx) - go s.lockerIndexer.Start(ctx) + for _, indexer := range s.indexers { + go indexer.Start(ctx) + } for { select { @@ -183,64 +183,21 @@ func (s *SolanaIndexer) ProcessRetryQueue(ctx context.Context) { } for _, item := range queue { - switch item.Indexer { - case token.NAME: - err := s.tokenIndexer.HandleUpdate(ctx, item.UpdateMessage.SubscribeUpdate) - if err != nil { - logger.Error("failed to retry", zap.String("indexer", token.NAME), zap.Error(err)) - offset++ - } else { - err = common.DeleteFromRetryQueue(ctx, s.pool, item.ID) - if err != nil { - logger.Error("failed to delete from retry queue", zap.Error(err)) - } - } - case damm_v2.NAME: - err := s.dammV2Indexer.HandleUpdate(ctx, item.UpdateMessage.SubscribeUpdate) - if err != nil { - logger.Error("failed to retry", zap.String("indexer", damm_v2.NAME), zap.Error(err)) - offset++ - } else { - err = common.DeleteFromRetryQueue(ctx, s.pool, item.ID) - if err != nil { - logger.Error("failed to delete from retry queue", zap.Error(err)) - } - } - case dbc.NAME: - err := s.dbcIndexer.HandleUpdate(ctx, item.UpdateMessage.SubscribeUpdate) - if err != nil { - logger.Error("failed to retry", zap.String("indexer", dbc.NAME), zap.Error(err)) - offset++ - } else { - err = common.DeleteFromRetryQueue(ctx, s.pool, item.ID) - if err != nil { - logger.Error("failed to delete from retry queue", zap.Error(err)) - } - } - case program.NAME: - err := s.programIndexer.HandleUpdate(ctx, item.UpdateMessage.SubscribeUpdate) - if err != nil { - logger.Error("failed to retry", zap.String("indexer", program.NAME), zap.Error(err)) - offset++ - } else { - err = common.DeleteFromRetryQueue(ctx, s.pool, item.ID) - if err != nil { - logger.Error("failed to delete from retry queue", zap.Error(err)) - } - } - case locker.NAME: - err := s.lockerIndexer.HandleUpdate(ctx, item.UpdateMessage.SubscribeUpdate) + indexer := s.indexers[item.Indexer] + if indexer == nil { + logger.Warn("unknown indexer in retry queue", zap.String("indexer", item.Indexer)) + offset++ + continue + } + err := indexer.HandleUpdate(ctx, item.UpdateMessage.SubscribeUpdate) + if err != nil { + logger.Error("failed to retry", zap.String("indexer", locker.NAME), zap.Error(err)) + offset++ + } else { + err = common.DeleteFromRetryQueue(ctx, s.pool, item.ID) if err != nil { - logger.Error("failed to retry", zap.String("indexer", locker.NAME), zap.Error(err)) - offset++ - } else { - err = common.DeleteFromRetryQueue(ctx, s.pool, item.ID) - if err != nil { - logger.Error("failed to delete from retry queue", zap.Error(err)) - } + logger.Error("failed to delete from retry queue", zap.Error(err)) } - default: - logger.Warn("unknown indexer in retry queue", zap.String("indexer", item.Indexer)) } count++ } @@ -258,6 +215,81 @@ func (s *SolanaIndexer) ProcessRetryQueue(ctx context.Context) { ) } +type solanaHealth struct { + ChainSlot uint64 `json:"chain_slot"` + Errors []string `json:"errors,omitempty"` + Indexers []indexerHealthRow `json:"indexers"` +} + +type indexerHealthRow struct { + Name string `json:"name"` + SlotDiff uint64 `json:"slot_diff"` + IndexedSlot uint64 `json:"indexed_slot"` + RetryQueueCount int `json:"retry_queue_count"` + UpdatedAt *time.Time `json:"updated_at"` +} + +func (s *SolanaIndexer) GetHealth(ctx context.Context, maxSlotDiff uint64, maxRetryQueue int) (*solanaHealth, error) { + chainSlot, err := s.rpcClient.GetSlot(ctx, rpc.CommitmentConfirmed) + if err != nil { + return nil, fmt.Errorf("failed to get chain slot: %w", err) + } + + names := make([]string, 0, len(s.indexers)) + for name := range s.indexers { + names = append(names, name) + } + + sql := ` + WITH retry_queue_by_indexer AS ( + SELECT + indexer, + COUNT(*) AS retry_queue_count + FROM sol_retry_queue + GROUP BY indexer + ) SELECT DISTINCT ON (indexers.name) + indexers.name, + to_slot AS indexed_slot, + GREATEST(@chain_slot - to_slot, 0) AS slot_diff, + COALESCE(retry_queue_count, 0) AS retry_queue_count, + updated_at + FROM UNNEST(@indexers::TEXT[]) AS indexers(name) + LEFT JOIN sol_slot_checkpoints ON sol_slot_checkpoints.name = indexers.name + LEFT JOIN retry_queue_by_indexer ON indexer = indexers.name + ORDER BY indexers.name, from_slot DESC NULLS LAST + ;` + + rows, err := s.pool.Query(ctx, sql, pgx.NamedArgs{ + "chain_slot": chainSlot, + "indexers": names, + }) + if err != nil { + return nil, fmt.Errorf("failed to query indexer health: %w", err) + } + + healths, err := pgx.CollectRows(rows, pgx.RowToStructByName[indexerHealthRow]) + if err != nil { + return nil, fmt.Errorf("failed to collect indexer health rows: %w", err) + } + + errors := make([]string, 0) + for _, h := range healths { + if h.RetryQueueCount > maxRetryQueue { + errors = append(errors, fmt.Sprintf("indexer %s has high retry queue count: %d", h.Name, h.RetryQueueCount)) + } + + if h.SlotDiff > maxSlotDiff { + errors = append(errors, fmt.Sprintf("indexer %s has high slot diff: %d", h.Name, h.SlotDiff)) + } + } + + return &solanaHealth{ + ChainSlot: chainSlot, + Errors: errors, + Indexers: healths, + }, nil +} + func (s *SolanaIndexer) Close() { s.pool.Close() }