diff --git a/cmd/trader/data/args.go b/cmd/trader/data/args.go index d78fe76..8eaec08 100644 --- a/cmd/trader/data/args.go +++ b/cmd/trader/data/args.go @@ -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 diff --git a/cmd/trader/data/command.go b/cmd/trader/data/command.go index 0d4390f..a12d6b7 100644 --- a/cmd/trader/data/command.go +++ b/cmd/trader/data/command.go @@ -63,6 +63,7 @@ func New() *cobra.Command { cmd.AddCommand(newSyncCmd()) cmd.AddCommand(newBuildCmd()) cmd.AddCommand(newUpdateCmd()) + cmd.AddCommand(newConvertCmd()) return cmd } diff --git a/cmd/trader/data/convert.go b/cmd/trader/data/convert.go new file mode 100644 index 0000000..762d03f --- /dev/null +++ b/cmd/trader/data/convert.go @@ -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 { + 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) +} diff --git a/cmd/trader/data/convert_integration_test.go b/cmd/trader/data/convert_integration_test.go new file mode 100644 index 0000000..0716185 --- /dev/null +++ b/cmd/trader/data/convert_integration_test.go @@ -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(",,,