Repository navigation
fix: add txes to redis with ttl #86
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Changes from all commits
92fba59
d8ab3cb
85a588b
2963aa2
6a86487
70a72e6
f7f33ae
dffcd5d
a078f7a
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,88 @@ | ||
| package collector | ||
|
|
||
| import ( | ||
| "context" | ||
| "fmt" | ||
| "sync" | ||
| "time" | ||
|
|
||
| "github.com/redis/go-redis/v9" | ||
| "go.uber.org/zap" | ||
| ) | ||
|
|
||
| const ( | ||
| redisKeyPrefix = "mempool-dumpster:" | ||
| redisTTL = 5 * time.Minute | ||
| redisPingTimeout = 30 * time.Second | ||
| redisAddTxTimeout = 10 * time.Second | ||
| redisQueueSize = 4096 | ||
| redisNumWorkers = 4 | ||
| ) | ||
|
|
||
| type Redis struct { | ||
| log *zap.SugaredLogger | ||
| client *redis.Client | ||
| queue chan string | ||
| wg sync.WaitGroup | ||
| } | ||
|
|
||
| func NewRedis(log *zap.SugaredLogger, endpoint string) (*Redis, error) { | ||
| opts, err := redis.ParseURL(endpoint) | ||
| if err != nil { | ||
| return nil, fmt.Errorf("failed to parse redis endpoint: %w", err) | ||
| } | ||
|
|
||
| client := redis.NewClient(opts) | ||
|
|
||
| ctx, cancel := context.WithTimeout(context.Background(), redisPingTimeout) | ||
| defer cancel() | ||
|
|
||
| if err := client.Ping(ctx).Err(); err != nil { | ||
| client.Close() | ||
| return nil, fmt.Errorf("failed to ping redis: %w", err) | ||
| } | ||
|
|
||
| rd := &Redis{ | ||
| log: log, | ||
| client: client, | ||
| queue: make(chan string, redisQueueSize), | ||
| } | ||
|
|
||
| rd.wg.Add(redisNumWorkers) | ||
| for range redisNumWorkers { | ||
| go rd.worker() | ||
| } | ||
|
|
||
| return rd, nil | ||
| } | ||
|
|
||
| // AddTx sends a tx hash to the background queue for async Redis write. | ||
| // Drops the hash if the queue is full to avoid blocking the caller. | ||
| func (r *Redis) AddTx(hash string) { | ||
| select { | ||
| case r.queue <- hash: | ||
| default: | ||
| r.log.Warnw("redis queue full, dropping tx", "tx", hash) | ||
| } | ||
|
Comment on lines
+61
to
+66
|
||
| } | ||
|
|
||
| func (r *Redis) worker() { | ||
| defer r.wg.Done() | ||
| for hash := range r.queue { | ||
| r.processHash(hash) | ||
| } | ||
| } | ||
|
|
||
| func (r *Redis) processHash(hash string) { | ||
| ctx, cancel := context.WithTimeout(context.Background(), redisAddTxTimeout) | ||
| defer cancel() | ||
| if err := r.client.Set(ctx, redisKeyPrefix+hash, "1", redisTTL).Err(); err != nil { | ||
| r.log.Errorw("failed to add tx to redis", "error", err, "tx", hash) | ||
| } | ||
| } | ||
|
|
||
| func (r *Redis) Close() error { | ||
| close(r.queue) | ||
| r.wg.Wait() | ||
| return r.client.Close() | ||
| } | ||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,60 @@ | ||
| package collector | ||
|
|
||
| import ( | ||
| "testing" | ||
| "time" | ||
|
|
||
| "github.com/alicebob/miniredis/v2" | ||
| "github.com/stretchr/testify/require" | ||
| "go.uber.org/zap" | ||
| ) | ||
|
|
||
| func TestRedis_AddTx(t *testing.T) { | ||
| mr := miniredis.RunT(t) | ||
| log := zap.NewNop().Sugar() | ||
|
|
||
| r, err := NewRedis(log, "redis://"+mr.Addr()) | ||
| require.NoError(t, err) | ||
|
|
||
| hash := "0xabc123" | ||
| r.AddTx(hash) | ||
|
|
||
| // Close flushes the queue and waits for workers to finish | ||
| require.NoError(t, r.Close()) | ||
|
|
||
| // key exists with correct prefix | ||
| require.True(t, mr.Exists(redisKeyPrefix+hash)) | ||
|
|
||
| // value is "1" | ||
| val, err := mr.Get(redisKeyPrefix + hash) | ||
| require.NoError(t, err) | ||
| require.Equal(t, "1", val) | ||
|
|
||
| // TTL is set | ||
| ttl := mr.TTL(redisKeyPrefix + hash) | ||
| require.Equal(t, redisTTL, ttl) | ||
| } | ||
|
|
||
| func TestRedis_TTLExpiry(t *testing.T) { | ||
| mr := miniredis.RunT(t) | ||
| log := zap.NewNop().Sugar() | ||
|
|
||
| r, err := NewRedis(log, "redis://"+mr.Addr()) | ||
| require.NoError(t, err) | ||
|
|
||
| hash := "0xdef456" | ||
| r.AddTx(hash) | ||
| require.NoError(t, r.Close()) | ||
|
|
||
| require.True(t, mr.Exists(redisKeyPrefix+hash)) | ||
|
|
||
| // fast-forward past TTL | ||
| mr.FastForward(redisTTL + time.Second) | ||
| require.False(t, mr.Exists(redisKeyPrefix+hash)) | ||
| } | ||
|
|
||
| func TestRedis_BadEndpoint(t *testing.T) { | ||
| log := zap.NewNop().Sugar() | ||
| _, err := NewRedis(log, "redis://localhost:1") | ||
| require.Error(t, err) | ||
| } |
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -40,6 +40,7 @@ type TxProcessorOpts struct { | |
| Location string // location of the collector, will be stored in sourcelogs | ||
| CheckNodeURI string | ||
| ClickhouseDSN string | ||
| RedisEndpoint string | ||
| HTTPReceivers []string | ||
| ReceiversAllowedSources []string | ||
| APIServer *api.Server | ||
|
|
@@ -51,8 +52,9 @@ type TxProcessor struct { | |
| uid string | ||
| location string | ||
|
|
||
| outDir string | ||
| txC chan common.TxIn // note: it's important that the value is sent in here instead of a pointer, otherwise there are memory race conditions | ||
| outDir string | ||
| txC chan common.TxIn // note: it's important that the value is sent in here instead of a pointer, otherwise there are memory race conditions | ||
| txCDone chan struct{} // closed when startTransactionReceiverLoop exits | ||
|
|
||
| outFilesLock sync.RWMutex | ||
| outFiles map[int64]OutFiles | ||
|
|
@@ -74,6 +76,9 @@ type TxProcessor struct { | |
|
|
||
| clickhouseDSN string | ||
| clickhouse *Clickhouse | ||
|
|
||
| redisEndpoint string | ||
| redis *Redis | ||
| } | ||
|
|
||
| type OutFiles struct { | ||
|
|
@@ -95,8 +100,9 @@ func NewTxProcessor(opts TxProcessorOpts) *TxProcessor { | |
| } | ||
|
|
||
| return &TxProcessor{ //nolint:exhaustruct | ||
| log: opts.Log, | ||
| txC: make(chan common.TxIn, 100), | ||
| log: opts.Log, | ||
| txC: make(chan common.TxIn, 100), | ||
| txCDone: make(chan struct{}), | ||
|
|
||
| uid: opts.UID, | ||
| location: opts.Location, | ||
|
|
@@ -109,6 +115,7 @@ func NewTxProcessor(opts TxProcessorOpts) *TxProcessor { | |
|
|
||
| checkNodeURI: opts.CheckNodeURI, | ||
| clickhouseDSN: opts.ClickhouseDSN, | ||
| redisEndpoint: opts.RedisEndpoint, | ||
|
|
||
| receivers: receivers, | ||
| receiversAllowedSources: opts.ReceiversAllowedSources, | ||
|
|
@@ -118,6 +125,12 @@ func NewTxProcessor(opts TxProcessorOpts) *TxProcessor { | |
|
|
||
| func (p *TxProcessor) Shutdown() { | ||
| p.log.Info("Shutting down TxProcessor ...") | ||
| p.stopTransactionReceiverLoop() | ||
| if p.redis != nil { | ||
| if err := p.redis.Close(); err != nil { | ||
| p.log.Errorw("failed to close Redis", "error", err) | ||
| } | ||
| } | ||
| if p.clickhouse != nil { | ||
| p.clickhouse.FlushCurrentBatches() | ||
| } | ||
|
|
@@ -139,8 +152,17 @@ func (p *TxProcessor) Start() { | |
| p.log.Info("Connected to Clickhouse!") | ||
| } | ||
|
|
||
| if p.redisEndpoint != "" { | ||
| p.log.Info("Connecting to Redis...") | ||
| p.redis, err = NewRedis(p.log, p.redisEndpoint) | ||
| if err != nil { | ||
| p.log.Fatalw("failed to connect to Redis", "error", err) | ||
| } | ||
| p.log.Info("Connected to Redis!") | ||
| } | ||
|
shanejonas marked this conversation as resolved.
|
||
|
|
||
| if p.checkNodeURI != "" { | ||
| p.log.Infof("Conecting to check-node at %s ...", p.checkNodeURI) | ||
| p.log.Infof("Connecting to check-node at %s ...", p.checkNodeURI) | ||
| p.ethClient, err = ethclient.Dial(p.checkNodeURI) | ||
| if err != nil { | ||
| p.log.Fatal(err) | ||
|
|
@@ -165,7 +187,13 @@ func (p *TxProcessor) Start() { | |
| p.log.Info("TxProcessor started successfully") | ||
| } | ||
|
|
||
| func (p *TxProcessor) stopTransactionReceiverLoop() { | ||
| close(p.txC) | ||
| <-p.txCDone | ||
| } | ||
|
Comment on lines
+190
to
+193
|
||
|
|
||
| func (p *TxProcessor) startTransactionReceiverLoop() { | ||
| defer close(p.txCDone) | ||
| p.log.Info("Waiting for transactions...") | ||
| for txIn := range p.txC { | ||
| if txIn.Tx == nil { | ||
|
|
@@ -275,6 +303,12 @@ func (p *TxProcessor) processTx(txIn common.TxIn) { | |
| } | ||
| } | ||
|
|
||
| // Add tx hash to Redis asynchronously via background workers. | ||
| // Protect will check this to filter txs coming back through. | ||
| if p.redis != nil { | ||
| p.redis.AddTx(txHashLower) | ||
| } | ||
|
|
||
| // Add transaction to Clickhouse | ||
| if p.clickhouse != nil { | ||
| err = p.clickhouse.AddTransaction(txIn) // send to Clickhouse | ||
|
|
||
Uh oh!
There was an error while loading. Please reload this page.