MCPcopy Create free account
hub / github.com/ExtendDB/extenddb / writer_thread

Function writer_thread

samples/stream_consumer.py:47–91  ·  view source on GitHub ↗
()

Source from the content-addressed store, hash-verified

45 time.sleep(0.5)
46 raise TimeoutError(f"Table {table_name} did not become ACTIVE within {timeout}s")
47def writer_thread() -> None:
48 client = make_client()
49
50 # Phase 1: Insert items
51 for i in range(1, ITEM_COUNT + 1):
52 client.put_item(
53 TableName=TABLE_NAME,
54 Item={
55 "pk": {"S": f"item-{i}"},
56 "version": {"N": "1"},
57 "data": {"S": f"initial-value-{i}"},
58 },
59 )
60 log("WRITER", f"INSERT pk=item-{i} version=1 data=initial-value-{i}")
61 time.sleep(0.3)
62
63 time.sleep(1)
64
65 # Phase 2: Update items
66 for i in range(1, ITEM_COUNT + 1):
67 client.update_item(
68 TableName=TABLE_NAME,
69 Key={"pk": {"S": f"item-{i}"}},
70 UpdateExpression="SET version = :v, #d = :d",
71 ExpressionAttributeNames={"#d": "data"},
72 ExpressionAttributeValues={
73 ":v": {"N": "2"},
74 ":d": {"S": f"updated-value-{i}"},
75 },
76 )
77 log("WRITER", f"UPDATE pk=item-{i} version=2 data=updated-value-{i}")
78 time.sleep(0.3)
79
80 time.sleep(1)
81
82 # Phase 3: Delete items
83 for i in range(1, ITEM_COUNT + 1):
84 client.delete_item(
85 TableName=TABLE_NAME,
86 Key={"pk": {"S": f"item-{i}"}},
87 )
88 log("WRITER", f"DELETE pk=item-{i}")
89 time.sleep(0.3)
90
91 log("WRITER", "All writes complete. Waiting for poller to catch up...")
92def format_image(image: dict | None) -> str:
93 if not image:
94 return "{}"

Callers

nothing calls this directly

Calls 5

logFunction · 0.85
put_itemMethod · 0.80
update_itemMethod · 0.80
delete_itemMethod · 0.80
make_clientFunction · 0.70

Tested by

no test coverage detected