diff --git a/packages/google-cloud-spanner/tests/unit/_async/test_streamed.py b/packages/google-cloud-spanner/tests/unit/_async/test_streamed.py index 15114616fc26..27aebd048016 100644 --- a/packages/google-cloud-spanner/tests/unit/_async/test_streamed.py +++ b/packages/google-cloud-spanner/tests/unit/_async/test_streamed.py @@ -297,6 +297,105 @@ async def test__merge_chunk_float64_w_float64(self): with pytest.raises(Unmergeable): streamed._merge_chunk(chunk) + @CrossSync.pytest + async def test__merge_chunk_float32_nan_string(self): + from google.cloud.spanner_v1 import TypeCode + + iterator = _MockCancellableIterator() + streamed = self._make_one(iterator) + fields = [self._make_scalar_field("weight", TypeCode.FLOAT32)] + streamed._metadata = self._make_result_set_metadata(fields) + streamed._pending_chunk = self._make_value("Na") + chunk = self._make_value("N") + + merged = streamed._merge_chunk(chunk) + self.assertEqual(merged.string_value, "NaN") + + @CrossSync.pytest + async def test__merge_chunk_float32_w_empty(self): + from google.cloud.spanner_v1 import TypeCode + + iterator = _MockCancellableIterator() + streamed = self._make_one(iterator) + fields = [self._make_scalar_field("weight", TypeCode.FLOAT32)] + streamed._metadata = self._make_result_set_metadata(fields) + streamed._pending_chunk = self._make_value(1.25) + chunk = self._make_value("") + + merged = streamed._merge_chunk(chunk) + self.assertEqual(merged.number_value, 1.25) + + @CrossSync.pytest + async def test__merge_chunk_float32_w_float32(self): + from google.cloud.spanner_v1 import TypeCode + from google.cloud.spanner_v1._async.streamed import Unmergeable + + iterator = _MockCancellableIterator() + streamed = self._make_one(iterator) + fields = [self._make_scalar_field("weight", TypeCode.FLOAT32)] + streamed._metadata = self._make_result_set_metadata(fields) + streamed._pending_chunk = self._make_value(1.25) + chunk = self._make_value(2.5) + + with pytest.raises(Unmergeable): + streamed._merge_chunk(chunk) + + @CrossSync.pytest + async def test__merge_chunk_date(self): + from google.cloud.spanner_v1 import TypeCode + + iterator = _MockCancellableIterator() + streamed = self._make_one(iterator) + fields = [self._make_scalar_field("anniversary", TypeCode.DATE)] + streamed._metadata = self._make_result_set_metadata(fields) + streamed._pending_chunk = self._make_value("2026-") + chunk = self._make_value("10-08") + + merged = streamed._merge_chunk(chunk) + self.assertEqual(merged.string_value, "2026-10-08") + + @CrossSync.pytest + async def test__merge_chunk_timestamp(self): + from google.cloud.spanner_v1 import TypeCode + + iterator = _MockCancellableIterator() + streamed = self._make_one(iterator) + fields = [self._make_scalar_field("created_at", TypeCode.TIMESTAMP)] + streamed._metadata = self._make_result_set_metadata(fields) + streamed._pending_chunk = self._make_value("2026-10-08T13:") + chunk = self._make_value("15:00Z") + + merged = streamed._merge_chunk(chunk) + self.assertEqual(merged.string_value, "2026-10-08T13:15:00Z") + + @CrossSync.pytest + async def test__merge_chunk_interval(self): + from google.cloud.spanner_v1 import TypeCode + + iterator = _MockCancellableIterator() + streamed = self._make_one(iterator) + fields = [self._make_scalar_field("duration", TypeCode.INTERVAL)] + streamed._metadata = self._make_result_set_metadata(fields) + streamed._pending_chunk = self._make_value("P1Y2M") + chunk = self._make_value("3DT4H") + + merged = streamed._merge_chunk(chunk) + self.assertEqual(merged.string_value, "P1Y2M3DT4H") + + @CrossSync.pytest + async def test__merge_chunk_uuid(self): + from google.cloud.spanner_v1 import TypeCode + + iterator = _MockCancellableIterator() + streamed = self._make_one(iterator) + fields = [self._make_scalar_field("session_id", TypeCode.UUID)] + streamed._metadata = self._make_result_set_metadata(fields) + streamed._pending_chunk = self._make_value("12345678-1234-") + chunk = self._make_value("5678-1234-567812345678") + + merged = streamed._merge_chunk(chunk) + self.assertEqual(merged.string_value, "12345678-1234-5678-1234-567812345678") + @CrossSync.pytest async def test__merge_chunk_string(self): from google.cloud.spanner_v1 import TypeCode @@ -1547,6 +1646,828 @@ async def test_merge_values_zero_fields(self): streamed._merge_values([self._make_value(1)]) self.assertEqual(streamed._rows, []) + @CrossSync.pytest + async def test___iter___chunked_values_empty_string_start(self): + from google.cloud.spanner_v1 import TypeCode + + fields = [ + self._make_scalar_field("column_one", TypeCode.STRING), + self._make_scalar_field("column_two", TypeCode.STRING), + ] + metadata = self._make_result_set_metadata(fields) + chunk_one = self._make_partial_result_set( + [self._make_value("")], + metadata=metadata, + chunked_value=True, + ) + chunk_two = self._make_partial_result_set( + [ + self._make_value("tail-after-empty"), + self._make_value("value_two"), + self._make_value("value_three"), + self._make_value(""), + ], + chunked_value=True, + ) + chunk_three = self._make_partial_result_set( + [self._make_value("tail-two")], + last=True, + ) + iterator = _MockCancellableIterator(chunk_one, chunk_two, chunk_three) + streamed = self._make_one(iterator) + found = [row async for row in streamed] + self.assertEqual( + found, + [ + ["tail-after-empty", "value_two"], + ["value_three", "tail-two"], + ], + ) + self.assertIsNone(streamed._pending_chunk) + + @CrossSync.pytest + async def test___iter___chunked_values_empty_list_head_and_tail(self): + from google.cloud.spanner_v1 import TypeCode + + fields = [ + self._make_array_field("column_one", TypeCode.STRING), + self._make_array_field("column_two", TypeCode.STRING), + ] + metadata = self._make_result_set_metadata(fields) + chunk_one = self._make_partial_result_set( + [self._make_list_value([])], + metadata=metadata, + chunked_value=True, + ) + chunk_two = self._make_partial_result_set( + [ + self._make_list_value(["a", "b"]), + self._make_list_value([]), + ], + chunked_value=True, + ) + chunk_three = self._make_partial_result_set( + [ + self._make_list_value(["v1", "v2"]), + self._make_list_value(["x", "y"]), + ], + chunked_value=True, + ) + chunk_four = self._make_partial_result_set( + [ + self._make_list_value([]), + self._make_list_value(["w1", "w2"]), + ], + chunked_value=True, + ) + chunk_five = self._make_partial_result_set( + [ + self._make_list_value([]), + self._make_list_value([]), + ], + chunked_value=True, + ) + chunk_six = self._make_partial_result_set( + [ + self._make_list_value([]), + self._make_list_value(["z1", "z2"]), + ], + last=True, + ) + iterator = _MockCancellableIterator( + chunk_one, chunk_two, chunk_three, chunk_four, chunk_five, chunk_six + ) + streamed = self._make_one(iterator) + found = [row async for row in streamed] + self.assertEqual( + found, + [ + [["a", "b"], ["v1", "v2"]], + [["x", "y"], ["w1", "w2"]], + [[], ["z1", "z2"]], + ], + ) + self.assertIsNone(streamed._pending_chunk) + + @CrossSync.pytest + async def test___iter___chunked_values_bytes_and_bytes_array(self): + from google.protobuf.struct_pb2 import Value + + from google.cloud.spanner_v1 import TypeCode + + fields = [ + self._make_scalar_field("raw_bytes", TypeCode.BYTES), + self._make_array_field("bytes_array", TypeCode.BYTES), + ] + metadata = self._make_result_set_metadata(fields) + + full_bytes = b"hello-chunked-bytes-world" + bytes_part_one = "hello-chunked-" + bytes_part_two = "bytes-world" + + element_one = b"first-element" + element_two = b"second-chunked-element" + element_two_part_one = "second-chunked-" + element_two_part_two = "element" + element_three = b"third-element" + + chunk_one = self._make_partial_result_set( + [self._make_value(bytes_part_one)], + metadata=metadata, + chunked_value=True, + ) + chunk_two = self._make_partial_result_set( + [ + self._make_value(bytes_part_two), + self._make_list_value( + value_pbs=[ + self._make_value("first-element"), + Value(null_value="NULL_VALUE"), + self._make_value(element_two_part_one), + ] + ), + ], + chunked_value=True, + ) + chunk_three = self._make_partial_result_set( + [ + self._make_list_value( + value_pbs=[ + self._make_value(element_two_part_two), + self._make_value("third-element"), + ] + ) + ], + last=True, + ) + iterator = _MockCancellableIterator(chunk_one, chunk_two, chunk_three) + streamed = self._make_one(iterator) + found = [row async for row in streamed] + self.assertEqual( + found, + [ + [ + full_bytes, + [element_one, None, element_two, element_three], + ] + ], + ) + self.assertIsNone(streamed._pending_chunk) + + @CrossSync.pytest + async def test___iter___chunked_values_intermediate_single_value_chunk(self): + from google.cloud.spanner_v1 import TypeCode + + fields = [ + self._make_scalar_field("identifier", TypeCode.STRING), + self._make_scalar_field("text_payload", TypeCode.STRING), + ] + metadata = self._make_result_set_metadata(fields) + chunk_one = self._make_partial_result_set( + [self._make_value("id1"), self._make_value("hello-")], + metadata=metadata, + chunked_value=True, + ) + chunk_two = self._make_partial_result_set( + [self._make_value("world-")], + chunked_value=True, + ) + chunk_three = self._make_partial_result_set( + [ + self._make_value("again"), + self._make_value("id2"), + self._make_value("text2"), + ], + last=True, + ) + iterator = _MockCancellableIterator(chunk_one, chunk_two, chunk_three) + streamed = self._make_one(iterator) + found = [row async for row in streamed] + self.assertEqual( + found, + [ + ["id1", "hello-world-again"], + ["id2", "text2"], + ], + ) + self.assertIsNone(streamed._pending_chunk) + + @CrossSync.pytest + async def test___iter___multi_row_multi_chunk_interleaved_splits(self): + from google.cloud.spanner_v1 import TypeCode + + fields = [ + self._make_scalar_field("identifier", TypeCode.INT64), + self._make_scalar_field("text_payload", TypeCode.STRING), + ] + metadata = self._make_result_set_metadata(fields) + chunk_one = self._make_partial_result_set( + [self._make_value(1), self._make_value("spl")], + metadata=metadata, + chunked_value=True, + ) + chunk_two = self._make_partial_result_set( + [ + self._make_value("it_text_row1"), + self._make_value(2), + self._make_value("split_"), + ], + chunked_value=True, + ) + chunk_three = self._make_partial_result_set( + [self._make_value("text_row2")], + last=True, + ) + iterator = _MockCancellableIterator(chunk_one, chunk_two, chunk_three) + streamed = self._make_one(iterator) + found = [row async for row in streamed] + self.assertEqual( + found, + [ + [1, "split_text_row1"], + [2, "split_text_row2"], + ], + ) + self.assertIsNone(streamed._pending_chunk) + + @CrossSync.pytest + async def test___iter___subsequent_metadata_ignored(self): + from google.cloud.spanner_v1 import TypeCode + + fields_one = [self._make_scalar_field("primary_column", TypeCode.STRING)] + metadata_one = self._make_result_set_metadata(fields_one) + fields_two = [self._make_scalar_field("ignored_column", TypeCode.STRING)] + metadata_two = self._make_result_set_metadata(fields_two) + + chunk_one = self._make_partial_result_set( + [self._make_value("alpha")], + metadata=metadata_one, + ) + chunk_two = self._make_partial_result_set( + [self._make_value("beta")], + metadata=metadata_two, + last=True, + ) + iterator = _MockCancellableIterator(chunk_one, chunk_two) + streamed = self._make_one(iterator) + found = [row async for row in streamed] + self.assertEqual(found, [["alpha"], ["beta"]]) + self.assertEqual(streamed.fields[0].name, "primary_column") + + @CrossSync.pytest + async def test___iter___empty_chunk_in_middle_of_stream(self): + from google.cloud.spanner_v1 import TypeCode + + fields = [self._make_scalar_field("column_name", TypeCode.STRING)] + metadata = self._make_result_set_metadata(fields) + chunk_one = self._make_partial_result_set( + [self._make_value("value_one")], + metadata=metadata, + ) + chunk_two = self._make_partial_result_set([]) + chunk_three = self._make_partial_result_set( + [self._make_value("value_two")], + last=True, + ) + iterator = _MockCancellableIterator(chunk_one, chunk_two, chunk_three) + streamed = self._make_one(iterator) + found = [row async for row in streamed] + self.assertEqual(found, [["value_one"], ["value_two"]]) + + @CrossSync.pytest + async def test___iter___chunked_values_incompatible_types_error(self): + from google.cloud.spanner_v1 import TypeCode + from google.cloud.spanner_v1._async.streamed import Unmergeable + + fields = [self._make_scalar_field("registered", TypeCode.BOOL)] + metadata = self._make_result_set_metadata(fields) + chunk_one = self._make_partial_result_set( + [self._make_value(True)], + metadata=metadata, + chunked_value=True, + ) + chunk_two = self._make_partial_result_set( + [self._make_value(False)], + last=True, + ) + iterator = _MockCancellableIterator(chunk_one, chunk_two) + streamed = self._make_one(iterator) + with self.assertRaises(Unmergeable): + _ = [row async for row in streamed] + + @CrossSync.pytest + async def test___iter___resumption_across_chunked_value_retry(self): + from google.api_core.exceptions import ServiceUnavailable + + from google.cloud.spanner_v1 import TypeCode + + fields = [ + self._make_scalar_field("identifier", TypeCode.STRING), + self._make_scalar_field("data_payload", TypeCode.STRING), + ] + metadata = self._make_result_set_metadata(fields) + + attempt_count = 0 + last_resume_token = None + + async def create_stream_attempt(resume_token=None): + nonlocal attempt_count + attempt_count += 1 + if attempt_count == 1: + self.assertIsNone(resume_token) + chunk_one = self._make_partial_result_set( + [self._make_value("id1"), self._make_value("part1-")], + metadata=metadata, + chunked_value=True, + ) + chunk_one.resume_token = b"token-1" + yield chunk_one + raise ServiceUnavailable("transient error during chunked scalar") + else: + self.assertEqual(resume_token, b"token-1") + chunk_two = self._make_partial_result_set( + [self._make_value("part2-actual-")], + chunked_value=True, + ) + chunk_two.resume_token = b"token-2" + yield chunk_two + chunk_three = self._make_partial_result_set( + [self._make_value("part3-")], + chunked_value=True, + ) + yield chunk_three + chunk_four = self._make_partial_result_set( + [ + self._make_value("part4"), + self._make_value("id2"), + self._make_value("data2"), + ], + last=True, + ) + yield chunk_four + + class ResumableResponseIterator: + def __init__(self): + self._current_generator = create_stream_attempt(None) + + def __aiter__(self): + return self + + async def __anext__(self): + nonlocal last_resume_token + while True: + try: + item = await anext(self._current_generator) + if item.resume_token: + last_resume_token = item.resume_token + return item + except ServiceUnavailable: + await self._current_generator.aclose() + self._current_generator = create_stream_attempt( + last_resume_token + ) + + async def aclose(self): + if self._current_generator is not None: + await self._current_generator.aclose() + + iterator = ResumableResponseIterator() + streamed = self._make_one(iterator) + found = [row async for row in streamed] + await iterator.aclose() + self.assertEqual( + found, + [ + ["id1", "part1-part2-actual-part3-part4"], + ["id2", "data2"], + ], + ) + self.assertEqual(attempt_count, 2) + self.assertIsNone(streamed._pending_chunk) + + @CrossSync.pytest + async def test___iter___resumption_across_chunked_array_retry(self): + from google.api_core.exceptions import ServiceUnavailable + + from google.cloud.spanner_v1 import TypeCode + + fields = [ + self._make_scalar_field("identifier", TypeCode.STRING), + self._make_array_field("tags", TypeCode.STRING), + ] + metadata = self._make_result_set_metadata(fields) + + attempt_count = 0 + last_resume_token = None + + async def create_stream_attempt(resume_token=None): + nonlocal attempt_count + attempt_count += 1 + if attempt_count == 1: + self.assertIsNone(resume_token) + chunk_one = self._make_partial_result_set( + [ + self._make_value("id1"), + self._make_list_value(["tag1", "tag2-"]), + ], + metadata=metadata, + chunked_value=True, + ) + chunk_one.resume_token = b"token-array" + yield chunk_one + raise ServiceUnavailable("transient error during chunked array") + else: + self.assertEqual(resume_token, b"token-array") + chunk_two = self._make_partial_result_set( + [self._make_list_value(["tag2-tail", "tag3"])], + last=True, + ) + yield chunk_two + + class ResumableResponseIterator: + def __init__(self): + self._current_generator = create_stream_attempt(None) + + def __aiter__(self): + return self + + async def __anext__(self): + nonlocal last_resume_token + while True: + try: + item = await anext(self._current_generator) + if item.resume_token: + last_resume_token = item.resume_token + return item + except ServiceUnavailable: + await self._current_generator.aclose() + self._current_generator = create_stream_attempt( + last_resume_token + ) + + async def aclose(self): + if self._current_generator is not None: + await self._current_generator.aclose() + + iterator = ResumableResponseIterator() + streamed = self._make_one(iterator) + found = [row async for row in streamed] + await iterator.aclose() + self.assertEqual( + found, + [ + ["id1", ["tag1", "tag2-tag2-tail", "tag3"]], + ], + ) + self.assertEqual(attempt_count, 2) + self.assertIsNone(streamed._pending_chunk) + + @CrossSync.pytest + async def test___iter___chunked_values_numeric(self): + import decimal + + from google.cloud.spanner_v1 import TypeCode + + fields = [self._make_scalar_field("amount", TypeCode.NUMERIC)] + metadata = self._make_result_set_metadata(fields) + chunk_one = self._make_partial_result_set( + [self._make_value("1234.")], + metadata=metadata, + chunked_value=True, + ) + chunk_two = self._make_partial_result_set( + [self._make_value("5678")], + last=True, + ) + iterator = _MockCancellableIterator(chunk_one, chunk_two) + streamed = self._make_one(iterator) + found = [row async for row in streamed] + self.assertEqual(found, [[decimal.Decimal("1234.5678")]]) + self.assertIsNone(streamed._pending_chunk) + + @CrossSync.pytest + async def test___iter___chunked_values_json(self): + from google.cloud.spanner_v1 import JsonObject, TypeCode + + fields = [self._make_scalar_field("payload", TypeCode.JSON)] + metadata = self._make_result_set_metadata(fields) + chunk_one = self._make_partial_result_set( + [self._make_value('{"user": "')], + metadata=metadata, + chunked_value=True, + ) + chunk_two = self._make_partial_result_set( + [self._make_value('alice", "score": 10}')], + last=True, + ) + iterator = _MockCancellableIterator(chunk_one, chunk_two) + streamed = self._make_one(iterator) + found = [row async for row in streamed] + self.assertEqual(found, [[JsonObject({"user": "alice", "score": 10})]]) + self.assertIsNone(streamed._pending_chunk) + + @CrossSync.pytest + async def test___iter___chunked_values_lazy_decode(self): + from google.protobuf.struct_pb2 import Value + + from google.cloud.spanner_v1 import TypeCode + + fields = [self._make_scalar_field("name", TypeCode.STRING)] + metadata = self._make_result_set_metadata(fields) + chunk_one = self._make_partial_result_set( + [self._make_value("hello-")], + metadata=metadata, + chunked_value=True, + ) + chunk_two = self._make_partial_result_set( + [self._make_value("world")], + last=True, + ) + iterator = _MockCancellableIterator(chunk_one, chunk_two) + streamed = self._make_one(iterator, lazy_decode=True) + found = [row async for row in streamed] + self.assertEqual(found, [[Value(string_value="hello-world")]]) + self.assertEqual(streamed.decode_row(found[0]), ["hello-world"]) + self.assertIsNone(streamed._pending_chunk) + + @CrossSync.pytest + async def test___iter___chunked_values_float32(self): + from google.cloud.spanner_v1 import TypeCode + + fields = [self._make_scalar_field("ratio", TypeCode.FLOAT32)] + metadata = self._make_result_set_metadata(fields) + chunk_one = self._make_partial_result_set( + [self._make_value("Infi")], + metadata=metadata, + chunked_value=True, + ) + chunk_two = self._make_partial_result_set( + [self._make_value("nity")], + last=True, + ) + iterator = _MockCancellableIterator(chunk_one, chunk_two) + streamed = self._make_one(iterator) + found = [row async for row in streamed] + self.assertEqual(found, [[float("inf")]]) + self.assertIsNone(streamed._pending_chunk) + + @CrossSync.pytest + async def test___iter___chunked_values_float32_array(self): + from google.cloud.spanner_v1 import TypeCode + + fields = [self._make_array_field("ratios", TypeCode.FLOAT32)] + metadata = self._make_result_set_metadata(fields) + chunk_one = self._make_partial_result_set( + [self._make_list_value([1.5, 2.5])], + metadata=metadata, + chunked_value=True, + ) + chunk_two = self._make_partial_result_set( + [self._make_list_value(["", 3.5])], + last=True, + ) + iterator = _MockCancellableIterator(chunk_one, chunk_two) + streamed = self._make_one(iterator) + found = [row async for row in streamed] + self.assertEqual(found, [[[1.5, 2.5, 3.5]]]) + self.assertIsNone(streamed._pending_chunk) + + @CrossSync.pytest + async def test___iter___chunked_values_date(self): + import datetime + + from google.cloud.spanner_v1 import TypeCode + + fields = [self._make_scalar_field("anniversary", TypeCode.DATE)] + metadata = self._make_result_set_metadata(fields) + chunk_one = self._make_partial_result_set( + [self._make_value("2026-")], + metadata=metadata, + chunked_value=True, + ) + chunk_two = self._make_partial_result_set( + [self._make_value("10-08")], + last=True, + ) + iterator = _MockCancellableIterator(chunk_one, chunk_two) + streamed = self._make_one(iterator) + found = [row async for row in streamed] + self.assertEqual(found, [[datetime.date(2026, 10, 8)]]) + self.assertIsNone(streamed._pending_chunk) + + @CrossSync.pytest + async def test___iter___chunked_values_date_array(self): + import datetime + + from google.cloud.spanner_v1 import TypeCode + + fields = [self._make_array_field("dates", TypeCode.DATE)] + metadata = self._make_result_set_metadata(fields) + chunk_one = self._make_partial_result_set( + [self._make_list_value(["2026-01-01", "2026-02-"])], + metadata=metadata, + chunked_value=True, + ) + chunk_two = self._make_partial_result_set( + [self._make_list_value(["01", "2026-03-01"])], + last=True, + ) + iterator = _MockCancellableIterator(chunk_one, chunk_two) + streamed = self._make_one(iterator) + found = [row async for row in streamed] + self.assertEqual( + found, + [ + [ + [ + datetime.date(2026, 1, 1), + datetime.date(2026, 2, 1), + datetime.date(2026, 3, 1), + ] + ] + ], + ) + self.assertIsNone(streamed._pending_chunk) + + @CrossSync.pytest + async def test___iter___chunked_values_timestamp(self): + import datetime + + from google.cloud.spanner_v1 import TypeCode + + fields = [self._make_scalar_field("created_at", TypeCode.TIMESTAMP)] + metadata = self._make_result_set_metadata(fields) + chunk_one = self._make_partial_result_set( + [self._make_value("2026-10-08T13:")], + metadata=metadata, + chunked_value=True, + ) + chunk_two = self._make_partial_result_set( + [self._make_value("15:00Z")], + last=True, + ) + iterator = _MockCancellableIterator(chunk_one, chunk_two) + streamed = self._make_one(iterator) + found = [row async for row in streamed] + self.assertEqual( + found, + [[datetime.datetime(2026, 10, 8, 13, 15, 0, tzinfo=datetime.timezone.utc)]], + ) + self.assertIsNone(streamed._pending_chunk) + + @CrossSync.pytest + async def test___iter___chunked_values_timestamp_array(self): + import datetime + + from google.cloud.spanner_v1 import TypeCode + + fields = [self._make_array_field("timestamps", TypeCode.TIMESTAMP)] + metadata = self._make_result_set_metadata(fields) + chunk_one = self._make_partial_result_set( + [self._make_list_value(["2026-01-01T00:00:00Z", "2026-02-01T00:"])], + metadata=metadata, + chunked_value=True, + ) + chunk_two = self._make_partial_result_set( + [self._make_list_value(["00:00Z", "2026-03-01T00:00:00Z"])], + last=True, + ) + iterator = _MockCancellableIterator(chunk_one, chunk_two) + streamed = self._make_one(iterator) + found = [row async for row in streamed] + self.assertEqual( + found, + [ + [ + [ + datetime.datetime( + 2026, 1, 1, 0, 0, 0, tzinfo=datetime.timezone.utc + ), + datetime.datetime( + 2026, 2, 1, 0, 0, 0, tzinfo=datetime.timezone.utc + ), + datetime.datetime( + 2026, 3, 1, 0, 0, 0, tzinfo=datetime.timezone.utc + ), + ] + ] + ], + ) + self.assertIsNone(streamed._pending_chunk) + + @CrossSync.pytest + async def test___iter___chunked_values_interval(self): + from google.cloud.spanner_v1 import Interval, TypeCode + + fields = [self._make_scalar_field("duration", TypeCode.INTERVAL)] + metadata = self._make_result_set_metadata(fields) + chunk_one = self._make_partial_result_set( + [self._make_value("P1Y2M")], + metadata=metadata, + chunked_value=True, + ) + chunk_two = self._make_partial_result_set( + [self._make_value("3DT4H")], + last=True, + ) + iterator = _MockCancellableIterator(chunk_one, chunk_two) + streamed = self._make_one(iterator) + found = [row async for row in streamed] + self.assertEqual(found, [[Interval.from_str("P1Y2M3DT4H")]]) + self.assertIsNone(streamed._pending_chunk) + + @CrossSync.pytest + async def test___iter___chunked_values_interval_array(self): + from google.cloud.spanner_v1 import Interval, TypeCode + + fields = [self._make_array_field("durations", TypeCode.INTERVAL)] + metadata = self._make_result_set_metadata(fields) + chunk_one = self._make_partial_result_set( + [self._make_list_value(["P1Y", "P2"])], + metadata=metadata, + chunked_value=True, + ) + chunk_two = self._make_partial_result_set( + [self._make_list_value(["M", "P3D"])], + last=True, + ) + iterator = _MockCancellableIterator(chunk_one, chunk_two) + streamed = self._make_one(iterator) + found = [row async for row in streamed] + self.assertEqual( + found, + [ + [ + [ + Interval.from_str("P1Y"), + Interval.from_str("P2M"), + Interval.from_str("P3D"), + ] + ] + ], + ) + self.assertIsNone(streamed._pending_chunk) + + @CrossSync.pytest + async def test___iter___chunked_values_uuid(self): + import uuid + + from google.cloud.spanner_v1 import TypeCode + + fields = [self._make_scalar_field("session_id", TypeCode.UUID)] + metadata = self._make_result_set_metadata(fields) + chunk_one = self._make_partial_result_set( + [self._make_value("12345678-1234-")], + metadata=metadata, + chunked_value=True, + ) + chunk_two = self._make_partial_result_set( + [self._make_value("5678-1234-567812345678")], + last=True, + ) + iterator = _MockCancellableIterator(chunk_one, chunk_two) + streamed = self._make_one(iterator) + found = [row async for row in streamed] + self.assertEqual(found, [[uuid.UUID("12345678-1234-5678-1234-567812345678")]]) + self.assertIsNone(streamed._pending_chunk) + + @CrossSync.pytest + async def test___iter___chunked_values_uuid_array(self): + import uuid + + from google.cloud.spanner_v1 import TypeCode + + fields = [self._make_array_field("ids", TypeCode.UUID)] + metadata = self._make_result_set_metadata(fields) + chunk_one = self._make_partial_result_set( + [ + self._make_list_value( + ["12345678-1234-5678-1234-567812345678", "00000000-0000-"] + ) + ], + metadata=metadata, + chunked_value=True, + ) + chunk_two = self._make_partial_result_set( + [self._make_list_value(["0000-0000-000000000000"])], + last=True, + ) + iterator = _MockCancellableIterator(chunk_one, chunk_two) + streamed = self._make_one(iterator) + found = [row async for row in streamed] + self.assertEqual( + found, + [ + [ + [ + uuid.UUID("12345678-1234-5678-1234-567812345678"), + uuid.UUID("00000000-0000-0000-0000-000000000000"), + ] + ] + ], + ) + self.assertIsNone(streamed._pending_chunk) + class _MockCancellableIterator(object): cancel_calls = 0 @@ -1692,6 +2613,18 @@ async def test_multiple_row_multiple_chunks(self): async def test_multiple_row_chunks_non_chunks_interleaved(self): await self._match_results("Multiple Row Chunks/Non Chunks Interleaved") + @CrossSync.pytest + async def test_json_array_chunking(self): + await self._match_results("JSON Array Chunking Test") + + @CrossSync.pytest + async def test_numeric_array_chunking(self): + await self._match_results("NUMERIC Array Chunking Test") + + @CrossSync.pytest + async def test_nested_struct_array_chunking_with_null(self): + await self._match_results("Nested Struct Array Chunking Test With null") + def _generate_partial_result_sets(prs_text_pbs): from google.cloud.spanner_v1 import PartialResultSet @@ -1727,7 +2660,9 @@ def _normalize_float(cell): def _normalize_results(rows_data, fields): """Helper for _parse_streaming_read_acceptance_tests""" - from google.cloud.spanner_v1 import TypeCode + import decimal + + from google.cloud.spanner_v1 import JsonObject, TypeCode normalized = [] for row_data in rows_data: @@ -1735,16 +2670,32 @@ def _normalize_results(rows_data, fields): assert len(row_data) == len(fields) for cell, field in zip(row_data, fields): if field.type_.code == TypeCode.INT64: - cell = int(cell) - if field.type_.code == TypeCode.FLOAT64: + cell = int(cell) if cell is not None else None + elif field.type_.code == TypeCode.FLOAT64: cell = _normalize_float(cell) + elif field.type_.code == TypeCode.NUMERIC: + cell = decimal.Decimal(cell) if cell is not None else None + elif field.type_.code == TypeCode.JSON: + cell = JsonObject.from_str(cell) if cell is not None else None elif field.type_.code == TypeCode.BYTES: - cell = cell.encode("utf8") + cell = cell.encode("utf8") if cell is not None else None elif field.type_.code == TypeCode.ARRAY: - if field.type_.array_element_type.code == TypeCode.INT64: + if cell is None: + pass + elif field.type_.array_element_type.code == TypeCode.INT64: cell = _normalize_int_array(cell) elif field.type_.array_element_type.code == TypeCode.FLOAT64: cell = [_normalize_float(subcell) for subcell in cell] + elif field.type_.array_element_type.code == TypeCode.NUMERIC: + cell = [ + decimal.Decimal(subcell) if subcell is not None else None + for subcell in cell + ] + elif field.type_.array_element_type.code == TypeCode.JSON: + cell = [ + JsonObject.from_str(subcell) if subcell is not None else None + for subcell in cell + ] row.append(cell) normalized.append(row) return normalized diff --git a/packages/google-cloud-spanner/tests/unit/streaming-read-acceptance-test.json b/packages/google-cloud-spanner/tests/unit/streaming-read-acceptance-test.json index 9b44b4077812..557e787a36f2 100644 --- a/packages/google-cloud-spanner/tests/unit/streaming-read-acceptance-test.json +++ b/packages/google-cloud-spanner/tests/unit/streaming-read-acceptance-test.json @@ -213,5 +213,43 @@ "{\n \"values\": [\"f\"]\n}" ], "name": "Multiple Row Chunks/Non Chunks Interleaved" + }, + { + "result": {"value": [[[ + "{}", + "{\"color\":\"red\",\"value\":\"#f00\"}" + ]]]}, + "chunks": [ + "{\n \"metadata\": {\n \"rowType\": {\n \"fields\": [{\n \"name\": \"f1\",\n \"type\": {\n \"code\": \"ARRAY\",\n \"arrayElementType\": {\n \"code\": \"JSON\"\n }\n }\n }]\n }\n },\n \"values\": [[\"{}\", \"{\\\"color\\\":\\\"red\\\"\"]],\n \"chunkedValue\": true\n}", + "{\n \"values\": [[\",\\\"value\\\":\\\"#f00\\\"}\"]],\n \"chunkedValue\": false\n}" + ], + "name": "JSON Array Chunking Test" + }, + { + "result": {"value": [[[ + "1.1", + "2.3", + "4.0", + null, + "5.5" + ]]]}, + "chunks": [ + "{\n \"metadata\": {\n \"rowType\": {\n \"fields\": [{\n \"name\": \"f1\",\n \"type\": {\n \"code\": \"ARRAY\",\n \"arrayElementType\": {\n \"code\": \"NUMERIC\"\n }\n }\n }]\n }\n },\n \"values\": [[\"1.1\", \"2.\"]],\n \"chunkedValue\": true\n}", + "{\n \"values\": [[\"3\", \"4.0\"]],\n \"chunkedValue\": true\n}", + "{\n \"values\": [[\"\", null, \"5.5\"]]\n}" + ], + "name": "NUMERIC Array Chunking Test" + }, + { + "result": {"value": [[[[[ + ["abc"], + null, + ["def"] + ]]]]]}, + "chunks": [ + "{\n \"metadata\": {\n \"rowType\": {\n \"fields\": [{\n \"name\": \"f1\",\n \"type\": {\n \"code\": \"ARRAY\",\n \"arrayElementType\": {\n \"code\": \"STRUCT\",\n \"structType\": {\n \"fields\": [{\n \"name\": \"f11\",\n \"type\": {\n \"code\": \"ARRAY\",\n \"arrayElementType\": {\n \"code\": \"STRUCT\",\n \"structType\": {\n \"fields\": [{\n \"name\": \"f12\",\n \"type\": {\n \"code\": \"STRING\"\n }\n }]\n }\n }\n }\n }]\n }\n }\n }\n }]\n }\n },\n \"values\": [[[[[\"abc\"], null]]]],\n \"chunkedValue\": true\n}", + "{\n \"values\": [[[[[\"def\"]]]]]\n}" + ], + "name": "Nested Struct Array Chunking Test With null" } ]} diff --git a/packages/google-cloud-spanner/tests/unit/test_streamed.py b/packages/google-cloud-spanner/tests/unit/test_streamed.py index 4c8ff04d3193..2959b49568a6 100644 --- a/packages/google-cloud-spanner/tests/unit/test_streamed.py +++ b/packages/google-cloud-spanner/tests/unit/test_streamed.py @@ -258,6 +258,98 @@ def test__merge_chunk_float64_w_float64(self): with self.assertRaises(Unmergeable): streamed._merge_chunk(chunk) + def test__merge_chunk_float32_nan_string(self): + from google.cloud.spanner_v1 import TypeCode + + iterator = _MockCancellableIterator() + streamed = self._make_one(iterator) + fields = [self._make_scalar_field("weight", TypeCode.FLOAT32)] + streamed._metadata = self._make_result_set_metadata(fields) + streamed._pending_chunk = self._make_value("Na") + chunk = self._make_value("N") + + merged = streamed._merge_chunk(chunk) + self.assertEqual(merged.string_value, "NaN") + + def test__merge_chunk_float32_w_empty(self): + from google.cloud.spanner_v1 import TypeCode + + iterator = _MockCancellableIterator() + streamed = self._make_one(iterator) + fields = [self._make_scalar_field("weight", TypeCode.FLOAT32)] + streamed._metadata = self._make_result_set_metadata(fields) + streamed._pending_chunk = self._make_value(1.25) + chunk = self._make_value("") + + merged = streamed._merge_chunk(chunk) + self.assertEqual(merged.number_value, 1.25) + + def test__merge_chunk_float32_w_float32(self): + from google.cloud.spanner_v1 import TypeCode + from google.cloud.spanner_v1.streamed import Unmergeable + + iterator = _MockCancellableIterator() + streamed = self._make_one(iterator) + fields = [self._make_scalar_field("weight", TypeCode.FLOAT32)] + streamed._metadata = self._make_result_set_metadata(fields) + streamed._pending_chunk = self._make_value(1.25) + chunk = self._make_value(2.5) + + with self.assertRaises(Unmergeable): + streamed._merge_chunk(chunk) + + def test__merge_chunk_date(self): + from google.cloud.spanner_v1 import TypeCode + + iterator = _MockCancellableIterator() + streamed = self._make_one(iterator) + fields = [self._make_scalar_field("anniversary", TypeCode.DATE)] + streamed._metadata = self._make_result_set_metadata(fields) + streamed._pending_chunk = self._make_value("2026-") + chunk = self._make_value("10-08") + + merged = streamed._merge_chunk(chunk) + self.assertEqual(merged.string_value, "2026-10-08") + + def test__merge_chunk_timestamp(self): + from google.cloud.spanner_v1 import TypeCode + + iterator = _MockCancellableIterator() + streamed = self._make_one(iterator) + fields = [self._make_scalar_field("created_at", TypeCode.TIMESTAMP)] + streamed._metadata = self._make_result_set_metadata(fields) + streamed._pending_chunk = self._make_value("2026-10-08T13:") + chunk = self._make_value("15:00Z") + + merged = streamed._merge_chunk(chunk) + self.assertEqual(merged.string_value, "2026-10-08T13:15:00Z") + + def test__merge_chunk_interval(self): + from google.cloud.spanner_v1 import TypeCode + + iterator = _MockCancellableIterator() + streamed = self._make_one(iterator) + fields = [self._make_scalar_field("duration", TypeCode.INTERVAL)] + streamed._metadata = self._make_result_set_metadata(fields) + streamed._pending_chunk = self._make_value("P1Y2M") + chunk = self._make_value("3DT4H") + + merged = streamed._merge_chunk(chunk) + self.assertEqual(merged.string_value, "P1Y2M3DT4H") + + def test__merge_chunk_uuid(self): + from google.cloud.spanner_v1 import TypeCode + + iterator = _MockCancellableIterator() + streamed = self._make_one(iterator) + fields = [self._make_scalar_field("session_id", TypeCode.UUID)] + streamed._metadata = self._make_result_set_metadata(fields) + streamed._pending_chunk = self._make_value("12345678-1234-") + chunk = self._make_value("5678-1234-567812345678") + + merged = streamed._merge_chunk(chunk) + self.assertEqual(merged.string_value, "12345678-1234-5678-1234-567812345678") + def test__merge_chunk_string(self): from google.cloud.spanner_v1 import TypeCode @@ -1448,6 +1540,791 @@ def test_merge_values_zero_fields(self): streamed._merge_values([self._make_value(1)]) self.assertEqual(streamed._rows, []) + def test___iter___chunked_values_empty_string_start(self): + from google.cloud.spanner_v1 import TypeCode + + fields = [ + self._make_scalar_field("column_one", TypeCode.STRING), + self._make_scalar_field("column_two", TypeCode.STRING), + ] + metadata = self._make_result_set_metadata(fields) + chunk_one = self._make_partial_result_set( + [self._make_value("")], + metadata=metadata, + chunked_value=True, + ) + chunk_two = self._make_partial_result_set( + [ + self._make_value("tail-after-empty"), + self._make_value("value_two"), + self._make_value("value_three"), + self._make_value(""), + ], + chunked_value=True, + ) + chunk_three = self._make_partial_result_set( + [self._make_value("tail-two")], + last=True, + ) + iterator = _MockCancellableIterator(chunk_one, chunk_two, chunk_three) + streamed = self._make_one(iterator) + found = list(streamed) + self.assertEqual( + found, + [ + ["tail-after-empty", "value_two"], + ["value_three", "tail-two"], + ], + ) + self.assertIsNone(streamed._pending_chunk) + + def test___iter___chunked_values_empty_list_head_and_tail(self): + from google.cloud.spanner_v1 import TypeCode + + fields = [ + self._make_array_field("column_one", TypeCode.STRING), + self._make_array_field("column_two", TypeCode.STRING), + ] + metadata = self._make_result_set_metadata(fields) + chunk_one = self._make_partial_result_set( + [self._make_list_value([])], + metadata=metadata, + chunked_value=True, + ) + chunk_two = self._make_partial_result_set( + [ + self._make_list_value(["a", "b"]), + self._make_list_value([]), + ], + chunked_value=True, + ) + chunk_three = self._make_partial_result_set( + [ + self._make_list_value(["v1", "v2"]), + self._make_list_value(["x", "y"]), + ], + chunked_value=True, + ) + chunk_four = self._make_partial_result_set( + [ + self._make_list_value([]), + self._make_list_value(["w1", "w2"]), + ], + chunked_value=True, + ) + chunk_five = self._make_partial_result_set( + [ + self._make_list_value([]), + self._make_list_value([]), + ], + chunked_value=True, + ) + chunk_six = self._make_partial_result_set( + [ + self._make_list_value([]), + self._make_list_value(["z1", "z2"]), + ], + last=True, + ) + iterator = _MockCancellableIterator( + chunk_one, chunk_two, chunk_three, chunk_four, chunk_five, chunk_six + ) + streamed = self._make_one(iterator) + found = list(streamed) + self.assertEqual( + found, + [ + [["a", "b"], ["v1", "v2"]], + [["x", "y"], ["w1", "w2"]], + [[], ["z1", "z2"]], + ], + ) + self.assertIsNone(streamed._pending_chunk) + + def test___iter___chunked_values_bytes_and_bytes_array(self): + from google.protobuf.struct_pb2 import Value + + from google.cloud.spanner_v1 import TypeCode + + fields = [ + self._make_scalar_field("raw_bytes", TypeCode.BYTES), + self._make_array_field("bytes_array", TypeCode.BYTES), + ] + metadata = self._make_result_set_metadata(fields) + + full_bytes = b"hello-chunked-bytes-world" + bytes_part_one = "hello-chunked-" + bytes_part_two = "bytes-world" + + element_one = b"first-element" + element_two = b"second-chunked-element" + element_two_part_one = "second-chunked-" + element_two_part_two = "element" + element_three = b"third-element" + + chunk_one = self._make_partial_result_set( + [self._make_value(bytes_part_one)], + metadata=metadata, + chunked_value=True, + ) + chunk_two = self._make_partial_result_set( + [ + self._make_value(bytes_part_two), + self._make_list_value( + value_pbs=[ + self._make_value("first-element"), + Value(null_value="NULL_VALUE"), + self._make_value(element_two_part_one), + ] + ), + ], + chunked_value=True, + ) + chunk_three = self._make_partial_result_set( + [ + self._make_list_value( + value_pbs=[ + self._make_value(element_two_part_two), + self._make_value("third-element"), + ] + ) + ], + last=True, + ) + iterator = _MockCancellableIterator(chunk_one, chunk_two, chunk_three) + streamed = self._make_one(iterator) + found = list(streamed) + self.assertEqual( + found, + [ + [ + full_bytes, + [element_one, None, element_two, element_three], + ] + ], + ) + self.assertIsNone(streamed._pending_chunk) + + def test___iter___chunked_values_intermediate_single_value_chunk(self): + from google.cloud.spanner_v1 import TypeCode + + fields = [ + self._make_scalar_field("identifier", TypeCode.STRING), + self._make_scalar_field("text_payload", TypeCode.STRING), + ] + metadata = self._make_result_set_metadata(fields) + chunk_one = self._make_partial_result_set( + [self._make_value("id1"), self._make_value("hello-")], + metadata=metadata, + chunked_value=True, + ) + chunk_two = self._make_partial_result_set( + [self._make_value("world-")], + chunked_value=True, + ) + chunk_three = self._make_partial_result_set( + [ + self._make_value("again"), + self._make_value("id2"), + self._make_value("text2"), + ], + last=True, + ) + iterator = _MockCancellableIterator(chunk_one, chunk_two, chunk_three) + streamed = self._make_one(iterator) + found = list(streamed) + self.assertEqual( + found, + [ + ["id1", "hello-world-again"], + ["id2", "text2"], + ], + ) + self.assertIsNone(streamed._pending_chunk) + + def test___iter___multi_row_multi_chunk_interleaved_splits(self): + from google.cloud.spanner_v1 import TypeCode + + fields = [ + self._make_scalar_field("identifier", TypeCode.INT64), + self._make_scalar_field("text_payload", TypeCode.STRING), + ] + metadata = self._make_result_set_metadata(fields) + chunk_one = self._make_partial_result_set( + [self._make_value(1), self._make_value("spl")], + metadata=metadata, + chunked_value=True, + ) + chunk_two = self._make_partial_result_set( + [ + self._make_value("it_text_row1"), + self._make_value(2), + self._make_value("split_"), + ], + chunked_value=True, + ) + chunk_three = self._make_partial_result_set( + [self._make_value("text_row2")], + last=True, + ) + iterator = _MockCancellableIterator(chunk_one, chunk_two, chunk_three) + streamed = self._make_one(iterator) + found = list(streamed) + self.assertEqual( + found, + [ + [1, "split_text_row1"], + [2, "split_text_row2"], + ], + ) + self.assertIsNone(streamed._pending_chunk) + + def test___iter___subsequent_metadata_ignored(self): + from google.cloud.spanner_v1 import TypeCode + + fields_one = [self._make_scalar_field("primary_column", TypeCode.STRING)] + metadata_one = self._make_result_set_metadata(fields_one) + fields_two = [self._make_scalar_field("ignored_column", TypeCode.STRING)] + metadata_two = self._make_result_set_metadata(fields_two) + + chunk_one = self._make_partial_result_set( + [self._make_value("alpha")], + metadata=metadata_one, + ) + chunk_two = self._make_partial_result_set( + [self._make_value("beta")], + metadata=metadata_two, + last=True, + ) + iterator = _MockCancellableIterator(chunk_one, chunk_two) + streamed = self._make_one(iterator) + found = list(streamed) + self.assertEqual(found, [["alpha"], ["beta"]]) + self.assertEqual(streamed.fields[0].name, "primary_column") + + def test___iter___empty_chunk_in_middle_of_stream(self): + from google.cloud.spanner_v1 import TypeCode + + fields = [self._make_scalar_field("column_name", TypeCode.STRING)] + metadata = self._make_result_set_metadata(fields) + chunk_one = self._make_partial_result_set( + [self._make_value("value_one")], + metadata=metadata, + ) + chunk_two = self._make_partial_result_set([]) + chunk_three = self._make_partial_result_set( + [self._make_value("value_two")], + last=True, + ) + iterator = _MockCancellableIterator(chunk_one, chunk_two, chunk_three) + streamed = self._make_one(iterator) + found = list(streamed) + self.assertEqual(found, [["value_one"], ["value_two"]]) + + def test___iter___chunked_values_incompatible_types_error(self): + from google.cloud.spanner_v1 import TypeCode + from google.cloud.spanner_v1.streamed import Unmergeable + + fields = [self._make_scalar_field("registered", TypeCode.BOOL)] + metadata = self._make_result_set_metadata(fields) + chunk_one = self._make_partial_result_set( + [self._make_value(True)], + metadata=metadata, + chunked_value=True, + ) + chunk_two = self._make_partial_result_set( + [self._make_value(False)], + last=True, + ) + iterator = _MockCancellableIterator(chunk_one, chunk_two) + streamed = self._make_one(iterator) + with self.assertRaises(Unmergeable): + list(streamed) + + def test___iter___resumption_across_chunked_value_retry(self): + from google.api_core.exceptions import ServiceUnavailable + + from google.cloud.spanner_v1 import TypeCode + + fields = [ + self._make_scalar_field("identifier", TypeCode.STRING), + self._make_scalar_field("data_payload", TypeCode.STRING), + ] + metadata = self._make_result_set_metadata(fields) + + attempt_count = 0 + last_resume_token = None + + def create_stream_attempt(resume_token=None): + nonlocal attempt_count + attempt_count += 1 + if attempt_count == 1: + self.assertIsNone(resume_token) + chunk_one = self._make_partial_result_set( + [self._make_value("id1"), self._make_value("part1-")], + metadata=metadata, + chunked_value=True, + ) + chunk_one.resume_token = b"token-1" + yield chunk_one + raise ServiceUnavailable("transient error during chunked scalar") + else: + self.assertEqual(resume_token, b"token-1") + chunk_two = self._make_partial_result_set( + [self._make_value("part2-actual-")], + chunked_value=True, + ) + chunk_two.resume_token = b"token-2" + yield chunk_two + chunk_three = self._make_partial_result_set( + [self._make_value("part3-")], + chunked_value=True, + ) + yield chunk_three + chunk_four = self._make_partial_result_set( + [ + self._make_value("part4"), + self._make_value("id2"), + self._make_value("data2"), + ], + last=True, + ) + yield chunk_four + + class ResumableResponseIterator: + def __init__(self): + self._current_generator = create_stream_attempt(None) + + def __iter__(self): + return self + + def __next__(self): + nonlocal last_resume_token + while True: + try: + item = next(self._current_generator) + if item.resume_token: + last_resume_token = item.resume_token + return item + except ServiceUnavailable: + self._current_generator = create_stream_attempt( + last_resume_token + ) + + streamed = self._make_one(ResumableResponseIterator()) + found = list(streamed) + self.assertEqual( + found, + [ + ["id1", "part1-part2-actual-part3-part4"], + ["id2", "data2"], + ], + ) + self.assertEqual(attempt_count, 2) + self.assertIsNone(streamed._pending_chunk) + + def test___iter___resumption_across_chunked_array_retry(self): + from google.api_core.exceptions import ServiceUnavailable + + from google.cloud.spanner_v1 import TypeCode + + fields = [ + self._make_scalar_field("identifier", TypeCode.STRING), + self._make_array_field("tags", TypeCode.STRING), + ] + metadata = self._make_result_set_metadata(fields) + + attempt_count = 0 + last_resume_token = None + + def create_stream_attempt(resume_token=None): + nonlocal attempt_count + attempt_count += 1 + if attempt_count == 1: + self.assertIsNone(resume_token) + chunk_one = self._make_partial_result_set( + [ + self._make_value("id1"), + self._make_list_value(["tag1", "tag2-"]), + ], + metadata=metadata, + chunked_value=True, + ) + chunk_one.resume_token = b"token-array" + yield chunk_one + raise ServiceUnavailable("transient error during chunked array") + else: + self.assertEqual(resume_token, b"token-array") + chunk_two = self._make_partial_result_set( + [self._make_list_value(["tag2-tail", "tag3"])], + last=True, + ) + yield chunk_two + + class ResumableResponseIterator: + def __init__(self): + self._current_generator = create_stream_attempt(None) + + def __iter__(self): + return self + + def __next__(self): + nonlocal last_resume_token + while True: + try: + item = next(self._current_generator) + if item.resume_token: + last_resume_token = item.resume_token + return item + except ServiceUnavailable: + self._current_generator = create_stream_attempt( + last_resume_token + ) + + streamed = self._make_one(ResumableResponseIterator()) + found = list(streamed) + self.assertEqual( + found, + [ + ["id1", ["tag1", "tag2-tag2-tail", "tag3"]], + ], + ) + self.assertEqual(attempt_count, 2) + self.assertIsNone(streamed._pending_chunk) + + def test___iter___chunked_values_numeric(self): + import decimal + + from google.cloud.spanner_v1 import TypeCode + + fields = [self._make_scalar_field("amount", TypeCode.NUMERIC)] + metadata = self._make_result_set_metadata(fields) + chunk_one = self._make_partial_result_set( + [self._make_value("1234.")], + metadata=metadata, + chunked_value=True, + ) + chunk_two = self._make_partial_result_set( + [self._make_value("5678")], + last=True, + ) + iterator = _MockCancellableIterator(chunk_one, chunk_two) + streamed = self._make_one(iterator) + found = list(streamed) + self.assertEqual(found, [[decimal.Decimal("1234.5678")]]) + self.assertIsNone(streamed._pending_chunk) + + def test___iter___chunked_values_json(self): + from google.cloud.spanner_v1 import JsonObject, TypeCode + + fields = [self._make_scalar_field("payload", TypeCode.JSON)] + metadata = self._make_result_set_metadata(fields) + chunk_one = self._make_partial_result_set( + [self._make_value('{"user": "')], + metadata=metadata, + chunked_value=True, + ) + chunk_two = self._make_partial_result_set( + [self._make_value('alice", "score": 10}')], + last=True, + ) + iterator = _MockCancellableIterator(chunk_one, chunk_two) + streamed = self._make_one(iterator) + found = list(streamed) + self.assertEqual(found, [[JsonObject({"user": "alice", "score": 10})]]) + self.assertIsNone(streamed._pending_chunk) + + def test___iter___chunked_values_lazy_decode(self): + from google.protobuf.struct_pb2 import Value + + from google.cloud.spanner_v1 import TypeCode + + fields = [self._make_scalar_field("name", TypeCode.STRING)] + metadata = self._make_result_set_metadata(fields) + chunk_one = self._make_partial_result_set( + [self._make_value("hello-")], + metadata=metadata, + chunked_value=True, + ) + chunk_two = self._make_partial_result_set( + [self._make_value("world")], + last=True, + ) + iterator = _MockCancellableIterator(chunk_one, chunk_two) + streamed = self._make_one(iterator, lazy_decode=True) + found = list(streamed) + self.assertEqual(found, [[Value(string_value="hello-world")]]) + self.assertEqual(streamed.decode_row(found[0]), ["hello-world"]) + self.assertIsNone(streamed._pending_chunk) + + def test___iter___chunked_values_float32(self): + from google.cloud.spanner_v1 import TypeCode + + fields = [self._make_scalar_field("ratio", TypeCode.FLOAT32)] + metadata = self._make_result_set_metadata(fields) + chunk_one = self._make_partial_result_set( + [self._make_value("Infi")], + metadata=metadata, + chunked_value=True, + ) + chunk_two = self._make_partial_result_set( + [self._make_value("nity")], + last=True, + ) + iterator = _MockCancellableIterator(chunk_one, chunk_two) + streamed = self._make_one(iterator) + found = list(streamed) + self.assertEqual(found, [[float("inf")]]) + self.assertIsNone(streamed._pending_chunk) + + def test___iter___chunked_values_float32_array(self): + from google.cloud.spanner_v1 import TypeCode + + fields = [self._make_array_field("ratios", TypeCode.FLOAT32)] + metadata = self._make_result_set_metadata(fields) + chunk_one = self._make_partial_result_set( + [self._make_list_value([1.5, 2.5])], + metadata=metadata, + chunked_value=True, + ) + chunk_two = self._make_partial_result_set( + [self._make_list_value(["", 3.5])], + last=True, + ) + iterator = _MockCancellableIterator(chunk_one, chunk_two) + streamed = self._make_one(iterator) + found = list(streamed) + self.assertEqual(found, [[[1.5, 2.5, 3.5]]]) + self.assertIsNone(streamed._pending_chunk) + + def test___iter___chunked_values_date(self): + import datetime + + from google.cloud.spanner_v1 import TypeCode + + fields = [self._make_scalar_field("anniversary", TypeCode.DATE)] + metadata = self._make_result_set_metadata(fields) + chunk_one = self._make_partial_result_set( + [self._make_value("2026-")], + metadata=metadata, + chunked_value=True, + ) + chunk_two = self._make_partial_result_set( + [self._make_value("10-08")], + last=True, + ) + iterator = _MockCancellableIterator(chunk_one, chunk_two) + streamed = self._make_one(iterator) + found = list(streamed) + self.assertEqual(found, [[datetime.date(2026, 10, 8)]]) + self.assertIsNone(streamed._pending_chunk) + + def test___iter___chunked_values_date_array(self): + import datetime + + from google.cloud.spanner_v1 import TypeCode + + fields = [self._make_array_field("dates", TypeCode.DATE)] + metadata = self._make_result_set_metadata(fields) + chunk_one = self._make_partial_result_set( + [self._make_list_value(["2026-01-01", "2026-02-"])], + metadata=metadata, + chunked_value=True, + ) + chunk_two = self._make_partial_result_set( + [self._make_list_value(["01", "2026-03-01"])], + last=True, + ) + iterator = _MockCancellableIterator(chunk_one, chunk_two) + streamed = self._make_one(iterator) + found = list(streamed) + self.assertEqual( + found, + [ + [ + [ + datetime.date(2026, 1, 1), + datetime.date(2026, 2, 1), + datetime.date(2026, 3, 1), + ] + ] + ], + ) + self.assertIsNone(streamed._pending_chunk) + + def test___iter___chunked_values_timestamp(self): + import datetime + + from google.cloud.spanner_v1 import TypeCode + + fields = [self._make_scalar_field("created_at", TypeCode.TIMESTAMP)] + metadata = self._make_result_set_metadata(fields) + chunk_one = self._make_partial_result_set( + [self._make_value("2026-10-08T13:")], + metadata=metadata, + chunked_value=True, + ) + chunk_two = self._make_partial_result_set( + [self._make_value("15:00Z")], + last=True, + ) + iterator = _MockCancellableIterator(chunk_one, chunk_two) + streamed = self._make_one(iterator) + found = list(streamed) + self.assertEqual( + found, + [[datetime.datetime(2026, 10, 8, 13, 15, 0, tzinfo=datetime.timezone.utc)]], + ) + self.assertIsNone(streamed._pending_chunk) + + def test___iter___chunked_values_timestamp_array(self): + import datetime + + from google.cloud.spanner_v1 import TypeCode + + fields = [self._make_array_field("timestamps", TypeCode.TIMESTAMP)] + metadata = self._make_result_set_metadata(fields) + chunk_one = self._make_partial_result_set( + [self._make_list_value(["2026-01-01T00:00:00Z", "2026-02-01T00:"])], + metadata=metadata, + chunked_value=True, + ) + chunk_two = self._make_partial_result_set( + [self._make_list_value(["00:00Z", "2026-03-01T00:00:00Z"])], + last=True, + ) + iterator = _MockCancellableIterator(chunk_one, chunk_two) + streamed = self._make_one(iterator) + found = list(streamed) + self.assertEqual( + found, + [ + [ + [ + datetime.datetime( + 2026, 1, 1, 0, 0, 0, tzinfo=datetime.timezone.utc + ), + datetime.datetime( + 2026, 2, 1, 0, 0, 0, tzinfo=datetime.timezone.utc + ), + datetime.datetime( + 2026, 3, 1, 0, 0, 0, tzinfo=datetime.timezone.utc + ), + ] + ] + ], + ) + self.assertIsNone(streamed._pending_chunk) + + def test___iter___chunked_values_interval(self): + from google.cloud.spanner_v1 import Interval, TypeCode + + fields = [self._make_scalar_field("duration", TypeCode.INTERVAL)] + metadata = self._make_result_set_metadata(fields) + chunk_one = self._make_partial_result_set( + [self._make_value("P1Y2M")], + metadata=metadata, + chunked_value=True, + ) + chunk_two = self._make_partial_result_set( + [self._make_value("3DT4H")], + last=True, + ) + iterator = _MockCancellableIterator(chunk_one, chunk_two) + streamed = self._make_one(iterator) + found = list(streamed) + self.assertEqual(found, [[Interval.from_str("P1Y2M3DT4H")]]) + self.assertIsNone(streamed._pending_chunk) + + def test___iter___chunked_values_interval_array(self): + from google.cloud.spanner_v1 import Interval, TypeCode + + fields = [self._make_array_field("durations", TypeCode.INTERVAL)] + metadata = self._make_result_set_metadata(fields) + chunk_one = self._make_partial_result_set( + [self._make_list_value(["P1Y", "P2"])], + metadata=metadata, + chunked_value=True, + ) + chunk_two = self._make_partial_result_set( + [self._make_list_value(["M", "P3D"])], + last=True, + ) + iterator = _MockCancellableIterator(chunk_one, chunk_two) + streamed = self._make_one(iterator) + found = list(streamed) + self.assertEqual( + found, + [ + [ + [ + Interval.from_str("P1Y"), + Interval.from_str("P2M"), + Interval.from_str("P3D"), + ] + ] + ], + ) + self.assertIsNone(streamed._pending_chunk) + + def test___iter___chunked_values_uuid(self): + import uuid + + from google.cloud.spanner_v1 import TypeCode + + fields = [self._make_scalar_field("session_id", TypeCode.UUID)] + metadata = self._make_result_set_metadata(fields) + chunk_one = self._make_partial_result_set( + [self._make_value("12345678-1234-")], + metadata=metadata, + chunked_value=True, + ) + chunk_two = self._make_partial_result_set( + [self._make_value("5678-1234-567812345678")], + last=True, + ) + iterator = _MockCancellableIterator(chunk_one, chunk_two) + streamed = self._make_one(iterator) + found = list(streamed) + self.assertEqual(found, [[uuid.UUID("12345678-1234-5678-1234-567812345678")]]) + self.assertIsNone(streamed._pending_chunk) + + def test___iter___chunked_values_uuid_array(self): + import uuid + + from google.cloud.spanner_v1 import TypeCode + + fields = [self._make_array_field("ids", TypeCode.UUID)] + metadata = self._make_result_set_metadata(fields) + chunk_one = self._make_partial_result_set( + [ + self._make_list_value( + ["12345678-1234-5678-1234-567812345678", "00000000-0000-"] + ) + ], + metadata=metadata, + chunked_value=True, + ) + chunk_two = self._make_partial_result_set( + [self._make_list_value(["0000-0000-000000000000"])], + last=True, + ) + iterator = _MockCancellableIterator(chunk_one, chunk_two) + streamed = self._make_one(iterator) + found = list(streamed) + self.assertEqual( + found, + [ + [ + [ + uuid.UUID("12345678-1234-5678-1234-567812345678"), + uuid.UUID("00000000-0000-0000-0000-000000000000"), + ] + ] + ], + ) + self.assertIsNone(streamed._pending_chunk) + class _MockCancellableIterator(object): cancel_calls = 0 @@ -1565,6 +2442,15 @@ def test_multiple_row_multiple_chunks(self): def test_multiple_row_chunks_non_chunks_interleaved(self): self._match_results("Multiple Row Chunks/Non Chunks Interleaved") + def test_json_array_chunking(self): + self._match_results("JSON Array Chunking Test") + + def test_numeric_array_chunking(self): + self._match_results("NUMERIC Array Chunking Test") + + def test_nested_struct_array_chunking_with_null(self): + self._match_results("Nested Struct Array Chunking Test With null") + def _generate_partial_result_sets(prs_text_pbs): from google.cloud.spanner_v1 import PartialResultSet @@ -1600,7 +2486,9 @@ def _normalize_float(cell): def _normalize_results(rows_data, fields): """Helper for _parse_streaming_read_acceptance_tests""" - from google.cloud.spanner_v1 import TypeCode + import decimal + + from google.cloud.spanner_v1 import JsonObject, TypeCode normalized = [] for row_data in rows_data: @@ -1608,16 +2496,32 @@ def _normalize_results(rows_data, fields): assert len(row_data) == len(fields) for cell, field in zip(row_data, fields): if field.type_.code == TypeCode.INT64: - cell = int(cell) - if field.type_.code == TypeCode.FLOAT64: + cell = int(cell) if cell is not None else None + elif field.type_.code == TypeCode.FLOAT64: cell = _normalize_float(cell) + elif field.type_.code == TypeCode.NUMERIC: + cell = decimal.Decimal(cell) if cell is not None else None + elif field.type_.code == TypeCode.JSON: + cell = JsonObject.from_str(cell) if cell is not None else None elif field.type_.code == TypeCode.BYTES: - cell = cell.encode("utf8") + cell = cell.encode("utf8") if cell is not None else None elif field.type_.code == TypeCode.ARRAY: - if field.type_.array_element_type.code == TypeCode.INT64: + if cell is None: + pass + elif field.type_.array_element_type.code == TypeCode.INT64: cell = _normalize_int_array(cell) elif field.type_.array_element_type.code == TypeCode.FLOAT64: cell = [_normalize_float(subcell) for subcell in cell] + elif field.type_.array_element_type.code == TypeCode.NUMERIC: + cell = [ + decimal.Decimal(subcell) if subcell is not None else None + for subcell in cell + ] + elif field.type_.array_element_type.code == TypeCode.JSON: + cell = [ + JsonObject.from_str(subcell) if subcell is not None else None + for subcell in cell + ] row.append(cell) normalized.append(row) return normalized