Compare commits

...
8 Commits
Author SHA1 Message Date
natercio c08261c1fb minor error msg tweak
Generate check / check-changes (pull_request) Successful in 36s
Generate check / verify-generate (pull_request) Skipped
Quality / check-changes (pull_request) Successful in 36s
Quality / run-tests (pull_request) Failing after 1m19s
2026-08-02 11:07:29 +01:00
natercio 1df5bffbd6 simplify output and stop on invalid kind 2026-08-02 10:56:08 +01:00
natercio fa95552971 moved open-figi implementation into own package 2026-08-02 10:55:31 +01:00
natercio 8f6cefcd68 stronger kind tests and validation 2026-08-01 16:32:20 +01:00
natercio d2b2f934a0 initial file store
Generate check / check-changes (pull_request) Successful in 20s
Quality / check-changes (pull_request) Successful in 34s
Generate check / verify-generate (pull_request) Has been skipped
Quality / run-tests (pull_request) Successful in 31s
2026-07-12 07:56:52 +01:00
natercio 4fc5cd22ea make use of the store optional 2026-07-11 16:43:27 +01:00
natercio b70ac40e37 implement a store abstraction 2026-07-11 16:38:40 +01:00
natercio 0a657682f6 use options in the report 2026-07-11 16:28:35 +01:00
9 changed files with 216 additions and 91 deletions
+3 -2
View File
@@ -10,6 +10,7 @@ import (
"time" "time"
"github.com/nmoniz/any2anexoj/internal" "github.com/nmoniz/any2anexoj/internal"
"github.com/nmoniz/any2anexoj/internal/ofigi"
"github.com/nmoniz/any2anexoj/internal/trading212" "github.com/nmoniz/any2anexoj/internal/trading212"
"github.com/spf13/pflag" "github.com/spf13/pflag"
"golang.org/x/sync/errgroup" "golang.org/x/sync/errgroup"
@@ -71,7 +72,7 @@ func run(ctx context.Context) error {
} }
eg.Go(func() error { eg.Go(func() error {
return internal.BuildReport(ctx, reader, writer, selector) return internal.BuildReport(ctx, reader, writer, internal.WithSelector(selector))
}) })
err = eg.Wait() err = eg.Wait()
@@ -97,7 +98,7 @@ func run(ctx context.Context) error {
func getReader(platform string, ofAPIKey string) (internal.RecordReader, error) { func getReader(platform string, ofAPIKey string) (internal.RecordReader, error) {
switch platform { switch platform {
case "trading212": case "trading212":
return trading212.NewRecordReader(os.Stdin, internal.NewOpenFIGI(&http.Client{Timeout: 5 * time.Second}, ofAPIKey)), nil return trading212.NewRecordReader(os.Stdin, ofigi.NewClient(&http.Client{Timeout: 5 * time.Second}, ofAPIKey)), nil
default: default:
return nil, fmt.Errorf("unsupported platform: %s", platform) return nil, fmt.Errorf("unsupported platform: %s", platform)
} }
+12
View File
@@ -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
}
+29
View File
@@ -0,0 +1,29 @@
package internal
import (
"fmt"
"os"
)
// FileStore loads and saves the state into a file
type FileStore struct {
file *os.File
}
func NewFileStore(filename string) (*FileStore, error) {
f, err := os.Open(filename)
if err != nil {
return nil, err
}
return &FileStore{
file: f,
}, nil
}
func (fs *FileStore) Load() (map[string]*FillerQueue, error) {
return nil, fmt.Errorf("not implemented")
}
func (fs *FileStore) Save(map[string]*FillerQueue) error {
return fmt.Errorf("not implemented")
}
+11 -5
View File
@@ -7,11 +7,12 @@ const (
KindBuy KindBuy
KindSell KindSell
KindSplit KindSplit
sentinelKind
) )
// String returns a human readable value // String returns a unique string value for Kind k
func (d Kind) String() string { func (k Kind) String() string {
switch d { switch k {
case KindBuy: case KindBuy:
return "buy" return "buy"
case KindSell: case KindSell:
@@ -23,8 +24,13 @@ func (d Kind) String() string {
} }
} }
// Is returns true when k equals o // Valid returns true if k is an accepted value for the Kind type
func (k Kind) Valid() bool {
return k > 0 && k < sentinelKind
}
// Is returns true when k and o are valid and equal.
func (k Kind) Is(o any) bool { func (k Kind) Is(o any) bool {
other, ok := o.(Kind) other, ok := o.(Kind)
return ok && k == other return ok && k.Valid() && k == other
} }
+75 -40
View File
@@ -1,60 +1,95 @@
package internal package internal
import "testing" import (
"fmt"
"testing"
)
func TestSide_String(t *testing.T) { func TestSide_String(t *testing.T) {
tests := []struct { const unknown = "unknown"
name string
side Kind seen := make(map[string]Kind, sentinelKind)
want string for k := Kind(1); k < sentinelKind; k++ {
}{ t.Run(fmt.Sprintf("Kind %d", k), func(t *testing.T) {
{"buy", KindBuy, "buy"}, str := k.String()
{"sell", KindSell, "sell"},
{"unknown", KindUnknown, "unknown"}, if other, ok := seen[str]; ok {
t.Errorf("want Kind(%d).String to be unique but was a duplicate of Kind(%d)", k, other)
} else {
seen[str] = k
} }
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) { if len(str) == 0 {
if got := tt.side.String(); got != tt.want { t.Errorf("want Kind(%d).String to be non-empty", k)
t.Errorf("want Side.String() to be %v but got %v", tt.want, got) }
if str == unknown {
t.Errorf("want Kind(%d).String to be a known value", k)
} }
}) })
} }
if KindUnknown.String() != unknown {
t.Errorf("want Kind(0) to be unknown")
}
if Kind(sentinelKind).String() != unknown {
t.Errorf("want Kind(%d) to be unknown", sentinelKind)
}
} }
func TestSide_IsBuy(t *testing.T) { func TestSide_Valid(t *testing.T) {
tests := []struct { for k := Kind(1); k < sentinelKind; k++ {
name string if !k.Valid() {
side Kind t.Errorf("want %s(%d) to be valid", k, k)
want bool
}{
{"buy", KindBuy, true},
{"sell", KindSell, false},
{"unknown", KindUnknown, false},
} }
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
if got := tt.side.Is(KindBuy); got != tt.want {
t.Errorf("want Side.IsBuy() to be %v but got %v", tt.want, got)
} }
})
if KindUnknown.Valid() {
t.Errorf("want Kind(0) to be invalid")
}
if Kind(sentinelKind).Valid() {
t.Errorf("want Kind(%d) to be invalid", sentinelKind)
} }
} }
func TestSide_IsSell(t *testing.T) { func TestSide_Is(t *testing.T) {
tests := []struct { t.Run("valid is self", func(t *testing.T) {
name string for k := Kind(1); k < sentinelKind; k++ {
side Kind if !k.Is(k) {
want bool t.Errorf("want Kind(%d).Is(%d) to be true", k, k)
}{
{"buy", KindBuy, false},
{"sell", KindSell, true},
{"unknown", KindUnknown, false},
} }
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
if got := tt.side.Is(KindSell); got != tt.want {
t.Errorf("want Side.IsSell() to be %v but got %v", tt.want, got)
} }
}) })
t.Run("valid is unknown", func(t *testing.T) {
for k := Kind(1); k < sentinelKind; k++ {
if k.Is(KindUnknown) {
t.Errorf("want Kind(%d).Is(0) to be false", k)
} }
if k.Is(sentinelKind) {
t.Errorf("want Kind(%d).Is(%d) to be false", k, sentinelKind)
}
if k.Is(struct{}{}) {
t.Errorf("want Kind(%d).Is(other type) to be false", k)
}
}
})
t.Run("unknown is unknown", func(t *testing.T) {
if KindUnknown.Is(KindUnknown) {
t.Errorf("want Kind(0).Is(0) to be false")
}
if KindUnknown.Is(sentinelKind) {
t.Errorf("want Kind(0).Is(%d) to be false", sentinelKind)
}
if sentinelKind.Is(sentinelKind) {
t.Errorf("want Kind(%d).Is(%d) to be false", sentinelKind, sentinelKind)
}
})
} }
@@ -1,4 +1,4 @@
package internal package ofigi
import ( import (
"bytes" "bytes"
@@ -16,8 +16,8 @@ import (
var OpenFIGIAPIKeyHeader = http.CanonicalHeaderKey("X-OPENFIGI-APIKEY") var OpenFIGIAPIKeyHeader = http.CanonicalHeaderKey("X-OPENFIGI-APIKEY")
// OpenFIGI is a small adapter for the openfigi.com api. // Client is a thin adapter for the openfigi.com api.
type OpenFIGI struct { type Client struct {
client *http.Client client *http.Client
apiKey string apiKey string
mappingLimiter *rate.Limiter mappingLimiter *rate.Limiter
@@ -25,12 +25,13 @@ type OpenFIGI struct {
mu sync.RWMutex mu sync.RWMutex
// TODO: there's no eviction policy at the moment as this is only used by short-lived application // TODO: there's no eviction policy at the moment as this is only used by short-lived application
// which processes a relatively small amount of records. We need to consider using an external // which processes a relatively small amount of records. We need to consider using an external
// cache lib (like golang-lru or go-cache) if this becomes a problem or implement this ourselves. // cache lib (like golang-lru or go-cache) if this becomes a problem or implement eviction
// ourselves.
securityTypeCache map[string]string securityTypeCache map[string]string
} }
// NewOpenFIGI creates an OpenFIGI client that uses the API key if provided // NewOpenFIGI creates an OpenFIGI client that uses the API key if provided
func NewOpenFIGI(c *http.Client, apiKey string) *OpenFIGI { func NewOpenFIGI(c *http.Client, apiKey string) *Client {
// Rate limits as per https://www.openfigi.com/api/documentation#rate-limits // Rate limits as per https://www.openfigi.com/api/documentation#rate-limits
limiter := rate.NewLimiter(rate.Every(time.Minute), 25) limiter := rate.NewLimiter(rate.Every(time.Minute), 25)
if len(apiKey) > 0 { if len(apiKey) > 0 {
@@ -40,7 +41,7 @@ func NewOpenFIGI(c *http.Client, apiKey string) *OpenFIGI {
slog.Debug("OpenFIGI client: created with puplic rate limits") slog.Debug("OpenFIGI client: created with puplic rate limits")
} }
return &OpenFIGI{ return &Client{
client: c, client: c,
apiKey: apiKey, apiKey: apiKey,
mappingLimiter: limiter, mappingLimiter: limiter,
@@ -48,7 +49,7 @@ func NewOpenFIGI(c *http.Client, apiKey string) *OpenFIGI {
} }
} }
func (of *OpenFIGI) SecurityTypeByISIN(ctx context.Context, isin string) (string, error) { func (of *Client) SecurityTypeByISIN(ctx context.Context, isin string) (string, error) {
of.mu.RLock() of.mu.RLock()
if secType, ok := of.securityTypeCache[isin]; ok { if secType, ok := of.securityTypeCache[isin]; ok {
of.mu.RUnlock() of.mu.RUnlock()
@@ -1,4 +1,4 @@
package internal_test package ofigi_test
import ( import (
"bytes" "bytes"
@@ -9,7 +9,7 @@ import (
"testing" "testing"
"time" "time"
"github.com/nmoniz/any2anexoj/internal" "github.com/nmoniz/any2anexoj/internal/ofigi"
) )
func TestOpenFIGI_SecurityTypeByISIN(t *testing.T) { func TestOpenFIGI_SecurityTypeByISIN(t *testing.T) {
@@ -110,7 +110,7 @@ func TestOpenFIGI_SecurityTypeByISIN(t *testing.T) {
} }
for _, tt := range tests { for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) { t.Run(tt.name, func(t *testing.T) {
of := internal.NewOpenFIGI(tt.client, "") of := ofigi.NewOpenFIGI(tt.client, "")
got, gotErr := of.SecurityTypeByISIN(context.Background(), tt.isin) got, gotErr := of.SecurityTypeByISIN(context.Background(), tt.isin)
if gotErr != nil { if gotErr != nil {
@@ -145,7 +145,7 @@ func TestOpenFIGI_SecurityTypeByISIN_Cache(t *testing.T) {
}, nil }, nil
}) })
of := internal.NewOpenFIGI(c, "") of := ofigi.NewOpenFIGI(c, "")
got, gotErr := of.SecurityTypeByISIN(t.Context(), "NL0000235190") got, gotErr := of.SecurityTypeByISIN(t.Context(), "NL0000235190")
if gotErr != nil { if gotErr != nil {
@@ -171,15 +171,15 @@ func TestOpenFIGI_SecurityTypeByISIN_APIKey(t *testing.T) {
wantAPIKey := "123abc-456xyz" wantAPIKey := "123abc-456xyz"
c := NewTestClient(t, func(req *http.Request) (*http.Response, error) { c := NewTestClient(t, func(req *http.Request) (*http.Response, error) {
value, ok := req.Header[internal.OpenFIGIAPIKeyHeader] value, ok := req.Header[ofigi.OpenFIGIAPIKeyHeader]
if !ok { if !ok {
t.Fatalf("want %q header but got none: %v", internal.OpenFIGIAPIKeyHeader, req.Header) t.Fatalf("want %q header but got none: %v", ofigi.OpenFIGIAPIKeyHeader, req.Header)
} }
if len(value) != 1 { if len(value) != 1 {
t.Fatalf("want exactly one %q header value but got %d", internal.OpenFIGIAPIKeyHeader, len(value)) t.Fatalf("want exactly one %q header value but got %d", ofigi.OpenFIGIAPIKeyHeader, len(value))
} }
if value[0] != wantAPIKey { if value[0] != wantAPIKey {
t.Fatalf("want %q header value %q but got %q", internal.OpenFIGIAPIKeyHeader, wantAPIKey, value[0]) t.Fatalf("want %q header value %q but got %q", ofigi.OpenFIGIAPIKeyHeader, wantAPIKey, value[0])
} }
return &http.Response{ return &http.Response{
Status: http.StatusText(http.StatusOK), Status: http.StatusText(http.StatusOK),
@@ -187,7 +187,7 @@ func TestOpenFIGI_SecurityTypeByISIN_APIKey(t *testing.T) {
Body: io.NopCloser(bytes.NewBufferString(`[{"data":[{"securityType":"Common Stock"}]}]`)), Body: io.NopCloser(bytes.NewBufferString(`[{"data":[{"securityType":"Common Stock"}]}]`)),
}, nil }, nil
}) })
of := internal.NewOpenFIGI(c, wantAPIKey) of := ofigi.NewOpenFIGI(c, wantAPIKey)
_, err := of.SecurityTypeByISIN(t.Context(), "US1234567890") _, err := of.SecurityTypeByISIN(t.Context(), "US1234567890")
if err != nil { if err != nil {
@@ -197,9 +197,9 @@ func TestOpenFIGI_SecurityTypeByISIN_APIKey(t *testing.T) {
t.Run("without API key", func(t *testing.T) { t.Run("without API key", func(t *testing.T) {
c := NewTestClient(t, func(req *http.Request) (*http.Response, error) { c := NewTestClient(t, func(req *http.Request) (*http.Response, error) {
_, ok := req.Header[internal.OpenFIGIAPIKeyHeader] _, ok := req.Header[ofigi.OpenFIGIAPIKeyHeader]
if ok { if ok {
t.Fatalf("want no %s header but got one", internal.OpenFIGIAPIKeyHeader) t.Fatalf("want no %s header but got one", ofigi.OpenFIGIAPIKeyHeader)
} }
return &http.Response{ return &http.Response{
Status: http.StatusText(http.StatusOK), Status: http.StatusText(http.StatusOK),
@@ -207,7 +207,7 @@ func TestOpenFIGI_SecurityTypeByISIN_APIKey(t *testing.T) {
Body: io.NopCloser(bytes.NewBufferString(`[{"data":[{"securityType":"Common Stock"}]}]`)), Body: io.NopCloser(bytes.NewBufferString(`[{"data":[{"securityType":"Common Stock"}]}]`)),
}, nil }, nil
}) })
of := internal.NewOpenFIGI(c, "") of := ofigi.NewOpenFIGI(c, "")
_, err := of.SecurityTypeByISIN(t.Context(), "US1234567890") _, err := of.SecurityTypeByISIN(t.Context(), "US1234567890")
if err != nil { if err != nil {
t.Fatalf("want success but got an error: %s", err) t.Fatalf("want success but got an error: %s", err)
+60 -19
View File
@@ -51,18 +51,58 @@ type ReportWriter interface {
Write(context.Context, ReportItem) error Write(context.Context, ReportItem) error
} }
type optionals struct {
selector Selector
store Store
}
func applyOptions(defaults optionals, opts []Option) optionals {
for _, opt := range opts {
opt(&defaults)
}
return defaults
}
type Option func(*optionals)
// Selector returns true if a record should be selected for processing, false otherwise. // Selector returns true if a record should be selected for processing, false otherwise.
type Selector func(Record) bool type Selector func(Record) bool
func WithSelector(s Selector) Option {
return func(o *optionals) {
o.selector = s
}
}
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 // BuildReport reads records from a RecordReader and, if the record passes the Selector, it is
// processed into the ReportWriter. // processed into the ReportWriter.
func BuildReport(ctx context.Context, reader RecordReader, writer ReportWriter, sel Selector) error { func BuildReport(ctx context.Context, reader RecordReader, writer ReportWriter, options ...Option) error {
buys := make(map[string]*FillerQueue) optionals := applyOptions(optionals{
selector: Any(),
store: EphemeralStore{},
}, options)
var buysCount, sellsCount int64 buys, err := optionals.store.Load()
var lastTimestamp time.Time if err != nil {
progTicker := time.NewTicker(10 * time.Second) return fmt.Errorf("loading state: %w", err)
}
var (
recordsCount int64
lastTimestamp time.Time
progTicker = time.NewTicker(10 * time.Second)
)
for { for {
select { select {
case <-ctx.Done(): case <-ctx.Done():
@@ -70,27 +110,29 @@ func BuildReport(ctx context.Context, reader RecordReader, writer ReportWriter,
case <-progTicker.C: case <-progTicker.C:
slog.InfoContext( slog.InfoContext(
ctx, "Progress update", ctx, "Progress update",
slog.Int64("total_records", buysCount+sellsCount), slog.Int64("records_count", recordsCount),
slog.Int64("sell_records", sellsCount),
slog.Int64("buy_records", buysCount),
slog.Time("last_record_timestamp", lastTimestamp), slog.Time("last_record_timestamp", lastTimestamp),
) )
default: default:
rec, err := reader.ReadRecord(ctx) rec, err := reader.ReadRecord(ctx)
if err != nil { if err != nil {
if errors.Is(err, io.EOF) { if errors.Is(err, io.EOF) {
err = optionals.store.Save(buys)
if err != nil {
return fmt.Errorf("saving state: %w", err)
}
return nil return nil
} }
return err return err
} }
if rec.Kind().Is(KindBuy) { if !rec.Kind().Valid() {
buysCount++ return fmt.Errorf("found invalid Kind(%d)", rec.Kind())
} else if rec.Kind().Is(KindSell) {
sellsCount++
} }
lastTimestamp = rec.Timestamp() lastTimestamp = rec.Timestamp()
recordsCount++
buyQueue, ok := buys[rec.Symbol()] buyQueue, ok := buys[rec.Symbol()]
if !ok { if !ok {
@@ -98,21 +140,16 @@ func BuildReport(ctx context.Context, reader RecordReader, writer ReportWriter,
buys[rec.Symbol()] = buyQueue buys[rec.Symbol()] = buyQueue
} }
err = processRecord(ctx, buyQueue, rec, sel, writer) err = processRecord(ctx, buyQueue, rec, optionals.selector, writer)
if err != nil { if err != nil {
return fmt.Errorf("processing record: %w", err) return fmt.Errorf("processing record: %w", err)
} }
} }
} }
} }
// processRecord either adds buys to the queue or consumes buys from the queue when processing a // processRecord either adds buys to the queue or consumes buys from the queue when processing a
// sell record. // sell record.
//
// NOTE: Selectors are only applied when processing sell records for performance reasons. It's much
// cheaper to just accumulate buys and only actually inspect any records once a sell happens. This
// avoids potential network requests to for every single record.
func processRecord(ctx context.Context, q *FillerQueue, rec Record, sel Selector, writer ReportWriter) error { func processRecord(ctx context.Context, q *FillerQueue, rec Record, sel Selector, writer ReportWriter) error {
slog.Debug( slog.Debug(
"Report: processing record", "Report: processing record",
@@ -122,6 +159,9 @@ func processRecord(ctx context.Context, q *FillerQueue, rec Record, sel Selector
switch rec.Kind() { switch rec.Kind() {
case KindBuy: case KindBuy:
// Selectors are only applied when processing sell records for performance reasons. It's much
// cheaper to just accumulate buys and only actually inspect any records once a sell happens. This
// avoids potential network requests to for every single record.
q.Push(NewFiller(rec)) q.Push(NewFiller(rec))
case KindSell: case KindSell:
@@ -142,7 +182,8 @@ func processRecord(ctx context.Context, q *FillerQueue, rec Record, sel Selector
return ErrInsufficientBoughtVolume return ErrInsufficientBoughtVolume
} }
// Since we don't apply selectors while processing buys we need to apply them here. // Since we don't apply selectors while processing buys we need to apply them here, before we
// actually use them.
if !sel(buy) { if !sel(buy) {
continue continue
} }
+1 -1
View File
@@ -43,7 +43,7 @@ func TestBuildReport(t *testing.T) {
Taxes: decimal.Decimal{}, Taxes: decimal.Decimal{},
})).Times(1) })).Times(1)
gotErr := internal.BuildReport(t.Context(), reader, writer, internal.Any()) gotErr := internal.BuildReport(t.Context(), reader, writer)
if gotErr != nil { if gotErr != nil {
t.Fatalf("got unexpected err: %v", gotErr) t.Fatalf("got unexpected err: %v", gotErr)
} }