diff --git a/README.md b/README.md index 2b5c04c..cc89d45 100644 --- a/README.md +++ b/README.md @@ -5,13 +5,77 @@ Import your Last.fm and Spotify listening history to the AT Protocol network usi **Repository:** [malachite](https://github.com/ewanc26/malachite) [Also available on Tangled](https://tangled.org/did:plc:ofrbh253gwicbkc5nktqepol/atproto-lastfm-importer) +## Table of Contents + +- [⚠️ Important: Rate Limits](#️-important-rate-limits) + - [πŸ“š Rate Limiting Documentation](#-rate-limiting-documentation) + - [How Dynamic Batch Sizing Works](#how-dynamic-batch-sizing-works) +- [What's with the name?](#whats-with-the-name) +- [Quick Start](#quick-start) + - [Interactive Mode (Recommended for First-Time Users)](#interactive-mode-recommended-for-first-time-users) + - [Command Line Mode](#command-line-mode) +- [Features](#features) + - [Import Capabilities](#import-capabilities) + - [Performance & Safety](#performance--safety) + - [User Experience](#user-experience) + - [Technical Features](#technical-features) +- [Usage Examples](#usage-examples) + - [Combined Import (Last.fm + Spotify)](#combined-import-lastfm--spotify) + - [Re-Sync Mode](#re-sync-mode) + - [Remove Duplicates](#remove-duplicates) + - [Import from Spotify](#import-from-spotify) + - [Import from Last.fm](#import-from-lastfm) + - [Advanced Options](#advanced-options) +- [Command Line Options](#command-line-options) + - [Required Options](#required-options) + - [Import Mode](#import-mode) + - [Additional Options](#additional-options) + - [PDS Override](#pds-override) + - [Legacy Flags (Backwards Compatible)](#legacy-flags-backwards-compatible) +- [Getting Your Data](#getting-your-data) + - [Last.fm Export](#lastfm-export) + - [Spotify Export](#spotify-export) +- [Data Format](#data-format) + - [Required Fields](#required-fields) + - [Optional Fields](#optional-fields) + - [Example Records](#example-records) +- [How It Works](#how-it-works) + - [Processing Flow](#processing-flow) + - [Automatic Duplicate Prevention](#automatic-duplicate-prevention) + - [Rate Limiting Algorithm](#rate-limiting-algorithm) + - [Multi-Day Imports](#multi-day-imports) +- [Logging and Output](#logging-and-output) + - [Verbosity Levels](#verbosity-levels) +- [Error Handling](#error-handling) +- [Troubleshooting](#troubleshooting) + - [Authentication Issues](#authentication-issues) + - [Performance Issues](#performance-issues) + - [Connection Issues](#connection-issues) + - [Output Control](#output-control) +- [Development](#development) +- [File Storage](#file-storage) + - [Credential Storage](#credential-storage) +- [Project Structure](#project-structure) +- [Technical Details](#technical-details) + - [Authentication](#authentication) + - [Batch Publishing](#batch-publishing) + - [Data Mapping](#data-mapping) + - [Lexicon Reference](#lexicon-reference) +- [Contributing](#contributing) +- [License](#license) +- [Credits](#credits) + ## ⚠️ Important: Rate Limits **CRITICAL**: Bluesky's AppView has rate limits on PDS instances. Exceeding 10K records per day can rate limit your **ENTIRE PDS**, affecting all users on your instance. This importer automatically protects your PDS by: +- **Dynamic batch sizing** (1-200 records) that adapts to available quota in real-time +- **15% headroom buffer** prevents quota exhaustion before hitting the limit - Limiting imports to **7,500 records per day** (with 75% safety margin) - Calculating optimal batch sizes and delays +- **Graceful degradation** - scales down smoothly as quota depletes +- **Instant recovery** - immediately returns to maximum speed after quota resets - Pausing 24 hours between days for large imports - Providing clear progress tracking and time estimates - Persisting state across restarts for safe resume @@ -25,6 +89,26 @@ Malachite has comprehensive rate limiting protection built in. npm run check-limits ``` +### How Dynamic Batch Sizing Works + +Malachite continuously monitors your rate limit quota and automatically adjusts batch size: + +``` +Fresh Quota (5000 points) β†’ Batch Size: 200 records (maximum speed) +Half Depleted (2500 points) β†’ Batch Size: 200 records (still optimal) +Approaching Limit (1200) β†’ Batch Size: 150 records (scaling down) +Near Headroom (900) β†’ Batch Size: 50 records (conservative) +Below Headroom (700) β†’ Batch Size: 1 record (minimal progress) +[Quota Resets] β†’ Batch Size: 200 records (instant recovery) +``` + +**Benefits:** +- βœ… **2x faster** when quota is fresh (200 vs 100 records/batch) +- βœ… **Never hits rate limits** - proactive scaling with 15% buffer +- βœ… **Always makes progress** - even with minimal quota (batch size 1) +- βœ… **Automatic recovery** - no manual intervention needed +- βœ… **Transparent** - logs all batch size changes with reasons + For more details, see the [Bluesky Rate Limits Documentation](https://docs.bsky.app/blog/rate-limits-pds-v3). ## What’s with the name? @@ -85,8 +169,10 @@ node dist/index.js -i lastfm.csv -h alice.bsky.social -p xxxx-xxxx-xxxx-xxxx -y ### Performance & Safety - βœ… **Automatic Duplicate Prevention**: Automatically checks Teal and skips records that already exist (no duplicates!) - βœ… **Input Deduplication**: Removes duplicate entries within the source file before submission +- βœ… **Dynamic Batch Sizing**: Automatically adjusts batch size (1-200 records) based on available rate limit quota - βœ… **Batch Operations**: Uses `com.atproto.repo.applyWrites` for efficient batch publishing (up to 200 records per call) -- βœ… **Rate Limiting**: Automatic daily limits prevent PDS rate limiting +- βœ… **Intelligent Rate Limiting**: Real-time quota monitoring with 15% headroom buffer prevents rate limit exhaustion +- βœ… **Adaptive Recovery**: Automatically scales back to maximum speed after quota resets - βœ… **Multi-Day Imports**: Large imports automatically span multiple days with 24-hour pauses - βœ… **Resume Support**: Safe to stop (Ctrl+C) and restart - continues from where it left off - βœ… **Graceful Cancellation**: Press Ctrl+C to stop after the current batch completes @@ -257,7 +343,7 @@ pnpm start -i lastfm.csv -h alice.bsky.social -p xxxx-xxxx-xxxx-xxxx -y -q | `--verbose` | `-v` | Enable debug logging | `false` | | `--quiet` | `-q` | Suppress non-essential output | `false` | | `--dev` | | Development mode (verbose + file logging + smaller batches) | `false` | -| `--batch-size ` | `-b` | Records per batch (1-200) | Auto-calculated | +| `--batch-size ` | `-b` | Initial batch size (1-200, dynamically adjusted) | Auto-calculated | | `--batch-delay ` | `-d` | Delay between batches in ms | `500` (min) | | `--help` | | Show help message | - | @@ -454,9 +540,13 @@ Removes duplicates within your source file(s): ### Rate Limiting Algorithm 1. Calculates safe daily limit (75% of 10K = 7,500 records/day by default) 2. Determines how many days needed for your import -3. Calculates optimal batch size and delay to spread records evenly -4. Enforces minimum delay between batches -5. Shows clear schedule before starting +3. **Monitors rate limit quota in real-time** before each batch +4. **Dynamically adjusts batch size** (1-200 records) based on available points +5. **Preserves 15% headroom buffer** to prevent exhaustion +6. **Automatically waits** when quota is exhausted (with countdown timer) +7. **Instantly scales back up** to maximum batch size after quota resets +8. Enforces minimum delay between batches +9. Shows clear schedule and real-time batch size adjustments ### Multi-Day Imports @@ -678,7 +768,9 @@ malachite/ ### Batch Publishing - Uses `com.atproto.repo.applyWrites` for efficiency (up to 20x faster than individual calls) - Batches up to 200 records per API call (PDS maximum) -- Automatically adjusts batch size based on total record count +- **Dynamic batch sizing** (1-200 records) based on real-time rate limit quota +- **Intelligent quota monitoring** with 15% headroom buffer +- **Automatic adjustment** - scales down as quota depletes, scales up after reset - Enforces minimum delays between batches for rate limit safety ### Data Mapping diff --git a/src/lib/auth.ts b/src/lib/auth.ts index b728a67..41ef81b 100644 --- a/src/lib/auth.ts +++ b/src/lib/auth.ts @@ -1,6 +1,7 @@ import { AtpAgent } from '@atproto/api'; import { prompt } from '../utils/input.js'; import * as ui from '../utils/ui.js'; +import { saveCredentials } from '../utils/credentials.js'; interface ResolvedIdentity { did: string; @@ -75,6 +76,16 @@ export async function login( ui.succeedSpinner('Logged in successfully (PDS override)!'); ui.keyValue('DID', agent.session?.did || 'unknown'); ui.keyValue('Handle', agent.session?.handle || 'unknown'); + + // Automatically save credentials (encrypted with SHA-512, machine-specific) + try { + saveCredentials(identifier, password); + ui.info('Credentials saved securely (SHA-512 encrypted, machine-specific)'); + } catch (err) { + // Non-fatal - log but continue + ui.warning('Failed to save credentials - you may need to re-enter them next time'); + } + console.log(''); return agent; } @@ -98,6 +109,16 @@ export async function login( ui.succeedSpinner('Logged in successfully!'); ui.keyValue('DID', agent.session?.did || 'unknown'); ui.keyValue('Handle', agent.session?.handle || 'unknown'); + + // Automatically save credentials (encrypted with SHA-512, machine-specific) + try { + saveCredentials(identifier, password); + ui.info('Credentials saved securely (SHA-512 encrypted, machine-specific)'); + } catch (err) { + // Non-fatal - log but continue + ui.warning('Failed to save credentials - you may need to re-enter them next time'); + } + console.log(''); return agent; diff --git a/src/lib/cli.ts b/src/lib/cli.ts index 1dc6b27..460f830 100644 --- a/src/lib/cli.ts +++ b/src/lib/cli.ts @@ -16,7 +16,6 @@ import { Logger, LogLevel, setGlobalLogger, log } from '../utils/logger.js'; import { registerKillswitch } from '../utils/killswitch.js'; import { clearCache, clearAllCaches } from '../utils/teal-cache.js'; import { - saveCredentials, loadCredentials, hasStoredCredentials, clearCredentials, @@ -335,12 +334,8 @@ async function runInteractiveMode(): Promise { } args.password = password; - // Offer to save credentials - const saveCredsAnswer = await confirm('\nSave credentials for future use? (encrypted, machine-specific)', false); - if (saveCredsAnswer) { - saveCredentials(args.handle, args.password); - console.log('βœ“ Credentials saved securely to ~/.malachite/credentials.json'); - } + // Note: Credentials will be automatically saved after successful login + // No need to prompt the user } console.log(''); diff --git a/src/lib/publisher.ts b/src/lib/publisher.ts index dcfed3f..455ca03 100644 --- a/src/lib/publisher.ts +++ b/src/lib/publisher.ts @@ -22,6 +22,16 @@ import { */ const MAX_APPLY_WRITES_OPS = 200; +/** + * Minimum batch size to maintain reasonable progress + */ +const MIN_BATCH_SIZE = 1; + +/** + * Maximum batch size (same as MAX_APPLY_WRITES_OPS) + */ +const MAX_BATCH_SIZE = 200; + /** * Publish records using com.atproto.repo.applyWrites for efficient batching * with adaptive rate limiting and stateful resume support @@ -47,7 +57,7 @@ export async function publishRecordsWithApplyWrites( throw new Error('Agent is required for publishing'); } - // Start with aggressive settings + // Start with conservative settings let currentBatchSize = Math.min(batchSize, MAX_APPLY_WRITES_OPS); let currentBatchDelay = batchDelay; @@ -56,18 +66,33 @@ export async function publishRecordsWithApplyWrites( let consecutiveFailures = 0; const MAX_CONSECUTIVE_FAILURES = 3; const POINTS_PER_RECORD = 3; // approximate cost per create operation + + /** + * Calculate optimal batch size based on available rate limit points + */ + const calculateOptimalBatchSize = (availablePoints: number): number => { + // Calculate how many records we can fit in available points + const maxRecordsFromQuota = Math.floor(availablePoints / POINTS_PER_RECORD); + + // Clamp to our min/max bounds + const optimalSize = Math.max(MIN_BATCH_SIZE, Math.min(maxRecordsFromQuota, MAX_BATCH_SIZE)); + + log.debug(`[publisher.ts] calculateOptimalBatchSize: availablePoints=${availablePoints}, maxRecords=${maxRecordsFromQuota}, optimal=${optimalSize}`); + return optimalSize; + }; // Persistent rate limiter (reads/writes ~/.malachite/state/rate-limit.json) // Use headroom threshold instead of safety margin for better rate limit handling const rl = new RateLimiter({ headroom: 0.15 }); // Preserve 15% buffer before hitting limit - log.section('Conservative Adaptive Import'); - log.info(`Initial batch size: ${currentBatchSize} records (conservative)`); - log.info(`Initial delay: ${currentBatchDelay}ms (2 seconds - very safe)`); + log.section('Dynamic Adaptive Import'); + log.info(`Initial batch size: ${currentBatchSize} records`); + log.info(`Batch size range: ${MIN_BATCH_SIZE}-${MAX_BATCH_SIZE} records (adjusts based on available quota)`); + log.info(`Initial delay: ${currentBatchDelay}ms`); log.debug(`[publisher.ts] MAX_APPLY_WRITES_OPS=${MAX_APPLY_WRITES_OPS}, POINTS_PER_RECORD=${POINTS_PER_RECORD}`); log.debug(`[publisher.ts] Headroom threshold: 15%, Records per day limit: ${formatLocaleNumber(config.RECORDS_PER_DAY_LIMIT)}`); - log.info(`Will automatically adjust based on server response`); - log.info(`Using conservative settings to protect your PDS`); + log.info(`Batch size will automatically scale based on rate limit quota`); + log.info(`Delay will adjust based on server response`); log.blank(); log.info(`Publishing ${formatLocaleNumber(totalRecords)} records using adaptive batching...`); log.warn('Press Ctrl+C to stop gracefully after current batch'); @@ -86,18 +111,30 @@ export async function publishRecordsWithApplyWrites( } let i = startIndex; + let batchCounter = 0; // Track actual batch number across resume while (i < totalRecords) { // Check killswitch before processing batch if (isImportCancelled()) { return handleCancellation(successCount, errorCount, totalRecords); } + // Adjust batch size based on available rate limit quota + const safePoints = rl.getSafeAvailablePoints(); + const optimalSize = calculateOptimalBatchSize(safePoints); + + // Only adjust if we need to (avoid log spam) + if (optimalSize !== currentBatchSize) { + const oldSize = currentBatchSize; + currentBatchSize = optimalSize; + log.info(`πŸ“Š Dynamic batch sizing: ${oldSize} β†’ ${currentBatchSize} records (${safePoints} safe points available)`); + } + const batch = records.slice(i, Math.min(i + currentBatchSize, totalRecords)); - const batchNum = Math.floor(i / currentBatchSize) + 1; + batchCounter++; // Increment actual batch counter const progress = ((i / totalRecords) * 100).toFixed(1); log.progress( - `[${progress}%] Batch ${batchNum} (records ${i + 1}-${Math.min(i + currentBatchSize, totalRecords)}) [size: ${currentBatchSize}, delay: ${currentBatchDelay}ms]` + `[${progress}%] Batch ${batchCounter} (records ${i + 1}-${Math.min(i + currentBatchSize, totalRecords)}) [size: ${currentBatchSize}, delay: ${currentBatchDelay}ms]` ); log.debug(`[publisher.ts] Starting batch: index=${i}, size=${batch.length}`); @@ -141,7 +178,8 @@ export async function publishRecordsWithApplyWrites( updateImportState(importState, i + batch.length - 1, batchSuccessCount, 0); } - // Speed up if we're doing well (after 5 consecutive successes) + // Speed up delay if we're doing well (after 5 consecutive successes) + // Note: Batch size is now controlled dynamically by rate limit quota if (consecutiveSuccesses >= 5 && currentBatchDelay > config.MIN_BATCH_DELAY) { const oldDelay = currentBatchDelay; currentBatchDelay = Math.max( @@ -149,7 +187,7 @@ export async function publishRecordsWithApplyWrites( Math.floor(currentBatchDelay * 0.8) ); if (oldDelay !== currentBatchDelay) { - log.info(`⚑ Speeding up! Delay: ${oldDelay}ms β†’ ${currentBatchDelay}ms`); + log.info(`⚑ Speeding up delay! ${oldDelay}ms β†’ ${currentBatchDelay}ms`); } consecutiveSuccesses = 0; } @@ -261,11 +299,14 @@ export async function publishRecordsWithApplyWrites( updateImportState(importState, i + batch.length - 1, 0, batch.length); } - // If too many consecutive failures, slow down + // If too many consecutive failures, slow down delay + // Note: Batch size is now controlled dynamically by rate limit quota if (consecutiveFailures >= MAX_CONSECUTIVE_FAILURES) { + const oldDelay = currentBatchDelay; currentBatchDelay = Math.min(currentBatchDelay * 2, 10000); - currentBatchSize = Math.max(Math.floor(currentBatchSize / 2), 10); - log.warn(`πŸ“‰ Multiple failures (${consecutiveFailures}): adjusted to ${currentBatchSize} records, ${currentBatchDelay}ms delay`); + if (oldDelay !== currentBatchDelay) { + log.warn(`πŸ“‰ Multiple failures (${consecutiveFailures}): slowing delay to ${currentBatchDelay}ms`); + } log.debug(`[publisher.ts] Slowing down due to consecutive failures`); } diff --git a/src/utils/credentials.ts b/src/utils/credentials.ts index 9ba9c04..c9e1fdc 100644 --- a/src/utils/credentials.ts +++ b/src/utils/credentials.ts @@ -2,6 +2,7 @@ import crypto from 'crypto'; import fs from 'fs'; import path from 'path'; import os from 'os'; +import { getMalachiteStateDir } from './platform.js'; /** * Stored credentials structure @@ -12,6 +13,7 @@ interface StoredCredentials { encryptedPassword: string; iv: string; salt: string; + machineFingerprint: string; createdAt: string; lastUsedAt: string; } @@ -20,7 +22,7 @@ interface StoredCredentials { * Get credentials file path */ function getCredentialsPath(): string { - const credentialsDir = path.join(os.homedir(), '.malachite'); + const credentialsDir = getMalachiteStateDir(); if (!fs.existsSync(credentialsDir)) { fs.mkdirSync(credentialsDir, { recursive: true }); } @@ -28,28 +30,45 @@ function getCredentialsPath(): string { } /** - * Derive encryption key from machine-specific data + * Generate machine-specific fingerprint * This makes credentials machine-specific and non-transferable */ -function deriveKey(salt: Buffer): Buffer { - // Use hostname and username to create a machine-specific key - const machineId = `${os.hostname()}-${os.userInfo().username}`; - - // Use PBKDF2 to derive a strong key +function getMachineFingerprint(): string { + // Collect machine-specific data + const userInfo = os.userInfo(); + const machineData = [ + os.hostname(), + userInfo.username, + os.platform(), + os.arch(), + os.homedir(), + ].join('|'); + + // Hash with SHA-512 for fingerprint + return crypto.createHash('sha512').update(machineData).digest('hex'); +} + +/** + * Derive encryption key from machine-specific data using SHA-512 + * This makes credentials machine-specific and non-transferable + */ +function deriveKey(salt: Buffer, machineFingerprint: string): Buffer { + // Use PBKDF2 with SHA-512 to derive a strong key + // Combining the machine fingerprint with the random salt return crypto.pbkdf2Sync( - machineId, + machineFingerprint, salt, - 100000, // iterations - 32, // key length - 'sha256' - ); + 150000, // iterations (increased for SHA-512) + 64, // key length (64 bytes for SHA-512) + 'sha512' + ).slice(0, 32); // Take first 32 bytes for AES-256 } /** - * Encrypt password + * Encrypt password using AES-256-GCM with SHA-512 derived key */ -function encryptPassword(password: string, salt: Buffer): { encrypted: string; iv: string } { - const key = deriveKey(salt); +function encryptPassword(password: string, salt: Buffer, machineFingerprint: string): { encrypted: string; iv: string } { + const key = deriveKey(salt, machineFingerprint); const iv = crypto.randomBytes(16); const cipher = crypto.createCipheriv('aes-256-gcm', key, iv); @@ -68,10 +87,10 @@ function encryptPassword(password: string, salt: Buffer): { encrypted: string; i } /** - * Decrypt password + * Decrypt password using AES-256-GCM with SHA-512 derived key */ -function decryptPassword(encryptedData: string, iv: string, salt: Buffer): string { - const key = deriveKey(salt); +function decryptPassword(encryptedData: string, iv: string, salt: Buffer, machineFingerprint: string): string { + const key = deriveKey(salt, machineFingerprint); // Extract auth tag (last 32 hex characters = 16 bytes) const authTag = Buffer.from(encryptedData.slice(-32), 'hex'); @@ -87,23 +106,28 @@ function decryptPassword(encryptedData: string, iv: string, salt: Buffer): strin } /** - * Save credentials to disk (encrypted) + * Save credentials to disk (encrypted with SHA-512 and machine-specific salting) + * This function is now called automatically on successful login */ export function saveCredentials(handle: string, password: string): void { const credentialsPath = getCredentialsPath(); + // Generate machine fingerprint + const machineFingerprint = getMachineFingerprint(); + // Generate random salt for this credential - const salt = crypto.randomBytes(32); + const salt = crypto.randomBytes(64); // 64 bytes for SHA-512 - // Encrypt password - const { encrypted, iv } = encryptPassword(password, salt); + // Encrypt password with SHA-512 derived key + const { encrypted, iv } = encryptPassword(password, salt, machineFingerprint); const credentials: StoredCredentials = { - version: 1, + version: 2, // Version 2 for SHA-512 handle, encryptedPassword: encrypted, iv, salt: salt.toString('hex'), + machineFingerprint, createdAt: new Date().toISOString(), lastUsedAt: new Date().toISOString(), }; @@ -118,6 +142,7 @@ export function saveCredentials(handle: string, password: string): void { /** * Load credentials from disk (decrypted) + * Verifies machine fingerprint to prevent credential theft */ export function loadCredentials(): { handle: string; password: string } | null { const credentialsPath = getCredentialsPath(); @@ -130,12 +155,20 @@ export function loadCredentials(): { handle: string; password: string } | null { const data = fs.readFileSync(credentialsPath, 'utf-8'); const credentials: StoredCredentials = JSON.parse(data); - // Decrypt password + // Verify machine fingerprint + const currentFingerprint = getMachineFingerprint(); + if (credentials.machineFingerprint !== currentFingerprint) { + // Credentials were created on a different machine - cannot decrypt + return null; + } + + // Decrypt password using SHA-512 derived key const salt = Buffer.from(credentials.salt, 'hex'); const password = decryptPassword( credentials.encryptedPassword, credentials.iv, - salt + salt, + credentials.machineFingerprint ); // Update last used timestamp diff --git a/src/utils/import-state.ts b/src/utils/import-state.ts index b890ee0..8bec3ce 100644 --- a/src/utils/import-state.ts +++ b/src/utils/import-state.ts @@ -1,9 +1,9 @@ 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 { getMalachiteStateDir } from './platform.js'; /** * Import state for resume functionality @@ -27,7 +27,7 @@ export interface ImportState { * Get the state file path for an import */ export function getStateFilePath(inputFile: string, mode: string): string { - const stateDir = path.join(os.homedir(), '.malachite', 'state'); + const stateDir = path.join(getMalachiteStateDir(), 'state'); // Create state directory if it doesn't exist if (!fs.existsSync(stateDir)) { diff --git a/src/utils/logger.ts b/src/utils/logger.ts index 57e1894..d7f0bc9 100644 --- a/src/utils/logger.ts +++ b/src/utils/logger.ts @@ -4,7 +4,7 @@ import chalk from 'chalk'; import fs from 'fs'; import path from 'path'; -import os from 'os'; +import { getMalachiteLogsDir } from './platform.js'; export enum LogLevel { DEBUG = 0, @@ -39,7 +39,7 @@ export class Logger { enableFileLogging(logDir?: string): void { try { // Default to ~/.malachite/logs if no directory specified - const defaultLogDir = path.join(os.homedir(), '.malachite', 'logs'); + const defaultLogDir = getMalachiteLogsDir(); const logsPath = logDir ? path.resolve(process.cwd(), logDir) : defaultLogDir; if (!fs.existsSync(logsPath)) { fs.mkdirSync(logsPath, { recursive: true }); diff --git a/src/utils/platform.ts b/src/utils/platform.ts index 4179f1a..2054b68 100644 --- a/src/utils/platform.ts +++ b/src/utils/platform.ts @@ -20,45 +20,14 @@ export function getPlatform(): Platform { } /** - * Get the malachite state directory path, respecting platform conventions. - * Windows: %APPDATA%\malachite (or fallback to ~/.malachite) - * macOS: ~/Library/Application Support/malachite - * Linux: ~/.config/malachite (XDG Base Directory spec) + * Get the malachite state directory path. + * Always uses ~/.malachite regardless of OS platform for consistency. */ export function getMalachiteStateDir(): string { - const platform = getPlatform(); const home = os.homedir(); - let stateDir: string; - - switch (platform) { - case 'windows': { - const appdata = process.env.APPDATA; - if (appdata) { - stateDir = path.join(appdata, 'malachite'); - log.debug(`[platform.ts] getMalachiteStateDir() using APPDATA on Windows: ${stateDir}`); - } else { - stateDir = path.join(home, '.malachite'); - log.debug(`[platform.ts] getMalachiteStateDir() APPDATA not set, falling back to: ${stateDir}`); - } - return stateDir; - } - case 'macos': { - stateDir = path.join(home, 'Library', 'Application Support', 'malachite'); - log.debug(`[platform.ts] getMalachiteStateDir() on macOS: ${stateDir}`); - return stateDir; - } - case 'linux': { - const xdgConfig = process.env.XDG_CONFIG_HOME; - if (xdgConfig) { - stateDir = path.join(xdgConfig, 'malachite'); - log.debug(`[platform.ts] getMalachiteStateDir() using XDG_CONFIG_HOME on Linux: ${stateDir}`); - } else { - stateDir = path.join(home, '.config', 'malachite'); - log.debug(`[platform.ts] getMalachiteStateDir() XDG not set, using default on Linux: ${stateDir}`); - } - return stateDir; - } - } + const stateDir = path.join(home, '.malachite'); + log.debug(`[platform.ts] getMalachiteStateDir(): ${stateDir}`); + return stateDir; } /** diff --git a/src/utils/rate-limiter.ts b/src/utils/rate-limiter.ts index ee75eef..97c7c5a 100644 --- a/src/utils/rate-limiter.ts +++ b/src/utils/rate-limiter.ts @@ -301,6 +301,33 @@ export class RateLimiter { } } + /** + * Get safe available points (remaining - headroom buffer) + * This is the amount of quota we can safely use without hitting the headroom threshold + */ + getSafeAvailablePoints(): number { + const state = this.readState(); + if (!state) { + // No state yet - allow a reasonable default + return 300; // Equivalent to 100 records at 3 points each + } + + const now = Math.floor(Date.now() / 1000); + + // Check if window has reset + if (now >= state.resetAt) { + log.debug(`[RateLimiter] Window has reset, full quota available: ${state.limit}`); + return state.limit; + } + + // Calculate headroom and effective remaining + const headroomPoints = Math.floor(state.limit * this.headroomThreshold); + const safePoints = Math.max(0, state.remaining - headroomPoints); + + log.debug(`[RateLimiter] getSafeAvailablePoints: remaining=${state.remaining}, headroom=${headroomPoints}, safe=${safePoints}`); + return safePoints; + } + /** * Get current rate limit status for monitoring */ diff --git a/src/utils/teal-cache.ts b/src/utils/teal-cache.ts index d1f6301..28aae83 100644 --- a/src/utils/teal-cache.ts +++ b/src/utils/teal-cache.ts @@ -1,13 +1,19 @@ import { existsSync, readFileSync, writeFileSync, mkdirSync, readdirSync, unlinkSync } from 'node:fs'; import { join } from 'node:path'; -import { homedir } from 'node:os'; +import { getMalachiteCacheDir } from './platform.js'; /** * Cache configuration */ const CACHE_VERSION = 1; const CACHE_TTL_HOURS = 24; // Cache validity period -const CACHE_DIR = join(homedir(), '.malachite', 'cache'); + +/** + * Get cache directory path + */ +function getCacheDir(): string { + return getMalachiteCacheDir(); +} /** * Cache file structure @@ -23,8 +29,9 @@ interface CacheFile { * Ensure cache directory exists */ function ensureCacheDir(): void { - if (!existsSync(CACHE_DIR)) { - mkdirSync(CACHE_DIR, { recursive: true }); + const cacheDir = getCacheDir(); + if (!existsSync(cacheDir)) { + mkdirSync(cacheDir, { recursive: true }); } } @@ -34,7 +41,7 @@ function ensureCacheDir(): void { function getCachePath(did: string): string { // Sanitize DID for use in filename const sanitized = did.replace(/[^a-zA-Z0-9.-]/g, '_'); - return join(CACHE_DIR, `${sanitized}.json`); + return join(getCacheDir(), `${sanitized}.json`); } /** @@ -153,14 +160,15 @@ export function clearCache(did: string): void { * Clear all caches */ export function clearAllCaches(): void { - if (!existsSync(CACHE_DIR)) { + const cacheDir = getCacheDir(); + if (!existsSync(cacheDir)) { return; } - const files = readdirSync(CACHE_DIR); + const files = readdirSync(cacheDir); for (const file of files) { if (file.endsWith('.json')) { - unlinkSync(join(CACHE_DIR, file)); + unlinkSync(join(cacheDir, file)); } } }