From fffd60ea2f08fa5557b2a11b9e98354fc8df751b Mon Sep 17 00:00:00 2001 From: Vaibhav Dewangan Date: Fri, 31 Jul 2026 21:52:43 +0530 Subject: [PATCH] fix(batch): set has_errors on per-object/reference errors in stream recv In _BatchBaseSync/_BatchBaseAsync.__recv, each failed object or reference from the server was wrapped in a fresh BatchObjectReturn/BatchReferenceReturn with only `errors` set. has_errors defaults to False on the dataclass and __add__ only ORs it forward (`self.has_errors or other.has_errors`), so it never became True even though `errors` was non-empty. This meant `if result.has_errors:` never fired on data.ingest() (and on batch.stream(), though failed_objects/failed_references were populated correctly there since those are tracked separately). Fixes #2107 --- mock_tests/test_batch.py | 41 ++++++++++++++++++++++++++++ weaviate/collections/batch/async_.py | 2 ++ weaviate/collections/batch/sync.py | 2 ++ 3 files changed, 45 insertions(+) diff --git a/mock_tests/test_batch.py b/mock_tests/test_batch.py index c00bd08db..f8b114534 100644 --- a/mock_tests/test_batch.py +++ b/mock_tests/test_batch.py @@ -60,3 +60,44 @@ def test_ssb_canceled_stream( for i in range(HOW_MANY): batch.add_object({"name": f"Object {i}"}) assert len(service.uuids) == HOW_MANY + + +class MockFailedObjectWeaviateService(weaviate_pb2_grpc.WeaviateServicer): + def BatchStream( + self, + request_iterator: Generator[batch_pb2.BatchStreamRequest, None, None], + context: grpc.ServicerContext, + ) -> Generator[batch_pb2.BatchStreamReply, None, None]: + yield batch_pb2.BatchStreamReply(started=batch_pb2.BatchStreamReply.Started()) + for request in request_iterator: + if request.HasField("data"): + uuids = [obj.uuid for obj in request.data.objects.values] + yield batch_pb2.BatchStreamReply( + results=batch_pb2.BatchStreamReply.Results( + errors=[ + batch_pb2.BatchStreamReply.Results.Error( + uuid=uuid, error="mock failure" + ) + for uuid in uuids + ] + ) + ) + if request.HasField("stop"): + return + + +@pytest.fixture(scope="function") +def failed_object_stream( + canceled_stream_client: weaviate.WeaviateClient, start_grpc_server: grpc.Server +): + service = MockFailedObjectWeaviateService() + weaviate_pb2_grpc.add_WeaviateServicer_to_server(service, start_grpc_server) + return canceled_stream_client.collections.use(mock_class["class"]) + + +def test_ingest_has_errors_on_failed_object( + failed_object_stream: weaviate.collections.Collection, +): + result = failed_object_stream.data.ingest([{"name": "Object 1"}]) + assert result.has_errors is True + assert len(result.errors) == 1 diff --git a/weaviate/collections/batch/async_.py b/weaviate/collections/batch/async_.py index 8b997586c..2e7aa9279 100644 --- a/weaviate/collections/batch/async_.py +++ b/weaviate/collections/batch/async_.py @@ -407,6 +407,7 @@ async def __recv(self) -> None: result_objs += BatchObjectReturn( _all_responses=[err], errors={cached.index: err}, + has_errors=True, ) failed_objs.append(err) logger.warning( @@ -428,6 +429,7 @@ async def __recv(self) -> None: ) result_refs += BatchReferenceReturn( errors={cached.index: err}, + has_errors=True, ) failed_refs.append(err) logger.warning( diff --git a/weaviate/collections/batch/sync.py b/weaviate/collections/batch/sync.py index 6627f8911..cff631e0e 100644 --- a/weaviate/collections/batch/sync.py +++ b/weaviate/collections/batch/sync.py @@ -365,6 +365,7 @@ def __recv(self) -> None: result_objs += BatchObjectReturn( _all_responses=[err], errors={cached.index: err}, + has_errors=True, ) failed_objs.append(err) logger.warning( @@ -387,6 +388,7 @@ def __recv(self) -> None: failed_refs.append(err) result_refs += BatchReferenceReturn( errors={cached.index: err}, + has_errors=True, ) logger.warning( {