Skip to content

Commit 6e35dd0

Browse files
authored
fix race condition that could cause events to be dropped (#36)
1 parent 7c3736e commit 6e35dd0

1 file changed

Lines changed: 105 additions & 100 deletions

File tree

src/schematic/event_buffer.py

Lines changed: 105 additions & 100 deletions
Original file line numberDiff line numberDiff line change
@@ -31,8 +31,7 @@ def __init__(
3131
self.max_events = max_events
3232
self.max_retries = max_retries
3333
self.initial_retry_delay = initial_retry_delay
34-
self.flush_lock = threading.Lock()
35-
self.push_lock = threading.Lock()
34+
self.lock = threading.Lock() # Single lock for all buffer operations
3635
self.shutdown = threading.Event()
3736
self.stopped = False
3837

@@ -42,55 +41,53 @@ def __init__(
4241
self.flush_thread.start()
4342

4443
def _flush(self):
45-
with self.flush_lock:
44+
with self.lock:
4645
if not self.events:
4746
return
47+
events_to_process = [event for event in self.events if event is not None]
48+
self.events.clear()
4849

49-
events = [event for event in self.events if event is not None]
50-
51-
# Initialize retry counter and success flag
52-
retry_count = 0
53-
success = False
54-
last_exception = None
55-
56-
# Try with retries and exponential backoff
57-
while retry_count <= self.max_retries and not success:
58-
try:
59-
if retry_count > 0:
60-
# Log retry attempt
61-
self.logger.info(f"Retrying event batch submission (attempt {retry_count} of {self.max_retries})")
62-
63-
# Attempt to send events
64-
self.events_api.create_event_batch(events=events)
65-
success = True
66-
67-
except Exception as e:
68-
last_exception = e
69-
retry_count += 1
70-
71-
if retry_count <= self.max_retries:
72-
# Calculate backoff with jitter
73-
delay = self.initial_retry_delay * (2 ** (retry_count - 1))
74-
jitter = random.uniform(0, 0.1 * delay) # 10% jitter
75-
wait_time = delay + jitter
76-
77-
self.logger.warning(
78-
f"Event batch submission failed: {e}. "
79-
f"Retrying in {wait_time:.2f} seconds..."
80-
)
81-
82-
# Wait before retry
83-
time.sleep(wait_time)
84-
85-
# After all retries, if still not successful, log the error
86-
if not success:
87-
self.logger.error(
88-
f"Event batch submission failed after {self.max_retries} retries: {last_exception}"
89-
)
90-
elif retry_count > 0:
91-
self.logger.info(f"Event batch submission succeeded after {retry_count} retries")
50+
if events_to_process:
51+
self._process_events(events_to_process)
9252

93-
self.events.clear()
53+
def _process_events(self, events_to_process):
54+
"""Process events with retry logic - called without holding lock"""
55+
retry_count = 0
56+
success = False
57+
last_exception = None
58+
59+
while retry_count <= self.max_retries and not success:
60+
try:
61+
if retry_count > 0:
62+
self.logger.info(f"Retrying event batch submission (attempt {retry_count} of {self.max_retries})")
63+
64+
self.events_api.create_event_batch(events=events_to_process)
65+
success = True
66+
67+
except Exception as e:
68+
last_exception = e
69+
retry_count += 1
70+
71+
if retry_count <= self.max_retries:
72+
# Calculate backoff with jitter
73+
delay = self.initial_retry_delay * (2 ** (retry_count - 1))
74+
jitter = random.uniform(0, 0.1 * delay) # 10% jitter
75+
wait_time = delay + jitter
76+
77+
self.logger.warning(
78+
f"Event batch submission failed: {e}. "
79+
f"Retrying in {wait_time:.2f} seconds..."
80+
)
81+
82+
# Wait before retry
83+
time.sleep(wait_time)
84+
85+
if not success:
86+
self.logger.error(
87+
f"Event batch submission failed after {self.max_retries} retries: {last_exception}"
88+
)
89+
elif retry_count > 0:
90+
self.logger.info(f"Event batch submission succeeded after {retry_count} retries")
9491

9592
def _periodic_flush(self):
9693
while not self.shutdown.is_set():
@@ -102,11 +99,17 @@ def push(self, event: CreateEventRequestBody):
10299
self.logger.error("Event buffer is stopped, not accepting new events")
103100
return
104101

105-
with self.push_lock:
102+
should_flush = False
103+
with self.lock:
106104
if len(self.events) >= self.max_events:
107-
self._flush()
105+
should_flush = True
106+
else:
107+
self.events.append(event)
108108

109-
self.events.append(event)
109+
if should_flush:
110+
self._flush()
111+
with self.lock:
112+
self.events.append(event)
110113

111114
def stop(self):
112115
try:
@@ -136,62 +139,58 @@ def __init__(
136139
self.initial_retry_delay = initial_retry_delay
137140
self.shutdown_event = asyncio.Event()
138141
self.stopped = False
139-
self.flush_lock = asyncio.Lock()
140-
self.push_lock = asyncio.Lock()
142+
self.lock = asyncio.Lock() # Single lock for all buffer operations
141143

142144
# Start periodic flushing task
143145
self.flush_task = asyncio.create_task(self._periodic_flush())
144146

145147
async def _flush(self):
146-
async with self.flush_lock:
148+
async with self.lock:
147149
if not self.events:
148150
return
151+
events_to_process = [event for event in self.events if event is not None]
152+
self.events.clear()
149153

150-
events = [event for event in self.events if event is not None]
151-
152-
# Initialize retry counter and success flag
153-
retry_count = 0
154-
success = False
155-
last_exception = None
156-
157-
# Try with retries and exponential backoff
158-
while retry_count <= self.max_retries and not success:
159-
try:
160-
if retry_count > 0:
161-
# Log retry attempt
162-
self.logger.info(f"Retrying event batch submission (attempt {retry_count} of {self.max_retries})")
163-
164-
# Attempt to send events
165-
await self.events_api.create_event_batch(events=events)
166-
success = True
167-
168-
except Exception as e:
169-
last_exception = e
170-
retry_count += 1
171-
172-
if retry_count <= self.max_retries:
173-
# Calculate backoff with jitter
174-
delay = self.initial_retry_delay * (2 ** (retry_count - 1))
175-
jitter = random.uniform(0, 0.1 * delay) # 10% jitter
176-
wait_time = delay + jitter
177-
178-
self.logger.warning(
179-
f"Event batch submission failed: {e}. "
180-
f"Retrying in {wait_time:.2f} seconds..."
181-
)
182-
183-
# Wait before retry (asyncio sleep)
184-
await asyncio.sleep(wait_time)
185-
186-
# After all retries, if still not successful, log the error
187-
if not success:
188-
self.logger.error(
189-
f"Event batch submission failed after {self.max_retries} retries: {last_exception}"
190-
)
191-
elif retry_count > 0:
192-
self.logger.info(f"Event batch submission succeeded after {retry_count} retries")
154+
if events_to_process:
155+
await self._process_events_async(events_to_process)
193156

194-
self.events.clear()
157+
async def _process_events_async(self, events_to_process):
158+
"""Process events with retry logic - called without holding lock"""
159+
# Initialize retry counter and success flag
160+
retry_count = 0
161+
success = False
162+
last_exception = None
163+
164+
while retry_count <= self.max_retries and not success:
165+
try:
166+
if retry_count > 0:
167+
self.logger.info(f"Retrying event batch submission (attempt {retry_count} of {self.max_retries})")
168+
169+
await self.events_api.create_event_batch(events=events_to_process)
170+
success = True
171+
172+
except Exception as e:
173+
last_exception = e
174+
retry_count += 1
175+
176+
if retry_count <= self.max_retries:
177+
delay = self.initial_retry_delay * (2 ** (retry_count - 1))
178+
jitter = random.uniform(0, 0.1 * delay) # 10% jitter
179+
wait_time = delay + jitter
180+
181+
self.logger.warning(
182+
f"Event batch submission failed: {e}. "
183+
f"Retrying in {wait_time:.2f} seconds..."
184+
)
185+
186+
await asyncio.sleep(wait_time)
187+
188+
if not success:
189+
self.logger.error(
190+
f"Event batch submission failed after {self.max_retries} retries: {last_exception}"
191+
)
192+
elif retry_count > 0:
193+
self.logger.info(f"Event batch submission succeeded after {retry_count} retries")
195194

196195
async def _periodic_flush(self):
197196
while not self.shutdown_event.is_set():
@@ -208,11 +207,17 @@ async def push(self, event: CreateEventRequestBody):
208207
self.logger.error("Event buffer is stopped, not accepting new events")
209208
return
210209

211-
async with self.push_lock:
210+
should_flush = False
211+
async with self.lock:
212212
if len(self.events) >= self.max_events:
213-
await self._flush()
213+
should_flush = True
214+
else:
215+
self.events.append(event)
214216

215-
self.events.append(event)
217+
if should_flush:
218+
await self._flush()
219+
async with self.lock:
220+
self.events.append(event)
216221

217222
async def stop(self):
218223
try:

0 commit comments

Comments
 (0)