Move the event emitter and lifecycle events from the engine into the reducer, making the reducer the single point of orchestration for reduction. This eliminates the engine package entirely. - Add events.go to pkg/reducer with Start, Step, and Stop events. - Extend Reducer interface to embed Emitter and add Expression() method. - Update NormalOrderReducer to embed BaseEmitter and emit lifecycle events. - Update all plugins to attach to Reducer instead of Engine. - Remove internal/engine package. - Add Off() method to BaseEmitter to complete Emitter interface. - Fix Emitter.On signature to use generic type E instead of string.
45 lines
821 B
Go
45 lines
821 B
Go
package plugins
|
|
|
|
import (
|
|
"fmt"
|
|
"os"
|
|
"time"
|
|
|
|
"git.maximhutz.com/max/lambda/internal/statistics"
|
|
"git.maximhutz.com/max/lambda/pkg/reducer"
|
|
)
|
|
|
|
// An observer, to track reduction performance.
|
|
type Statistics struct {
|
|
start time.Time
|
|
steps uint64
|
|
}
|
|
|
|
// Create a new reduction performance Statistics.
|
|
func NewStatistics(r reducer.Reducer) *Statistics {
|
|
plugin := &Statistics{}
|
|
r.On(reducer.StartEvent, plugin.Start)
|
|
r.On(reducer.StepEvent, plugin.Step)
|
|
r.On(reducer.StopEvent, plugin.Stop)
|
|
|
|
return plugin
|
|
}
|
|
|
|
func (t *Statistics) Start() {
|
|
t.start = time.Now()
|
|
t.steps = 0
|
|
}
|
|
|
|
func (t *Statistics) Step() {
|
|
t.steps++
|
|
}
|
|
|
|
func (t *Statistics) Stop() {
|
|
results := statistics.Results{
|
|
StepsTaken: t.steps,
|
|
TimeElapsed: uint64(time.Since(t.start).Milliseconds()),
|
|
}
|
|
|
|
fmt.Fprint(os.Stderr, results.String())
|
|
}
|