-
Notifications
You must be signed in to change notification settings - Fork 137
Expand file tree
/
Copy pathtest_batch.py
More file actions
103 lines (86 loc) · 3.84 KB
/
Copy pathtest_batch.py
File metadata and controls
103 lines (86 loc) · 3.84 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
from typing import Generator
import grpc
import pytest
import weaviate
from weaviate.proto.v1 import batch_pb2, weaviate_pb2_grpc
from .conftest import MOCK_IP, MOCK_PORT, MOCK_PORT_GRPC, mock_class, HTTPServer
HOW_MANY = 1000
class MockCanceledStreamWeaviateService(weaviate_pb2_grpc.WeaviateServicer):
called = False
uuids = set[str]()
def BatchStream(
self,
request_iterator: Generator[batch_pb2.BatchStreamRequest, None, None],
context: grpc.ServicerContext,
) -> Generator[batch_pb2.BatchStreamReply, None, None]:
if not self.called:
self.called = True
context.set_code(grpc.StatusCode.CANCELLED)
context.set_details("context canceled")
return
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]
self.uuids.update(uuids)
yield batch_pb2.BatchStreamReply(acks=batch_pb2.BatchStreamReply.Acks(uuids=uuids))
if request.HasField("stop"):
return
@pytest.fixture(scope="function")
def canceled_stream_client(
weaviate_mock: HTTPServer, start_grpc_server: grpc.Server
) -> Generator[weaviate.WeaviateClient, None, None]:
weaviate_mock.expect_request(f"/v1/schema/{mock_class['class']}").respond_with_json(mock_class)
client = weaviate.connect_to_local(port=MOCK_PORT, host=MOCK_IP, grpc_port=MOCK_PORT_GRPC)
yield client
client.close()
@pytest.fixture(scope="function")
def canceled_stream(
canceled_stream_client: weaviate.WeaviateClient, start_grpc_server: grpc.Server
):
service = MockCanceledStreamWeaviateService()
weaviate_pb2_grpc.add_WeaviateServicer_to_server(service, start_grpc_server)
return canceled_stream_client.collections.use(mock_class["class"]), service
def test_ssb_canceled_stream(
canceled_stream: tuple[weaviate.collections.Collection, MockCanceledStreamWeaviateService],
):
collection, service = canceled_stream
with collection.batch.stream() as batch:
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