Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
10 changes: 7 additions & 3 deletions cmd/trader/data/args.go
Original file line number Diff line number Diff line change
Expand Up @@ -65,17 +65,21 @@ type datasetArgFlags struct {
// addDatasetArgFlags registers --from, --to (both required), --format
// (issue #111; defaults to "table"), and --exchange/--kind (issue
// #331; required only for a non-FX provider) on cmd.
func addDatasetArgFlags(cmd *cobra.Command, flags *datasetArgFlags) {
func addDatasetConversionFlags(cmd *cobra.Command, flags *datasetArgFlags) {
cmd.Flags().StringVar(&flags.from, "from", "", "range start (YYYY-MM-DD or RFC3339), required")
cmd.Flags().StringVar(&flags.to, "to", "", "range end (YYYY-MM-DD or RFC3339), required")
cmd.Flags().StringVar(&flags.format, "format", formatTable,
"output format: "+formatTable+" or "+formatJSON)
cmd.Flags().StringVar(&flags.exchange, "exchange", "",
"listing exchange (for example ARCA or NASDAQ); required for a non-FX provider such as alpaca")
cmd.Flags().StringVar(&flags.kind, "kind", "",
`instrument kind, "equity" or "etf"; required for a non-FX provider such as alpaca`)
}

func addDatasetArgFlags(cmd *cobra.Command, flags *datasetArgFlags) {
addDatasetConversionFlags(cmd, flags)
cmd.Flags().StringVar(&flags.format, "format", formatTable,
"output format: "+formatTable+" or "+formatJSON)
}

// fxProviders names the providers whose bare INSTRUMENT argument is a
// 6-letter FX pair symbol, resolved via svc.RegisterFXInstrument.
// "oanda" is Trader's only FX provider today; every other configured
Expand Down
1 change: 1 addition & 0 deletions cmd/trader/data/command.go
Original file line number Diff line number Diff line change
Expand Up @@ -63,6 +63,7 @@ func New() *cobra.Command {
cmd.AddCommand(newSyncCmd())
cmd.AddCommand(newBuildCmd())
cmd.AddCommand(newUpdateCmd())
cmd.AddCommand(newConvertCmd())

return cmd
}
112 changes: 112 additions & 0 deletions cmd/trader/data/convert.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,112 @@
package data

import (
"archive/zip"
"context"
"fmt"
"io"
"os"
"path/filepath"
"strings"

"github.com/spf13/cobra"

svc "github.com/rustyeddy/trader/internal/service/marketdata"
)

// newConvertCmd imports one native Stooq archive member into managed raw
// partitions and builds the requested canonical range. Extraction is always
// temporary; the original ZIP remains the source of truth.
func newConvertCmd() *cobra.Command {
var flags datasetArgFlags
var archivePath string

cmd := &cobra.Command{
Use: "convert INSTRUMENT INTERVAL",
Short: "Convert one downloaded provider archive into canonical data.",
Args: cobra.ExactArgs(2),
RunE: func(cmd *cobra.Command, args []string) error {
if archivePath == "" {
return fmt.Errorf("--archive is required")
}
dc, ok := dataContextFrom(cmd.Context())
if !ok {
return fmt.Errorf("data service is not configured on this command's context")
}
if strings.ToLower(dc.Provider) != "stooq" {
return fmt.Errorf("convert currently supports only provider stooq")
}
req, err := resolveDatasetRequest(cmd, args, flags)
if err != nil {
return err
}
if strings.ToUpper(args[1]) != "D1" {
return fmt.Errorf("stooq conversion currently supports only D1")
}
tmp, err := os.MkdirTemp("", "trader-stooq-convert-")
if err != nil {
return fmt.Errorf("create temporary extraction directory: %w", err)
}
defer func() { _ = os.RemoveAll(tmp) }()
extracted, err := extractStooqMember(cmd.Context(), archivePath, args[0], tmp)
if err != nil {
return err
}
resp, err := dc.Service.Convert(cmd.Context(), svc.ConvertRequest{
DatasetRequest: req, ArchivePath: extracted,
})
if err != nil {
return err
}
if _, err := fmt.Fprintf(cmd.OutOrStdout(), "imported %d rows across %d raw months; published %d canonical partitions\n",
resp.Import.RowsImported, resp.Import.MonthsWritten, len(resp.Build.Result.Published)); err != nil {
Comment on lines +61 to +62
Comment on lines +55 to +62
return err
}
return nil
},
}
cmd.Flags().StringVar(&archivePath, "archive", "", "native Stooq ZIP archive path")
addDatasetConversionFlags(cmd, &flags)
return cmd
}

func extractStooqMember(ctx context.Context, archivePath, symbol, destination string) (string, error) {
zr, err := zip.OpenReader(archivePath)
if err != nil {
return "", fmt.Errorf("open archive: %w", err)
}
defer func() { _ = zr.Close() }()
want := strings.ToLower(strings.TrimSpace(symbol)) + ".us.txt"
for _, entry := range zr.File {
if ctx.Err() != nil {
return "", ctx.Err()
}
if strings.ToLower(filepath.Base(entry.Name)) != want {
continue
}
if entry.FileInfo().IsDir() {
continue
}
in, err := entry.Open()
if err != nil {
return "", fmt.Errorf("open archive member %q: %w", entry.Name, err)
}
path := filepath.Join(destination, filepath.Base(entry.Name))
out, err := os.Create(path)
if err != nil {
_ = in.Close()
return "", fmt.Errorf("create extracted archive member: %w", err)
}
_, copyErr := io.Copy(out, in)
closeInErr := in.Close()
closeOutErr := out.Close()
if copyErr != nil {
return "", fmt.Errorf("extract archive member %q: %w", entry.Name, copyErr)
}
if closeInErr != nil || closeOutErr != nil {
return "", fmt.Errorf("close extracted archive member %q: %v %v", entry.Name, closeInErr, closeOutErr)
}
return path, nil
}
return "", fmt.Errorf("archive %q contains no %s member", archivePath, want)
}
52 changes: 52 additions & 0 deletions cmd/trader/data/convert_integration_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,52 @@
package data_test

import (
"archive/zip"
"crypto/sha256"
"os"
"path/filepath"
"testing"

"github.com/stretchr/testify/require"
)

func TestDataConvertStooqZIPToCanonical(t *testing.T) {
archivePath := filepath.Join(t.TempDir(), "daily.zip")
content := []byte("<TICKER>,<PER>,<DATE>,<TIME>,<OPEN>,<HIGH>,<LOW>,<CLOSE>,<VOL>,<OPENINT>\n" +
"SPY.US,D,20200131,000000,100,101,99,100.5,1000,0\n" +
"SPY.US,D,20200203,000000,100.5,102,100,101.5,2000,0\n")
writeZIPMember(t, archivePath, "data/daily/us/nyse etfs/2/spy.us.txt", content)
before := sha256.Sum256(mustReadFile(t, archivePath))
rawRoot, storeRoot := t.TempDir(), t.TempDir()

out, err := runData(t, storeRoot, rawRoot,
"convert", "SPY", "D1", "--provider", "stooq", "--archive", archivePath,
"--exchange", "ARCA", "--kind", "etf", "--from", "2020-01-01", "--to", "2020-03-01")
require.NoError(t, err)
require.Contains(t, out, "imported 2 rows across 2 raw months; published 2 canonical partitions")
require.Equal(t, before, sha256.Sum256(mustReadFile(t, archivePath)))
require.FileExists(t, filepath.Join(rawRoot, "SPY", "2020", "01", "SPY-2020-01-d1.csv"))
canonical := filepath.Join(storeRoot, "stooq", "SPY", "2020", "01", "SPY-2020-01-d1.csv")
require.FileExists(t, canonical)
require.Contains(t, string(mustReadFile(t, canonical)), `"provider":"stooq"`)
}

func writeZIPMember(t *testing.T, path, name string, content []byte) {
t.Helper()
f, err := os.Create(path)
require.NoError(t, err)
zw := zip.NewWriter(f)
w, err := zw.Create(name)
require.NoError(t, err)
_, err = w.Write(content)
require.NoError(t, err)
require.NoError(t, zw.Close())
require.NoError(t, f.Close())
}

func mustReadFile(t *testing.T, path string) []byte {
t.Helper()
data, err := os.ReadFile(path)
require.NoError(t, err)
return data
}
45 changes: 45 additions & 0 deletions cmd/trader/data/convert_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,45 @@
package data

import (
"archive/zip"
"context"
"os"
"path/filepath"
"testing"

"github.com/stretchr/testify/require"
)

func TestExtractStooqMember_SelectsSPYAndPreservesContent(t *testing.T) {
archivePath := filepath.Join(t.TempDir(), "daily.zip")
f, err := os.Create(archivePath)
require.NoError(t, err)
zw := zip.NewWriter(f)
w, err := zw.Create("data/daily/us/nyse etfs/2/spy.us.txt")
require.NoError(t, err)
_, err = w.Write([]byte("<TICKER>,<PER>\nSPY.US,D\n"))
require.NoError(t, err)
require.NoError(t, zw.Close())
require.NoError(t, f.Close())

destination := t.TempDir()
path, err := extractStooqMember(context.Background(), archivePath, "SPY", destination)
require.NoError(t, err)
contents, err := os.ReadFile(path)
require.NoError(t, err)
require.Equal(t, "<TICKER>,<PER>\nSPY.US,D\n", string(contents))
}

func TestExtractStooqMemberReportsMissingSymbol(t *testing.T) {
archivePath := filepath.Join(t.TempDir(), "daily.zip")
f, err := os.Create(archivePath)
require.NoError(t, err)
zw := zip.NewWriter(f)
_, err = zw.Create("data/daily/us/nyse etfs/2/qqq.us.txt")
require.NoError(t, err)
require.NoError(t, zw.Close())
require.NoError(t, f.Close())

_, err = extractStooqMember(context.Background(), archivePath, "SPY", t.TempDir())
require.ErrorContains(t, err, "contains no spy.us.txt member")
}
40 changes: 40 additions & 0 deletions internal/marketdata/convert.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,40 @@
package marketdata

import (
"context"
"fmt"
"time"

"github.com/rustyeddy/trader/instrument"
"github.com/rustyeddy/trader/internal/marketdata/internal/provider/stooq"
)

// StooqImportResult summarizes a native Stooq archive import without
// exposing provider-internal record types across the Manager boundary.
type StooqImportResult struct {
MonthsWritten int
RowsImported int
FirstDate time.Time
LastDate time.Time
}

// ImportStooqArchive imports one native Stooq archive file into the
// manager's configured managed raw root. The caller owns archive extraction;
// this method only adapts the source into Trader's existing raw partitions.
func (m *Manager) ImportStooqArchive(ctx context.Context, path string, id instrument.ID) (StooqImportResult, error) {
if !m.configured() {
return StooqImportResult{}, fmt.Errorf("marketdata: Stooq archive import: %w: manager is not configured", ErrInvalidConfig)
}
if m.providerName != "stooq" {
return StooqImportResult{}, fmt.Errorf("marketdata: Stooq archive import requires provider stooq, got %q", m.providerName)
}
if m.rawRoot == "" {
return StooqImportResult{}, fmt.Errorf("marketdata: Stooq archive import requires a raw root")
}
listing, err := m.resolver.ResolveInstrument(id, m.providerName, "")
if err != nil {
return StooqImportResult{}, fmt.Errorf("marketdata: Stooq archive import: resolve listing: %w", err)
}
result, err := stooq.ImportArchive(ctx, path, m.rawRoot, listing.Symbol())
return StooqImportResult{MonthsWritten: result.MonthsWritten, RowsImported: result.RowsImported, FirstDate: result.FirstDate, LastDate: result.LastDate}, err
}
34 changes: 34 additions & 0 deletions internal/service/marketdata/convert.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,34 @@
package marketdata

import (
"context"
"fmt"
"log/slog"

"github.com/rustyeddy/trader/marketdata"
)

// Convert imports one already-downloaded provider archive and then builds
// canonical data through the same Manager Plan/Build path used elsewhere.
func (s *Service) Convert(ctx context.Context, req ConvertRequest) (resp ConvertResponse, err error) {
if err := req.Validate(); err != nil {
return ConvertResponse{}, err
}
if req.ArchivePath == "" {
return ConvertResponse{}, fmt.Errorf("%w: archive path is required", ErrInvalidRequest)
}
if req.Interval != marketdata.D1 {
return ConvertResponse{}, fmt.Errorf("%w: stooq conversion supports only D1", ErrInvalidRequest)
}
defer func() {
s.logOutcome(ctx, slog.LevelInfo, "convert completed", "convert failed", req.DatasetRequest, err,
"rows_imported", resp.Import.RowsImported, "published_partitions", len(resp.Build.Result.Published))
}()
imported, err := s.manager.ImportStooqArchive(ctx, req.ArchivePath, req.Instrument)
if err != nil {
return ConvertResponse{}, err
}
build, err := s.Build(ctx, BuildRequest{DatasetRequest: req.DatasetRequest})
resp = ConvertResponse{Import: imported, Build: build}
return resp, err
}
59 changes: 59 additions & 0 deletions internal/service/marketdata/convert_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,59 @@
package marketdata_test

import (
"context"
"os"
"path/filepath"
"testing"
"time"

"github.com/stretchr/testify/require"

"github.com/rustyeddy/trader/instrument"
"github.com/rustyeddy/trader/internal/clock"
marketruntime "github.com/rustyeddy/trader/internal/marketdata"
svc "github.com/rustyeddy/trader/internal/service/marketdata"
"github.com/rustyeddy/trader/marketdata"
"github.com/rustyeddy/trader/num"
)

func TestConvertImportsAndBuildsStooqArchive(t *testing.T) {
resolver := instrument.NewMemoryResolver()
id, err := svc.RegisterETFInstrument(resolver, svc.EquityRegistration{
Provider: "stooq", Exchange: "ARCA", Ticker: "SPY", Currency: num.MustParseCurrency("USD"),
})
require.NoError(t, err)
rawRoot, storeRoot := t.TempDir(), t.TempDir()
manager, err := marketruntime.New(marketruntime.Config{
Clock: clock.NewSimulated(time.Date(2020, 3, 1, 0, 0, 0, 0, time.UTC)), StoreRoot: storeRoot,
RawRoot: rawRoot, Resolver: resolver, ProviderName: "stooq",
Calendar: marketdata.NewUSEquityCalendar(marketdata.StandardUSEquityHolidays(2020)),
})
require.NoError(t, err)
service, err := svc.New(manager, nil)
require.NoError(t, err)
path := filepath.Join(t.TempDir(), "spy.us.txt")
require.NoError(t, os.WriteFile(path, []byte(
"<TICKER>,<PER>,<DATE>,<TIME>,<OPEN>,<HIGH>,<LOW>,<CLOSE>,<VOL>,<OPENINT>\n"+
"SPY.US,D,20200131,000000,100,101,99,100.5,1000,0\n"), 0o644))
span, err := marketdata.NewTimeRange(time.Date(2020, 1, 1, 0, 0, 0, 0, time.UTC), time.Date(2020, 2, 1, 0, 0, 0, 0, time.UTC))
require.NoError(t, err)

resp, err := service.Convert(context.Background(), svc.ConvertRequest{
DatasetRequest: svc.DatasetRequest{Instrument: id, Interval: marketdata.D1, Range: span}, ArchivePath: path,
})
require.NoError(t, err)
require.Equal(t, 1, resp.Import.RowsImported)
require.Len(t, resp.Build.Result.Published, 1)
}

func TestConvertRejectsNonDailyIntervalBeforeImport(t *testing.T) {
service := newTestService(t)
span, err := marketdata.NewTimeRange(time.Date(2020, 1, 1, 0, 0, 0, 0, time.UTC), time.Date(2020, 2, 1, 0, 0, 0, 0, time.UTC))
require.NoError(t, err)
_, err = service.Convert(context.Background(), svc.ConvertRequest{
DatasetRequest: svc.DatasetRequest{Instrument: instrument.ETFID("ARCA", "SPY"), Interval: marketdata.H1, Range: span}, ArchivePath: "unused",
})
require.ErrorIs(t, err, svc.ErrInvalidRequest)
require.ErrorContains(t, err, "supports only D1")
}
7 changes: 7 additions & 0 deletions internal/service/marketdata/request.go
Original file line number Diff line number Diff line change
Expand Up @@ -111,3 +111,10 @@ type BuildRequest struct {
type UpdateRequest struct {
DatasetRequest
}

// ConvertRequest imports one provider-native archive and builds the requested
// canonical dataset from the resulting managed raw partitions.
type ConvertRequest struct {
DatasetRequest
ArchivePath string
}
Loading
Loading