diff --git a/abciapp/app.go b/abciapp/app.go index 172f5f4..e2491b5 100644 --- a/abciapp/app.go +++ b/abciapp/app.go @@ -2,7 +2,6 @@ package abciapp import ( "context" - "fmt" "os" "sync" "time" @@ -10,6 +9,7 @@ import ( dbm "github.com/cometbft/cometbft-db" abcitypes "github.com/cometbft/cometbft/abci/types" "github.com/cometbft/cometbft/crypto" + cmtlog "github.com/cometbft/cometbft/libs/log" "github.com/cometbft/cometbft/privval" bftstore "github.com/cometbft/cometbft/store" "github.com/cosmos/iavl" @@ -26,6 +26,7 @@ import ( type DIDPLCApplication struct { runnerContext context.Context + logger cmtlog.Logger plc plc.PLC txFactory *transaction.Factory indexDB dbm.DB @@ -54,7 +55,7 @@ type DIDPLCApplication struct { } // store and plc must be able to share transaction objects -func NewDIDPLCApplication(appContext context.Context, pv *privval.FilePV, treeDB dbm.DB, indexDB transaction.ExtendedDB, clearData func(), snapshotDirectory, didBloomFilterPath string, mempoolSubmitter types.MempoolSubmitter) (*DIDPLCApplication, *transaction.Factory, plc.PLC, func(), error) { +func NewDIDPLCApplication(appContext context.Context, logger cmtlog.Logger, pv *privval.FilePV, treeDB dbm.DB, indexDB transaction.ExtendedDB, clearData func(), snapshotDirectory, didBloomFilterPath string, mempoolSubmitter types.MempoolSubmitter) (*DIDPLCApplication, *transaction.Factory, plc.PLC, func(), error) { mkTree := func() *iavl.MutableTree { // Using SpeedDefault appears to cause the processing time for ExecuteOperation to double on average // Using SpeedBetterCompression appears to cause the processing time to double again @@ -80,6 +81,7 @@ func NewDIDPLCApplication(appContext context.Context, pv *privval.FilePV, treeDB d := &DIDPLCApplication{ runnerContext: runnerContext, + logger: logger.With("module", "plcapp"), tree: tree, indexDB: indexDB, mempoolSubmitter: mempoolSubmitter, @@ -92,7 +94,7 @@ func NewDIDPLCApplication(appContext context.Context, pv *privval.FilePV, treeDB d.validatorPrivKey = pv.Key.PrivKey } - d.txFactory, err = transaction.NewFactory(tree, indexDB, store.Consensus.CountOperations, store.NewDIDBloomFilterStore(didBloomFilterPath)) + d.txFactory, err = transaction.NewFactory(tree, indexDB, store.Consensus.CountOperations, store.NewDIDBloomFilterStore(d.logger, didBloomFilterPath)) if err != nil { return nil, nil, nil, cancelRunnerContext, stacktrace.Propagate(err, "") } @@ -109,7 +111,7 @@ func NewDIDPLCApplication(appContext context.Context, pv *privval.FilePV, treeDB *d.tree = *mkTree() - d.txFactory, err = transaction.NewFactory(tree, indexDB, store.Consensus.CountOperations, store.NewDIDBloomFilterStore(didBloomFilterPath)) + d.txFactory, err = transaction.NewFactory(tree, indexDB, store.Consensus.CountOperations, store.NewDIDBloomFilterStore(d.logger, didBloomFilterPath)) if err != nil { return stacktrace.Propagate(err, "") } @@ -131,9 +133,9 @@ func NewDIDPLCApplication(appContext context.Context, pv *privval.FilePV, treeDB st := time.Now() err := d.txFactory.SaveDIDBloomFilter() if err != nil { - fmt.Println("FAILED TO SAVE BLOOM FILTER:", stacktrace.Propagate(err, "")) + d.logger.Error("failed to save bloom filter", "error", stacktrace.Propagate(err, "")) } - fmt.Println("SAVED BLOOM FILTER IN", time.Since(st)) + d.logger.Debug("saved bloom filter", "took", time.Since(st)) } }) @@ -195,11 +197,23 @@ func NewDIDPLCApplication(appContext context.Context, pv *privval.FilePV, treeDB }, nil } +func (d *DIDPLCApplication) logMethod(method string, keyvals ...any) func(...any) { + st := time.Now() + d.logger.Debug(method+" start", keyvals...) + return func(extra ...any) { + args := make([]any, 0, len(keyvals)+len(extra)+2) + args = append(args, keyvals...) + args = append(args, extra...) + args = append(args, "took", time.Since(st)) + d.logger.Debug(method+" done", args...) + } +} + func (d *DIDPLCApplication) FinishInitializing(blockStore *bftstore.BlockStore) error { d.blockStore = blockStore var err error - d.blockChallengeCoordinator, err = newBlockChallengeCoordinator(d.runnerContext, d.txFactory, blockStore, d.validatorPubKey) + d.blockChallengeCoordinator, err = newBlockChallengeCoordinator(d.runnerContext, d.logger, d.txFactory, blockStore, d.validatorPubKey) if err != nil { return stacktrace.Propagate(err, "") } diff --git a/abciapp/app_test.go b/abciapp/app_test.go index ce0f586..8c8ced8 100644 --- a/abciapp/app_test.go +++ b/abciapp/app_test.go @@ -6,6 +6,7 @@ import ( dbm "github.com/cometbft/cometbft-db" "github.com/cometbft/cometbft/abci/types" + cmtlog "github.com/cometbft/cometbft/libs/log" "github.com/dgraph-io/badger/v4" cbornode "github.com/ipfs/go-ipld-cbor" "github.com/stretchr/testify/require" @@ -22,7 +23,8 @@ func txJSONToCBOR(t *testing.T, jsonBytes []byte) []byte { } func TestCheckTx(t *testing.T) { - app, _, _, cleanup, err := abciapp.NewDIDPLCApplication(t.Context(), nil, dbm.NewMemDB(), memDBWrapper{dbm.NewMemDB()}, nil, "", "", nil) + logger := cmtlog.NewNopLogger() + app, _, _, cleanup, err := abciapp.NewDIDPLCApplication(t.Context(), logger, nil, dbm.NewMemDB(), memDBWrapper{dbm.NewMemDB()}, nil, "", "", nil) require.NoError(t, err) t.Cleanup(cleanup) diff --git a/abciapp/block_challenge.go b/abciapp/block_challenge.go index 3b6a3ff..6047266 100644 --- a/abciapp/block_challenge.go +++ b/abciapp/block_challenge.go @@ -4,22 +4,21 @@ import ( "bytes" "context" "embed" - "fmt" "math/big" "time" "github.com/Yiling-J/theine-go" "github.com/cometbft/cometbft/crypto" + cmtlog "github.com/cometbft/cometbft/libs/log" bftstore "github.com/cometbft/cometbft/store" "github.com/consensys/gnark-crypto/ecc" "github.com/consensys/gnark-crypto/ecc/bn254" "github.com/consensys/gnark-crypto/ecc/bn254/fr/mimc" - "github.com/consensys/gnark/backend" "github.com/consensys/gnark/backend/groth16" "github.com/consensys/gnark/backend/witness" "github.com/consensys/gnark/constraint" - "github.com/consensys/gnark/constraint/solver" "github.com/consensys/gnark/frontend" + gnarklogger "github.com/consensys/gnark/logger" "github.com/palantir/stacktrace" "github.com/rs/zerolog" "github.com/samber/lo" @@ -51,12 +50,15 @@ func init() { vkFile := lo.Must(blockChallengeCircuitFS.Open("proofcircuit/BlockChallenge_VerifyingKey")) defer vkFile.Close() lo.Must(blockChallengeVerifyingKey.ReadFrom(vkFile)) + + gnarklogger.Set(zerolog.Nop()) } type blockChallengeCoordinator struct { g singleflight.Group[int64, []byte] runnerContext context.Context + logger cmtlog.Logger isConfiguredToBeValidator bool validatorAddress []byte @@ -66,9 +68,10 @@ type blockChallengeCoordinator struct { sharedWitnessDataCache *theine.LoadingCache[int64, proof.BlockChallengeCircuit] } -func newBlockChallengeCoordinator(runnerContext context.Context, txFactory *transaction.Factory, blockStore *bftstore.BlockStore, pubKey crypto.PubKey) (*blockChallengeCoordinator, error) { +func newBlockChallengeCoordinator(runnerContext context.Context, logger cmtlog.Logger, txFactory *transaction.Factory, blockStore *bftstore.BlockStore, pubKey crypto.PubKey) (*blockChallengeCoordinator, error) { c := &blockChallengeCoordinator{ runnerContext: runnerContext, + logger: logger, txFactory: txFactory, nodeBlockStore: blockStore, isConfiguredToBeValidator: pubKey != nil, @@ -126,7 +129,7 @@ func (c *blockChallengeCoordinator) notifyOfIncomingBlockHeight(height int64) { go func() { _, err := c.loadOrComputeBlockChallengeProof(c.runnerContext, height) if err != nil { - fmt.Printf("FAILED TO COMPUTE CHALLENGE FOR BLOCK %d: %v\n", height, stacktrace.Propagate(err, "")) + c.logger.Error("failed to compute block challenge", "height", height, "error", stacktrace.Propagate(err, "")) } }() } @@ -148,6 +151,7 @@ func (c *blockChallengeCoordinator) loadOrComputeBlockChallengeProof(ctx context return nil, stacktrace.Propagate(err, "") } if proof == nil { + st := time.Now() // compute and store proof, err = c.computeBlockChallengeProof(tx, height) if err != nil { @@ -169,6 +173,8 @@ func (c *blockChallengeCoordinator) loadOrComputeBlockChallengeProof(ctx context if err != nil { return nil, stacktrace.Propagate(err, "") } + + c.logger.Debug("computed and stored block challenge", "height", height, "took", time.Since(st)) } return proof, nil }) @@ -180,10 +186,7 @@ func (c *blockChallengeCoordinator) computeBlockChallengeProof(tx transaction.Re if err != nil { return nil, stacktrace.Propagate(err, "") } - - // TODO consider using a different logger once we clean up our logging act - // TODO open an issue in the gnark repo because backend.WithSolverOptions(solver.WithLogger(zerolog.Nop())) has no effect... - proof, err := groth16.Prove(blockChallengeConstraintSystem, blockChallengeProvingKey, witness, backend.WithSolverOptions(solver.WithLogger(zerolog.Nop()))) + proof, err := groth16.Prove(blockChallengeConstraintSystem, blockChallengeProvingKey, witness) if err != nil { return nil, stacktrace.Propagate(err, "") } diff --git a/abciapp/execution.go b/abciapp/execution.go index 4fa0f58..d6d4529 100644 --- a/abciapp/execution.go +++ b/abciapp/execution.go @@ -3,7 +3,7 @@ package abciapp import ( "bytes" "context" - "fmt" + "encoding/hex" "slices" "time" @@ -25,6 +25,7 @@ func (d *DIDPLCApplication) InitChain(_ context.Context, req *abcitypes.RequestI // PrepareProposal implements [types.Application]. func (d *DIDPLCApplication) PrepareProposal(ctx context.Context, req *abcitypes.RequestPrepareProposal) (*abcitypes.ResponsePrepareProposal, error) { + defer (d.logMethod("PrepareProposal", "height", req.Height, "txs", len(req.Txs)))() defer d.DiscardChanges() if req.Height == 2 { @@ -105,6 +106,8 @@ func (d *DIDPLCApplication) PrepareProposal(ctx context.Context, req *abcitypes. // ProcessProposal implements [types.Application]. func (d *DIDPLCApplication) ProcessProposal(ctx context.Context, req *abcitypes.RequestProcessProposal) (*abcitypes.ResponseProcessProposal, error) { + defer (d.logMethod("ProcessProposal", "height", req.Height, "hash", req.Hash, "txs", len(req.Txs)))() + // always reset state before processing a new proposal d.DiscardChanges() // do not unconditionally defer DiscardChanges because we want to re-use the results in FinalizeBlock when we vote accept @@ -140,12 +143,10 @@ func (d *DIDPLCApplication) ProcessProposal(ctx context.Context, req *abcitypes. return &abcitypes.ResponseProcessProposal{Status: abcitypes.ResponseProcessProposal_REJECT}, nil } - st := time.Now() result, err = finishProcessTx(ctx, d.transactionProcessorDependenciesForOngoingProcessing(true, req.Time), processor, tx) if err != nil { return nil, stacktrace.Propagate(err, "") } - fmt.Println("FINISHPROCESSTX TOOK", time.Since(st)) } // when preparing a proposal, invalid transactions should have been discarded @@ -165,6 +166,8 @@ func (d *DIDPLCApplication) ProcessProposal(ctx context.Context, req *abcitypes. // ExtendVote implements [types.Application]. func (d *DIDPLCApplication) ExtendVote(ctx context.Context, req *abcitypes.RequestExtendVote) (*abcitypes.ResponseExtendVote, error) { + defer (d.logMethod("ExtendVote", "height", req.Height, "hash", req.Hash))() + proof, err := d.blockChallengeCoordinator.loadOrComputeBlockChallengeProof(ctx, req.Height) if err != nil { return nil, stacktrace.Propagate(err, "") @@ -176,6 +179,8 @@ func (d *DIDPLCApplication) ExtendVote(ctx context.Context, req *abcitypes.Reque // VerifyVoteExtension implements [types.Application]. func (d *DIDPLCApplication) VerifyVoteExtension(_ context.Context, req *abcitypes.RequestVerifyVoteExtension) (*abcitypes.ResponseVerifyVoteExtension, error) { + defer (d.logMethod("VerifyVoteExtension", "height", req.Height, "hash", req.Hash, "validator", hex.EncodeToString(req.ValidatorAddress)))() + if len(req.VoteExtension) > 200 { // that definitely ain't right return &abcitypes.ResponseVerifyVoteExtension{ @@ -195,6 +200,8 @@ func (d *DIDPLCApplication) VerifyVoteExtension(_ context.Context, req *abcitype // FinalizeBlock implements [types.Application]. func (d *DIDPLCApplication) FinalizeBlock(ctx context.Context, req *abcitypes.RequestFinalizeBlock) (*abcitypes.ResponseFinalizeBlock, error) { + defer (d.logMethod("FinalizeBlock", "height", req.Height, "hash", req.Hash))() + if bytes.Equal(req.Hash, d.lastProcessedProposalHash) && d.lastProcessedProposalExecTxResults != nil { // the block that was decided was the one we processed in ProcessProposal, and ProcessProposal processed successfully // reuse the uncommitted results @@ -231,6 +238,8 @@ func (d *DIDPLCApplication) FinalizeBlock(ctx context.Context, req *abcitypes.Re // Commit implements [types.Application]. func (d *DIDPLCApplication) Commit(context.Context, *abcitypes.RequestCommit) (*abcitypes.ResponseCommit, error) { + defer (d.logMethod("Commit"))() + // ensure we always advance tree version by creating ongoingWrite if it hasn't been created already d.createOngoingTxIfNeeded(time.Now()) diff --git a/abciapp/range_challenge.go b/abciapp/range_challenge.go index 645532f..cd7d0ff 100644 --- a/abciapp/range_challenge.go +++ b/abciapp/range_challenge.go @@ -3,8 +3,8 @@ package abciapp import ( "context" "encoding/binary" + "encoding/hex" "errors" - "fmt" "math/big" "slices" "sync" @@ -12,7 +12,7 @@ import ( "github.com/Yiling-J/theine-go" "github.com/cometbft/cometbft/crypto" - "github.com/cometbft/cometbft/mempool" + cmtlog "github.com/cometbft/cometbft/libs/log" "github.com/cometbft/cometbft/privval" "github.com/cometbft/cometbft/rpc/core" bftstore "github.com/cometbft/cometbft/store" @@ -30,6 +30,7 @@ import ( type RangeChallengeCoordinator struct { runnerContext context.Context + logger cmtlog.Logger isConfiguredToBeValidator bool validatorPubKey crypto.PubKey @@ -58,15 +59,17 @@ type consensusReactor interface { func NewRangeChallengeCoordinator( runnerContext context.Context, + logger cmtlog.Logger, + pv *privval.FilePV, txFactory *transaction.Factory, blockStore *bftstore.BlockStore, nodeEventBus *cmttypes.EventBus, mempoolSubmitter types.MempoolSubmitter, - consensusReactor consensusReactor, - pv *privval.FilePV) (*RangeChallengeCoordinator, error) { + consensusReactor consensusReactor) (*RangeChallengeCoordinator, error) { c := &RangeChallengeCoordinator{ txFactory: txFactory, runnerContext: runnerContext, + logger: logger, nodeBlockStore: blockStore, nodeEventBus: nodeEventBus, mempoolSubmitter: mempoolSubmitter, @@ -98,7 +101,7 @@ func (c *RangeChallengeCoordinator) Start() error { c.wg.Go(func() { err := c.newBlocksSubscriber() if err != nil { - fmt.Println("newBlocksSubscriber FAILED:", err) + c.logger.Error("blocks subscriber failed", "error", stacktrace.Propagate(err, "")) } }) c.wg.Go(func() { @@ -114,7 +117,7 @@ func (c *RangeChallengeCoordinator) Start() error { if err != nil { // note: this is expected in certain circumstances, such as the proof for the toHeight block not being ready yet as the block was just finalized // (and the block may have been finalized without our votes) - fmt.Println("onNewBlock FAILED:", err) + c.logger.Error("range challenge block handler error", "error", stacktrace.Propagate(err, "")) } }() } @@ -277,15 +280,15 @@ func (c *RangeChallengeCoordinator) onNewBlock(ctx context.Context, newBlockHeig } } - fmt.Println("RANGE CHALLENGE EVAL", shouldCommitToChallenge, shouldCompleteChallenge) - - var transactionBytes []byte + var transactionBytes cmttypes.Tx if shouldCompleteChallenge { + c.logger.Info("Creating challenge completion transaction", "fromHeight", fromHeight, "toHeight", toHeight, "provenHeight", provenHeight, "includedOnHeight", includedOnHeight) transactionBytes, err = c.createCompleteChallengeTx(ctx, tx, int64(fromHeight), int64(toHeight), int64(provenHeight), int64(includedOnHeight)) if err != nil { return stacktrace.Propagate(err, "") } } else if shouldCommitToChallenge { + c.logger.Info("Creating challenge commitment transaction", "toHeight", toHeight) transactionBytes, err = c.createCommitToChallengeTx(ctx, tx, newBlockHeight) if err != nil { if errors.Is(err, errMissingProofs) { @@ -301,17 +304,16 @@ func (c *RangeChallengeCoordinator) onNewBlock(ctx context.Context, newBlockHeig return nil } + txHashHex := hex.EncodeToString(transactionBytes.Hash()) + c.logger.Debug("broadcasting range challenge transaction", "hash", txHashHex) result, err := c.mempoolSubmitter.BroadcastTx(ctx, transactionBytes, true) if err != nil { - if errors.Is(err, mempool.ErrTxInCache) { - // expected, as we don't wait for broadcast and therefore will try to repeatedly commit/complete - return nil - } return stacktrace.Propagate(err, "") } if result.CheckTx.Code == 0 && shouldCompleteChallenge { c.hasSubmittedChallengeCompletion = true } + c.logger.Debug("range challenge transaction included", "hash", txHashHex, "txResult", result.TxResult.Code) c.cachedNextProofFromHeight = mo.None[int64]() return nil } diff --git a/abciapp/snapshots.go b/abciapp/snapshots.go index e8bdcc2..5e1644f 100644 --- a/abciapp/snapshots.go +++ b/abciapp/snapshots.go @@ -15,7 +15,6 @@ import ( "strconv" "strings" "sync" - "time" dbm "github.com/cometbft/cometbft-db" abcitypes "github.com/cometbft/cometbft/abci/types" @@ -230,6 +229,8 @@ func (d *DIDPLCApplication) OfferSnapshot(_ context.Context, req *abcitypes.Requ } func (d *DIDPLCApplication) createSnapshot(treeVersion int64, tempFilename string) error { + defer (d.logMethod("createSnapshot", "treeVersion", treeVersion, "tempFilename", tempFilename))() + it, err := d.tree.GetImmutable(treeVersion) if err != nil { return stacktrace.Propagate(err, "") @@ -244,8 +245,6 @@ func (d *DIDPLCApplication) createSnapshot(treeVersion int64, tempFilename strin } defer f.Close() - st := time.Now() - err = writeSnapshot(f, d.indexDB, it) if err != nil { return stacktrace.Propagate(err, "") @@ -279,8 +278,6 @@ func (d *DIDPLCApplication) createSnapshot(treeVersion int64, tempFilename strin os.Rename(tempFilename, filepath.Join(d.snapshotDirectory, fmt.Sprintf("%020d.snapshot", treeVersion))) - fmt.Println("Took", time.Since(st), "to export") - return nil } diff --git a/main.go b/main.go index b8bda4f..941b809 100644 --- a/main.go +++ b/main.go @@ -97,7 +97,22 @@ func main() { appContext, cancelAppContext := context.WithCancel(context.Background()) defer cancelAppContext() - app, txFactory, plc, cleanup, err := abciapp.NewDIDPLCApplication(appContext, pv, treeDB, indexDB, recreateDatabases, filepath.Join(homeDir, "snapshots"), didBloomFilterPath, mempoolSubmitter) + logger := cmtlog.NewTMLogger(cmtlog.NewSyncWriter(os.Stdout)) + logger, err = cmtflags.ParseLogLevel(config.LogLevel, logger, bftconfig.DefaultLogLevel) + if err != nil { + log.Fatalf("failed to parse log level: %v", err) + } + + app, txFactory, plc, cleanup, err := abciapp.NewDIDPLCApplication( + appContext, + logger, + pv, + treeDB, + indexDB, + recreateDatabases, + filepath.Join(homeDir, "snapshots"), + didBloomFilterPath, + mempoolSubmitter) if err != nil { log.Fatalf("failed to create DIDPLC application: %v", err) } @@ -108,13 +123,6 @@ func main() { log.Fatalf("failed to load node's key: %v", err) } - logger := cmtlog.NewTMLogger(cmtlog.NewSyncWriter(os.Stdout)) - logger, err = cmtflags.ParseLogLevel(config.LogLevel, logger, bftconfig.DefaultLogLevel) - - if err != nil { - log.Fatalf("failed to parse log level: %v", err) - } - node, err := nm.NewNode( config.Config, pv, @@ -137,7 +145,15 @@ func main() { log.Fatalf("Finishing ABCI app initialization: %v", err) } - rangeChallengeCoordinator, err := abciapp.NewRangeChallengeCoordinator(appContext, txFactory, node.BlockStore(), node.EventBus(), mempoolSubmitter, node.ConsensusReactor(), pv) + rangeChallengeCoordinator, err := abciapp.NewRangeChallengeCoordinator( + appContext, + logger.With("module", "plcapp"), + pv, + txFactory, + node.BlockStore(), + node.EventBus(), + mempoolSubmitter, + node.ConsensusReactor()) if err != nil { log.Fatalf("Creating RangeChallengeCoordinator: %v", err) } diff --git a/startfresh.sh b/startfresh.sh index b23c13d..359f1e0 100755 --- a/startfresh.sh +++ b/startfresh.sh @@ -3,4 +3,5 @@ rm -r didplcbft-data/ go build -trimpath go run github.com/cometbft/cometbft/cmd/cometbft@v0.38.19 init --home didplcbft-data sed -i 's/^create_empty_blocks = true$/create_empty_blocks = false/g' didplcbft-data/config/config.toml +sed -i 's/^log_level = "info"$/log_level = "plcapp:debug,*:info"/g' didplcbft-data/config/config.toml ./didplcbft \ No newline at end of file diff --git a/store/did_bloom.go b/store/did_bloom.go index 44f5eca..251f21c 100644 --- a/store/did_bloom.go +++ b/store/did_bloom.go @@ -3,27 +3,31 @@ package store import ( "encoding/binary" "errors" - "fmt" "io" "math" "os" "slices" "github.com/bits-and-blooms/bloom/v3" + cmtlog "github.com/cometbft/cometbft/libs/log" "github.com/palantir/stacktrace" "tangled.org/gbl08ma.com/didplcbft/transaction" ) type DIDBloomFilterStore struct { + logger cmtlog.Logger filePath string } -func NewInMemoryDIDBloomFilterStore() *DIDBloomFilterStore { - return &DIDBloomFilterStore{} +func NewInMemoryDIDBloomFilterStore(logger cmtlog.Logger) *DIDBloomFilterStore { + return &DIDBloomFilterStore{ + logger: logger, + } } -func NewDIDBloomFilterStore(filePath string) *DIDBloomFilterStore { +func NewDIDBloomFilterStore(logger cmtlog.Logger, filePath string) *DIDBloomFilterStore { return &DIDBloomFilterStore{ + logger: logger, filePath: filePath, } } @@ -93,13 +97,13 @@ func (s *DIDBloomFilterStore) BuildDIDBloomFilter(tx transaction.Read) (*bloom.B return filter, nil } - fmt.Println("(RE)BUILDING DID BLOOM FILTER") - filterEstimatedItems := uint(100000000) // we know there are like 80M DIDs at the time of writing if estimatedDIDCount != 0 { filterEstimatedItems = max(filterEstimatedItems, uint(estimatedDIDCount*3)) } + s.logger.Info("Rebuilding DID bloom filter", "itemCapacity", filterEstimatedItems) + filter = bloom.NewWithEstimates(filterEstimatedItems, 0.01) didRangeStart := marshalDIDLogKey(make([]byte, 15), 0) @@ -112,16 +116,20 @@ func (s *DIDBloomFilterStore) BuildDIDBloomFilter(tx transaction.Read) (*bloom.B defer iterator.Close() + itemCount := 0 for iterator.Valid() { filter.Add(iterator.Key()[1:16]) iterator.Next() + itemCount++ } err = iterator.Error() if err != nil { return nil, stacktrace.Propagate(err, "") } + s.logger.Debug("rebuilt DID bloom filter", "itemCapacity", filterEstimatedItems, "itemCount", itemCount) + return filter, nil } diff --git a/testutil/testutil.go b/testutil/testutil.go index 2b75a25..07aa58d 100644 --- a/testutil/testutil.go +++ b/testutil/testutil.go @@ -7,6 +7,7 @@ import ( "github.com/klauspost/compress/zstd" "github.com/stretchr/testify/require" + cmtlog "github.com/cometbft/cometbft/libs/log" "tangled.org/gbl08ma.com/didplcbft/badgertodbm" "tangled.org/gbl08ma.com/didplcbft/dbmtoiavldb" "tangled.org/gbl08ma.com/didplcbft/dbmtoiavldb/zstddict" @@ -23,7 +24,7 @@ func NewTestTxFactory(t *testing.T) (*transaction.Factory, *iavl.MutableTree, tr _, indexDB, err := badgertodbm.NewBadgerInMemoryDB() require.NoError(t, err) - factory, err := transaction.NewFactory(tree, indexDB, store.Consensus.CountOperations, store.NewInMemoryDIDBloomFilterStore()) + factory, err := transaction.NewFactory(tree, indexDB, store.Consensus.CountOperations, store.NewInMemoryDIDBloomFilterStore(cmtlog.NewNopLogger())) require.NoError(t, err) return factory, tree, indexDB