price pipeline
Signed-off-by: GnomeZworc <nicolas.boufidjeline@g3e.fr>
This commit is contained in:
parent
d7ec1fb355
commit
d9581d7f7e
16 changed files with 790 additions and 10 deletions
133
internal/pipeline/pipeline.go
Normal file
133
internal/pipeline/pipeline.go
Normal file
|
|
@ -0,0 +1,133 @@
|
|||
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
|
||||
}
|
||||
Loading…
Add table
Add a link
Reference in a new issue