Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 4 additions & 0 deletions main.go
Original file line number Diff line number Diff line change
Expand Up @@ -68,10 +68,14 @@ func main() {
solanaIndexer := solana_indexer.New(config.Cfg)
defer solanaIndexer.Close()

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)
defer stop()

go healthServer.Start(ctx)

if err := solanaIndexer.Start(ctx); err != nil {
if !errors.Is(err, context.Canceled) {
panic(err)
Expand Down
80 changes: 80 additions & 0 deletions solana/indexer/server.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,80 @@
package indexer

import (
"context"
"encoding/json"
"net"
"net/http"

"github.com/gofiber/fiber/v2"
"github.com/mcuadros/go-defaults"
"go.uber.org/zap"
)

type Server struct {
*fiber.App
logger *zap.Logger
}

func NewServer(indexer *SolanaIndexer) *Server {
logger := indexer.logger.Named("Server")
app := fiber.New(fiber.Config{
JSONEncoder: json.Marshal,
JSONDecoder: json.Unmarshal,
UnescapePath: true,
})

app.Get("/solana/health", func(c *fiber.Ctx) error {
Comment thread
rickyrombo marked this conversation as resolved.
type queryParams struct {
MaxSlotDiff int64 `query:"max_slot_diff" default:"100"`
MaxRetryQueue int `query:"max_retry_queue" default:"10"`
}
var q queryParams
err := c.QueryParser(&q)
if err != nil {
return c.Status(fiber.StatusBadRequest).JSON(fiber.Map{
"error": err.Error(),
})
}

defaults.SetDefaults(&q)

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(),
})
}
if len(health.Errors) > 0 {
c.Status(fiber.StatusInternalServerError)
}
return c.JSON(fiber.Map{
"data": health,
})
})
return &Server{
App: app,
logger: logger,
}
}

func (s *Server) Start(ctx context.Context) {
go func() {
<-ctx.Done()
s.logger.Info("received shutdown signal, stopping server")
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
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))
}
}
184 changes: 108 additions & 76 deletions solana/indexer/solana_indexer.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
}
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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 {
Expand Down Expand Up @@ -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++
}
Expand All @@ -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()
}
Loading