diff --git a/cache/bbolt.go b/cache/bbolt.go index e7e9507..bd07833 100644 --- a/cache/bbolt.go +++ b/cache/bbolt.go @@ -339,8 +339,8 @@ func (s *BoltStorage) Stats() (DBStats, error) { return tx.ForEach(func(name []byte, b *bbolt.Bucket) error { bucketName := string(name) - if strings.HasPrefix(bucketName, "records:") { - did := strings.TrimPrefix(bucketName, "records:") + if after, ok := strings.CutPrefix(bucketName, "records:"); ok { + did := after count := 0 b.ForEach(func(k, v []byte) error { count++ @@ -357,8 +357,8 @@ func (s *BoltStorage) Stats() (DBStats, error) { stats.UserStats[did] = m } - if strings.HasPrefix(bucketName, "processed:") { - did := strings.TrimPrefix(bucketName, "processed:") + if after, ok := strings.CutPrefix(bucketName, "processed:"); ok { + did := after count := 0 b.ForEach(func(k, v []byte) error { count++ @@ -375,8 +375,8 @@ func (s *BoltStorage) Stats() (DBStats, error) { stats.UserStats[did] = m } - if strings.HasPrefix(bucketName, "failed:") { - did := strings.TrimPrefix(bucketName, "failed:") + if after, ok := strings.CutPrefix(bucketName, "failed:"); ok { + did := after count := 0 b.ForEach(func(k, v []byte) error { count++ diff --git a/flags.go b/flags.go index 0e4d821..6995e3e 100644 --- a/flags.go +++ b/flags.go @@ -1,16 +1,13 @@ package main import ( - "time" - "github.com/urfave/cli/v3" "tangled.org/karitham.dev/lazuli/sync" ) const ( - DefaultBatchSize = 20 - DefaultBatchDelay = 2000 * time.Millisecond + DefaultBatchSize = 20 ) const ( @@ -131,12 +128,6 @@ var importFlags = []cli.Flag{ Value: DefaultBatchSize, Sources: cli.EnvVars("LAZULI_BATCH_SIZE"), }, - &cli.IntFlag{ - Name: "batch-delay", - Usage: "MS between batches (default: 2000)", - Value: int(DefaultBatchDelay.Milliseconds()), - Sources: cli.EnvVars("LAZULI_BATCH_DELAY"), - }, &cli.DurationFlag{ Name: "tolerance", Usage: "Time tolerance for cross-source deduplication (e.g., 5m, 10m)", diff --git a/main.go b/main.go index 66bfdf2..a74da1d 100644 --- a/main.go +++ b/main.go @@ -505,7 +505,6 @@ func (a *App) runImport(ctx context.Context, cmd *cli.Command) error { fresh := cmd.Bool("fresh") clearCache := cmd.Bool("clear-cache") batchSize := int(cmd.Int("batch-size")) - batchDelay := int(cmd.Int("batch-delay")) tolerance := cmd.Duration("tolerance") if clearCache { @@ -592,13 +591,11 @@ func (a *App) runImport(ctx context.Context, cmd *cli.Command) error { cfg := sync.DefaultConfig cfg.BatchSize = batchSize - cfg.BatchDelay = time.Duration(batchDelay) * time.Millisecond progressLog := a.createProgressLogger() publishOpts := sync.PublishOptions{ BatchSize: cfg.BatchSize, - BatchDelay: cfg.BatchDelay, DryRun: dryRun, ATProtoClient: repoClient, ProgressLog: progressLog, @@ -777,46 +774,6 @@ func (a *App) runDedupe(ctx context.Context, cmd *cli.Command) error { return nil } -func (a *App) loadRecordsForImport(ctx context.Context, lastfmPath, spotifyPath string, mode sync.ImportMode, tolerance time.Duration) ([]sync.PlayRecord, int, error) { - var lastfmRecords, spotifyRecords []sync.PlayRecord - var err error - - if mode == sync.ImportModeLastFM || mode == sync.ImportModeCombined { - lastfmRecords, err = sync.ParseInput(ctx, lastfmPath, lastfm.Parser{}) - if err != nil { - return nil, 0, fmt.Errorf("parse lastfm: %w", err) - } - } - - if mode == sync.ImportModeSpotify || mode == sync.ImportModeCombined { - spotifyRecords, err = sync.ParseInput(ctx, spotifyPath, spotify.Parser{}) - if err != nil { - return nil, 0, fmt.Errorf("parse spotify: %w", err) - } - } - - totalInput := len(lastfmRecords) + len(spotifyRecords) - - var ( - mergedRecords []sync.PlayRecord - stats sync.MergeStats - ) - - switch mode { - case sync.ImportModeCombined: - mergedRecords, stats = sync.MergeRecords(lastfmRecords, spotifyRecords, tolerance) - a.log.Debug("Merged records", - slog.Int("merged_total", stats.MergedTotal), - slog.Int("duplicates_removed", stats.DuplicatesRemoved)) - case sync.ImportModeLastFM: - mergedRecords = lastfmRecords - default: - mergedRecords = spotifyRecords - } - - return mergedRecords, totalInput, nil -} - func (a *App) outputRecords(records []sync.PlayRecord, outputPath string) error { var output io.Writer = os.Stdout if outputPath != "" { diff --git a/sync/atproto_auth.go b/sync/atproto_auth.go index 2c07ca2..f9da0b6 100644 --- a/sync/atproto_auth.go +++ b/sync/atproto_auth.go @@ -77,7 +77,7 @@ func (a *FixedPasswordAuth) Refresh(ctx context.Context, c *http.Client, priorRe } defer resp.Body.Close() - if !(resp.StatusCode >= 200 && resp.StatusCode < 300) { + if resp.StatusCode < 200 || resp.StatusCode >= 300 { var eb atclient.ErrorBody if err := json.NewDecoder(resp.Body).Decode(&eb); err != nil { return &atclient.APIError{StatusCode: resp.StatusCode} diff --git a/sync/atproto_auth_test.go b/sync/atproto_auth_test.go index f61bc6b..6d55d27 100644 --- a/sync/atproto_auth_test.go +++ b/sync/atproto_auth_test.go @@ -79,9 +79,10 @@ func TestLibraryPasswordAuth_Refresh_Method(t *testing.T) { // We call the library's method directly _ = pa.Refresh(context.Background(), http.DefaultClient, "old-refresh") - if methodUsed == http.MethodGet { + switch methodUsed { + case http.MethodGet: t.Log("Confirmed: Library uses GET for refreshSession (Buggy)") - } else if methodUsed == http.MethodPost { + case http.MethodPost: t.Log("Library uses POST for refreshSession") } } diff --git a/sync/config.go b/sync/config.go index 11c3045..11d4546 100644 --- a/sync/config.go +++ b/sync/config.go @@ -7,8 +7,6 @@ import ( const ( RecordType = "fm.teal.alpha.feed.play" DefaultBatchSize = 20 - DefaultBatchDelay = 2000 * time.Millisecond - MinBatchDelay = 1000 * time.Millisecond DefaultCrossSourceTolerance = 5 * time.Minute CrossSourceTolerance = DefaultCrossSourceTolerance CacheTTL = 24 * time.Hour @@ -32,7 +30,6 @@ type Config struct { RecordType string `json:"recordType"` ClientAgent string `json:"clientAgent"` BatchSize int `json:"batchSize"` - BatchDelay time.Duration `json:"batchDelay"` CrossSourceTolerance time.Duration `json:"crossSourceTolerance"` CacheTTL time.Duration `json:"cacheTTL"` CacheVersion int `json:"cacheVersion"` @@ -45,7 +42,6 @@ var DefaultConfig = Config{ RecordType: RecordType, ClientAgent: ClientAgent, BatchSize: DefaultBatchSize, - BatchDelay: DefaultBatchDelay, CrossSourceTolerance: CrossSourceTolerance, CacheTTL: CacheTTL, CacheVersion: CacheVersion, diff --git a/sync/publish.go b/sync/publish.go index 9123991..befed71 100644 --- a/sync/publish.go +++ b/sync/publish.go @@ -18,7 +18,6 @@ import ( type PublishOptions struct { BatchSize int - BatchDelay time.Duration DryRun bool ATProtoClient ATProtoClient ProgressLog func(ProgressReport) @@ -29,7 +28,6 @@ func Publish(ctx context.Context, client AuthClient, opts PublishOptions, limite startTime := time.Now() batchSize := defaultBatchSize(opts.BatchSize) - batchDelay := defaultBatchDelay(opts.BatchDelay) atprotoClient, err := buildClient(client, opts.ATProtoClient) if err != nil { @@ -57,7 +55,6 @@ func Publish(ctx context.Context, client AuthClient, opts PublishOptions, limite slog.Info("starting iterative import", slog.Int("total_records", totalRecords), slog.Int("batch_size", batchSize), - slog.Duration("batch_delay", batchDelay), slog.Int("daily_write_limit", WriteLimitDay), slog.Int("daily_token_limit", GlobalLimitDay), slog.String("rate_limit", fmt.Sprintf("1 write per %.1fs", 86400.0/WriteLimitDay))) @@ -243,13 +240,6 @@ func defaultBatchSize(size int) int { return DefaultBatchSize } -func defaultBatchDelay(delay time.Duration) time.Duration { - if delay > 0 { - return delay - } - return DefaultBatchDelay -} - func buildClient(client AuthClient, customClient ATProtoClient) (ATProtoClient, error) { if customClient != nil { return customClient, nil @@ -275,10 +265,6 @@ func newPublishResult(success, errors, total int, start time.Time, cancelled boo } } -func makeRecordKeys(records []PlayRecord) []string { - return CreateRecordKeys(records) -} - func logResult(success, errors int, startTime time.Time) { if errors > 0 { slog.Warn("import completed with errors",