Compare commits
2
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
4fc5cd22ea | ||
|
|
b70ac40e37 |
@@ -0,0 +1,12 @@
|
||||
package internal
|
||||
|
||||
// EphemeralStore loads an empty state and discards everything on save
|
||||
type EphemeralStore struct{}
|
||||
|
||||
func (EphemeralStore) Load() (map[string]*FillerQueue, error) {
|
||||
return make(map[string]*FillerQueue), nil
|
||||
}
|
||||
|
||||
func (EphemeralStore) Save(map[string]*FillerQueue) error {
|
||||
return nil
|
||||
}
|
||||
+27
-5
@@ -53,6 +53,7 @@ type ReportWriter interface {
|
||||
|
||||
type optionals struct {
|
||||
selector Selector
|
||||
store Store
|
||||
}
|
||||
|
||||
func applyOptions(defaults optionals, opts []Option) optionals {
|
||||
@@ -64,23 +65,38 @@ func applyOptions(defaults optionals, opts []Option) optionals {
|
||||
|
||||
type Option func(*optionals)
|
||||
|
||||
// Selector returns true if a record should be selected for processing, false otherwise.
|
||||
type Selector func(Record) bool
|
||||
|
||||
func WithSelector(s Selector) Option {
|
||||
return func(o *optionals) {
|
||||
o.selector = s
|
||||
}
|
||||
}
|
||||
|
||||
// Selector returns true if a record should be selected for processing, false otherwise.
|
||||
type Selector func(Record) bool
|
||||
type Store interface {
|
||||
Load() (map[string]*FillerQueue, error)
|
||||
Save(map[string]*FillerQueue) error
|
||||
}
|
||||
|
||||
func WithStore(s Store) Option {
|
||||
return func(o *optionals) {
|
||||
o.store = s
|
||||
}
|
||||
}
|
||||
|
||||
// BuildReport reads records from a RecordReader and, if the record passes the Selector, it is
|
||||
// processed into the ReportWriter.
|
||||
func BuildReport(ctx context.Context, reader RecordReader, writer ReportWriter, options ...Option) error {
|
||||
optionals := applyOptions(optionals{
|
||||
selector: Any(),
|
||||
store: EphemeralStore{},
|
||||
}, options)
|
||||
|
||||
buys := make(map[string]*FillerQueue)
|
||||
buys, err := optionals.store.Load()
|
||||
if err != nil {
|
||||
return fmt.Errorf("loading state: %w", err)
|
||||
}
|
||||
|
||||
var buysCount, sellsCount int64
|
||||
var lastTimestamp time.Time
|
||||
@@ -102,6 +118,11 @@ func BuildReport(ctx context.Context, reader RecordReader, writer ReportWriter,
|
||||
rec, err := reader.ReadRecord(ctx)
|
||||
if err != nil {
|
||||
if errors.Is(err, io.EOF) {
|
||||
err = optionals.store.Save(buys)
|
||||
if err != nil {
|
||||
return fmt.Errorf("saving state: %w", err)
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
return err
|
||||
@@ -125,9 +146,10 @@ func BuildReport(ctx context.Context, reader RecordReader, writer ReportWriter,
|
||||
if err != nil {
|
||||
return fmt.Errorf("processing record: %w", err)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// processRecord either adds buys to the queue or consumes buys from the queue when processing a
|
||||
|
||||
Reference in New Issue
Block a user