-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathchunk_stress.py
More file actions
114 lines (99 loc) · 3.91 KB
/
Copy pathchunk_stress.py
File metadata and controls
114 lines (99 loc) · 3.91 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
104
105
106
107
108
109
110
111
112
113
114
"""Long-session append path — segment rollover under sustained writes."""
from __future__ import annotations
import itertools
from typing import Any
from adk_aerospike import AerospikeSessionService
from ai_ecosystem_benchmark import BaseBenchmarkWorkload
from ._async_bridge import run_async
from ._fixtures import filler_text, make_event, new_session_id
from ._redis_backend import (
RedisSessionService,
close_session,
init_session,
purge_user_sessions,
session_service,
)
class ChunkStressWorkload(BaseBenchmarkWorkload):
"""Appends to one long-running session, exercising segment rollover.
``aerospike_session_append_chunked`` — sequential appends into append-only
K_ORDERED segment records that roll over on overflow (production long
conversations). The method name is retained for benchmark-profile stability.
"""
APP = "bench_eco_chunk"
def __init__(
self,
aerospike_connection_string: str | None = None,
redis_connection_string: str | None = None,
**params: Any,
) -> None:
super().__init__(
aerospike_connection_string=aerospike_connection_string,
redis_connection_string=redis_connection_string,
)
self._event_size_bytes = int(params.get("event_size_bytes", 600))
self._svc: AerospikeSessionService | RedisSessionService | None = None
self._session: Any = None
self._seq = itertools.count()
def setup(self) -> None:
if self.is_aerospike_enabled():
assert self.aerospike_connection_string is not None
self._svc = AerospikeSessionService.from_uri(self.aerospike_connection_string)
run_async(self._reset_session())
elif self.is_redis_enabled():
assert self.redis_connection_string is not None
self._svc = session_service(self.redis_connection_string)
run_async(self._reset_session_redis())
else:
raise RuntimeError("no backend connection string configured")
def between_benchmarks(self) -> None:
return None
def teardown(self) -> None:
if self._svc is not None and self._session is not None:
run_async(
self._svc.delete_session(
app_name=self.APP,
user_id="u0",
session_id=self._session.id,
)
)
if isinstance(self._svc, AerospikeSessionService):
self._svc.close()
else:
run_async(close_session(self._svc))
self._svc = None
self._session = None
def redis_session_append_chunked(self) -> None:
self.aerospike_session_append_chunked()
def aerospike_session_append_chunked(self) -> None:
svc = self._svc
session = self._session
assert svc is not None and session is not None
n = next(self._seq)
text = filler_text(self._event_size_bytes, seed=n)
run_async(svc.append_event(session, make_event(text, n)))
async def _reset_session(self) -> None:
svc = self._svc
assert svc is not None
resp = await svc.list_sessions(app_name=self.APP, user_id="u0")
for s in resp.sessions:
try:
await svc.delete_session(
app_name=self.APP, user_id="u0", session_id=s.id
)
except Exception:
pass
self._session = await svc.create_session(
app_name=self.APP,
user_id="u0",
session_id=new_session_id(),
)
async def _reset_session_redis(self) -> None:
svc = self._svc
assert isinstance(svc, RedisSessionService)
await init_session(svc)
await purge_user_sessions(svc, self.APP, "u0")
self._session = await svc.create_session(
app_name=self.APP,
user_id="u0",
session_id=new_session_id(),
)