diff --git a/params/version.go b/params/version.go index 67b04ff4ae..06272984cf 100644 --- a/params/version.go +++ b/params/version.go @@ -24,7 +24,7 @@ import ( const ( VersionMajor = 5 // Major version component of the current release VersionMinor = 4 // Minor version component of the current release - VersionPatch = 0 // Patch version component of the current release + VersionPatch = 1 // Patch version component of the current release VersionMeta = "mainnet" // Version metadata to append to the version string ) diff --git a/rollup/circuitcapacitychecker/impl.go b/rollup/circuitcapacitychecker/impl.go index 93fdf9a345..fe10a02b06 100644 --- a/rollup/circuitcapacitychecker/impl.go +++ b/rollup/circuitcapacitychecker/impl.go @@ -73,19 +73,7 @@ func (ccc *CircuitCapacityChecker) ApplyTransaction(traces *types.BlockTrace) (* } encodeStart := time.Now() - ccc.jsonBuffer.Reset() - err := json.NewEncoder(&ccc.jsonBuffer).Encode(traces) - if err != nil { - log.Error("fail to json marshal traces in ApplyTransaction", "id", ccc.ID, "TxHash", traces.Transactions[0].TxHash, "err", err) - return nil, ErrUnknown - } - - tracesStr := C.CString(string(ccc.jsonBuffer.Bytes())) - defer func() { - C.free(unsafe.Pointer(tracesStr)) - }() - - rustTrace := C.parse_json_to_rust_trace(tracesStr) + rustTrace := MakeRustTrace(traces, &ccc.jsonBuffer) if rustTrace == nil { log.Error("fail to parse json in to rust trace", "id", ccc.ID, "TxHash", traces.Transactions[0].TxHash) return nil, ErrUnknown @@ -93,26 +81,34 @@ func (ccc *CircuitCapacityChecker) ApplyTransaction(traces *types.BlockTrace) (* encodeTimer.UpdateSince(encodeStart) log.Debug("start to check circuit capacity for tx", "id", ccc.ID, "TxHash", traces.Transactions[0].TxHash) + return ccc.applyTransactionRustTrace(rustTrace) +} + +func (ccc *CircuitCapacityChecker) ApplyTransactionRustTrace(rustTrace unsafe.Pointer) (*types.RowConsumption, error) { + ccc.Lock() + defer ccc.Unlock() + return ccc.applyTransactionRustTrace(rustTrace) +} + +func (ccc *CircuitCapacityChecker) applyTransactionRustTrace(rustTrace unsafe.Pointer) (*types.RowConsumption, error) { rawResult := C.apply_tx(C.uint64_t(ccc.ID), rustTrace) defer func() { C.free_c_chars(rawResult) }() - log.Debug("check circuit capacity for tx done", "id", ccc.ID, "TxHash", traces.Transactions[0].TxHash) result := &WrappedRowUsage{} - if err = json.Unmarshal([]byte(C.GoString(rawResult)), result); err != nil { - log.Error("fail to json unmarshal apply_tx result", "id", ccc.ID, "TxHash", traces.Transactions[0].TxHash, "err", err) + if err := json.Unmarshal([]byte(C.GoString(rawResult)), result); err != nil { + log.Error("fail to json unmarshal apply_tx result", "id", ccc.ID, "err", err) return nil, ErrUnknown } if result.Error != "" { - log.Error("fail to apply_tx in CircuitCapacityChecker", "id", ccc.ID, "TxHash", traces.Transactions[0].TxHash, "err", result.Error) + log.Error("fail to apply_tx in CircuitCapacityChecker", "id", ccc.ID, "err", result.Error) return nil, ErrUnknown } if result.AccRowUsage == nil { log.Error("fail to apply_tx in CircuitCapacityChecker", - "id", ccc.ID, "TxHash", traces.Transactions[0].TxHash, - "result.AccRowUsage == nil", result.AccRowUsage == nil, + "id", ccc.ID, "result.AccRowUsage == nil", result.AccRowUsage == nil, "err", "AccRowUsage is empty unexpectedly") return nil, ErrUnknown } @@ -128,19 +124,7 @@ func (ccc *CircuitCapacityChecker) ApplyBlock(traces *types.BlockTrace) (*types. defer ccc.Unlock() encodeStart := time.Now() - ccc.jsonBuffer.Reset() - err := json.NewEncoder(&ccc.jsonBuffer).Encode(traces) - if err != nil { - log.Error("fail to json marshal traces in ApplyBlock", "id", ccc.ID, "blockNumber", traces.Header.Number, "blockHash", traces.Header.Hash(), "err", err) - return nil, ErrUnknown - } - - tracesStr := C.CString(string(ccc.jsonBuffer.Bytes())) - defer func() { - C.free(unsafe.Pointer(tracesStr)) - }() - - rustTrace := C.parse_json_to_rust_trace(tracesStr) + rustTrace := MakeRustTrace(traces, &ccc.jsonBuffer) if rustTrace == nil { log.Error("fail to parse json in to rust trace", "id", ccc.ID, "TxHash", traces.Transactions[0].TxHash) return nil, ErrUnknown @@ -155,7 +139,7 @@ func (ccc *CircuitCapacityChecker) ApplyBlock(traces *types.BlockTrace) (*types. log.Debug("check circuit capacity for block done", "id", ccc.ID, "blockNumber", traces.Header.Number, "blockHash", traces.Header.Hash()) result := &WrappedRowUsage{} - if err = json.Unmarshal([]byte(C.GoString(rawResult)), result); err != nil { + if err := json.Unmarshal([]byte(C.GoString(rawResult)), result); err != nil { log.Error("fail to json unmarshal apply_block result", "id", ccc.ID, "blockNumber", traces.Header.Number, "blockHash", traces.Header.Hash(), "err", err) return nil, ErrUnknown } @@ -219,3 +203,23 @@ func (ccc *CircuitCapacityChecker) SetLightMode(lightMode bool) error { return nil } + +func MakeRustTrace(trace *types.BlockTrace, buffer *bytes.Buffer) unsafe.Pointer { + if buffer == nil { + buffer = new(bytes.Buffer) + } + buffer.Reset() + + err := json.NewEncoder(buffer).Encode(trace) + if err != nil { + log.Error("fail to json marshal traces in MakeRustTrace", "err", err) + return nil + } + + tracesStr := C.CString(string(buffer.Bytes())) + defer func() { + C.free(unsafe.Pointer(tracesStr)) + }() + + return C.parse_json_to_rust_trace(tracesStr) +} diff --git a/rollup/circuitcapacitychecker/mock.go b/rollup/circuitcapacitychecker/mock.go index c7b96002f8..4f89464124 100644 --- a/rollup/circuitcapacitychecker/mock.go +++ b/rollup/circuitcapacitychecker/mock.go @@ -3,7 +3,9 @@ package circuitcapacitychecker import ( + "bytes" "math/rand" + "unsafe" "github.com/scroll-tech/go-ethereum/common" "github.com/scroll-tech/go-ethereum/core/types" @@ -51,6 +53,10 @@ func (ccc *CircuitCapacityChecker) ApplyTransaction(traces *types.BlockTrace) (* }}, nil } +func (ccc *CircuitCapacityChecker) ApplyTransactionRustTrace(rustTrace unsafe.Pointer) (*types.RowConsumption, error) { + return ccc.ApplyTransaction(goTraces[rustTrace]) +} + // ApplyBlock gets a block's RowConsumption. // Will only return a dummy value in mock_ccc. func (ccc *CircuitCapacityChecker) ApplyBlock(traces *types.BlockTrace) (*types.RowConsumption, error) { @@ -82,3 +88,11 @@ func (ccc *CircuitCapacityChecker) Skip(txnHash common.Hash, err error) { ccc.skipHash = txnHash.String() ccc.skipError = err } + +var goTraces = make(map[unsafe.Pointer]*types.BlockTrace) + +func MakeRustTrace(trace *types.BlockTrace, buffer *bytes.Buffer) unsafe.Pointer { + rustTrace := new(struct{}) + goTraces[unsafe.Pointer(rustTrace)] = trace + return unsafe.Pointer(rustTrace) +} diff --git a/rollup/pipeline/pipeline.go b/rollup/pipeline/pipeline.go index c3c7b996fe..6aeac5b5a1 100644 --- a/rollup/pipeline/pipeline.go +++ b/rollup/pipeline/pipeline.go @@ -1,8 +1,12 @@ package pipeline import ( + "bytes" + "context" "errors" + "sync" "time" + "unsafe" "github.com/scroll-tech/go-ethereum/common" "github.com/scroll-tech/go-ethereum/core" @@ -30,7 +34,7 @@ func (e *ErrorWithTrace) Unwrap() error { } var ( - ErrApplyStageDone = errors.New("apply stage is done") + ErrPipelineDone = errors.New("pipeline is done") ErrUnexpectedL1MessageIndex = errors.New("unexpected L1 message index") lifetimeTimer = func() metrics.Timer { @@ -38,18 +42,24 @@ var ( metrics.DefaultRegistry.Register("pipeline/lifetime", t) return t }() - applyTimer = metrics.NewRegisteredTimer("pipeline/apply", nil) - applyIdleTimer = metrics.NewRegisteredTimer("pipeline/apply_idle", nil) - applyStallTimer = metrics.NewRegisteredTimer("pipeline/apply_stall", nil) - cccTimer = metrics.NewRegisteredTimer("pipeline/ccc", nil) - cccIdleTimer = metrics.NewRegisteredTimer("pipeline/ccc_idle", nil) + applyTimer = metrics.NewRegisteredTimer("pipeline/apply", nil) + applyIdleTimer = metrics.NewRegisteredTimer("pipeline/apply_idle", nil) + applyStallTimer = metrics.NewRegisteredTimer("pipeline/apply_stall", nil) + encodeTimer = metrics.NewRegisteredTimer("pipeline/encode", nil) + encodeIdleTimer = metrics.NewRegisteredTimer("pipeline/encode_idle", nil) + encodeStallTimer = metrics.NewRegisteredTimer("pipeline/encode_stall", nil) + cccTimer = metrics.NewRegisteredTimer("pipeline/ccc", nil) + cccIdleTimer = metrics.NewRegisteredTimer("pipeline/ccc_idle", nil) ) type Pipeline struct { - chain *core.BlockChain - vmConfig vm.Config - parent *types.Block - start time.Time + chain *core.BlockChain + vmConfig vm.Config + parent *types.Block + start time.Time + wg sync.WaitGroup + ctx context.Context + cancelCtx context.CancelFunc // accumulators ccc *circuitcapacitychecker.CircuitCapacityChecker @@ -83,6 +93,8 @@ func NewPipeline( // make sure we are not sharing a tracer with the caller and not in debug mode vmConfig.Tracer = nil vmConfig.Debug = false + + ctx, cancel := context.WithCancel(context.Background()) return &Pipeline{ chain: chain, vmConfig: vmConfig, @@ -92,6 +104,8 @@ func NewPipeline( ccc: ccc, state: state, gasPool: new(core.GasPool).AddGas(header.GasLimit), + ctx: ctx, + cancelCtx: cancel, } } @@ -103,22 +117,20 @@ func (p *Pipeline) WithBeforeTxHook(beforeTxHook func()) *Pipeline { func (p *Pipeline) Start(deadline time.Time) error { p.start = time.Now() p.txnQueue = make(chan *types.Transaction) - applyStageRespCh, candidateCh, err := p.traceAndApplyStage(p.txnQueue) + applyStageRespCh, applyToEncodeCh, err := p.traceAndApplyStage(p.txnQueue) if err != nil { log.Error("Failed starting traceAndApplyStage", "err", err) return err } p.applyStageRespCh = applyStageRespCh - p.ResultCh = p.cccStage(candidateCh, deadline) + encodeToCccCh := p.encodeStage(applyToEncodeCh) + p.ResultCh = p.cccStage(encodeToCccCh, deadline) return nil } // Stop forces pipeline to stop its operation and return whatever progress it has so far func (p *Pipeline) Stop() { - if p.txnQueue != nil { - close(p.txnQueue) - p.txnQueue = nil - } + close(p.txnQueue) } func (p *Pipeline) TryPushTxns(txs types.OrderedTransactionSet, onFailingTxn func(txnIndex int, tx *types.Transaction, err error) bool) *Result { @@ -137,7 +149,7 @@ func (p *Pipeline) TryPushTxns(txs types.OrderedTransactionSet, onFailingTxn fun case err == nil, errors.Is(err, core.ErrNonceTooLow): txs.Shift() default: - if errors.Is(err, ErrApplyStageDone) || onFailingTxn(p.txs.Len(), tx, err) { + if errors.Is(err, ErrPipelineDone) || onFailingTxn(p.txs.Len(), tx, err) { p.Stop() return nil } @@ -154,12 +166,11 @@ func (p *Pipeline) TryPushTxns(txs types.OrderedTransactionSet, onFailingTxn fun } func (p *Pipeline) TryPushTxn(tx *types.Transaction) (*Result, error) { - if p.txnQueue == nil { - return nil, ErrApplyStageDone - } select { case p.txnQueue <- tx: + case <-p.ctx.Done(): + return nil, ErrPipelineDone case res := <-p.ResultCh: return res, nil } @@ -167,7 +178,7 @@ func (p *Pipeline) TryPushTxn(tx *types.Transaction) (*Result, error) { select { case err, valid := <-p.applyStageRespCh: if !valid { - return nil, ErrApplyStageDone + return nil, ErrPipelineDone } return nil, err case res := <-p.ResultCh: @@ -177,18 +188,13 @@ func (p *Pipeline) TryPushTxn(tx *types.Transaction) (*Result, error) { // Release releases all resources related to the pipeline func (p *Pipeline) Release() { - p.Stop() - - select { - case <-p.applyStageRespCh: - <-p.ResultCh - case <-p.ResultCh: - <-p.applyStageRespCh - } + p.cancelCtx() + p.wg.Wait() } type BlockCandidate struct { LastTrace *types.BlockTrace + RustTrace unsafe.Pointer NextL1MsgIndex uint64 // accumulated state @@ -199,25 +205,47 @@ type BlockCandidate struct { CoalescedLogs []*types.Log } +// sendCancellable tries to send msg to resCh but allows send operation to be cancelled +// by closing cancelCh. Returns true if cancelled. +func sendCancellable[T any, C comparable](resCh chan T, msg T, cancelCh <-chan C) bool { + var zeroC C + + select { + case resCh <- msg: + return false + case cancelSignal := <-cancelCh: + if cancelSignal != zeroC { + panic("shouldn't have happened") + } + return true + } +} + func (p *Pipeline) traceAndApplyStage(txsIn <-chan *types.Transaction) (<-chan error, <-chan *BlockCandidate, error) { p.state.StartPrefetcher("miner") - newCandidateCh := make(chan *BlockCandidate) + downstreamCh := make(chan *BlockCandidate, p.downstreamChCapacity()) resCh := make(chan error) + p.wg.Add(1) go func() { defer func() { - close(newCandidateCh) + close(downstreamCh) close(resCh) p.state.StopPrefetcher() + p.wg.Done() }() var tx *types.Transaction for { - applyIdleTimer.Time(func() { - tx = <-txsIn - }) - if tx == nil { + idleStart := time.Now() + select { + case tx = <-txsIn: + if tx == nil { + return + } + case <-p.ctx.Done(): return } + applyIdleTimer.UpdateSince(idleStart) applyStart := time.Now() @@ -234,13 +262,13 @@ func (p *Pipeline) traceAndApplyStage(txsIn <-chan *types.Transaction) (<-chan e if tx.IsL1MessageTx() && tx.AsL1MessageTx().QueueIndex != p.nextL1MsgIndex { // Continue, we might still be able to include some L2 messages - resCh <- ErrUnexpectedL1MessageIndex + sendCancellable(resCh, ErrUnexpectedL1MessageIndex, p.ctx.Done()) continue } if !tx.IsL1MessageTx() && !p.chain.Config().Scroll.IsValidBlockSize(p.blockSize+tx.Size()) { // can't fit this txn in this block, silently ignore and continue looking for more txns - resCh <- nil + sendCancellable(resCh, nil, p.ctx.Done()) continue } @@ -267,21 +295,16 @@ func (p *Pipeline) traceAndApplyStage(txsIn <-chan *types.Transaction) (<-chan e } stallStart := time.Now() - select { - case newCandidateCh <- &BlockCandidate{ - LastTrace: trace, + if sendCancellable(downstreamCh, &BlockCandidate{ NextL1MsgIndex: p.nextL1MsgIndex, + LastTrace: trace, Header: types.CopyHeader(&p.Header), State: p.state.Copy(), Txs: p.txs, Receipts: p.receipts, CoalescedLogs: p.coalescedLogs, - }: - case tx = <-txsIn: - if tx != nil { - panic("shouldn't have happened") - } + }, p.ctx.Done()) { // next stage terminated and caller terminated us as well return } @@ -294,10 +317,10 @@ func (p *Pipeline) traceAndApplyStage(txsIn <-chan *types.Transaction) (<-chan e } } applyTimer.UpdateSince(applyStart) - resCh <- err + sendCancellable(resCh, err, p.ctx.Done()) } }() - return resCh, newCandidateCh, nil + return resCh, downstreamCh, nil } type Result struct { @@ -309,25 +332,70 @@ type Result struct { FinalBlock *BlockCandidate } +func (p *Pipeline) encodeStage(traces <-chan *BlockCandidate) <-chan *BlockCandidate { + downstreamCh := make(chan *BlockCandidate, p.downstreamChCapacity()) + p.wg.Add(1) + + go func() { + defer func() { + close(downstreamCh) + p.wg.Done() + }() + buffer := new(bytes.Buffer) + for { + idleStart := time.Now() + select { + case trace := <-traces: + if trace == nil { + return + } + encodeIdleTimer.UpdateSince(idleStart) + + encodeStart := time.Now() + if p.ccc != nil { + trace.RustTrace = circuitcapacitychecker.MakeRustTrace(trace.LastTrace, buffer) + if trace.RustTrace == nil { + log.Error("making rust trace", "txHash", trace.LastTrace.Transactions[0].TxHash) + return + } + } + encodeTimer.UpdateSince(encodeStart) + + stallStart := time.Now() + sendCancellable(downstreamCh, trace, p.ctx.Done()) + encodeStallTimer.UpdateSince(stallStart) + case <-p.ctx.Done(): + return + } + + } + }() + return downstreamCh +} + func (p *Pipeline) cccStage(candidates <-chan *BlockCandidate, deadline time.Time) <-chan *Result { if p.ccc != nil { p.ccc.Reset() } - resultCh := make(chan *Result) + resultCh := make(chan *Result, 1) var lastCandidate *BlockCandidate var lastAccRows *types.RowConsumption var deadlineReached bool + p.wg.Add(1) go func() { deadlineTimer := time.NewTimer(time.Until(deadline)) defer func() { close(resultCh) deadlineTimer.Stop() lifetimeTimer.UpdateSince(p.start) + p.wg.Done() }() for { idleStart := time.Now() select { + case <-p.ctx.Done(): + return case <-deadlineTimer.C: cccIdleTimer.UpdateSince(idleStart) // note: currently we don't allow empty blocks, but if we ever do; make sure to CCC check it first @@ -345,7 +413,7 @@ func (p *Pipeline) cccStage(candidates <-chan *BlockCandidate, deadline time.Tim var accRows *types.RowConsumption var err error if candidate != nil && p.ccc != nil { - accRows, err = p.ccc.ApplyTransaction(candidate.LastTrace) + accRows, err = p.ccc.ApplyTransactionRustTrace(candidate.RustTrace) lastTxn := candidate.Txs[candidate.Txs.Len()-1] cccTimer.UpdateSince(cccStart) if err != nil { @@ -427,3 +495,13 @@ func (p *Pipeline) traceAndApply(tx *types.Transaction) (*types.Receipt, *types. } return receipt, trace, nil } + +// downstreamChCapacity returns the channel capacity that should be used for downstream channels. +// It aims to minimize stalls caused by different computational costs of different transactions +func (p *Pipeline) downstreamChCapacity() int { + cap := 1 + if p.chain.Config().Scroll.MaxTxPerBlock != nil { + cap = *p.chain.Config().Scroll.MaxTxPerBlock + } + return cap +}