"""Tests for Manifest Parser Lambda handler (dispatcher mode).""" import gzip import json import os from io import BytesIO from unittest.mock import MagicMock, patch, call import pytest from handler import ( handler, _parse_event, _parse_manifest, _reconstruct_key, _deserialize_ddb_item, _stream_presave_fans, _batch_iterator, _send_batches, ) # --------------------------------------------------------------------------- # Sample data # --------------------------------------------------------------------------- EXPORT_PREFIX = "resonance-engine/ddb-export/01771804222395-8dd0c6a7/" BUCKET = "dev-mymac80" MANIFEST_ENTRIES = [ { "itemCount": 52680, "md5Checksum": "abc123==", "etag": "abc-1", "dataFileS3Key": "dynamodb-exports/resonance-engine/AWSDynamoDB/" "01771804222395-8dd0c6a7/data/aaa111.json.gz", }, { "itemCount": 52681, "md5Checksum": "def456==", "etag": "def-1", "dataFileS3Key": "dynamodb-exports/resonance-engine/AWSDynamoDB/" "01771804222395-8dd0c6a7/data/bbb222.json.gz", }, { "itemCount": 52682, "md5Checksum": "ghi789==", "etag": "ghi-1", "dataFileS3Key": "dynamodb-exports/resonance-engine/AWSDynamoDB/" "01771804222395-8dd0c6a7/data/ccc333.json.gz", }, ] # DDB JSON records (with type markers) PRESAVE_RECORD_DDB = { "Item": { "partitionKey": {"S": "group:album12345"}, "sortKey": {"S": "task:spotify-presave:user123"}, "refreshToken": {"S": "refresh-tok-abc"}, "spotifyUserId": {"S": "spotify-user-123"}, } } PRESAVE_RECORD_DDB_2 = { "Item": { "partitionKey": {"S": "group:album67890"}, "sortKey": {"S": "task:spotify-presave:user456"}, "refreshToken": {"S": "refresh-tok-def"}, "spotifyUserId": {"S": "spotify-user-456"}, } } EMAIL_RECORD_DDB = { "Item": { "partitionKey": {"S": "group:album12345"}, "sortKey": {"S": "task:email-notification:abc"}, "email": {"S": "fan@example.com"}, } } GROUP_RECORD_DDB = { "Item": { "partitionKey": {"S": "group:album12345"}, "sortKey": {"S": "group:album12345"}, "count": {"N": "42"}, } } APPLE_RECORD_DDB = { "Item": { "partitionKey": {"S": "group:album12345"}, "sortKey": {"S": "task:apple-music-presave:xyz"}, "appleToken": {"S": "apple-tok"}, } } def _manifest_body(entries=None): """Build manifest-files.json body (JSON Lines).""" if entries is None: entries = MANIFEST_ENTRIES return "\n".join(json.dumps(e) for e in entries) def _make_gzip_body(records): """Build a gzip-compressed body from a list of DDB JSON records.""" lines = [json.dumps(r) for r in records] raw = "\n".join(lines).encode("utf-8") buf = BytesIO() with gzip.GzipFile(fileobj=buf, mode="wb") as gz: gz.write(raw) return buf.getvalue() def _mock_s3_gzip_response(body_bytes): """Create a mock S3 get_object response with BytesIO body for gzip.""" return {"Body": BytesIO(body_bytes)} # --------------------------------------------------------------------------- # _parse_event # --------------------------------------------------------------------------- class TestParseEvent: def test_s3_event(self): event = { "Records": [ { "s3": { "bucket": {"name": BUCKET}, "object": { "key": f"{EXPORT_PREFIX}manifest-summary.json" }, } } ] } bucket, prefix = _parse_event(event) assert bucket == BUCKET assert prefix == EXPORT_PREFIX def test_s3_event_url_encoded_key(self): encoded_key = ( "resonance-engine/ddb-export/" "01771804222395-8dd0c6a7/manifest-summary.json" ).replace("/", "%2F") event = { "Records": [ { "s3": { "bucket": {"name": BUCKET}, "object": {"key": encoded_key}, } } ] } bucket, prefix = _parse_event(event) assert bucket == BUCKET assert prefix == EXPORT_PREFIX def test_manual_invocation(self): event = {"bucket": BUCKET, "export_prefix": EXPORT_PREFIX} bucket, prefix = _parse_event(event) assert bucket == BUCKET assert prefix == EXPORT_PREFIX def test_manual_invocation_adds_trailing_slash(self): event = {"bucket": BUCKET, "export_prefix": EXPORT_PREFIX.rstrip("/")} _, prefix = _parse_event(event) assert prefix.endswith("/") def test_manual_invocation_default_bucket(self): with patch.dict(os.environ, {"DDB_EXPORT_BUCKET": "other-bucket"}): event = {"export_prefix": EXPORT_PREFIX} bucket, _ = _parse_event(event) assert bucket == "other-bucket" def test_manual_invocation_fallback_bucket(self): env = {k: v for k, v in os.environ.items() if k != "DDB_EXPORT_BUCKET"} with patch.dict(os.environ, env, clear=True): event = {"export_prefix": EXPORT_PREFIX} bucket, _ = _parse_event(event) assert bucket == "dev-mymac80" def test_missing_export_prefix_raises_valueerror(self): with pytest.raises(ValueError, match="export_prefix"): _parse_event({"bucket": BUCKET}) def test_empty_export_prefix_raises_valueerror(self): with pytest.raises(ValueError, match="export_prefix"): _parse_event({"bucket": BUCKET, "export_prefix": ""}) # --------------------------------------------------------------------------- # _parse_manifest # --------------------------------------------------------------------------- class TestParseManifest: def test_parses_json_lines(self): files = _parse_manifest(_manifest_body()) assert len(files) == 3 assert files[0]["itemCount"] == 52680 assert files[1]["dataFileS3Key"].endswith("bbb222.json.gz") def test_skips_blank_lines(self): body = "\n" + _manifest_body() + "\n\n" assert len(_parse_manifest(body)) == 3 def test_empty_manifest(self): assert _parse_manifest("") == [] def test_skips_entries_without_required_fields(self): incomplete = json.dumps({"itemCount": 100}) complete = json.dumps(MANIFEST_ENTRIES[0]) body = f"{incomplete}\n{complete}" files = _parse_manifest(body) assert len(files) == 1 assert files[0]["itemCount"] == 52680 # --------------------------------------------------------------------------- # _reconstruct_key # --------------------------------------------------------------------------- class TestReconstructKey: def test_extracts_data_relative_path(self): original = ( "dynamodb-exports/resonance-engine/AWSDynamoDB/" "01771804222395-8dd0c6a7/data/aaa111.json.gz" ) result = _reconstruct_key(EXPORT_PREFIX, original) assert result == f"{EXPORT_PREFIX}data/aaa111.json.gz" def test_simple_aws_data_key(self): original = "AWSDynamoDB/data/file.json.gz" result = _reconstruct_key(EXPORT_PREFIX, original) assert result == f"{EXPORT_PREFIX}data/file.json.gz" def test_nested_data_directory(self): original = "prefix/AWSDynamoDB/exportid/data/nested/file.json.gz" result = _reconstruct_key(EXPORT_PREFIX, original) assert result == f"{EXPORT_PREFIX}data/nested/file.json.gz" def test_fallback_no_data_prefix(self): original = "some/weird/path/file.json.gz" result = _reconstruct_key(EXPORT_PREFIX, original) assert result == f"{EXPORT_PREFIX}data/file.json.gz" def test_does_not_match_data_as_substring(self): original = "my-database/export-data/file.json.gz" result = _reconstruct_key(EXPORT_PREFIX, original) assert result == f"{EXPORT_PREFIX}data/file.json.gz" def test_key_starting_with_data(self): original = "data/file.json.gz" result = _reconstruct_key(EXPORT_PREFIX, original) assert result == f"{EXPORT_PREFIX}data/file.json.gz" # --------------------------------------------------------------------------- # _deserialize_ddb_item # --------------------------------------------------------------------------- class TestDeserializeDdbItem: def test_string_type(self): assert _deserialize_ddb_item({"S": "hello"}) == "hello" def test_number_int(self): assert _deserialize_ddb_item({"N": "42"}) == 42 def test_number_float(self): assert _deserialize_ddb_item({"N": "3.14"}) == 3.14 def test_boolean_true(self): assert _deserialize_ddb_item({"BOOL": True}) is True def test_boolean_false(self): assert _deserialize_ddb_item({"BOOL": False}) is False def test_null(self): assert _deserialize_ddb_item({"NULL": True}) is None def test_list(self): raw = {"L": [{"S": "a"}, {"N": "1"}, {"BOOL": True}]} assert _deserialize_ddb_item(raw) == ["a", 1, True] def test_nested_map(self): raw = {"M": {"name": {"S": "Alice"}, "age": {"N": "30"}}} assert _deserialize_ddb_item(raw) == {"name": "Alice", "age": 30} def test_full_presave_record(self): result = _deserialize_ddb_item(PRESAVE_RECORD_DDB["Item"]) assert result == { "partitionKey": "group:album12345", "sortKey": "task:spotify-presave:user123", "refreshToken": "refresh-tok-abc", "spotifyUserId": "spotify-user-123", } def test_passthrough_non_dict(self): assert _deserialize_ddb_item("plain") == "plain" assert _deserialize_ddb_item(123) == 123 def test_multi_key_dict_is_not_type_descriptor(self): raw = { "partitionKey": {"S": "pk"}, "sortKey": {"S": "sk"}, } result = _deserialize_ddb_item(raw) assert result == {"partitionKey": "pk", "sortKey": "sk"} # --------------------------------------------------------------------------- # _stream_presave_fans # --------------------------------------------------------------------------- class TestStreamPresaveFans: def test_yields_presave_records(self): body = _make_gzip_body([PRESAVE_RECORD_DDB]) mock_s3 = MagicMock() mock_s3.get_object.return_value = _mock_s3_gzip_response(body) fans = list(_stream_presave_fans(mock_s3, BUCKET, "data/f.json.gz")) assert len(fans) == 1 assert fans[0] == { "spotifyUserId": "spotify-user-123", "refreshToken": "refresh-tok-abc", "partitionKey": "group:album12345", "sortKey": "task:spotify-presave:user123", } def test_filters_non_presave_records(self): body = _make_gzip_body([ EMAIL_RECORD_DDB, PRESAVE_RECORD_DDB, GROUP_RECORD_DDB, APPLE_RECORD_DDB, ]) mock_s3 = MagicMock() mock_s3.get_object.return_value = _mock_s3_gzip_response(body) fans = list(_stream_presave_fans(mock_s3, BUCKET, "data/f.json.gz")) assert len(fans) == 1 assert fans[0]["spotifyUserId"] == "spotify-user-123" def test_extracts_only_four_fields(self): """Even if the DDB record has extra fields, only 4 are yielded.""" record = { "Item": { "partitionKey": {"S": "group:album99"}, "sortKey": {"S": "task:spotify-presave:userX"}, "refreshToken": {"S": "tok-X"}, "spotifyUserId": {"S": "user-X"}, "extraField": {"S": "should-not-appear"}, "anotherField": {"N": "999"}, } } body = _make_gzip_body([record]) mock_s3 = MagicMock() mock_s3.get_object.return_value = _mock_s3_gzip_response(body) fans = list(_stream_presave_fans(mock_s3, BUCKET, "data/f.json.gz")) assert len(fans) == 1 assert set(fans[0].keys()) == { "spotifyUserId", "refreshToken", "partitionKey", "sortKey", } def test_empty_file_yields_nothing(self): body = _make_gzip_body([]) mock_s3 = MagicMock() mock_s3.get_object.return_value = _mock_s3_gzip_response(body) fans = list(_stream_presave_fans(mock_s3, BUCKET, "data/f.json.gz")) assert fans == [] def test_skips_invalid_json_lines(self): """Invalid JSON lines are skipped, valid presave records are yielded.""" lines = "not-json\n" + json.dumps(PRESAVE_RECORD_DDB) buf = BytesIO() with gzip.GzipFile(fileobj=buf, mode="wb") as gz: gz.write(lines.encode("utf-8")) mock_s3 = MagicMock() mock_s3.get_object.return_value = _mock_s3_gzip_response(buf.getvalue()) fans = list(_stream_presave_fans(mock_s3, BUCKET, "data/f.json.gz")) assert len(fans) == 1 assert fans[0]["spotifyUserId"] == "spotify-user-123" def test_multiple_presave_records(self): body = _make_gzip_body([PRESAVE_RECORD_DDB, PRESAVE_RECORD_DDB_2]) mock_s3 = MagicMock() mock_s3.get_object.return_value = _mock_s3_gzip_response(body) fans = list(_stream_presave_fans(mock_s3, BUCKET, "data/f.json.gz")) assert len(fans) == 2 user_ids = {f["spotifyUserId"] for f in fans} assert user_ids == {"spotify-user-123", "spotify-user-456"} def test_missing_required_fields_skipped(self): """A presave record missing spotifyUserId is skipped.""" record = { "Item": { "partitionKey": {"S": "group:album1"}, "sortKey": {"S": "task:spotify-presave:nouser"}, "refreshToken": {"S": "tok"}, } } body = _make_gzip_body([record]) mock_s3 = MagicMock() mock_s3.get_object.return_value = _mock_s3_gzip_response(body) fans = list(_stream_presave_fans(mock_s3, BUCKET, "data/f.json.gz")) assert len(fans) == 0 def test_missing_refresh_token_skipped(self): """A presave record missing refreshToken is skipped.""" record = { "Item": { "partitionKey": {"S": "group:album1"}, "sortKey": {"S": "task:spotify-presave:user1"}, "spotifyUserId": {"S": "user-1"}, } } body = _make_gzip_body([record]) mock_s3 = MagicMock() mock_s3.get_object.return_value = _mock_s3_gzip_response(body) fans = list(_stream_presave_fans(mock_s3, BUCKET, "data/f.json.gz")) assert len(fans) == 0 def test_missing_partition_key_skipped(self): """A presave record missing partitionKey is skipped.""" record = { "Item": { "sortKey": {"S": "task:spotify-presave:user1"}, "spotifyUserId": {"S": "user-1"}, "refreshToken": {"S": "tok-1"}, } } body = _make_gzip_body([record]) mock_s3 = MagicMock() mock_s3.get_object.return_value = _mock_s3_gzip_response(body) fans = list(_stream_presave_fans(mock_s3, BUCKET, "data/f.json.gz")) assert len(fans) == 0 # --------------------------------------------------------------------------- # _batch_iterator # --------------------------------------------------------------------------- class TestBatchIterator: def test_exact_batch_size(self): items = list(range(10)) batches = list(_batch_iterator(items, 5)) assert len(batches) == 2 assert batches[0] == [0, 1, 2, 3, 4] assert batches[1] == [5, 6, 7, 8, 9] def test_remainder_batch(self): items = list(range(7)) batches = list(_batch_iterator(items, 3)) assert len(batches) == 3 assert batches[0] == [0, 1, 2] assert batches[1] == [3, 4, 5] assert batches[2] == [6] def test_empty_input(self): batches = list(_batch_iterator([], 5)) assert batches == [] def test_single_item(self): batches = list(_batch_iterator([42], 5)) assert batches == [[42]] def test_items_fewer_than_batch_size(self): batches = list(_batch_iterator([1, 2], 500)) assert len(batches) == 1 assert batches[0] == [1, 2] def test_works_with_generator(self): def gen(): for i in range(5): yield i batches = list(_batch_iterator(gen(), 2)) assert batches == [[0, 1], [2, 3], [4]] # --------------------------------------------------------------------------- # _send_batches # --------------------------------------------------------------------------- class TestSendBatches: def test_sends_single_message(self): mock_sqs = MagicMock() mock_sqs.send_message_batch.return_value = { "Successful": [{"Id": "0"}], "Failed": [], } msg = {"source": "backfill", "fans": [{"spotifyUserId": "u1"}], "fan_count": 1} result = _send_batches( mock_sqs, "https://sqs.example.com/queue", [msg], ) assert result == 1 mock_sqs.send_message_batch.assert_called_once() entries = mock_sqs.send_message_batch.call_args.kwargs["Entries"] body = json.loads(entries[0]["MessageBody"]) assert body["source"] == "backfill" assert body["fan_count"] == 1 def test_batches_in_groups_of_10(self): mock_sqs = MagicMock() mock_sqs.send_message_batch.return_value = { "Successful": [], "Failed": [], } messages = [{"id": i} for i in range(12)] result = _send_batches(mock_sqs, "https://sqs/q", messages) assert result == 12 assert mock_sqs.send_message_batch.call_count == 2 first = mock_sqs.send_message_batch.call_args_list[0].kwargs["Entries"] second = mock_sqs.send_message_batch.call_args_list[1].kwargs["Entries"] assert len(first) == 10 assert len(second) == 2 def test_raises_on_failure(self): mock_sqs = MagicMock() mock_sqs.send_message_batch.return_value = { "Successful": [], "Failed": [{"Id": "0", "Code": "InternalError", "Message": "oops"}], } with pytest.raises(RuntimeError, match="SQS batch send had 1 failure"): _send_batches(mock_sqs, "https://sqs/q", [{"test": True}]) # --------------------------------------------------------------------------- # handler (integration) # --------------------------------------------------------------------------- def _build_mock_s3_for_handler(manifest_body, data_file_records_map=None): """Build a mock S3 client that returns manifest + gzip data files. Args: manifest_body: String body for manifest-files.json. data_file_records_map: Dict mapping S3 key suffix to list of DDB records. If None, all data file requests return empty gzip. Returns: MagicMock: S3 client mock. """ if data_file_records_map is None: data_file_records_map = {} mock_s3 = MagicMock() def get_object_side_effect(Bucket, Key): if Key.endswith("manifest-files.json"): return { "Body": MagicMock( read=MagicMock( return_value=manifest_body.encode("utf-8") ) ) } # It is a data file (.json.gz) — find matching records for suffix, records in data_file_records_map.items(): if Key.endswith(suffix): body = _make_gzip_body(records) return _mock_s3_gzip_response(body) # Default: empty gzip body = _make_gzip_body([]) return _mock_s3_gzip_response(body) mock_s3.get_object.side_effect = get_object_side_effect return mock_s3 class TestHandler: @patch.dict( os.environ, {"SQS_QUEUE_URL": "https://sqs.us-east-1.amazonaws.com/123/queue"}, ) @patch("handler.boto3") def test_dispatches_presave_fans_as_batches(self, mock_boto3): """Presave records are extracted and sent as fan batches.""" records = [PRESAVE_RECORD_DDB, EMAIL_RECORD_DDB, PRESAVE_RECORD_DDB_2] mock_s3 = _build_mock_s3_for_handler( _manifest_body(MANIFEST_ENTRIES[:1]), {"aaa111.json.gz": records}, ) mock_sqs = MagicMock() mock_sqs.send_message_batch.return_value = { "Successful": [], "Failed": [], } mock_boto3.client.side_effect = lambda svc: { "s3": mock_s3, "sqs": mock_sqs, }[svc] result = handler( {"bucket": BUCKET, "export_prefix": EXPORT_PREFIX}, None ) assert result["statusCode"] == 200 assert result["fansDispatched"] == 2 assert result["batchesSent"] == 1 # Verify SQS message format entries = mock_sqs.send_message_batch.call_args.kwargs["Entries"] assert len(entries) == 1 body = json.loads(entries[0]["MessageBody"]) assert body["source"] == "backfill" assert body["fan_count"] == 2 assert len(body["fans"]) == 2 user_ids = {f["spotifyUserId"] for f in body["fans"]} assert user_ids == {"spotify-user-123", "spotify-user-456"} @patch.dict( os.environ, {"SQS_QUEUE_URL": "https://sqs.us-east-1.amazonaws.com/123/queue"}, ) @patch("handler.boto3") def test_empty_manifest_returns_zero(self, mock_boto3): mock_s3 = _build_mock_s3_for_handler("") mock_sqs = MagicMock() mock_boto3.client.side_effect = lambda svc: { "s3": mock_s3, "sqs": mock_sqs, }[svc] result = handler( {"bucket": BUCKET, "export_prefix": EXPORT_PREFIX}, None ) assert result == { "statusCode": 200, "fansDispatched": 0, "batchesSent": 0, } mock_sqs.send_message_batch.assert_not_called() @patch.dict( os.environ, {"SQS_QUEUE_URL": "https://sqs.us-east-1.amazonaws.com/123/queue"}, ) @patch("handler.boto3") def test_s3_event_trigger(self, mock_boto3): mock_s3 = _build_mock_s3_for_handler( _manifest_body(MANIFEST_ENTRIES[:1]), {"aaa111.json.gz": [PRESAVE_RECORD_DDB]}, ) mock_sqs = MagicMock() mock_sqs.send_message_batch.return_value = { "Successful": [], "Failed": [], } mock_boto3.client.side_effect = lambda svc: { "s3": mock_s3, "sqs": mock_sqs, }[svc] event = { "Records": [ { "s3": { "bucket": {"name": BUCKET}, "object": { "key": f"{EXPORT_PREFIX}manifest-summary.json" }, } } ] } result = handler(event, None) assert result["statusCode"] == 200 assert result["fansDispatched"] == 1 # Verify the manifest was read from the right key manifest_call = mock_s3.get_object.call_args_list[0] assert manifest_call.kwargs["Key"] == f"{EXPORT_PREFIX}manifest-files.json" @patch.dict( os.environ, {"SQS_QUEUE_URL": "https://sqs.us-east-1.amazonaws.com/123/queue"}, ) @patch("handler.boto3") def test_s3_read_failure_propagates(self, mock_boto3): mock_s3 = MagicMock() mock_s3.get_object.side_effect = Exception("NoSuchKey") mock_sqs = MagicMock() mock_boto3.client.side_effect = lambda svc: { "s3": mock_s3, "sqs": mock_sqs, }[svc] with pytest.raises(Exception, match="NoSuchKey"): handler( {"bucket": BUCKET, "export_prefix": EXPORT_PREFIX}, None ) @patch.dict( os.environ, {"SQS_QUEUE_URL": "https://sqs.us-east-1.amazonaws.com/123/queue"}, ) @patch("handler.boto3") def test_sqs_message_schema(self, mock_boto3): """Verify every SQS message has the expected backfill schema.""" mock_s3 = _build_mock_s3_for_handler( _manifest_body(MANIFEST_ENTRIES[:1]), {"aaa111.json.gz": [PRESAVE_RECORD_DDB]}, ) mock_sqs = MagicMock() mock_sqs.send_message_batch.return_value = { "Successful": [], "Failed": [], } mock_boto3.client.side_effect = lambda svc: { "s3": mock_s3, "sqs": mock_sqs, }[svc] handler({"bucket": BUCKET, "export_prefix": EXPORT_PREFIX}, None) entries = mock_sqs.send_message_batch.call_args.kwargs["Entries"] for entry in entries: body = json.loads(entry["MessageBody"]) assert set(body.keys()) == {"source", "fans", "fan_count"} assert body["source"] == "backfill" assert isinstance(body["fans"], list) assert isinstance(body["fan_count"], int) assert body["fan_count"] == len(body["fans"]) for fan in body["fans"]: assert set(fan.keys()) == { "spotifyUserId", "refreshToken", "partitionKey", "sortKey", } @patch.dict( os.environ, {"SQS_QUEUE_URL": "https://sqs.us-east-1.amazonaws.com/123/queue"}, ) @patch("handler.boto3") def test_sqs_batch_failure_raises(self, mock_boto3): mock_s3 = _build_mock_s3_for_handler( _manifest_body(MANIFEST_ENTRIES[:1]), {"aaa111.json.gz": [PRESAVE_RECORD_DDB]}, ) mock_sqs = MagicMock() mock_sqs.send_message_batch.return_value = { "Successful": [], "Failed": [ {"Id": "0", "Code": "InternalError", "Message": "oops"} ], } mock_boto3.client.side_effect = lambda svc: { "s3": mock_s3, "sqs": mock_sqs, }[svc] with pytest.raises(RuntimeError, match="SQS batch send had 1 failure"): handler( {"bucket": BUCKET, "export_prefix": EXPORT_PREFIX}, None ) @patch.dict( os.environ, { "SQS_QUEUE_URL": "https://sqs.us-east-1.amazonaws.com/123/queue", "FAN_BATCH_SIZE": "500", }, ) @patch("handler.boto3") def test_batching_1200_fans(self, mock_boto3): """1200 fans with batch_size=500 produces 3 batches (500+500+200).""" # Build 1200 presave records with unique user IDs records = [] for i in range(1200): records.append({ "Item": { "partitionKey": {"S": f"group:album{i}"}, "sortKey": {"S": f"task:spotify-presave:user{i}"}, "refreshToken": {"S": f"tok-{i}"}, "spotifyUserId": {"S": f"user-{i}"}, } }) mock_s3 = _build_mock_s3_for_handler( _manifest_body(MANIFEST_ENTRIES[:1]), {"aaa111.json.gz": records}, ) mock_sqs = MagicMock() mock_sqs.send_message_batch.return_value = { "Successful": [], "Failed": [], } mock_boto3.client.side_effect = lambda svc: { "s3": mock_s3, "sqs": mock_sqs, }[svc] result = handler( {"bucket": BUCKET, "export_prefix": EXPORT_PREFIX}, None ) assert result["fansDispatched"] == 1200 assert result["batchesSent"] == 3 # 3 fan batches are buffered into a single _send_batches call (< 10) assert mock_sqs.send_message_batch.call_count == 1 # Verify batch sizes: 500, 500, 200 batch_sizes = [] for sqs_call in mock_sqs.send_message_batch.call_args_list: entries = sqs_call.kwargs["Entries"] for entry in entries: body = json.loads(entry["MessageBody"]) batch_sizes.append(body["fan_count"]) assert batch_sizes == [500, 500, 200] @patch.dict( os.environ, {"SQS_QUEUE_URL": "https://sqs.us-east-1.amazonaws.com/123/queue"}, ) @patch("handler.boto3") def test_mixed_record_types_only_presave_dispatched(self, mock_boto3): """Only spotify-presave records are dispatched, others are filtered.""" records = [ EMAIL_RECORD_DDB, PRESAVE_RECORD_DDB, GROUP_RECORD_DDB, PRESAVE_RECORD_DDB_2, APPLE_RECORD_DDB, ] mock_s3 = _build_mock_s3_for_handler( _manifest_body(MANIFEST_ENTRIES[:1]), {"aaa111.json.gz": records}, ) mock_sqs = MagicMock() mock_sqs.send_message_batch.return_value = { "Successful": [], "Failed": [], } mock_boto3.client.side_effect = lambda svc: { "s3": mock_s3, "sqs": mock_sqs, }[svc] result = handler( {"bucket": BUCKET, "export_prefix": EXPORT_PREFIX}, None ) assert result["fansDispatched"] == 2 entries = mock_sqs.send_message_batch.call_args.kwargs["Entries"] body = json.loads(entries[0]["MessageBody"]) user_ids = {f["spotifyUserId"] for f in body["fans"]} assert user_ids == {"spotify-user-123", "spotify-user-456"} @patch.dict( os.environ, {"SQS_QUEUE_URL": "https://sqs.us-east-1.amazonaws.com/123/queue"}, ) @patch("handler.boto3") def test_multiple_data_files(self, mock_boto3): """Multiple files in manifest are each streamed and dispatched.""" mock_s3 = _build_mock_s3_for_handler( _manifest_body(MANIFEST_ENTRIES[:2]), { "aaa111.json.gz": [PRESAVE_RECORD_DDB], "bbb222.json.gz": [PRESAVE_RECORD_DDB_2], }, ) mock_sqs = MagicMock() mock_sqs.send_message_batch.return_value = { "Successful": [], "Failed": [], } mock_boto3.client.side_effect = lambda svc: { "s3": mock_s3, "sqs": mock_sqs, }[svc] result = handler( {"bucket": BUCKET, "export_prefix": EXPORT_PREFIX}, None ) assert result["fansDispatched"] == 2 assert result["batchesSent"] == 2 # Two data file reads + one manifest read = 3 get_object calls assert mock_s3.get_object.call_count == 3 @patch.dict( os.environ, {"SQS_QUEUE_URL": "https://sqs.us-east-1.amazonaws.com/123/queue"}, ) @patch("handler.boto3") def test_file_with_no_presave_records(self, mock_boto3): """A data file with no presave records produces zero batches.""" mock_s3 = _build_mock_s3_for_handler( _manifest_body(MANIFEST_ENTRIES[:1]), {"aaa111.json.gz": [EMAIL_RECORD_DDB, GROUP_RECORD_DDB]}, ) mock_sqs = MagicMock() mock_boto3.client.side_effect = lambda svc: { "s3": mock_s3, "sqs": mock_sqs, }[svc] result = handler( {"bucket": BUCKET, "export_prefix": EXPORT_PREFIX}, None ) assert result["fansDispatched"] == 0 assert result["batchesSent"] == 0 mock_sqs.send_message_batch.assert_not_called() @patch.dict( os.environ, { "SQS_QUEUE_URL": "https://sqs.us-east-1.amazonaws.com/123/queue", "FAN_BATCH_SIZE": "2", }, ) @patch("handler.boto3") def test_custom_batch_size_from_env(self, mock_boto3): """FAN_BATCH_SIZE env var controls batch grouping.""" records = [PRESAVE_RECORD_DDB, PRESAVE_RECORD_DDB_2, { "Item": { "partitionKey": {"S": "group:album999"}, "sortKey": {"S": "task:spotify-presave:user999"}, "refreshToken": {"S": "tok-999"}, "spotifyUserId": {"S": "user-999"}, } }] mock_s3 = _build_mock_s3_for_handler( _manifest_body(MANIFEST_ENTRIES[:1]), {"aaa111.json.gz": records}, ) mock_sqs = MagicMock() mock_sqs.send_message_batch.return_value = { "Successful": [], "Failed": [], } mock_boto3.client.side_effect = lambda svc: { "s3": mock_s3, "sqs": mock_sqs, }[svc] result = handler( {"bucket": BUCKET, "export_prefix": EXPORT_PREFIX}, None ) assert result["fansDispatched"] == 3 assert result["batchesSent"] == 2 # batch_size=2 → 2+1 # 2 fan batches buffered into 1 SQS call assert mock_sqs.send_message_batch.call_count == 1 entries = mock_sqs.send_message_batch.call_args.kwargs["Entries"] batch_sizes = [] for entry in entries: body = json.loads(entry["MessageBody"]) batch_sizes.append(body["fan_count"]) assert batch_sizes == [2, 1] @patch.dict( os.environ, { "SQS_QUEUE_URL": "https://sqs.us-east-1.amazonaws.com/123/queue", "FAN_BATCH_SIZE": "1", }, ) @patch("handler.boto3") def test_message_buffer_flushes_at_10(self, mock_boto3): """12 fan batches (batch_size=1) produce 2 SQS calls: 10 + 2.""" records = [] for i in range(12): records.append({ "Item": { "partitionKey": {"S": f"group:album{i}"}, "sortKey": {"S": f"task:spotify-presave:u{i}"}, "refreshToken": {"S": f"tok-{i}"}, "spotifyUserId": {"S": f"user-{i}"}, } }) mock_s3 = _build_mock_s3_for_handler( _manifest_body(MANIFEST_ENTRIES[:1]), {"aaa111.json.gz": records}, ) mock_sqs = MagicMock() mock_sqs.send_message_batch.return_value = { "Successful": [], "Failed": [], } mock_boto3.client.side_effect = lambda svc: { "s3": mock_s3, "sqs": mock_sqs, }[svc] result = handler( {"bucket": BUCKET, "export_prefix": EXPORT_PREFIX}, None ) assert result["fansDispatched"] == 12 assert result["batchesSent"] == 12 # Buffer flushes at 10, remainder of 2 flushes at end assert mock_sqs.send_message_batch.call_count == 2 first_entries = mock_sqs.send_message_batch.call_args_list[0].kwargs["Entries"] second_entries = mock_sqs.send_message_batch.call_args_list[1].kwargs["Entries"] assert len(first_entries) == 10 assert len(second_entries) == 2 @patch.dict( os.environ, { "SQS_QUEUE_URL": "https://sqs.us-east-1.amazonaws.com/123/queue", "FAN_BATCH_SIZE": "500", }, ) @patch("handler.boto3") def test_message_buffer_flushes_on_byte_size(self, mock_boto3): """Buffer flushes before exceeding 1MB cumulative payload limit. With batch_size=500 and ~1KB refresh tokens, each fan-batch message is ~600KB. Two messages would total ~1.2MB, so the handler must flush after the first before adding the second. """ # Long refresh tokens (~1KB) so each 500-fan message is ~600KB long_token = "tok-" + "x" * 1000 records = [] for i in range(1500): records.append({ "Item": { "partitionKey": {"S": f"group:album{i}"}, "sortKey": {"S": f"task:spotify-presave:u{i}"}, "refreshToken": {"S": f"{long_token}-{i}"}, "spotifyUserId": {"S": f"spotify-user-{i:06d}"}, } }) mock_s3 = _build_mock_s3_for_handler( _manifest_body(MANIFEST_ENTRIES[:1]), {"aaa111.json.gz": records}, ) mock_sqs = MagicMock() mock_sqs.send_message_batch.return_value = { "Successful": [], "Failed": [], } mock_boto3.client.side_effect = lambda svc: { "s3": mock_s3, "sqs": mock_sqs, }[svc] result = handler( {"bucket": BUCKET, "export_prefix": EXPORT_PREFIX}, None ) assert result["fansDispatched"] == 1500 assert result["batchesSent"] == 3 # 500 + 500 + 500 # 3 messages each ~600KB — all under the 10-message limit, but # byte-size guard forces flush after each message. Without the # byte-size tracking, all 3 would go in one call. assert mock_sqs.send_message_batch.call_count == 3 # Verify no single SQS call exceeds 1MB total payload for sqs_call in mock_sqs.send_message_batch.call_args_list: entries = sqs_call.kwargs["Entries"] total_bytes = sum( len(e["MessageBody"].encode("utf-8")) for e in entries ) assert total_bytes <= 1_000_000, ( f"SQS batch payload {total_bytes} bytes exceeds 1MB limit" ) @patch.dict( os.environ, { "SQS_QUEUE_URL": "https://sqs.us-east-1.amazonaws.com/123/queue", "FAN_BATCH_SIZE": "not-a-number", }, ) @patch("handler.boto3") def test_invalid_batch_size_falls_back_to_default(self, mock_boto3): """Non-numeric FAN_BATCH_SIZE falls back to 500.""" mock_s3 = _build_mock_s3_for_handler( _manifest_body(MANIFEST_ENTRIES[:1]), {"aaa111.json.gz": [PRESAVE_RECORD_DDB]}, ) mock_sqs = MagicMock() mock_sqs.send_message_batch.return_value = { "Successful": [], "Failed": [], } mock_boto3.client.side_effect = lambda svc: { "s3": mock_s3, "sqs": mock_sqs, }[svc] result = handler( {"bucket": BUCKET, "export_prefix": EXPORT_PREFIX}, None ) # Should not crash — falls back to 500 assert result["statusCode"] == 200 assert result["fansDispatched"] == 1 @patch.dict( os.environ, { "SQS_QUEUE_URL": "https://sqs.us-east-1.amazonaws.com/123/queue", "FAN_BATCH_SIZE": "0", }, ) @patch("handler.boto3") def test_zero_batch_size_clamped_to_one(self, mock_boto3): """FAN_BATCH_SIZE=0 is clamped to 1.""" mock_s3 = _build_mock_s3_for_handler( _manifest_body(MANIFEST_ENTRIES[:1]), {"aaa111.json.gz": [PRESAVE_RECORD_DDB, PRESAVE_RECORD_DDB_2]}, ) mock_sqs = MagicMock() mock_sqs.send_message_batch.return_value = { "Successful": [], "Failed": [], } mock_boto3.client.side_effect = lambda svc: { "s3": mock_s3, "sqs": mock_sqs, }[svc] result = handler( {"bucket": BUCKET, "export_prefix": EXPORT_PREFIX}, None ) assert result["fansDispatched"] == 2 # batch_size=1 → 2 batches (1 fan each) assert result["batchesSent"] == 2