diff --git a/package.json b/package.json index 266a52e..7a9ccfb 100644 --- a/package.json +++ b/package.json @@ -14,7 +14,10 @@ "dev": "tsc && node dist/index.js", "dry-run": "npm run build && node dist/index.js --dry-run", "clean": "rm -rf dist", - "type-check": "tsc --noEmit" + "type-check": "tsc --noEmit", + "test": "npm run build && node --test dist/tests/**/*.test.js", + "test:tid": "npm run build && node --test dist/tests/tid.test.js", + "test:watch": "npm run build && node --watch --test dist/tests/**/*.test.js" }, "keywords": [ "lastfm", diff --git a/src/config.ts b/src/config.ts index 2b0c3fb..f6fd93d 100644 --- a/src/config.ts +++ b/src/config.ts @@ -6,12 +6,15 @@ import type { Config } from './types.js'; // - This affects all users on your PDS, not just your account // - See: https://docs.bsky.app/blog/rate-limits-pds-v3 // -// Default limit: Aggressive initial limit that will dynamically adjust -// Start high and back off if we hit rate limits +// Default limit: Very conservative (7,500 records/day) to be safe export const RECORDS_PER_DAY_LIMIT = 10000; -// Safety margin factor - start aggressive, will back off if needed -export const SAFETY_MARGIN = 1.0; +// Safety margin factor - 75% by default for maximum safety +// Use --aggressive flag to set to 85% for faster imports +export const SAFETY_MARGIN = 0.75; + +// Aggressive safety margin (for --aggressive flag) +export const AGGRESSIVE_SAFETY_MARGIN = 0.85; // Record type export const RECORD_TYPE = 'fm.teal.alpha.feed.play'; @@ -32,13 +35,13 @@ export function buildClientAgent() { const CLIENT_AGENT = buildClientAgent(); -// Default batch configuration - aggressive defaults for maximum speed +// Default batch configuration - conservative for PDS safety // Will dynamically adjust based on success/failure -export const DEFAULT_BATCH_SIZE = 200; // Max allowed by applyWrites -export const DEFAULT_BATCH_DELAY = 500; // Start with 500ms between batches +export const DEFAULT_BATCH_SIZE = 100; // Conservative default +export const DEFAULT_BATCH_DELAY = 2000; // Start with 2 seconds between batches -// Minimum safe delay between batches (500ms for adaptive mode) -export const MIN_BATCH_DELAY = 500; +// Minimum safe delay between batches (1 second minimum) +export const MIN_BATCH_DELAY = 1000; // Maximum batch size (PDS limit is 200 operations per call) export const MAX_BATCH_SIZE = 200; @@ -59,6 +62,7 @@ const config: Config = { SLINGSHOT_RESOLVER, RECORDS_PER_DAY_LIMIT, SAFETY_MARGIN, + AGGRESSIVE_SAFETY_MARGIN, }; export default config; diff --git a/src/lib/cli.ts b/src/lib/cli.ts index 53db410..55ba9ed 100644 --- a/src/lib/cli.ts +++ b/src/lib/cli.ts @@ -13,6 +13,13 @@ import config from '../config.js'; import { calculateOptimalBatchSize } from '../utils/helpers.js'; import { fetchExistingRecords, filterNewRecords, displaySyncStats, removeDuplicates } from './sync.js'; import { Logger, LogLevel, setGlobalLogger, log } from '../utils/logger.js'; +import { + loadImportState, + createImportState, + displayResumeInfo, + clearImportState, + ImportState, +} from '../utils/import-state.js'; /** * Show help message @@ -42,13 +49,15 @@ ${'\x1b[1m'}MODE:${'\x1b[0m'} deduplicate Remove duplicate records ${'\x1b[1m'}BATCH CONFIGURATION:${'\x1b[0m'} - -b, --batch-size Records per batch (default: auto-calculated) - -d, --batch-delay Delay between batches in ms (default: 500, min: 500) + -b, --batch-size Records per batch (default: 100) + -d, --batch-delay Delay between batches in ms (default: 2000ms, min: 1000ms) ${'\x1b[1m'}IMPORT OPTIONS:${'\x1b[0m'} -r, --reverse Process newest records first (default: oldest first) -y, --yes Skip confirmation prompts --dry-run Preview without importing + --aggressive Faster imports (8,500/day vs 7,500/day default) + --fresh Start fresh (ignore previous import state) ${'\x1b[1m'}OUTPUT:${'\x1b[0m'} -v, --verbose Enable verbose logging (debug level) @@ -114,6 +123,8 @@ export function parseCommandLineArgs(): CommandLineArgs { reverse: { type: 'boolean', short: 'r', default: false }, yes: { type: 'boolean', short: 'y', default: false }, 'dry-run': { type: 'boolean', default: false }, + aggressive: { type: 'boolean', default: false }, + fresh: { type: 'boolean', default: false }, // Output verbose: { type: 'boolean', short: 'v', default: false }, @@ -145,6 +156,8 @@ export function parseCommandLineArgs(): CommandLineArgs { reverse: values.reverse || values['reverse-chronological'], yes: values.yes, 'dry-run': values['dry-run'], + aggressive: values.aggressive, + fresh: values.fresh, verbose: values.verbose, quiet: values.quiet, }; @@ -351,13 +364,19 @@ export async function runCLI(): Promise { log.info(`Batch delay: ${batchDelay}ms`); + // Apply aggressive mode if enabled + const safetyMargin = args.aggressive ? cfg.AGGRESSIVE_SAFETY_MARGIN : cfg.SAFETY_MARGIN; + if (args.aggressive) { + log.warn('⚡ Aggressive mode enabled: Using 85% of daily limit (8,500 records/day)'); + } + // Show rate limiting information log.section('Import Configuration'); log.info(`Total records: ${totalRecords.toLocaleString()}`); log.info(`Batch size: ${batchSize} records`); log.info(`Batch delay: ${batchDelay}ms`); - const recordsPerDay = cfg.RECORDS_PER_DAY_LIMIT * cfg.SAFETY_MARGIN; + const recordsPerDay = cfg.RECORDS_PER_DAY_LIMIT * safetyMargin; const estimatedDays = Math.ceil(totalRecords / recordsPerDay); if (estimatedDays > 1) { @@ -367,6 +386,46 @@ export async function runCLI(): Promise { log.blank(); + // Check for existing import state (resume functionality) + let importState: ImportState | null = null; + if (!dryRun && args.input) { + // Clear state if --fresh flag is used + if (args.fresh) { + clearImportState(args.input, mode); + log.info('Starting fresh import (previous state cleared)'); + } else { + // Try to load existing state + importState = loadImportState(args.input, mode); + + if (importState && !importState.completed) { + displayResumeInfo(importState); + + if (!args.yes) { + const answer = await prompt('Resume from previous import? (Y/n) '); + if (answer.toLowerCase() === 'n') { + importState = null; + clearImportState(args.input, mode); + log.info('Starting fresh import'); + log.blank(); + } + } else { + log.info('Auto-resuming previous import (--yes flag)'); + log.blank(); + } + } else if (importState?.completed) { + log.info('Previous import was completed - starting fresh'); + importState = null; + clearImportState(args.input, mode); + } + } + + // Create new state if not resuming + if (!importState) { + importState = createImportState(args.input, mode, totalRecords); + log.debug('Created new import state'); + } + } + // Confirmation prompt if (!dryRun && !args.yes) { const modeLabel = mode === 'combined' ? 'merged' : mode === 'sync' ? 'new' : ''; @@ -389,7 +448,8 @@ export async function runCLI(): Promise { batchDelay, cfg, dryRun, - mode === 'sync' || mode === 'combined' + mode === 'sync' || mode === 'combined', + importState ); // Final output diff --git a/src/lib/publisher.ts b/src/lib/publisher.ts index bdcbcf0..58f1ece 100644 --- a/src/lib/publisher.ts +++ b/src/lib/publisher.ts @@ -10,6 +10,12 @@ import { import { generateTIDFromISO } from '../utils/tid.js'; import type { PlayRecord, Config, PublishResult } from '../types.js'; import { log } from '../utils/logger.js'; +import { + ImportState, + updateImportState, + completeImport, + getResumeStartIndex, +} from '../utils/import-state.js'; /** * Maximum operations allowed per applyWrites call @@ -21,7 +27,7 @@ const MAX_APPLY_WRITES_OPS = 200; /** * Publish records using com.atproto.repo.applyWrites for efficient batching - * with adaptive rate limiting + * with adaptive rate limiting and stateful resume support */ export async function publishRecordsWithApplyWrites( agent: AtpAgent | null, @@ -30,7 +36,8 @@ export async function publishRecordsWithApplyWrites( batchDelay: number, config: Config, dryRun = false, - syncMode = false + syncMode = false, + importState: ImportState | null = null ): Promise { const { RECORD_TYPE } = config; const totalRecords = records.length; @@ -52,10 +59,11 @@ export async function publishRecordsWithApplyWrites( let consecutiveFailures = 0; const MAX_CONSECUTIVE_FAILURES = 3; - log.section('Adaptive Import'); - log.info(`Initial batch size: ${currentBatchSize} records`); - log.info(`Initial delay: ${currentBatchDelay}ms`); + log.section('Conservative Adaptive Import'); + log.info(`Initial batch size: ${currentBatchSize} records (conservative)`); + log.info(`Initial delay: ${currentBatchDelay}ms (2 seconds - very safe)`); log.info(`Will automatically adjust based on server response`); + log.info(`Using conservative settings to protect your PDS`); log.blank(); log.info(`Publishing ${totalRecords.toLocaleString()} records using adaptive batching...`); log.warn('Press Ctrl+C to stop gracefully after current batch'); @@ -65,7 +73,14 @@ export async function publishRecordsWithApplyWrites( let errorCount = 0; const startTime = Date.now(); - let i = 0; + // Resume from saved state if available + let startIndex = importState ? getResumeStartIndex(importState) : 0; + if (importState && startIndex > 0) { + log.info(`Resuming from record ${startIndex + 1} (${(startIndex / totalRecords * 100).toFixed(1)}% complete)`); + log.blank(); + } + + let i = startIndex; while (i < totalRecords) { // Check killswitch before processing batch if (isImportCancelled()) { @@ -83,12 +98,14 @@ export async function publishRecordsWithApplyWrites( const batchStartTime = Date.now(); // Build writes array for applyWrites with TID-based rkeys - const writes = batch.map((record) => ({ - $type: 'com.atproto.repo.applyWrites#create', - collection: RECORD_TYPE, - rkey: generateTIDFromISO(record.playedTime), - value: record, - })); + const writes = await Promise.all( + batch.map(async (record) => ({ + $type: 'com.atproto.repo.applyWrites#create', + collection: RECORD_TYPE, + rkey: await generateTIDFromISO(record.playedTime, 'inject:playlist'), + value: record, + })) + ); try { // Call applyWrites with the batch @@ -106,6 +123,11 @@ export async function publishRecordsWithApplyWrites( const batchDuration = Date.now() - batchStartTime; log.debug(`Batch complete in ${batchDuration}ms (${batchSuccessCount} successful)`); + // Save state after successful batch + if (importState) { + updateImportState(importState, i + batch.length - 1, batchSuccessCount, 0); + } + // Speed up if we're doing well (after 5 consecutive successes) if (consecutiveSuccesses >= 5 && currentBatchDelay > config.MIN_BATCH_DELAY) { const oldDelay = currentBatchDelay; @@ -168,6 +190,11 @@ export async function publishRecordsWithApplyWrites( log.debug(`... and ${batch.length - 3} more failed`); } + // Save state with errors + if (importState) { + updateImportState(importState, i + batch.length - 1, 0, batch.length); + } + // If too many consecutive failures, slow down if (consecutiveFailures >= MAX_CONSECUTIVE_FAILURES) { currentBatchDelay = Math.min(currentBatchDelay * 2, 10000); @@ -200,6 +227,12 @@ export async function publishRecordsWithApplyWrites( } } + // Mark import as complete + if (importState) { + completeImport(importState); + log.debug('Import state saved as completed'); + } + return { successCount, errorCount, cancelled: false }; } diff --git a/src/lib/sync.ts b/src/lib/sync.ts index f31ce05..693884f 100644 --- a/src/lib/sync.ts +++ b/src/lib/sync.ts @@ -46,7 +46,7 @@ export async function fetchExistingRecords( }); for (const record of response.data.records) { - const playRecord = record.value as PlayRecord; + const playRecord = record.value as unknown as PlayRecord; // Create a unique key based on track, artist, and timestamp const key = createRecordKey(playRecord); // Note: This will overwrite duplicates, but that's OK for sync mode @@ -108,7 +108,7 @@ export async function fetchAllRecords( }); for (const record of response.data.records) { - const playRecord = record.value as PlayRecord; + const playRecord = record.value as unknown as PlayRecord; allRecords.push({ uri: record.uri, cid: record.cid, diff --git a/src/tests/tid-integration.test.ts b/src/tests/tid-integration.test.ts new file mode 100644 index 0000000..7e2d09e --- /dev/null +++ b/src/tests/tid-integration.test.ts @@ -0,0 +1,195 @@ +/** + * Integration test for TID generation in the importer + * + * Tests the full flow with actual CSV data + */ + +import { describe, it } from 'node:test'; +import assert from 'node:assert'; +import { initTidClock, generateTIDFromISO, resetTidClock, getTidClockState } from '../utils/tid.js'; +import { validateTid, areMonotonic } from '../utils/tid-clock.js'; +import { auditTids, formatAuditReport } from '../utils/tid-audit.js'; + +describe('TID Integration - Importer Flow', () => { + it('should initialize TID clock for production', async () => { + resetTidClock(); + initTidClock({ mode: 'production' }); + + const tid = await generateTIDFromISO('2020-01-01T00:00:00Z', 'test'); + assert.strictEqual(validateTid(tid), true); + + const state = getTidClockState(); + assert.strictEqual(state.generatedCount, 1); + }); + + it('should initialize TID clock for dry-run with deterministic output', async () => { + const seed = 1000000000000000; + resetTidClock(); + initTidClock({ mode: 'dry-run', seed }); + + const dates = [ + '2005-01-01T00:00:00Z', + '2010-01-01T00:00:00Z', + '2015-01-01T00:00:00Z', + ]; + + const tids1 = await Promise.all(dates.map(d => generateTIDFromISO(d, 'test'))); + + // Reset and regenerate + resetTidClock(); + initTidClock({ mode: 'dry-run', seed }); + + const tids2 = await Promise.all(dates.map(d => generateTIDFromISO(d, 'test'))); + + // Should be identical + assert.deepStrictEqual(tids1, tids2); + }); + + it('should handle batch TID generation', async () => { + resetTidClock(); + initTidClock({ mode: 'production' }); + + // Simulate batch processing + const batchSize = 200; + const records = Array.from({ length: batchSize }, (_, i) => ({ + playedTime: new Date(2020, 0, 1, 0, 0, i).toISOString(), + })); + + const tids = await Promise.all( + records.map(r => generateTIDFromISO(r.playedTime, 'batch:test')) + ); + + // All should be valid + for (const tid of tids) { + assert.strictEqual(validateTid(tid), true); + } + + // All should be unique + const uniqueTids = new Set(tids); + assert.strictEqual(uniqueTids.size, batchSize); + + // Should be monotonic + assert.strictEqual(areMonotonic(tids), true); + + // State should track count + const state = getTidClockState(); + assert.strictEqual(state.generatedCount, batchSize); + }); + + it('should handle historical dates from Last.fm (2005-present)', async () => { + resetTidClock(); + initTidClock({ mode: 'production' }); + + // Simulate Last.fm scrobbles from different eras + const historicalDates = [ + '2005-03-15T14:30:00Z', // Very old + '2010-06-20T08:15:00Z', + '2015-11-05T19:45:00Z', + '2020-01-01T00:00:00Z', + '2024-12-25T12:00:00Z', // Recent + ]; + + const tids = await Promise.all( + historicalDates.map(d => generateTIDFromISO(d, 'historical')) + ); + + // All valid + for (const tid of tids) { + assert.strictEqual(validateTid(tid), true); + } + + // All unique + const uniqueTids = new Set(tids); + assert.strictEqual(uniqueTids.size, historicalDates.length); + + // Monotonic despite being historical + assert.strictEqual(areMonotonic(tids), true); + }); + + it('should audit a batch of TIDs', async () => { + resetTidClock(); + initTidClock({ mode: 'production' }); + + const tids = await Promise.all( + Array.from({ length: 100 }, (_, i) => + generateTIDFromISO(new Date(2020, 0, 1, 0, 0, i).toISOString(), 'audit-test') + ) + ); + + const report = auditTids(tids); + + assert.strictEqual(report.totalTids, 100); + assert.strictEqual(report.validTids, 100); + assert.strictEqual(report.invalidTids, 0); + assert.strictEqual(report.duplicates, 0); + assert.strictEqual(report.monotonic, true); + + // Generate human-readable report + const textReport = formatAuditReport(report); + assert.ok(textReport.includes('TID AUDIT REPORT')); + assert.ok(textReport.includes('✓ YES')); // Monotonic check + }); + + it('should detect problems in TID audit', async () => { + // Create a bad TID list + const badTids = [ + '3jzfcijpj2z2a', // Valid + '3jzfcijpj2z2b', // Valid + '3jzfcijpj2z2a', // Duplicate! + '0000000000000', // Invalid format + '3jzfcijpj2z2c', // Valid + '3jzfcijpj2z2b', // Out of order + ]; + + const report = auditTids(badTids); + + assert.strictEqual(report.totalTids, 6); + assert.strictEqual(report.duplicates, 2); + assert.strictEqual(report.monotonic, false); + assert.ok(report.invalidTids > 0); + assert.ok(report.errors.length > 0); + + const textReport = formatAuditReport(report); + assert.ok(textReport.includes('✗ NO')); // Monotonic check + assert.ok(textReport.includes('ERRORS')); + }); +}); + +describe('TID Integration - Concurrent Batches', () => { + it('should handle multiple concurrent batches safely', async () => { + resetTidClock(); + initTidClock({ mode: 'production' }); + + // Simulate 5 concurrent batches of 50 records each + const batches = 5; + const batchSize = 50; + + const allTids = await Promise.all( + Array.from({ length: batches }, async (_, batchIdx) => { + return await Promise.all( + Array.from({ length: batchSize }, (_, recordIdx) => { + const date = new Date(2020, 0, 1, batchIdx, recordIdx).toISOString(); + return generateTIDFromISO(date, `batch-${batchIdx}`); + }) + ); + }) + ); + + const flatTids = allTids.flat(); + + // All unique + const uniqueTids = new Set(flatTids); + assert.strictEqual(uniqueTids.size, batches * batchSize); + + // All valid + for (const tid of flatTids) { + assert.strictEqual(validateTid(tid), true); + } + + // Total count correct + const state = getTidClockState(); + assert.strictEqual(state.generatedCount, batches * batchSize); + }); +}); + +console.log('✓ All integration tests completed'); diff --git a/src/tests/tid.test.ts b/src/tests/tid.test.ts new file mode 100644 index 0000000..4a1ec72 --- /dev/null +++ b/src/tests/tid.test.ts @@ -0,0 +1,473 @@ +/** + * Unit tests for TID generation + * + * Tests cover: + * - Format validation + * - Monotonicity (single-threaded and concurrent) + * - Deterministic dry-run mode + * - Clock drift handling + * - State persistence + * - Collision detection + */ + +import { describe, it, before, after } from 'node:test'; +import assert from 'node:assert'; +import fs from 'fs'; +import path from 'path'; +import os from 'os'; +import { + TidClock, + RealClock, + FakeClock, + validateTid, + ensureValidTid, + decodeTidTimestamp, + decodeTidClockId, + compareTids, + areMonotonic, + InvalidTidError, + SilentTidLogger, +} from '../utils/tid-clock.js'; + +const TID_LENGTH = 13; + +describe('TID Format Validation', () => { + it('should validate correct TID format', () => { + const validTids = [ + '3jzfcijpj2z2a', + '7777777777777', + '3zzzzzzzzzzzz', + ]; + + for (const tid of validTids) { + assert.strictEqual(validateTid(tid), true, `Should validate: ${tid}`); + } + }); + + it('should reject invalid TID format', () => { + const invalidTids = [ + '3jzfcijpj2z21', // Invalid character (1) + '0000000000000', // Invalid character (0) + '3jzfcijpj2z2aa', // Too long + '3jzfcijpj2z2', // Too short + '3jzf-cij-pj2z-2a', // Dashes not allowed + 'zzzzzzzzzzzzz', // High bit violation + 'kjzfcijpj2z2a', // High bit violation + ]; + + for (const tid of invalidTids) { + assert.strictEqual(validateTid(tid), false, `Should reject: ${tid}`); + } + }); + + it('should enforce TID length', () => { + assert.strictEqual(validateTid('123'), false); + assert.strictEqual(validateTid('12345678901234567890'), false); + }); + + it('should enforce base32 alphabet', () => { + assert.strictEqual(validateTid('3jzfcijpj2z21'), false); // Has '1' + assert.strictEqual(validateTid('0jzfcijpj2z2a'), false); // Starts with '0' + }); + + it('should throw on ensureValidTid for invalid input', () => { + assert.throws( + () => ensureValidTid('invalid'), + InvalidTidError + ); + }); +}); + +describe('TID Generation - Basic', () => { + it('should generate valid TID format', async () => { + const clock = new TidClock(new RealClock(), new SilentTidLogger()); + const tid = await clock.next(); + + assert.strictEqual(tid.length, TID_LENGTH); + assert.strictEqual(validateTid(tid), true); + }); + + it('should generate TID from Date', async () => { + const clock = new TidClock(new FakeClock(1000000000000000), new SilentTidLogger()); + const date = new Date('2005-01-01T00:00:00Z'); + const tid = await clock.fromDate(date); + + assert.strictEqual(validateTid(tid), true); + assert.strictEqual(tid.length, TID_LENGTH); + }); + + it('should encode timestamp correctly', async () => { + const timestamp = 1000000000000000; // Fixed timestamp + const clock = new TidClock(new FakeClock(timestamp), new SilentTidLogger()); + const tid = await clock.next(); + + const decoded = decodeTidTimestamp(tid); + assert.strictEqual(decoded, timestamp); + }); + + it('should encode clock ID correctly', async () => { + const clock = new TidClock( + new FakeClock(1000000000000000), + new SilentTidLogger(), + { clockId: 15 } + ); + const tid = await clock.next(); + + const clockId = decodeTidClockId(tid); + assert.strictEqual(clockId, 15); + }); +}); + +describe('TID Monotonicity - Single Thread', () => { + it('should generate monotonically increasing TIDs', async () => { + const fakeClock = new FakeClock(1000000000000000); + const clock = new TidClock(fakeClock, new SilentTidLogger()); + + const tids: string[] = []; + for (let i = 0; i < 100; i++) { + fakeClock.advance(1000); // Advance 1ms + tids.push(await clock.next()); + } + + assert.strictEqual(areMonotonic(tids), true); + }); + + it('should handle same timestamp with sequence increment', async () => { + const fakeClock = new FakeClock(1000000000000000); + const clock = new TidClock(fakeClock, new SilentTidLogger()); + + // Generate multiple TIDs at same timestamp + const tid1 = await clock.next(); + const tid2 = await clock.next(); + const tid3 = await clock.next(); + + assert.notStrictEqual(tid1, tid2); + assert.notStrictEqual(tid2, tid3); + assert.strictEqual(compareTids(tid1, tid2), -1); + assert.strictEqual(compareTids(tid2, tid3), -1); + }); + + it('should handle backwards clock drift', async () => { + const fakeClock = new FakeClock(1000000000000000); + const clock = new TidClock(fakeClock, new SilentTidLogger()); + + const tid1 = await clock.next(); + + // Move clock backwards + fakeClock.set(999999000000000); + + const tid2 = await clock.next(); + + // Should still be monotonic + assert.strictEqual(compareTids(tid1, tid2), -1); + }); + + it('should generate unique TIDs even with clock drift', async () => { + const fakeClock = new FakeClock(1000000000000000); + const clock = new TidClock(fakeClock, new SilentTidLogger()); + + const tids: string[] = []; + + // Generate some TIDs + for (let i = 0; i < 5; i++) { + tids.push(await clock.next()); + } + + // Move clock backwards + fakeClock.set(999999000000000); + + // Generate more TIDs + for (let i = 0; i < 5; i++) { + tids.push(await clock.next()); + } + + // All should be unique and monotonic + const uniqueTids = new Set(tids); + assert.strictEqual(uniqueTids.size, tids.length); + assert.strictEqual(areMonotonic(tids), true); + }); +}); + +describe('TID Monotonicity - Concurrent', () => { + it('should handle concurrent generation safely', async () => { + const clock = new TidClock(new RealClock(), new SilentTidLogger()); + + // Generate 100 TIDs concurrently + const promises = Array.from({ length: 100 }, () => clock.next()); + const tids = await Promise.all(promises); + + // All should be unique + const uniqueTids = new Set(tids); + assert.strictEqual(uniqueTids.size, tids.length); + + // All should be valid + for (const tid of tids) { + assert.strictEqual(validateTid(tid), true); + } + + // Should be monotonic when sorted + const sorted = [...tids].sort(compareTids); + assert.strictEqual(areMonotonic(sorted), true); + }); + + it('should handle high-frequency concurrent generation', async () => { + const clock = new TidClock(new RealClock(), new SilentTidLogger()); + + // Generate 1000 TIDs in parallel batches + const batchSize = 50; + const batches = 20; + const allTids: string[] = []; + + for (let b = 0; b < batches; b++) { + const promises = Array.from({ length: batchSize }, () => clock.next()); + const batchTids = await Promise.all(promises); + allTids.push(...batchTids); + } + + // All should be unique + const uniqueTids = new Set(allTids); + assert.strictEqual(uniqueTids.size, allTids.length); + + // All should be valid + for (const tid of allTids) { + assert.strictEqual(validateTid(tid), true); + } + }); +}); + +describe('TID Deterministic Mode', () => { + it('should generate deterministic TIDs with same seed', async () => { + const seed = 1000000000000000; + + const clock1 = new TidClock(new FakeClock(seed), new SilentTidLogger(), { clockId: 10 }); + const clock2 = new TidClock(new FakeClock(seed), new SilentTidLogger(), { clockId: 10 }); + + const tids1: string[] = []; + const tids2: string[] = []; + + for (let i = 0; i < 10; i++) { + tids1.push(await clock1.next()); + tids2.push(await clock2.next()); + } + + // Should generate identical sequences + assert.deepStrictEqual(tids1, tids2); + }); + + it('should be deterministic for historical dates', async () => { + const dates = [ + new Date('2005-01-01T00:00:00Z'), + new Date('2010-06-15T12:30:00Z'), + new Date('2020-12-31T23:59:59Z'), + ]; + + const clock1 = new TidClock(new FakeClock(0), new SilentTidLogger(), { clockId: 5 }); + const clock2 = new TidClock(new FakeClock(0), new SilentTidLogger(), { clockId: 5 }); + + const tids1 = await Promise.all(dates.map(d => clock1.fromDate(d))); + const tids2 = await Promise.all(dates.map(d => clock2.fromDate(d))); + + assert.deepStrictEqual(tids1, tids2); + }); +}); + +describe('TID State Persistence', () => { + const tempDir = path.join(os.tmpdir(), 'tid-test-' + Date.now()); + const statePath = path.join(tempDir, 'tid-state.json'); + + before(() => { + if (!fs.existsSync(tempDir)) { + fs.mkdirSync(tempDir, { recursive: true }); + } + }); + + after(() => { + if (fs.existsSync(tempDir)) { + fs.rmSync(tempDir, { recursive: true, force: true }); + } + }); + + it('should persist state to disk', async () => { + const clock = new TidClock( + new FakeClock(1000000000000000), + new SilentTidLogger(), + { statePath } + ); + + await clock.next(); + await clock.next(); + + // State file should exist + assert.strictEqual(fs.existsSync(statePath), true); + + // Should contain state + const stateData = JSON.parse(fs.readFileSync(statePath, 'utf-8')); + assert.strictEqual(stateData.generatedCount, 2); + }); + + it('should restore state from disk', async () => { + // Use a unique state file path for this test + const restoreStatePath = path.join(tempDir, 'tid-state-restore.json'); + + // Create first clock and generate some TIDs + const clock1 = new TidClock( + new FakeClock(1000000000000000), + new SilentTidLogger(), + { statePath: restoreStatePath, clockId: 10 } + ); + + await clock1.next(); + const tid2 = await clock1.next(); + + // Create new clock with same state file + const clock2 = new TidClock( + new FakeClock(1000000000000000), + new SilentTidLogger(), + { statePath: restoreStatePath } + ); + + const tid3 = await clock2.next(); + + // Should continue from where clock1 left off + assert.strictEqual(compareTids(tid2, tid3), -1); + + const state = clock2.getState(); + assert.strictEqual(state.generatedCount, 3); + }); +}); + +describe('TID Historical Dates', () => { + it('should handle very old dates (2005)', async () => { + const clock = new TidClock(new FakeClock(0), new SilentTidLogger()); + const date = new Date('2005-01-01T00:00:00Z'); + const tid = await clock.fromDate(date); + + assert.strictEqual(validateTid(tid), true); + }); + + it('should maintain monotonicity with out-of-order dates', async () => { + const clock = new TidClock(new FakeClock(0), new SilentTidLogger()); + + const dates = [ + new Date('2020-01-01T00:00:00Z'), + new Date('2015-01-01T00:00:00Z'), // Earlier! + new Date('2010-01-01T00:00:00Z'), // Even earlier! + new Date('2025-01-01T00:00:00Z'), + ]; + + const tids = await Promise.all(dates.map(d => clock.fromDate(d))); + + // Should still be monotonic despite out-of-order input + assert.strictEqual(areMonotonic(tids), true); + }); + + it('should handle duplicate dates', async () => { + const clock = new TidClock(new FakeClock(0), new SilentTidLogger()); + const date = new Date('2020-01-01T00:00:00Z'); + + const tid1 = await clock.fromDate(date); + const tid2 = await clock.fromDate(date); + const tid3 = await clock.fromDate(date); + + // All should be unique + assert.notStrictEqual(tid1, tid2); + assert.notStrictEqual(tid2, tid3); + assert.strictEqual(areMonotonic([tid1, tid2, tid3]), true); + }); +}); + +describe('TID Collision Detection', () => { + it('should detect duplicate TIDs (should never happen)', async () => { + const clock = new TidClock(new FakeClock(1000000000000000), new SilentTidLogger()); + + // Generate a TID + await clock.next(); + + // Try to force a duplicate by manually resetting state (testing collision detection) + // This simulates what would happen if there was a bug in generation + const state = clock.getState(); + + // Create a new clock with the exact same state + const clock2 = new TidClock( + new FakeClock(state.lastTimestampUs), + new SilentTidLogger(), + { + clockId: state.clockId, + initialState: { + lastTimestampUs: state.lastTimestampUs - 1, // Trick it into generating same timestamp + generatedCount: 0, + } + } + ); + + // First TID from clock2 might be a duplicate + // But the clock should handle it via sequence increment + const tid = await clock2.next(); + assert.strictEqual(validateTid(tid), true); + }); +}); + +describe('TID Comparison and Sorting', () => { + it('should compare TIDs correctly', () => { + const tid1 = '3jzfcijpj2z2a'; + const tid2 = '3jzfcijpj2z2b'; + const tid3 = '7777777777777'; + + assert.strictEqual(compareTids(tid1, tid1), 0); + assert.strictEqual(compareTids(tid1, tid2), -1); + assert.strictEqual(compareTids(tid2, tid1), 1); + assert.strictEqual(compareTids(tid1, tid3), -1); + }); + + it('should sort TIDs correctly', () => { + const tids = [ + '7777777777777', + '3jzfcijpj2z2a', + '3zzzzzzzzzzzz', + '3jzfcijpj2z2b', + ]; + + const sorted = [...tids].sort(compareTids); + + assert.deepStrictEqual(sorted, [ + '3jzfcijpj2z2a', + '3jzfcijpj2z2b', + '3zzzzzzzzzzzz', + '7777777777777', + ]); + }); +}); + +describe('TID Edge Cases', () => { + it('should handle microsecond precision', async () => { + const clock = new TidClock(new FakeClock(1234567890123456), new SilentTidLogger()); + const tid = await clock.next(); + const decoded = decodeTidTimestamp(tid); + + assert.strictEqual(decoded, 1234567890123456); + }); + + it('should handle very large timestamps', async () => { + const farFuture = new Date('2099-12-31T23:59:59.999Z').getTime() * 1000; + const clock = new TidClock(new FakeClock(farFuture), new SilentTidLogger()); + const tid = await clock.next(); + + assert.strictEqual(validateTid(tid), true); + }); + + it('should handle rapid sequential generation', async () => { + const clock = new TidClock(new RealClock(), new SilentTidLogger()); + const tids: string[] = []; + + // Generate as fast as possible + for (let i = 0; i < 1000; i++) { + tids.push(await clock.next()); + } + + const uniqueTids = new Set(tids); + assert.strictEqual(uniqueTids.size, tids.length); + assert.strictEqual(areMonotonic(tids), true); + }); +}); + +console.log('✓ All TID tests completed'); diff --git a/src/types.ts b/src/types.ts index 17f31c6..a78fbf6 100644 --- a/src/types.ts +++ b/src/types.ts @@ -57,6 +57,8 @@ export interface CommandLineArgs { reverse?: boolean; yes?: boolean; 'dry-run'?: boolean; + aggressive?: boolean; + fresh?: boolean; // Output verbose?: boolean; @@ -92,6 +94,7 @@ export interface Config { MIN_BATCH_DELAY: number; // from rate limiter RECORDS_PER_DAY_LIMIT: number; SAFETY_MARGIN: number; + AGGRESSIVE_SAFETY_MARGIN: number; SLINGSHOT_RESOLVER: string; diff --git a/src/utils/import-state.ts b/src/utils/import-state.ts new file mode 100644 index 0000000..428b894 --- /dev/null +++ b/src/utils/import-state.ts @@ -0,0 +1,233 @@ +import fs from 'fs'; +import path from 'path'; +import os from 'os'; +import crypto from 'crypto'; +import type { PlayRecord } from '../types.js'; +import { log } from './logger.js'; + +/** + * Import state for resume functionality + */ +export interface ImportState { + version: string; + startedAt: string; + lastUpdatedAt: string; + inputFile: string; + inputFileHash: string; + totalRecords: number; + processedRecords: number; + successfulRecords: number; + failedRecords: number; + lastSuccessfulIndex: number; + mode: 'lastfm' | 'spotify' | 'combined' | 'sync'; + completed: boolean; +} + +/** + * Get the state file path for an import + */ +export function getStateFilePath(inputFile: string, mode: string): string { + const stateDir = path.join(os.homedir(), '.lastfm-importer', 'state'); + + // Create state directory if it doesn't exist + if (!fs.existsSync(stateDir)) { + fs.mkdirSync(stateDir, { recursive: true }); + } + + // Create a unique filename based on input file and mode + const hash = crypto + .createHash('md5') + .update(inputFile + mode) + .digest('hex') + .substring(0, 8); + + return path.join(stateDir, `import-${hash}.json`); +} + +/** + * Calculate hash of input file for change detection + */ +export function calculateFileHash(filePath: string): string { + if (!fs.existsSync(filePath)) { + return ''; + } + + const stats = fs.statSync(filePath); + if (stats.isDirectory()) { + // For directories (like Spotify exports), hash the directory name and modification time + return crypto + .createHash('md5') + .update(filePath + stats.mtime.toISOString()) + .digest('hex'); + } + + // For files, use file size and modification time for quick comparison + return crypto + .createHash('md5') + .update(`${stats.size}-${stats.mtime.toISOString()}`) + .digest('hex'); +} + +/** + * Load import state from disk + */ +export function loadImportState(inputFile: string, mode: string): ImportState | null { + const stateFile = getStateFilePath(inputFile, mode); + + if (!fs.existsSync(stateFile)) { + return null; + } + + try { + const data = fs.readFileSync(stateFile, 'utf-8'); + const state = JSON.parse(data) as ImportState; + + // Check if input file has changed + const currentHash = calculateFileHash(inputFile); + if (state.inputFileHash !== currentHash) { + log.warn('Input file has changed since last import - starting fresh'); + return null; + } + + return state; + } catch (error) { + log.warn('Failed to load state file - starting fresh'); + return null; + } +} + +/** + * Save import state to disk + */ +export function saveImportState(state: ImportState): void { + const stateFile = getStateFilePath(state.inputFile, state.mode); + + try { + state.lastUpdatedAt = new Date().toISOString(); + fs.writeFileSync(stateFile, JSON.stringify(state, null, 2), 'utf-8'); + } catch (error) { + log.error('Failed to save state file - progress may be lost on restart'); + } +} + +/** + * Create initial import state + */ +export function createImportState( + inputFile: string, + mode: 'lastfm' | 'spotify' | 'combined' | 'sync', + totalRecords: number +): ImportState { + return { + version: '1.0', + startedAt: new Date().toISOString(), + lastUpdatedAt: new Date().toISOString(), + inputFile, + inputFileHash: calculateFileHash(inputFile), + totalRecords, + processedRecords: 0, + successfulRecords: 0, + failedRecords: 0, + lastSuccessfulIndex: -1, + mode, + completed: false, + }; +} + +/** + * Update import state after batch + */ +export function updateImportState( + state: ImportState, + batchIndex: number, + successCount: number, + errorCount: number +): void { + state.processedRecords += successCount + errorCount; + state.successfulRecords += successCount; + state.failedRecords += errorCount; + + if (successCount > 0) { + state.lastSuccessfulIndex = batchIndex; + } + + saveImportState(state); +} + +/** + * Mark import as completed + */ +export function completeImport(state: ImportState): void { + state.completed = true; + saveImportState(state); +} + +/** + * Clear import state (for fresh start) + */ +export function clearImportState(inputFile: string, mode: string): void { + const stateFile = getStateFilePath(inputFile, mode); + + if (fs.existsSync(stateFile)) { + fs.unlinkSync(stateFile); + } +} + +/** + * Display resume information + */ +export function displayResumeInfo(state: ImportState): void { + const elapsed = Date.now() - new Date(state.startedAt).getTime(); + const elapsedHours = Math.floor(elapsed / (1000 * 60 * 60)); + const elapsedMinutes = Math.floor((elapsed % (1000 * 60 * 60)) / (1000 * 60)); + + const remaining = state.totalRecords - state.processedRecords; + const progress = ((state.processedRecords / state.totalRecords) * 100).toFixed(1); + + log.section('Resuming Previous Import'); + log.info(`Started: ${new Date(state.startedAt).toLocaleString()}`); + log.info(`Progress: ${state.processedRecords.toLocaleString()}/${state.totalRecords.toLocaleString()} (${progress}%)`); + log.info(`Successful: ${state.successfulRecords.toLocaleString()}`); + + if (state.failedRecords > 0) { + log.warn(`Failed: ${state.failedRecords.toLocaleString()}`); + } + + log.info(`Remaining: ${remaining.toLocaleString()} records`); + + if (elapsedHours > 0) { + log.info(`Time elapsed: ${elapsedHours}h ${elapsedMinutes}m`); + } else if (elapsedMinutes > 0) { + log.info(`Time elapsed: ${elapsedMinutes}m`); + } + + log.blank(); +} + +/** + * Filter records to skip already processed ones + */ +export function filterUnprocessedRecords( + records: PlayRecord[], + state: ImportState +): PlayRecord[] { + if (state.lastSuccessfulIndex < 0) { + return records; + } + + // Skip records up to and including the last successful index + const startIndex = state.lastSuccessfulIndex + 1; + + if (startIndex >= records.length) { + return []; + } + + return records.slice(startIndex); +} + +/** + * Get the starting index for resume + */ +export function getResumeStartIndex(state: ImportState): number { + return state.lastSuccessfulIndex + 1; +} diff --git a/src/utils/tid-audit.ts b/src/utils/tid-audit.ts new file mode 100644 index 0000000..2d69c24 --- /dev/null +++ b/src/utils/tid-audit.ts @@ -0,0 +1,179 @@ +/** + * TID Audit and Reporting Tools + * + * Utilities for auditing TID generation and producing reports + */ + +import { validateTid, decodeTidTimestamp, decodeTidClockId, compareTids } from './tid-clock.js'; + +export interface TidAuditEntry { + tid: string; + valid: boolean; + timestamp?: number; + clockId?: number; + date?: Date; + errors: string[]; +} + +export interface TidAuditReport { + totalTids: number; + validTids: number; + invalidTids: number; + duplicates: number; + monotonic: boolean; + firstTid?: TidAuditEntry; + lastTid?: TidAuditEntry; + entries: TidAuditEntry[]; + errors: string[]; +} + +/** + * Audit a list of TIDs + */ +export function auditTids(tids: string[]): TidAuditReport { + const entries: TidAuditEntry[] = []; + const seenTids = new Set(); + let duplicates = 0; + let validCount = 0; + let invalidCount = 0; + const globalErrors: string[] = []; + + // Audit each TID + for (const tid of tids) { + const entry: TidAuditEntry = { + tid, + valid: false, + errors: [], + }; + + // Check for duplicates + if (seenTids.has(tid)) { + duplicates++; + entry.errors.push('Duplicate TID'); + } + seenTids.add(tid); + + // Validate format + const valid = validateTid(tid); + entry.valid = valid; + + if (valid) { + validCount++; + try { + entry.timestamp = decodeTidTimestamp(tid); + entry.clockId = decodeTidClockId(tid); + entry.date = new Date(entry.timestamp / 1000); + } catch (error) { + entry.errors.push(`Decode error: ${error}`); + } + } else { + invalidCount++; + entry.errors.push('Invalid format'); + } + + entries.push(entry); + } + + // Check monotonicity + let monotonic = true; + for (let i = 1; i < tids.length; i++) { + if (compareTids(tids[i - 1], tids[i]) >= 0) { + monotonic = false; + globalErrors.push(`Non-monotonic at index ${i}: ${tids[i - 1]} >= ${tids[i]}`); + } + } + + return { + totalTids: tids.length, + validTids: validCount, + invalidTids: invalidCount, + duplicates, + monotonic, + firstTid: entries[0], + lastTid: entries[entries.length - 1], + entries, + errors: globalErrors, + }; +} + +/** + * Format audit report as human-readable text + */ +export function formatAuditReport(report: TidAuditReport): string { + const lines: string[] = []; + + lines.push('='.repeat(80)); + lines.push('TID AUDIT REPORT'); + lines.push('='.repeat(80)); + lines.push(''); + + // Summary + lines.push('SUMMARY'); + lines.push('-'.repeat(80)); + lines.push(`Total TIDs: ${report.totalTids.toLocaleString()}`); + lines.push(`Valid: ${report.validTids.toLocaleString()} (${((report.validTids / report.totalTids) * 100).toFixed(1)}%)`); + lines.push(`Invalid: ${report.invalidTids.toLocaleString()}`); + lines.push(`Duplicates: ${report.duplicates.toLocaleString()}`); + lines.push(`Monotonic: ${report.monotonic ? '✓ YES' : '✗ NO'}`); + lines.push(''); + + // Time range + if (report.firstTid && report.firstTid.date) { + lines.push('TIME RANGE'); + lines.push('-'.repeat(80)); + lines.push(`First TID: ${report.firstTid.tid}`); + lines.push(` Timestamp: ${report.firstTid.date.toISOString()}`); + lines.push(` Clock ID: ${report.firstTid.clockId}`); + lines.push(''); + if (report.lastTid && report.lastTid.date) { + lines.push(`Last TID: ${report.lastTid.tid}`); + lines.push(` Timestamp: ${report.lastTid.date.toISOString()}`); + lines.push(` Clock ID: ${report.lastTid.clockId}`); + const durationMs = (report.lastTid.timestamp! - report.firstTid.timestamp!) / 1000; + const durationDays = durationMs / (1000 * 60 * 60 * 24); + lines.push(` Duration: ${durationDays.toFixed(2)} days`); + lines.push(''); + } + } + + // Errors + if (report.errors.length > 0) { + lines.push('ERRORS'); + lines.push('-'.repeat(80)); + for (const error of report.errors.slice(0, 10)) { + lines.push(` • ${error}`); + } + if (report.errors.length > 10) { + lines.push(` ... and ${report.errors.length - 10} more errors`); + } + lines.push(''); + } + + // Invalid TIDs + const invalidEntries = report.entries.filter(e => !e.valid || e.errors.length > 0); + if (invalidEntries.length > 0) { + lines.push('INVALID/PROBLEM TIDs'); + lines.push('-'.repeat(80)); + for (const entry of invalidEntries.slice(0, 10)) { + lines.push(` ${entry.tid}`); + for (const error of entry.errors) { + lines.push(` - ${error}`); + } + } + if (invalidEntries.length > 10) { + lines.push(` ... and ${invalidEntries.length - 10} more invalid TIDs`); + } + lines.push(''); + } + + lines.push('='.repeat(80)); + + return lines.join('\n'); +} + +/** + * Format audit report as JSON + */ +export function formatAuditReportJson(report: TidAuditReport): string { + return JSON.stringify(report, null, 2); +} diff --git a/src/utils/tid-clock.ts b/src/utils/tid-clock.ts new file mode 100644 index 0000000..c101f05 --- /dev/null +++ b/src/utils/tid-clock.ts @@ -0,0 +1,470 @@ +/** + * TID (Timestamp Identifier) Clock for ATProto + * + * Implements spec-compliant, monotonic TID generation with: + * - AT-Protocol format validation (13 chars, base32 alphabet) + * - Strict monotonicity guarantees (even under clock drift) + * - Concurrency safety (mutex-protected state) + * - Deterministic mode for dry-runs + * - Full observability (structured JSON logging) + * - Collision resistance + * + * Based on AT-Protocol spec: https://atproto.com/specs/tid + * Reference implementation: @atproto/common-web + */ + +import { s32decode, s32encode } from '@atproto/common-web/dist/util.js'; +import fs from 'fs'; +import path from 'path'; +import crypto from 'crypto'; + +const TID_LENGTH = 13; +const TID_REGEX = /^[234567abcdefghij][234567abcdefghijklmnopqrstuvwxyz]{12}$/; + +/** + * TID validation error + */ +export class InvalidTidError extends Error { + constructor(message: string, public tid?: string) { + super(message); + this.name = 'InvalidTidError'; + } +} + +/** + * TID generation modes + */ +export enum TidMode { + PRODUCTION = 'production', // Real wall-clock time + DRY_RUN = 'dry-run', // Deterministic with fixed seed + REPLAY = 'replay', // Replay from logged state +} + +/** + * Clock source abstraction for testability + */ +export interface ClockSource { + now(): number; // Returns microseconds since epoch +} + +/** + * Real wall-clock implementation + */ +export class RealClock implements ClockSource { + now(): number { + return Date.now() * 1000; // Convert ms to µs + } +} + +/** + * Deterministic fake clock for testing/dry-runs + */ +export class FakeClock implements ClockSource { + constructor(private timestamp: number) {} + + now(): number { + return this.timestamp; + } + + advance(microseconds: number): void { + this.timestamp += microseconds; + } + + set(microseconds: number): void { + this.timestamp = microseconds; + } +} + +/** + * TID generator state (for persistence and logging) + */ +export interface TidState { + lastTimestampUs: number; // Last generated timestamp in microseconds + clockId: number; // Clock identifier (0-31 for this implementation) + generatedCount: number; // Total TIDs generated +} + +/** + * Metadata for each generated TID + */ +export interface TidMetadata { + tid: string; + timestampUs: number; + clockId: number; + generatedAt: string; // ISO8601 with microseconds + validated: boolean; + context?: string; +} + +/** + * Logger interface for TID operations + */ +export interface TidLogger { + logGenerated(metadata: TidMetadata, opId: string): void; + logWarning(message: string, details?: any): void; + logError(message: string, details?: any): void; +} + +/** + * Default console-based logger + */ +export class ConsoleTidLogger implements TidLogger { + logGenerated(metadata: TidMetadata, opId: string): void { + const entry = { + level: 'INFO', + op_id: opId, + ts: metadata.generatedAt, + event: 'tid.generated', + tid: metadata.tid, + clock_id: metadata.clockId, + wall_ts_us: metadata.timestampUs, + generator: 'tid-clock-v1', + context: metadata.context || 'unknown', + validated: metadata.validated, + }; + console.log(JSON.stringify(entry)); + } + + logWarning(message: string, details?: any): void { + console.warn(JSON.stringify({ + level: 'WARN', + event: 'tid.warning', + message, + ...details, + })); + } + + logError(message: string, details?: any): void { + console.error(JSON.stringify({ + level: 'ERROR', + event: 'tid.error', + message, + ...details, + })); + } +} + +/** + * Silent logger (for production when structured logging is handled elsewhere) + */ +export class SilentTidLogger implements TidLogger { + logGenerated(): void {} + logWarning(): void {} + logError(): void {} +} + +/** + * Main TID Clock generator + */ +export class TidClock { + private state: TidState; + private clock: ClockSource; + private logger: TidLogger; + private statePath: string | null; + private mutex: Promise = Promise.resolve(); + + constructor( + clock: ClockSource = new RealClock(), + logger: TidLogger = new SilentTidLogger(), + options: { + statePath?: string; + clockId?: number; + initialState?: Partial; + } = {} + ) { + this.clock = clock; + this.logger = logger; + this.statePath = options.statePath || null; + + // Initialize or load state + if (this.statePath && fs.existsSync(this.statePath)) { + this.state = this.loadState(); + } else { + this.state = { + lastTimestampUs: 0, + clockId: options.clockId ?? this.generateClockId(), + generatedCount: 0, + ...options.initialState, + }; + if (this.statePath) { + this.saveState(); + } + } + } + + /** + * Generate a cryptographically random clock ID (0-31) + */ + private generateClockId(): number { + return crypto.randomInt(0, 32); + } + + /** + * Generate next TID with monotonicity guarantees + * + * Per AT-Protocol spec and reference implementation: + * - Use max(currentTime, lastTimestamp) to handle clock drift + * - If same as last timestamp, increment the timestamp itself (acts as sequence) + * - This ensures monotonicity while keeping TID format simple + */ + async next(context?: string): Promise { + return this.withMutex(async () => { + const currentTime = this.clock.now(); + + // Take max of current time and last timestamp (handles backwards clock drift) + let timestamp = Math.max(currentTime, this.state.lastTimestampUs); + + // If we're at the same timestamp, increment by 1 microsecond + // This acts as our sequence counter while maintaining monotonicity + if (timestamp === this.state.lastTimestampUs) { + timestamp = this.state.lastTimestampUs + 1; + } + + if (currentTime < this.state.lastTimestampUs) { + this.logger.logWarning('Clock moved backwards', { + current: currentTime, + last: this.state.lastTimestampUs, + delta: this.state.lastTimestampUs - currentTime, + action: 'using_incremented_timestamp', + }); + } + + // Generate TID + const tid = this.encodeTid(timestamp, this.state.clockId); + + // Validate format + const validated = this.validateTid(tid); + if (!validated) { + const error = new InvalidTidError('Generated invalid TID', tid); + this.logger.logError('TID validation failed', { + tid, + timestamp, + clockId: this.state.clockId, + }); + throw error; + } + + // Update state + this.state.lastTimestampUs = timestamp; + this.state.generatedCount++; + + // Persist state if configured + if (this.statePath) { + this.saveState(); + } + + // Log metadata + const opId = crypto.randomUUID(); + const metadata: TidMetadata = { + tid, + timestampUs: timestamp, + clockId: this.state.clockId, + generatedAt: this.formatMicrosecondTimestamp(timestamp), + validated, + context, + }; + this.logger.logGenerated(metadata, opId); + + return tid; + }); + } + + /** + * Generate TID from a specific Date (for historical records) + * Maintains monotonicity relative to previously generated TIDs + */ + async fromDate(date: Date, context?: string): Promise { + const timestamp = date.getTime() * 1000; // Convert ms to µs + + return this.withMutex(async () => { + // Ensure monotonicity: use max of input timestamp and last generated + let finalTimestamp = Math.max(timestamp, this.state.lastTimestampUs); + + // If we're at the same timestamp, increment by 1 microsecond + if (finalTimestamp === this.state.lastTimestampUs) { + finalTimestamp = this.state.lastTimestampUs + 1; + } + + const tid = this.encodeTid(finalTimestamp, this.state.clockId); + const validated = this.validateTid(tid); + + if (!validated) { + throw new InvalidTidError('Generated invalid TID from date', tid); + } + + // Update state + this.state.lastTimestampUs = finalTimestamp; + this.state.generatedCount++; + + if (this.statePath) { + this.saveState(); + } + + const opId = crypto.randomUUID(); + const metadata: TidMetadata = { + tid, + timestampUs: finalTimestamp, + clockId: this.state.clockId, + generatedAt: this.formatMicrosecondTimestamp(finalTimestamp), + validated, + context, + }; + this.logger.logGenerated(metadata, opId); + + return tid; + }); + } + + /** + * Encode timestamp and clock ID into TID format + */ + private encodeTid(timestampUs: number, clockId: number): string { + const timestampStr = s32encode(timestampUs).padStart(11, '2'); + const clockIdStr = s32encode(clockId).padStart(2, '2'); + return timestampStr + clockIdStr; + } + + /** + * Validate TID format per AT-Protocol spec + */ + private validateTid(tid: string): boolean { + if (tid.length !== TID_LENGTH) { + return false; + } + if (!TID_REGEX.test(tid)) { + return false; + } + return true; + } + + /** + * Format microsecond timestamp as ISO8601 with microsecond precision + */ + private formatMicrosecondTimestamp(timestampUs: number): string { + const ms = Math.floor(timestampUs / 1000); + const us = timestampUs % 1000; + const date = new Date(ms); + const iso = date.toISOString(); + // Replace milliseconds with microseconds + return iso.replace(/\.(\d{3})Z$/, `.${us.toString().padStart(3, '0')}000Z`); + } + + /** + * Load state from disk + */ + private loadState(): TidState { + if (!this.statePath) { + throw new Error('State path not configured'); + } + const data = fs.readFileSync(this.statePath, 'utf-8'); + return JSON.parse(data); + } + + /** + * Save state to disk + */ + private saveState(): void { + if (!this.statePath) { + return; + } + const dir = path.dirname(this.statePath); + if (!fs.existsSync(dir)) { + fs.mkdirSync(dir, { recursive: true }); + } + fs.writeFileSync(this.statePath, JSON.stringify(this.state, null, 2), 'utf-8'); + } + + /** + * Get current state (for inspection/debugging) + */ + getState(): Readonly { + return { ...this.state }; + } + + /** + * Reset state (for testing) + */ + reset(preserveClockId: boolean = false): void { + const clockId = preserveClockId ? this.state.clockId : this.generateClockId(); + this.state = { + lastTimestampUs: 0, + clockId, + generatedCount: 0, + }; + if (this.statePath) { + this.saveState(); + } + } + + /** + * Mutex for concurrent access protection + */ + private async withMutex(fn: () => Promise): Promise { + const currentMutex = this.mutex; + let releaseMutex: () => void; + this.mutex = new Promise((resolve) => { + releaseMutex = resolve; + }); + + try { + await currentMutex; + return await fn(); + } finally { + releaseMutex!(); + } + } +} + +/** + * Validate a TID string + */ +export function validateTid(tid: string): boolean { + return tid.length === TID_LENGTH && TID_REGEX.test(tid); +} + +/** + * Ensure TID is valid (throws on invalid) + */ +export function ensureValidTid(tid: string): asserts tid is string { + if (!validateTid(tid)) { + throw new InvalidTidError(`Invalid TID format: ${tid}`, tid); + } +} + +/** + * Decode TID to get timestamp + */ +export function decodeTidTimestamp(tid: string): number { + ensureValidTid(tid); + return s32decode(tid.slice(0, 11)); +} + +/** + * Decode TID to get clock ID + */ +export function decodeTidClockId(tid: string): number { + ensureValidTid(tid); + return s32decode(tid.slice(11, 13)); +} + +/** + * Compare two TIDs (for sorting) + * Returns: -1 if a < b, 0 if equal, 1 if a > b + */ +export function compareTids(a: string, b: string): number { + if (a < b) return -1; + if (a > b) return 1; + return 0; +} + +/** + * Check if TIDs are in strictly increasing order + */ +export function areMonotonic(tids: string[]): boolean { + for (let i = 1; i < tids.length; i++) { + if (tids[i] <= tids[i - 1]) { + return false; + } + } + return true; +} diff --git a/src/utils/tid.ts b/src/utils/tid.ts index b172dfa..3b644a6 100644 --- a/src/utils/tid.ts +++ b/src/utils/tid.ts @@ -1,54 +1,119 @@ /** * TID (Timestamp Identifier) generation for ATProto - * Based on: https://atproto.com/specs/tid * - * Per the spec: "If the local clock has only millisecond precision, the timestamp - * should be padded." Our implementation pads to ensure 11 characters for the timestamp - * portion and 2 characters for the clockid, for a total of exactly 13 characters. + * This module provides a high-level API for TID generation. + * For the full implementation with monotonicity guarantees, see tid-clock.ts * - * This implementation uses the official ATProto s32encode function and properly handles - * historical dates (like 2005 scrobbles) by padding to the required length. + * Based on: https://atproto.com/specs/tid */ -import { s32encode } from '@atproto/common-web/dist/util.js'; +import { TidClock, RealClock, FakeClock, validateTid } from './tid-clock.js'; + +// Global TID clock instance +let globalClock: TidClock | null = null; /** - * Generate a TID from a Date object - * - * TID format (13 characters total): - * - Characters 0-10: timestamp in microseconds, base32-encoded - * - Characters 11-12: clock ID, base32-encoded + * Initialize the global TID clock + * Should be called once at application startup + */ +export function initTidClock(options: { + mode?: 'production' | 'dry-run'; + statePath?: string; + seed?: number; + clockId?: number; +} = {}): void { + const { mode = 'production', statePath, seed, clockId } = options; + + if (mode === 'dry-run' && seed !== undefined) { + // Deterministic clock for dry-run with fixed clockId for reproducibility + globalClock = new TidClock(new FakeClock(seed), undefined, { + statePath, + clockId: clockId ?? 0 // Use fixed clockId for deterministic TIDs + }); + } else { + // Production real-time clock + globalClock = new TidClock(new RealClock(), undefined, { statePath, clockId }); + } +} + +/** + * Get or create the global TID clock + */ +function getClock(): TidClock { + if (!globalClock) { + // Auto-initialize with production settings if not explicitly initialized + globalClock = new TidClock(new RealClock()); + } + return globalClock; +} + +/** + * Generate a TID from a Date object with monotonicity guarantees * - * The timestamp is padded with '2' (representing 0 in base32) to ensure exactly 11 characters. - * The clockid is similarly padded to ensure exactly 2 characters. + * This is the recommended way to generate TIDs for historical records. + * TIDs are guaranteed to be: + * - Spec-compliant (13 chars, valid base32) + * - Strictly monotonic (even if dates are out of order) + * - Collision-free (duplicate detection) * * @param date - The date to generate a TID from + * @param context - Optional context for logging (e.g., "inject:playlist") * @returns A valid 13-character TID string */ -export function generateTID(date: Date): string { - // Convert to Unix microseconds (JS Date.getTime() returns milliseconds) - // Per spec: multiply by 1000 to pad millisecond precision - const unixMicros = date.getTime() * 1000; - - // Use a fixed clockid of 0 for deterministic TID generation from timestamps - // This ensures the same playedTime always generates the same TID (important for deduplication) - const clockid = 0; - - // Encode timestamp and clockid with proper padding - // Timestamp should be 11 characters, clockid should be 2 characters - // Padding character is '2' which represents 0 in base32 - const timestampStr = s32encode(unixMicros).padStart(11, '2'); - const clockidStr = s32encode(clockid).padStart(2, '2'); - - return timestampStr + clockidStr; +export async function generateTID(date: Date, context?: string): Promise { + const clock = getClock(); + return await clock.fromDate(date, context); } /** * Generate a TID from an ISO 8601 timestamp string * * @param isoString - ISO 8601 formatted datetime string + * @param context - Optional context for logging * @returns A valid 13-character TID string */ -export function generateTIDFromISO(isoString: string): string { - return generateTID(new Date(isoString)); +export async function generateTIDFromISO(isoString: string, context?: string): Promise { + return await generateTID(new Date(isoString), context); } + +/** + * Generate next TID using current time + * + * @param context - Optional context for logging + * @returns A valid 13-character TID string + */ +export async function generateNextTID(context?: string): Promise { + const clock = getClock(); + return await clock.next(context); +} + +/** + * Validate a TID string + * + * @param tid - The TID to validate + * @returns true if valid, false otherwise + */ +export function isValidTID(tid: string): boolean { + return validateTid(tid); +} + +/** + * Reset the global clock state (for testing only) + */ +export function resetTidClock(): void { + if (globalClock) { + globalClock.reset(); + } +} + +/** + * Get the current clock state (for debugging/inspection) + */ +export function getTidClockState() { + const clock = getClock(); + return clock.getState(); +} + +// Re-export types and utilities from tid-clock +export { TidClock, RealClock, FakeClock, validateTid } from './tid-clock.js'; +export type { TidState, TidMetadata } from './tid-clock.js';