Skip to content

fix(batch): filter re-added errors by uuid and avoid the deprecated all_responses read - #2180

Open
inchang-ing wants to merge 2 commits into
weaviate:mainfrom
inchang-ing:fix-2179-batch-readd-stale-errors
Open

inchang-ing wants to merge 2 commits into
weaviate:mainfrom
inchang-ing:fix-2179-batch-readd-stale-errors

Conversation

@inchang-ing

Copy link
Copy Markdown

Summary

Fixes #2179.

The re-add path in _BatchBase.__send_batch had two defects:

  1. It read the deprecated BatchObjectReturn.all_responses property, so a spurious Dep020 warning fired inside internal code during every rate-limited batch.
  2. It filtered _all_responses by list position (enumerate) while readded_objects contains batch-wide indices — the same key space as the errors and uuids dictionaries. From the second chunk on the two differ, so the re-added object's ErrorObject stayed in _all_responses even after its retry succeeded (and a valid entry could be dropped instead).

Fix

Filter _all_responses entries by the errored objects' uuids (object_.uuid ∈ readded_uuids) instead of by list position, and read the private _all_responses field directly so internal code no longer triggers the deprecated property's warning.

Validation

  • Regression test mock_tests/test_readd_no_stale_errors.py (mock gRPC server, mirrors the issue's probe: 4 objects, fixed_size(batch_size=2), first object of the second chunk rate-limited once):
    • asserts no Dep020 warning during the batch — fails on main (1 warning fires);
    • asserts no stale ErrorObject remains and _all_responses holds exactly the 4 final entries — fails on main (5 entries, 1 stale);
  • Full mock_tests/ suite: 65 passed; the 2 test_auth.py failures are pre-existing in a minimal environment (verified with the fix stashed).

DCO

Signed-off-by: inchang-ing 197932532+inchang-ing@users.noreply.github.com

@orca-security-eu orca-security-eu Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Orca Security Scan Summary

Status Check Issues by priority
Passed Passed Infrastructure as Code high 0   medium 0   low 0   info 0 View in Orca
Passed Passed SAST high 0   medium 0   low 0   info 0 View in Orca
Passed Passed Secrets high 0   medium 0   low 0   info 0 View in Orca
Passed Passed Vulnerabilities high 0   medium 0   low 0   info 0 View in Orca

@weaviate-git-bot

Copy link
Copy Markdown

To avoid any confusion in the future about your contribution to Weaviate, we work with a Contributor License Agreement. If you agree, you can simply add a comment to this PR that you agree with the CLA so that we can merge.

beep boop - the Weaviate bot 👋🤖

PS:
Are you already a member of the Weaviate Forum?

@chrikrah chrikrah left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

@inchang-ing I would hold this short of merge, on the formatter gate and on the choice of filter key. Credit first: I filed #2179, and mock_tests/test_readd_no_stale_errors.py is my probe from it with asserts bolted on. The 111-character line 14 that the formatter rejects below came in with my paste. The diagnosis and the fix are yours, and both reproduce.

$ cd /tmp/wv2180 && PYTHONPATH=$PWD python -m pytest mock_tests/test_readd_no_stale_errors.py -q
1 passed, 1 warning in 2.95s                   # head 4118492, Python 3.12.3

$ cd /tmp/wv2180-base && PYTHONPATH=$PWD python -m pytest mock_tests/test_readd_no_stale_errors.py -q
E           AssertionError: deprecated all_responses read during batch: ['Dep020: The `all_responses` attribute in the `BatchResults` ']
stale      = [ErrorObject(message='failed with status: 503 error', ... index=2, retry_count=1)]
1 failed, 1 warning in 2.99s                   # base 38f30c1, your test copied in

blocking: the lint job fails on this branch and nothing on the thread shows it, because lint-and-format in .github/workflows/main.yaml has never run on this head. The check rollup carries only Orca and the unicode scan while a maintainer's approval is pending.

# ruff 0.14.7 and flake8 7.4.1 with the four plugins, all pinned in requirements-devel.txt
$ cd /tmp/wv2180 && ruff format --diff weaviate test mock_tests integration packages/web; echo "exit=$?"
2 files would be reformatted, 475 files already formatted   # base.py:725, the test at :14 and :35
exit=1
$ flake8 weaviate test mock_tests integration packages/web; echo "exit=$?"
mock_tests/test_readd_no_stale_errors.py:19:1: D205 1 blank line required between summary line and description
mock_tests/test_readd_no_stale_errors.py:19:1: D209 Multi-line docstring closing quotes should be on a separate line
mock_tests/test_readd_no_stale_errors.py:19:1: D415 First line should end with a period, question mark, or exclamation point
exit=1
# the same two commands on base 38f30c1: exit=0, and "476 files already formatted"
# ruff check weaviate test mock_tests integration packages/web: exit=0 at head and at base
# cd weaviate && pyright (1.1.399, as the type-checking job pins it): 224 errors at both, no delta
# not run: integration/, which wants a live Weaviate.

blocking: the uuid key reopens the same divergence for a batch that holds one uuid twice. An ingest loop that re-sends a key does exactly that. Six objects through fixed_size(batch_size=3), objects 3 and 4 sharing an explicit uuid, the second chunk answering 503 on the first of them and invalid property foo on the second.

$ cd /tmp/wv2180-dup && PYTHONPATH=$PWD python -m pytest mock_tests/test_dup_uuid.py -q -s
imported /tmp/wv2180-dup/weaviate/collections/batch/base.py
errors=[4] len(_all_responses)=5 errors_kept=[]
1 passed, 2 warnings in 3.10s                  # head 4118492 as written

$ cd /tmp/wv2180-idx && PYTHONPATH=$PWD python -m pytest mock_tests/test_dup_uuid.py mock_tests/test_readd_no_stale_errors.py -q -s
imported /tmp/wv2180-idx/weaviate/collections/batch/base.py
errors=[4] len(_all_responses)=6 errors_kept=['invalid property foo']
2 passed, 2 warnings in 5.14s                  # head, predicate swapped to r.object_.index in readded_objects

The permanent failure stays in errors and disappears from _all_responses, because readded_uuids matches both objects. r.object_.index in readded_objects is one token away, BatchObject.index is non-optional at classes/batch.py:51, and that key is the one #2179 asked for.

non-blocking: worth quoting in the description, classes/batch.py:204 says the errors keys "will always be equivalent to the original_index".

@g-despot you merged #2132, the rate-limited batching work in this same file. Is the uuid key deliberate over the index one, or would you take the swap? I will re-run both gates against whichever you pick.

mock_tests/test_dup_uuid.py
import uuid as uuidlib
import weaviate
from weaviate.collections.classes.batch import ErrorObject
from weaviate.proto.v1 import batch_pb2, weaviate_pb2_grpc
from .conftest import MOCK_IP, MOCK_PORT, MOCK_PORT_GRPC, mock_class

SHARED = uuidlib.UUID("11111111-1111-4111-8111-111111111111")


class Svc(weaviate_pb2_grpc.WeaviateServicer):
    calls = 0

    def BatchObjects(self, request, context):
        self.calls += 1
        if self.calls == 2:  # second chunk: index 0 retryable, index 1 permanent
            return batch_pb2.BatchObjectsReply(errors=[
                {"index": 0, "error": "failed with status: 503 error"},
                {"index": 1, "error": "invalid property foo"},
            ])
        return batch_pb2.BatchObjectsReply()


def test_dup_uuid(weaviate_mock, start_grpc_server):
    import weaviate.collections.batch.base as mod
    print("\nimported", mod.__file__)
    weaviate_mock.expect_request(f"/v1/schema/{mock_class['class']}").respond_with_json(mock_class)
    weaviate_pb2_grpc.add_WeaviateServicer_to_server(Svc(), start_grpc_server)
    client = weaviate.connect_to_local(port=MOCK_PORT, host=MOCK_IP, grpc_port=MOCK_PORT_GRPC)
    try:
        col = client.collections.use(mock_class["class"])
        with col.batch.fixed_size(batch_size=3, concurrent_requests=1) as b:
            for i in range(6):
                b.add_object({"name": f"o{i}"}, uuid=SHARED if i in (3, 4) else None)
        r = col.batch.results.objs
        errs = [x for x in r._all_responses if isinstance(x, ErrorObject)]
        print(f"errors={sorted(r.errors)} len(_all_responses)={len(r._all_responses)} "
              f"errors_kept={[e.message for e in errs]}")
    finally:
        client.close()

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Batch: re-added rate-limited objects stay in results.objs.all_responses as errors

3 participants