#!/usr/bin/env python3 """SMF Token Dispatcher — sends fan-batch messages to SQS for Lambda processing. Reads SMF refresh tokens from a local JSON file (flat array of strings), packages them into fan-batch messages matching the manifest parser format, and sends them to the collector worker SQS queue. Usage: eval "$(awsume aws_dev -s 2>/dev/null)" python3 scripts/smf_dispatcher.py \ --file ~/Downloads/spotify-smf/spotify.json \ --queue-url https://sqs.us-east-1.amazonaws.com/103233932089/dev-resonance-engine-manifest-files \ --limit 10000 \ --batch-size 100 """ import argparse import json import logging import sys from itertools import islice import boto3 logging.basicConfig( level=logging.INFO, format="%(asctime)s %(levelname)-7s %(message)s", datefmt="%H:%M:%S", ) logger = logging.getLogger(__name__) def load_tokens(path, limit, offset=0): """Load tokens from a local JSON file (flat array of strings).""" logger.info("Loading tokens from %s (offset=%d, limit=%s)", path, offset, limit) tokens = [] skipped = 0 with open(path) as f: for line in f: line = line.strip().strip(",") if line.startswith('"AQ'): if skipped < offset: skipped += 1 continue token = line.strip('"') tokens.append(token) if limit and len(tokens) >= limit: break logger.info("Loaded %d tokens", len(tokens)) return tokens def make_fan_batches(tokens, batch_size): """Package tokens into fan-batch messages matching manifest parser format.""" batches = [] it = iter(enumerate(tokens)) while chunk := list(islice(it, batch_size)): fans = [] for idx, token in chunk: fans.append({ "spotifyUserId": f"smf-{idx}", "refreshToken": token, "partitionKey": "smf-rate-test", "sortKey": f"task:smf-rate-test-{idx}", }) batches.append({ "source": "backfill", "fans": fans, "fan_count": len(fans), }) return batches def send_to_sqs(sqs, queue_url, batches): """Send fan-batch messages to SQS.""" enqueued = 0 # SQS send_message_batch max 10 messages, each max 256KB it = iter(batches) sqs_batch_num = 0 while sqs_batch := list(islice(it, 10)): sqs_batch_num += 1 entries = [] cumulative_bytes = 0 for i, msg in enumerate(sqs_batch): body = json.dumps(msg) msg_bytes = len(body.encode("utf-8")) # Guard against 256KB per-message limit if msg_bytes > 262144: logger.error( "Message %d exceeds 256KB (%d bytes), skipping", enqueued + i, msg_bytes, ) continue # Guard against 1MB total batch limit if cumulative_bytes + msg_bytes > 1_000_000: # Flush what we have and start a new SQS batch if entries: resp = sqs.send_message_batch(QueueUrl=queue_url, Entries=entries) failed = resp.get("Failed", []) if failed: logger.error("SQS batch failures: %s", failed) enqueued += len(entries) - len(failed) entries = [] cumulative_bytes = 0 entries.append({ "Id": str(i), "MessageBody": body, }) cumulative_bytes += msg_bytes if entries: resp = sqs.send_message_batch(QueueUrl=queue_url, Entries=entries) failed = resp.get("Failed", []) if failed: logger.error("SQS batch failures: %s", failed) enqueued += len(entries) - len(failed) if sqs_batch_num % 10 == 0: logger.info("Sent %d SQS batches (%d messages)", sqs_batch_num, enqueued) return enqueued def main(): parser = argparse.ArgumentParser(description="SMF token dispatcher to SQS") parser.add_argument("--file", required=True, help="Local path to token JSON file") parser.add_argument("--queue-url", required=True, help="SQS queue URL") parser.add_argument("--limit", type=int, default=10000, help="Max tokens to dispatch (default: 10000)") parser.add_argument("--offset", type=int, default=0, help="Skip first N tokens (default: 0)") parser.add_argument("--batch-size", type=int, default=100, help="Fans per SQS message (default: 100)") parser.add_argument("--dry-run", action="store_true", help="Show what would be sent without sending") args = parser.parse_args() tokens = load_tokens(args.file, args.limit, args.offset) if not tokens: print("No tokens loaded.") sys.exit(1) batches = make_fan_batches(tokens, args.batch_size) logger.info( "Created %d fan-batch messages (%d fans each, %d total fans)", len(batches), args.batch_size, len(tokens), ) if args.dry_run: logger.info("DRY RUN — would send %d messages to %s", len(batches), args.queue_url) # Show a sample sample = batches[0] logger.info("Sample message: fan_count=%d, first fan=%s", sample["fan_count"], sample["fans"][0]["spotifyUserId"]) msg_size = len(json.dumps(sample).encode("utf-8")) logger.info("Sample message size: %d bytes", msg_size) return sqs = boto3.client("sqs", region_name="us-east-1") enqueued = send_to_sqs(sqs, args.queue_url, batches) logger.info("Done — dispatched %d messages (%d total fans)", enqueued, len(tokens)) if __name__ == "__main__": main()