kafka_prefix = 'yclriuzlwn', kafka_server = '127.0.0.1:49293'
consumer = <cimpl.Consumer object at 0x7f8972852598>
def test_storage_direct_writer(kafka_prefix: str, kafka_server, consumer: Consumer):
writer_config = {
"cls": "kafka",
"brokers": [kafka_server],
"client_id": "kafka_writer",
"prefix": kafka_prefix,
"anonymize": False,
}
storage_config = {
"cls": "pipeline",
"steps": [{"cls": "memory", "journal_writer": writer_config},],
}
storage = get_storage(**storage_config)
expected_messages = 0
for obj_type, objs in TEST_OBJECTS.items():
method = getattr(storage, obj_type + "_add")
if obj_type in (
"content",
"skipped_content",
"directory",
"revision",
"release",
"snapshot",
"origin",
):
method(objs)
expected_messages += len(objs)
elif obj_type in ("origin_visit",):
for obj in objs:
assert isinstance(obj, OriginVisit)
storage.origin_add_one(Origin(url=obj.origin))
visit = method(obj.origin, date=obj.date, type=obj.type)
expected_messages += 1
obj_d = obj.to_dict()
for k in ("visit", "origin", "date", "type"):
del obj_d[k]
storage.origin_visit_update(obj.origin, visit.visit, **obj_d)
expected_messages += 1
else:
assert False, obj_type
existing_topics = {
topic.split(".", 1)[1]
for topic in consumer.list_topics(timeout=10).topics.keys()
}
> assert existing_topics == {
"content",
"directory",
"origin",
"origin_visit",
"release",
"revision",
"snapshot",
"skipped_content",
}
E AssertionError: assert {'content', '...evision', ...} == {'content', '...evision', ...}
E Extra items in the left set:
E 'swh.journal.objects.snapshot'
E Full diff:
E {
E 'content',
E 'directory',
E 'origin',...
E
E ...Full output truncated (8 lines hidden), use '-vv' to show
.tox/py3/lib/python3.7/site-packages/swh/storage/tests/test_kafka_writer.py:71: AssertionError
TEST RESULT
TEST RESULT
- Run At
- May 13 2020, 5:40 PM