From 4d488000e08450a35db5b687440e7181ad7b4fcf Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?L=C3=ADvia?= Date: Sun, 12 Apr 2026 12:43:24 -0300 Subject: [PATCH] feat: initial cli pipeline logic --- cli/pipeline.ts | 173 ++++++++++++++++++++++++++++++++++++++++++++++++ lib/csv.ts | 8 +-- utils/retry.ts | 13 ++++ 3 files changed, 190 insertions(+), 4 deletions(-) create mode 100644 cli/pipeline.ts diff --git a/cli/pipeline.ts b/cli/pipeline.ts new file mode 100644 index 0000000..6916a8a --- /dev/null +++ b/cli/pipeline.ts @@ -0,0 +1,173 @@ +import { type Config } from "~/lib/config.ts"; +import { countScrobblesInRange, scrobble, ScrobblePayload, ScrobbleResult } from "~/api/lastfm.ts"; +import { CsvDocument, DocumentTrack, markSkipped } from "~/lib/csv.ts"; +import { AppError, describe } from "~/lib/errors.ts"; +import { sleep, withRetryR } from "~/utils/retry.ts"; + +const DAILY_SCROBBLE_LIMIT: number = 2880; + +export interface PipelineOptions { + config: Omit, "password">; + sessionKey: string; + dryRun?: boolean; + delayMs?: number; +} + +export interface PipelineSummary { + total: number; + accepted: number; + ignored: number; + failed: number; + skipped: number; +} + +export interface PendingEntry { + track: DocumentTrack; + index: number; + payload: ScrobblePayload; +} + +export interface LimitCheckResult { + scrobblesToday: number; + remaining: number; + importCount: number; + wouldExceed: boolean; + excess: number; +} + +export async function checkDailyLimit( + apiKey: string, + username: string, + importCount: number, +): Promise { + const begin = Math.floor( + Date.UTC( + new Date().getUTCFullYear(), + new Date().getUTCMonth(), + new Date().getUTCDate(), + ) / 1000, + ); + const now = Math.floor(Date.now() / 1e3); + + const result = await countScrobblesInRange(apiKey, username, begin, now); + const scrobblesToday = result.ok ? result.value : 0; + + const remaining = Math.max( + 0, + DAILY_SCROBBLE_LIMIT - scrobblesToday, + ); + const wouldExceed = importCount > remaining; + const excess = wouldExceed ? importCount - remaining : 0; + + return { scrobblesToday, remaining, importCount, wouldExceed, excess }; +} + +/** + * determines if a Last.fm error is worth retrying. + * + * transient: rate limit (29), service offline (11), operation failed (8) + * fatal: bad auth (4, 9, 10), invalid params (6), suspended key (26) + */ +export function isRetryable(error: unknown): boolean { + if (typeof error === "object" && error !== null) { + const e = error as AppError; + if (e.kind === "network") return true; + if (e.kind === "lastfm") { + return [8, 11, 16, 29].includes(e.tag); + } + } + return false; +} + +export async function runPipeline( + pending: PendingEntry[], + file: CsvDocument, + opts: PipelineOptions, + onProgress: ( + current: number, + total: number, + track: DocumentTrack, + result: "ok" | "ignored" | "failed", + detail?: string, + ) => void, +): Promise { + const summary: PipelineSummary = { + total: pending.length, + accepted: 0, + ignored: 0, + failed: 0, + skipped: 0, + }; + + const delay = opts.delayMs ?? 100; + let currentFile = file; + + for (const [i, entry] of pending.entries()) { + const { track, index, payload } = entry; + + if (opts.dryRun) { + summary.accepted++; + onProgress(i + 1, pending.length, track, "ok"); + continue; + } + + try { + const result: ScrobbleResult = await withRetryR( + () => scrobble(opts.config.apiKey, opts.config.secret, opts.sessionKey, payload), + { + maxAttempts: 4, + baseDelayMs: 500, + retryIf: isRetryable, + onRetry: (attempt, delayMs) => { + console.error(` ↺ Retry ${attempt} for "${track.title}" in ${delayMs}ms...`); + }, + }, + ); + + if (result.ignored) { + summary.ignored += result.ignored; + onProgress(i + 1, pending.length, track, "ignored"); + } else { + summary.accepted += result.accepted; + onProgress(i + 1, pending.length, track, "ok"); + } + + const marked = markSkipped(currentFile, index); + if (marked.ok) currentFile = marked.value; + } catch (e) { + summary.failed++; + const msg = describe(e as AppError); + onProgress(i + 1, pending.length, track, "failed", msg); + } + + if (i < pending.length - 1) await sleep(delay); + } + + return summary; +} + +/** + * generate evenly-spaced timestamps for tracks that have no date. + * starts 13.9 days before now and steps 30s per track + */ +export function generateTimestamps(count: number): number[] { + const begin = Math.floor(Date.now() / 1_000) - 1_200_960; + return Array.from({ length: count }, (_, i) => begin + (i + 1) * 30); +} + +export function trackToPipelineEntry( + track: DocumentTrack, + index: number, + timestamp: number, +): PendingEntry { + return { + track, + index, + payload: { + artist: track.artist, + album: track.album || undefined, + title: track.title, + timestamp, + }, + }; +} diff --git a/lib/csv.ts b/lib/csv.ts index bc89780..dc20f93 100644 --- a/lib/csv.ts +++ b/lib/csv.ts @@ -5,7 +5,7 @@ export const SKIP_PREFIX = "#SKIP:"; const COLUMNS = ["artist", "album", "title", "date"] as const; -export interface Track { +export interface DocumentTrack { artist: string; album: string; title: string; @@ -15,7 +15,7 @@ export interface Track { export interface CsvDocument { path: string; rawLines: string[]; - pending: Array<{ track: Track; lineIndex: number }>; + pending: Array<{ track: DocumentTrack; lineIndex: number }>; skippedCount: number; } @@ -83,7 +83,7 @@ export function markSkipped( return Ok({ ...doc, rawLines }); } -export function saveCsvDocument(path: string, tracks: Track[]): Result { +export function saveCsvDocument(path: string, tracks: DocumentTrack[]): Result { const header = COLUMNS.join(","); const rows = tracks.map((t) => [escape(t.artist), escape(t.album), escape(t.title), escape(t.date)].join(",")); @@ -99,7 +99,7 @@ function parse( line: string, index: number, path: string, -): Result { +): Result { const fields = split(line); if (fields.length < 3) { diff --git a/utils/retry.ts b/utils/retry.ts index e86a577..0b00d42 100644 --- a/utils/retry.ts +++ b/utils/retry.ts @@ -1,3 +1,5 @@ +import { Result } from "~/lib/result.ts"; + export interface RetryOptions { maxAttempts?: number; baseDelayMs?: number; @@ -19,6 +21,17 @@ function jitteredDelay(attempt: number, base: number, cap: number): number { return Math.floor(Math.random() * ceiling); } +export function withRetryR( + fn: () => Promise>, + opts: RetryOptions = {}, +): Promise { + return withRetry(() => + fn().then((result) => { + if (!result.ok) throw result.error; + return result.value; + }), opts); +} + export async function withRetry( fn: () => Promise, opts: RetryOptions = {}, -- 2.51.2