| 
									
										
										
										
											2020-10-06 07:54:49 +02:00
										 |  |  | import pytest | 
					
						
							| 
									
										
										
										
											2018-11-08 13:22:24 -05:00
										 |  |  | 
 | 
					
						
							| 
									
										
										
										
											2018-11-07 17:10:00 -05:00
										 |  |  | import boto3 | 
					
						
							| 
									
										
										
										
											2022-03-09 16:57:25 -01:00
										 |  |  | from moto import mock_dynamodb, mock_dynamodbstreams | 
					
						
							| 
									
										
										
										
											2018-11-07 17:10:00 -05:00
										 |  |  | 
 | 
					
						
							|  |  |  | 
 | 
					
						
							| 
									
										
										
										
											2018-11-08 13:22:24 -05:00
										 |  |  | class TestCore: | 
					
						
							| 
									
										
										
										
											2018-11-07 17:10:00 -05:00
										 |  |  |     stream_arn = None | 
					
						
							|  |  |  |     mocks = [] | 
					
						
							| 
									
										
										
										
											2019-10-31 08:44:26 -07:00
										 |  |  | 
 | 
					
						
							| 
									
										
										
										
											2022-10-27 17:37:11 -07:00
										 |  |  |     def setup_method(self): | 
					
						
							| 
									
										
										
										
											2022-03-09 16:57:25 -01:00
										 |  |  |         self.mocks = [mock_dynamodb(), mock_dynamodbstreams()] | 
					
						
							| 
									
										
										
										
											2018-11-07 17:10:00 -05:00
										 |  |  |         for m in self.mocks: | 
					
						
							|  |  |  |             m.start() | 
					
						
							| 
									
										
										
										
											2019-10-31 08:44:26 -07:00
										 |  |  | 
 | 
					
						
							| 
									
										
										
										
											2018-11-07 17:10:00 -05:00
										 |  |  |         # create a table with a stream | 
					
						
							|  |  |  |         conn = boto3.client("dynamodb", region_name="us-east-1") | 
					
						
							|  |  |  | 
 | 
					
						
							|  |  |  |         resp = conn.create_table( | 
					
						
							|  |  |  |             TableName="test-streams", | 
					
						
							|  |  |  |             KeySchema=[{"AttributeName": "id", "KeyType": "HASH"}], | 
					
						
							|  |  |  |             AttributeDefinitions=[{"AttributeName": "id", "AttributeType": "S"}], | 
					
						
							|  |  |  |             ProvisionedThroughput={"ReadCapacityUnits": 1, "WriteCapacityUnits": 1}, | 
					
						
							|  |  |  |             StreamSpecification={ | 
					
						
							|  |  |  |                 "StreamEnabled": True, | 
					
						
							|  |  |  |                 "StreamViewType": "NEW_AND_OLD_IMAGES", | 
					
						
							|  |  |  |             }, | 
					
						
							|  |  |  |         ) | 
					
						
							|  |  |  |         self.stream_arn = resp["TableDescription"]["LatestStreamArn"] | 
					
						
							|  |  |  | 
 | 
					
						
							| 
									
										
										
										
											2022-10-27 17:37:11 -07:00
										 |  |  |     def teardown_method(self): | 
					
						
							| 
									
										
										
										
											2018-11-07 17:10:00 -05:00
										 |  |  |         conn = boto3.client("dynamodb", region_name="us-east-1") | 
					
						
							|  |  |  |         conn.delete_table(TableName="test-streams") | 
					
						
							|  |  |  |         self.stream_arn = None | 
					
						
							|  |  |  | 
 | 
					
						
							|  |  |  |         for m in self.mocks: | 
					
						
							| 
									
										
										
										
											2021-04-06 11:22:42 +02:00
										 |  |  |             try: | 
					
						
							|  |  |  |                 m.stop() | 
					
						
							|  |  |  |             except RuntimeError: | 
					
						
							|  |  |  |                 pass | 
					
						
							| 
									
										
										
										
											2018-11-07 17:10:00 -05:00
										 |  |  | 
 | 
					
						
							|  |  |  |     def test_verify_stream(self): | 
					
						
							|  |  |  |         conn = boto3.client("dynamodb", region_name="us-east-1") | 
					
						
							|  |  |  |         resp = conn.describe_table(TableName="test-streams") | 
					
						
							|  |  |  |         assert "LatestStreamArn" in resp["Table"] | 
					
						
							|  |  |  | 
 | 
					
						
							|  |  |  |     def test_describe_stream(self): | 
					
						
							|  |  |  |         conn = boto3.client("dynamodbstreams", region_name="us-east-1") | 
					
						
							|  |  |  | 
 | 
					
						
							|  |  |  |         resp = conn.describe_stream(StreamArn=self.stream_arn) | 
					
						
							|  |  |  |         assert "StreamDescription" in resp | 
					
						
							|  |  |  |         desc = resp["StreamDescription"] | 
					
						
							|  |  |  |         assert desc["StreamArn"] == self.stream_arn | 
					
						
							|  |  |  |         assert desc["TableName"] == "test-streams" | 
					
						
							|  |  |  | 
 | 
					
						
							|  |  |  |     def test_list_streams(self): | 
					
						
							|  |  |  |         conn = boto3.client("dynamodbstreams", region_name="us-east-1") | 
					
						
							|  |  |  | 
 | 
					
						
							|  |  |  |         resp = conn.list_streams() | 
					
						
							|  |  |  |         assert resp["Streams"][0]["StreamArn"] == self.stream_arn | 
					
						
							|  |  |  | 
 | 
					
						
							|  |  |  |         resp = conn.list_streams(TableName="no-stream") | 
					
						
							|  |  |  |         assert not resp["Streams"] | 
					
						
							|  |  |  | 
 | 
					
						
							|  |  |  |     def test_get_shard_iterator(self): | 
					
						
							|  |  |  |         conn = boto3.client("dynamodbstreams", region_name="us-east-1") | 
					
						
							|  |  |  | 
 | 
					
						
							|  |  |  |         resp = conn.describe_stream(StreamArn=self.stream_arn) | 
					
						
							|  |  |  |         shard_id = resp["StreamDescription"]["Shards"][0]["ShardId"] | 
					
						
							| 
									
										
										
										
											2019-10-31 08:44:26 -07:00
										 |  |  | 
 | 
					
						
							| 
									
										
										
										
											2018-11-07 17:10:00 -05:00
										 |  |  |         resp = conn.get_shard_iterator( | 
					
						
							|  |  |  |             StreamArn=self.stream_arn, | 
					
						
							|  |  |  |             ShardId=shard_id, | 
					
						
							|  |  |  |             ShardIteratorType="TRIM_HORIZON", | 
					
						
							|  |  |  |         ) | 
					
						
							|  |  |  |         assert "ShardIterator" in resp | 
					
						
							| 
									
										
										
										
											2019-07-29 21:13:58 -05:00
										 |  |  | 
 | 
					
						
							|  |  |  |     def test_get_shard_iterator_at_sequence_number(self): | 
					
						
							|  |  |  |         conn = boto3.client("dynamodbstreams", region_name="us-east-1") | 
					
						
							|  |  |  | 
 | 
					
						
							|  |  |  |         resp = conn.describe_stream(StreamArn=self.stream_arn) | 
					
						
							|  |  |  |         shard_id = resp["StreamDescription"]["Shards"][0]["ShardId"] | 
					
						
							| 
									
										
										
										
											2019-10-31 08:44:26 -07:00
										 |  |  | 
 | 
					
						
							| 
									
										
										
										
											2019-07-29 21:13:58 -05:00
										 |  |  |         resp = conn.get_shard_iterator( | 
					
						
							|  |  |  |             StreamArn=self.stream_arn, | 
					
						
							|  |  |  |             ShardId=shard_id, | 
					
						
							|  |  |  |             ShardIteratorType="AT_SEQUENCE_NUMBER", | 
					
						
							|  |  |  |             SequenceNumber=resp["StreamDescription"]["Shards"][0][ | 
					
						
							|  |  |  |                 "SequenceNumberRange" | 
					
						
							|  |  |  |             ]["StartingSequenceNumber"], | 
					
						
							|  |  |  |         ) | 
					
						
							|  |  |  |         assert "ShardIterator" in resp | 
					
						
							|  |  |  | 
 | 
					
						
							|  |  |  |     def test_get_shard_iterator_after_sequence_number(self): | 
					
						
							|  |  |  |         conn = boto3.client("dynamodbstreams", region_name="us-east-1") | 
					
						
							|  |  |  | 
 | 
					
						
							|  |  |  |         resp = conn.describe_stream(StreamArn=self.stream_arn) | 
					
						
							|  |  |  |         shard_id = resp["StreamDescription"]["Shards"][0]["ShardId"] | 
					
						
							| 
									
										
										
										
											2019-10-31 08:44:26 -07:00
										 |  |  | 
 | 
					
						
							| 
									
										
										
										
											2019-07-29 21:13:58 -05:00
										 |  |  |         resp = conn.get_shard_iterator( | 
					
						
							|  |  |  |             StreamArn=self.stream_arn, | 
					
						
							|  |  |  |             ShardId=shard_id, | 
					
						
							|  |  |  |             ShardIteratorType="AFTER_SEQUENCE_NUMBER", | 
					
						
							|  |  |  |             SequenceNumber=resp["StreamDescription"]["Shards"][0][ | 
					
						
							|  |  |  |                 "SequenceNumberRange" | 
					
						
							|  |  |  |             ]["StartingSequenceNumber"], | 
					
						
							|  |  |  |         ) | 
					
						
							|  |  |  |         assert "ShardIterator" in resp | 
					
						
							| 
									
										
										
										
											2019-10-31 08:44:26 -07:00
										 |  |  | 
 | 
					
						
							| 
									
										
										
										
											2018-11-08 10:54:54 -05:00
										 |  |  |     def test_get_records_empty(self): | 
					
						
							| 
									
										
										
										
											2018-11-07 17:10:00 -05:00
										 |  |  |         conn = boto3.client("dynamodbstreams", region_name="us-east-1") | 
					
						
							|  |  |  | 
 | 
					
						
							|  |  |  |         resp = conn.describe_stream(StreamArn=self.stream_arn) | 
					
						
							|  |  |  |         shard_id = resp["StreamDescription"]["Shards"][0]["ShardId"] | 
					
						
							| 
									
										
										
										
											2019-10-31 08:44:26 -07:00
										 |  |  | 
 | 
					
						
							| 
									
										
										
										
											2018-11-07 17:10:00 -05:00
										 |  |  |         resp = conn.get_shard_iterator( | 
					
						
							|  |  |  |             StreamArn=self.stream_arn, ShardId=shard_id, ShardIteratorType="LATEST" | 
					
						
							|  |  |  |         ) | 
					
						
							|  |  |  |         iterator_id = resp["ShardIterator"] | 
					
						
							|  |  |  | 
 | 
					
						
							|  |  |  |         resp = conn.get_records(ShardIterator=iterator_id) | 
					
						
							|  |  |  |         assert "Records" in resp | 
					
						
							| 
									
										
										
										
											2018-11-08 10:54:54 -05:00
										 |  |  |         assert len(resp["Records"]) == 0 | 
					
						
							| 
									
										
										
										
											2018-11-07 17:10:00 -05:00
										 |  |  | 
 | 
					
						
							| 
									
										
										
										
											2018-11-08 10:54:54 -05:00
										 |  |  |     def test_get_records_seq(self): | 
					
						
							|  |  |  |         conn = boto3.client("dynamodb", region_name="us-east-1") | 
					
						
							|  |  |  | 
 | 
					
						
							|  |  |  |         conn.put_item( | 
					
						
							|  |  |  |             TableName="test-streams", | 
					
						
							| 
									
										
										
										
											2020-04-06 19:55:54 +10:00
										 |  |  |             Item={"id": {"S": "entry1"}, "first_col": {"S": "foo"}}, | 
					
						
							| 
									
										
										
										
											2018-11-08 10:54:54 -05:00
										 |  |  |         ) | 
					
						
							|  |  |  |         conn.put_item( | 
					
						
							|  |  |  |             TableName="test-streams", | 
					
						
							|  |  |  |             Item={ | 
					
						
							|  |  |  |                 "id": {"S": "entry1"}, | 
					
						
							|  |  |  |                 "first_col": {"S": "bar"}, | 
					
						
							|  |  |  |                 "second_col": {"S": "baz"}, | 
					
						
							| 
									
										
										
										
											2020-04-06 19:55:54 +10:00
										 |  |  |                 "a": {"L": [{"M": {"b": {"S": "bar1"}}}]}, | 
					
						
							| 
									
										
										
										
											2018-11-08 10:54:54 -05:00
										 |  |  |             }, | 
					
						
							|  |  |  |         ) | 
					
						
							|  |  |  |         conn.delete_item(TableName="test-streams", Key={"id": {"S": "entry1"}}) | 
					
						
							| 
									
										
										
										
											2019-10-31 08:44:26 -07:00
										 |  |  | 
 | 
					
						
							| 
									
										
										
										
											2018-11-08 10:54:54 -05:00
										 |  |  |         conn = boto3.client("dynamodbstreams", region_name="us-east-1") | 
					
						
							| 
									
										
										
										
											2019-10-31 08:44:26 -07:00
										 |  |  | 
 | 
					
						
							| 
									
										
										
										
											2018-11-08 10:54:54 -05:00
										 |  |  |         resp = conn.describe_stream(StreamArn=self.stream_arn) | 
					
						
							|  |  |  |         shard_id = resp["StreamDescription"]["Shards"][0]["ShardId"] | 
					
						
							| 
									
										
										
										
											2019-10-31 08:44:26 -07:00
										 |  |  | 
 | 
					
						
							| 
									
										
										
										
											2018-11-08 10:54:54 -05:00
										 |  |  |         resp = conn.get_shard_iterator( | 
					
						
							|  |  |  |             StreamArn=self.stream_arn, | 
					
						
							|  |  |  |             ShardId=shard_id, | 
					
						
							|  |  |  |             ShardIteratorType="TRIM_HORIZON", | 
					
						
							|  |  |  |         ) | 
					
						
							|  |  |  |         iterator_id = resp["ShardIterator"] | 
					
						
							|  |  |  | 
 | 
					
						
							|  |  |  |         resp = conn.get_records(ShardIterator=iterator_id) | 
					
						
							|  |  |  |         assert len(resp["Records"]) == 3 | 
					
						
							|  |  |  |         assert resp["Records"][0]["eventName"] == "INSERT" | 
					
						
							|  |  |  |         assert resp["Records"][1]["eventName"] == "MODIFY" | 
					
						
							| 
									
										
										
										
											2020-11-02 09:21:09 -08:00
										 |  |  |         assert resp["Records"][2]["eventName"] == "REMOVE" | 
					
						
							| 
									
										
										
										
											2018-11-08 10:54:54 -05:00
										 |  |  | 
 | 
					
						
							| 
									
										
										
										
											2019-07-29 21:13:58 -05:00
										 |  |  |         sequence_number_modify = resp["Records"][1]["dynamodb"]["SequenceNumber"] | 
					
						
							|  |  |  | 
 | 
					
						
							| 
									
										
										
										
											2018-11-08 10:54:54 -05:00
										 |  |  |         # now try fetching from the next shard iterator, it should be | 
					
						
							|  |  |  |         # empty | 
					
						
							|  |  |  |         resp = conn.get_records(ShardIterator=resp["NextShardIterator"]) | 
					
						
							|  |  |  |         assert len(resp["Records"]) == 0 | 
					
						
							| 
									
										
										
										
											2018-11-08 13:22:24 -05:00
										 |  |  | 
 | 
					
						
							| 
									
										
										
										
											2019-07-29 21:21:02 -05:00
										 |  |  |         # check that if we get the shard iterator AT_SEQUENCE_NUMBER will get the MODIFY event | 
					
						
							| 
									
										
										
										
											2019-07-29 21:13:58 -05:00
										 |  |  |         resp = conn.get_shard_iterator( | 
					
						
							|  |  |  |             StreamArn=self.stream_arn, | 
					
						
							|  |  |  |             ShardId=shard_id, | 
					
						
							|  |  |  |             ShardIteratorType="AT_SEQUENCE_NUMBER", | 
					
						
							|  |  |  |             SequenceNumber=sequence_number_modify, | 
					
						
							|  |  |  |         ) | 
					
						
							|  |  |  |         iterator_id = resp["ShardIterator"] | 
					
						
							|  |  |  |         resp = conn.get_records(ShardIterator=iterator_id) | 
					
						
							|  |  |  |         assert len(resp["Records"]) == 2 | 
					
						
							|  |  |  |         assert resp["Records"][0]["eventName"] == "MODIFY" | 
					
						
							| 
									
										
										
										
											2020-11-02 09:21:09 -08:00
										 |  |  |         assert resp["Records"][1]["eventName"] == "REMOVE" | 
					
						
							| 
									
										
										
										
											2019-07-29 21:13:58 -05:00
										 |  |  | 
 | 
					
						
							| 
									
										
										
										
											2019-07-29 21:21:02 -05:00
										 |  |  |         # check that if we get the shard iterator AFTER_SEQUENCE_NUMBER will get the DELETE event | 
					
						
							| 
									
										
										
										
											2019-07-29 21:13:58 -05:00
										 |  |  |         resp = conn.get_shard_iterator( | 
					
						
							|  |  |  |             StreamArn=self.stream_arn, | 
					
						
							|  |  |  |             ShardId=shard_id, | 
					
						
							|  |  |  |             ShardIteratorType="AFTER_SEQUENCE_NUMBER", | 
					
						
							|  |  |  |             SequenceNumber=sequence_number_modify, | 
					
						
							|  |  |  |         ) | 
					
						
							|  |  |  |         iterator_id = resp["ShardIterator"] | 
					
						
							|  |  |  |         resp = conn.get_records(ShardIterator=iterator_id) | 
					
						
							|  |  |  |         assert len(resp["Records"]) == 1 | 
					
						
							| 
									
										
										
										
											2020-11-02 09:21:09 -08:00
										 |  |  |         assert resp["Records"][0]["eventName"] == "REMOVE" | 
					
						
							| 
									
										
										
										
											2019-07-29 21:13:58 -05:00
										 |  |  | 
 | 
					
						
							| 
									
										
										
										
											2018-11-08 13:22:24 -05:00
										 |  |  | 
 | 
					
						
							|  |  |  | class TestEdges: | 
					
						
							|  |  |  |     mocks = [] | 
					
						
							|  |  |  | 
 | 
					
						
							| 
									
										
										
										
											2022-10-27 17:37:11 -07:00
										 |  |  |     def setup_method(self): | 
					
						
							| 
									
										
										
										
											2022-03-09 16:57:25 -01:00
										 |  |  |         self.mocks = [mock_dynamodb(), mock_dynamodbstreams()] | 
					
						
							| 
									
										
										
										
											2018-11-08 13:22:24 -05:00
										 |  |  |         for m in self.mocks: | 
					
						
							|  |  |  |             m.start() | 
					
						
							|  |  |  | 
 | 
					
						
							| 
									
										
										
										
											2022-10-27 17:37:11 -07:00
										 |  |  |     def teardown_method(self): | 
					
						
							| 
									
										
										
										
											2018-11-08 13:22:24 -05:00
										 |  |  |         for m in self.mocks: | 
					
						
							| 
									
										
										
										
											2021-04-06 11:22:42 +02:00
										 |  |  |             try: | 
					
						
							|  |  |  |                 m.stop() | 
					
						
							|  |  |  |             except RuntimeError: | 
					
						
							|  |  |  |                 pass | 
					
						
							| 
									
										
										
										
											2018-11-08 13:22:24 -05:00
										 |  |  | 
 | 
					
						
							|  |  |  |     def test_enable_stream_on_table(self): | 
					
						
							|  |  |  |         conn = boto3.client("dynamodb", region_name="us-east-1") | 
					
						
							|  |  |  |         resp = conn.create_table( | 
					
						
							|  |  |  |             TableName="test-streams", | 
					
						
							|  |  |  |             KeySchema=[{"AttributeName": "id", "KeyType": "HASH"}], | 
					
						
							|  |  |  |             AttributeDefinitions=[{"AttributeName": "id", "AttributeType": "S"}], | 
					
						
							|  |  |  |             ProvisionedThroughput={"ReadCapacityUnits": 1, "WriteCapacityUnits": 1}, | 
					
						
							|  |  |  |         ) | 
					
						
							|  |  |  |         assert "StreamSpecification" not in resp["TableDescription"] | 
					
						
							| 
									
										
										
										
											2019-10-31 08:44:26 -07:00
										 |  |  | 
 | 
					
						
							| 
									
										
										
										
											2018-11-08 13:22:24 -05:00
										 |  |  |         resp = conn.update_table( | 
					
						
							|  |  |  |             TableName="test-streams", | 
					
						
							| 
									
										
										
										
											2019-11-21 17:00:18 -05:00
										 |  |  |             StreamSpecification={"StreamViewType": "KEYS_ONLY", "StreamEnabled": True}, | 
					
						
							| 
									
										
										
										
											2018-11-08 13:22:24 -05:00
										 |  |  |         ) | 
					
						
							|  |  |  |         assert "StreamSpecification" in resp["TableDescription"] | 
					
						
							|  |  |  |         assert resp["TableDescription"]["StreamSpecification"] == { | 
					
						
							|  |  |  |             "StreamEnabled": True, | 
					
						
							|  |  |  |             "StreamViewType": "KEYS_ONLY", | 
					
						
							|  |  |  |         } | 
					
						
							|  |  |  |         assert "LatestStreamLabel" in resp["TableDescription"] | 
					
						
							|  |  |  | 
 | 
					
						
							|  |  |  |         # now try to enable it again | 
					
						
							| 
									
										
										
										
											2020-10-06 07:54:49 +02:00
										 |  |  |         with pytest.raises(conn.exceptions.ResourceInUseException): | 
					
						
							| 
									
										
										
										
											2018-11-08 13:22:24 -05:00
										 |  |  |             resp = conn.update_table( | 
					
						
							|  |  |  |                 TableName="test-streams", | 
					
						
							| 
									
										
										
										
											2019-11-21 17:53:58 -05:00
										 |  |  |                 StreamSpecification={ | 
					
						
							|  |  |  |                     "StreamViewType": "OLD_IMAGES", | 
					
						
							|  |  |  |                     "StreamEnabled": True, | 
					
						
							|  |  |  |                 }, | 
					
						
							| 
									
										
										
										
											2018-11-08 13:22:24 -05:00
										 |  |  |             ) | 
					
						
							| 
									
										
										
										
											2019-10-31 08:44:26 -07:00
										 |  |  | 
 | 
					
						
							| 
									
										
										
										
											2018-11-08 13:22:24 -05:00
										 |  |  |     def test_stream_with_range_key(self): | 
					
						
							|  |  |  |         dyn = boto3.client("dynamodb", region_name="us-east-1") | 
					
						
							|  |  |  | 
 | 
					
						
							|  |  |  |         resp = dyn.create_table( | 
					
						
							|  |  |  |             TableName="test-streams", | 
					
						
							|  |  |  |             KeySchema=[ | 
					
						
							|  |  |  |                 {"AttributeName": "id", "KeyType": "HASH"}, | 
					
						
							|  |  |  |                 {"AttributeName": "color", "KeyType": "RANGE"}, | 
					
						
							|  |  |  |             ], | 
					
						
							|  |  |  |             AttributeDefinitions=[ | 
					
						
							|  |  |  |                 {"AttributeName": "id", "AttributeType": "S"}, | 
					
						
							|  |  |  |                 {"AttributeName": "color", "AttributeType": "S"}, | 
					
						
							|  |  |  |             ], | 
					
						
							|  |  |  |             ProvisionedThroughput={"ReadCapacityUnits": 1, "WriteCapacityUnits": 1}, | 
					
						
							| 
									
										
										
										
											2019-11-21 17:00:18 -05:00
										 |  |  |             StreamSpecification={"StreamViewType": "NEW_IMAGES", "StreamEnabled": True}, | 
					
						
							| 
									
										
										
										
											2018-11-08 13:22:24 -05:00
										 |  |  |         ) | 
					
						
							|  |  |  |         stream_arn = resp["TableDescription"]["LatestStreamArn"] | 
					
						
							|  |  |  | 
 | 
					
						
							|  |  |  |         streams = boto3.client("dynamodbstreams", region_name="us-east-1") | 
					
						
							|  |  |  |         resp = streams.describe_stream(StreamArn=stream_arn) | 
					
						
							|  |  |  |         shard_id = resp["StreamDescription"]["Shards"][0]["ShardId"] | 
					
						
							|  |  |  | 
 | 
					
						
							|  |  |  |         resp = streams.get_shard_iterator( | 
					
						
							|  |  |  |             StreamArn=stream_arn, ShardId=shard_id, ShardIteratorType="LATEST" | 
					
						
							|  |  |  |         ) | 
					
						
							|  |  |  |         iterator_id = resp["ShardIterator"] | 
					
						
							|  |  |  | 
 | 
					
						
							|  |  |  |         dyn.put_item( | 
					
						
							|  |  |  |             TableName="test-streams", Item={"id": {"S": "row1"}, "color": {"S": "blue"}} | 
					
						
							|  |  |  |         ) | 
					
						
							|  |  |  |         dyn.put_item( | 
					
						
							|  |  |  |             TableName="test-streams", | 
					
						
							|  |  |  |             Item={"id": {"S": "row2"}, "color": {"S": "green"}}, | 
					
						
							|  |  |  |         ) | 
					
						
							|  |  |  | 
 | 
					
						
							|  |  |  |         resp = streams.get_records(ShardIterator=iterator_id) | 
					
						
							|  |  |  |         assert len(resp["Records"]) == 2 | 
					
						
							|  |  |  |         assert resp["Records"][0]["eventName"] == "INSERT" | 
					
						
							|  |  |  |         assert resp["Records"][1]["eventName"] == "INSERT" |