account/internal/pipeline/pipeline.go
GnomeZworc d9581d7f7e
price pipeline
Signed-off-by: GnomeZworc <nicolas.boufidjeline@g3e.fr>
2026-06-13 12:33:39 +02:00

133 lines
3.6 KiB
Go

package pipeline
import (
"context"
"log/slog"
"time"
"git.g3e.fr/H6N/account/internal/store"
)
type Config struct {
OpenFIGIKey string
CoinGeckoKey string
}
type Pipeline struct {
store *store.Store
cfg Config
logger *slog.Logger
}
func New(s *store.Store, cfg Config, logger *slog.Logger) *Pipeline {
return &Pipeline{store: s, cfg: cfg, logger: logger}
}
// FetchAll récupère les prix de tous les instruments et les écrit dans price_history.
func (p *Pipeline) FetchAll(ctx context.Context) error {
instruments, err := p.store.ListInstruments(ctx)
if err != nil {
return err
}
now := time.Now().UTC()
// Séparer les instruments par type pour batch CoinGecko
var cryptoIDs []string
cryptoByID := make(map[string]store.Instrument)
for _, inst := range instruments {
switch inst.Type {
case "devise":
if err := p.store.UpsertPrice(ctx, inst.ID, now, 1.0); err != nil {
p.logger.Error("upsert devise price", "code", inst.Code, "error", err)
}
case "action", "etf":
ticker, err := p.resolveTicker(ctx, inst)
if err != nil {
p.logger.Warn("could not resolve ticker", "instrument", inst.Code, "error", err)
continue
}
price, err := FetchYahooPrice(ctx, ticker)
if err != nil {
p.logger.Warn("yahoo fetch failed", "ticker", ticker, "error", err)
continue
}
if err := p.store.UpsertPrice(ctx, inst.ID, now, price); err != nil {
p.logger.Error("upsert price", "instrument", inst.Code, "error", err)
}
p.logger.Info("price fetched", "instrument", inst.Code, "ticker", ticker, "price", price)
case "crypto":
coinID, err := p.resolveCoinID(ctx, inst)
if err != nil {
p.logger.Warn("could not resolve coingecko id", "instrument", inst.Code, "error", err)
continue
}
cryptoIDs = append(cryptoIDs, coinID)
cryptoByID[coinID] = inst
}
}
// Batch fetch crypto
if len(cryptoIDs) > 0 {
prices, err := FetchCoinGeckoPrices(ctx, p.cfg.CoinGeckoKey, cryptoIDs)
if err != nil {
p.logger.Error("coingecko fetch failed", "error", err)
} else {
for coinID, price := range prices {
inst := cryptoByID[coinID]
if err := p.store.UpsertPrice(ctx, inst.ID, now, price); err != nil {
p.logger.Error("upsert crypto price", "instrument", inst.Code, "error", err)
continue
}
p.logger.Info("price fetched", "instrument", inst.Code, "coin_id", coinID, "price", price)
}
}
}
return nil
}
// CleanHistory supprime les prix intraday des jours passés, ne gardant que le dernier par instrument.
func (p *Pipeline) CleanHistory(ctx context.Context) error {
if err := p.store.CleanPastDays(ctx); err != nil {
p.logger.Error("clean price history failed", "error", err)
return err
}
p.logger.Info("price history cleaned")
return nil
}
// resolveTicker retourne le ticker Yahoo Finance, en le résolvant via OpenFIGI si nécessaire.
func (p *Pipeline) resolveTicker(ctx context.Context, inst store.Instrument) (string, error) {
ticker, err := p.store.GetTicker(ctx, inst.ID)
if err == nil {
return ticker, nil
}
resolved, err := ResolveISIN(ctx, p.cfg.OpenFIGIKey, inst.Code)
if err != nil {
return "", err
}
_ = p.store.UpsertTicker(ctx, inst.ID, resolved)
return resolved, nil
}
// resolveCoinID retourne l'ID CoinGecko, en le résolvant via l'API search si nécessaire.
func (p *Pipeline) resolveCoinID(ctx context.Context, inst store.Instrument) (string, error) {
coinID, err := p.store.GetTicker(ctx, inst.ID)
if err == nil {
return coinID, nil
}
resolved, err := ResolveCryptoID(ctx, p.cfg.CoinGeckoKey, inst.Code)
if err != nil {
return "", err
}
_ = p.store.UpsertTicker(ctx, inst.ID, resolved)
return resolved, nil
}