forked from BasedHardware/omi
-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathconversation_finalization.py
More file actions
231 lines (212 loc) · 9.88 KB
/
Copy pathconversation_finalization.py
File metadata and controls
231 lines (212 loc) · 9.88 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
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
"""Protected Cloud Tasks worker for durable listen conversation finalization."""
from __future__ import annotations
import asyncio
import logging
from typing import Any
from fastapi import APIRouter, Depends, Request
from fastapi.responses import JSONResponse
from database import conversation_finalization_jobs as jobs_db
from database.sync_jobs import release_job_run_lock, try_acquire_job_run_lock
from services.conversation_finalization import (
final_attempt_failed,
get_listen_finalization_tasks_max_attempts_for_worker,
)
from utils.cloud_tasks import verify_listen_finalization_cloud_tasks_oidc
from utils.account_cutover.access import should_skip_background_account_mutation
from utils.conversations import lifecycle as lifecycle_service
from utils.conversations.finalizer import (
ConversationFinalizationDisposition,
ConversationFinalizationError,
finalize_persisted_conversation,
)
from utils.executors import db_executor, run_blocking
from utils.metrics import LISTEN_FINALIZATION_RETRIES_TOTAL
from utils.observability.journeys import (
record_capture_finalization_terminal,
record_conversation_finalization_client_terminal,
)
logger = logging.getLogger(__name__)
router = APIRouter()
def _parse_task_payload(payload: Any) -> tuple[str, int] | None:
"""Accept exactly the opaque durable task schema, never credential fields."""
if not isinstance(payload, dict) or set(payload) != {'job_id', 'dispatch_generation'}:
return None
job_id = payload.get('job_id')
generation = payload.get('dispatch_generation')
if not isinstance(job_id, str) or not job_id or len(job_id) > 128:
return None
if not isinstance(generation, int) or isinstance(generation, bool) or generation < 1:
return None
return job_id, generation
async def _retry_or_dead_letter(
job_id: str,
dispatch_generation: int,
lease_epoch: int,
task_retry_count: int,
reason: str,
) -> bool:
"""Record a task failure; return whether this was the terminal delivery."""
max_attempts = get_listen_finalization_tasks_max_attempts_for_worker()
if task_retry_count >= max_attempts - 1:
marked_dead_letter = await run_blocking(
db_executor,
final_attempt_failed,
job_id,
dispatch_generation,
lease_epoch,
task_retry_count + 1,
)
if not marked_dead_letter:
return False
return True
await run_blocking(
db_executor,
jobs_db.mark_finalization_retryable,
job_id,
dispatch_generation,
lease_epoch,
reason,
)
LISTEN_FINALIZATION_RETRIES_TOTAL.inc()
return False
@router.post('/v1/conversation-finalization-jobs/run', include_in_schema=False)
async def run_listen_finalization_job(
request: Request,
task_retry_count: int = Depends(verify_listen_finalization_cloud_tasks_oidc),
):
try:
parsed = _parse_task_payload(await request.json())
except Exception:
parsed = None
if parsed is None:
logger.warning('listen finalization handler dropped invalid opaque task payload')
return JSONResponse(status_code=200, content={'status': 'dropped', 'reason': 'invalid_payload'})
job_id, dispatch_generation = parsed
lock_key = f'listen-finalization:{job_id}'
lock_token = await run_blocking(db_executor, try_acquire_job_run_lock, lock_key)
if not lock_token:
return JSONResponse(status_code=409, content={'status': 'locked'})
release_lock = True
claimed_lease_epoch: int | None = None
job: dict[str, Any] | None = None
try:
claim = await run_blocking(
db_executor,
jobs_db.claim_finalization_job,
job_id,
dispatch_generation,
)
claim_status = claim['status']
if claim_status == 'completed':
return JSONResponse(status_code=200, content={'status': 'acked', 'job_status': 'completed'})
if claim_status == 'leased':
return JSONResponse(status_code=409, content={'status': claim_status})
if claim_status == 'stale_generation':
# The reconciler has already enqueued the newer generation. An old
# named task is no longer actionable and must be acknowledged so
# Cloud Tasks does not retry this permanently fenced payload.
logger.info(
'listen finalization stale generation task acknowledged job=%s dispatch_generation=%s',
job_id,
dispatch_generation,
)
return JSONResponse(status_code=200, content={'status': 'dropped', 'reason': claim_status})
if claim_status != 'claimed':
return JSONResponse(status_code=200, content={'status': 'dropped', 'reason': claim_status})
claimed_lease_epoch = claim['lease_epoch']
if claimed_lease_epoch is None:
logger.error('listen finalization claim returned no lease epoch job=%s', job_id)
return JSONResponse(status_code=500, content={'status': 'retry'})
job = await run_blocking(db_executor, jobs_db.get_finalization_job, job_id)
if not job or not isinstance(job.get('uid'), str) or not isinstance(job.get('conversation_id'), str):
terminal = await _retry_or_dead_letter(
job_id, dispatch_generation, claimed_lease_epoch, task_retry_count, 'invalid_job'
)
if terminal:
logger.error('listen finalization final attempt failed job=%s error=invalid_job', job_id)
return JSONResponse(status_code=200, content={'status': 'dead_letter'})
return JSONResponse(status_code=500, content={'status': 'retry'})
if await run_blocking(db_executor, should_skip_background_account_mutation, job['uid']):
# Prequeued finalization must not mutate migrating/new accounts.
completed = await run_blocking(
db_executor,
lifecycle_service.complete_fenced_finalization,
job_id,
dispatch_generation,
claimed_lease_epoch,
)
if not completed:
return JSONResponse(status_code=409, content={'status': 'completion_conflict'})
record_capture_finalization_terminal('stale', job.get('created_at'))
record_conversation_finalization_client_terminal('cancelled', job)
return JSONResponse(status_code=200, content={'status': 'skipped', 'reason': 'account_cutover'})
try:
disposition = await finalize_persisted_conversation(
job['uid'],
job['conversation_id'],
finalization_job_id=job_id,
dispatch_generation=dispatch_generation,
lease_epoch=claimed_lease_epoch,
force_process=bool(job.get('force_process')),
final_attempt=task_retry_count >= get_listen_finalization_tasks_max_attempts_for_worker() - 1,
)
except ConversationFinalizationError:
terminal = await _retry_or_dead_letter(
job_id, dispatch_generation, claimed_lease_epoch, task_retry_count, 'processing_failed'
)
if terminal:
logger.error('listen finalization final attempt failed job=%s failure=processing_failed', job_id)
return JSONResponse(status_code=200, content={'status': 'dead_letter'})
return JSONResponse(status_code=500, content={'status': 'retry'})
if disposition == ConversationFinalizationDisposition.fenced:
completed = await run_blocking(
db_executor,
lifecycle_service.complete_fenced_finalization,
job_id,
dispatch_generation,
claimed_lease_epoch,
)
else:
completed = await run_blocking(
db_executor,
jobs_db.mark_finalization_completed,
job_id,
dispatch_generation,
claimed_lease_epoch,
)
if not completed:
return JSONResponse(status_code=409, content={'status': 'completion_conflict'})
accepted_at = job.get('created_at') if job else None
if disposition == ConversationFinalizationDisposition.fenced:
record_capture_finalization_terminal('stale', accepted_at)
record_conversation_finalization_client_terminal('cancelled', job)
else:
record_capture_finalization_terminal('success', accepted_at)
record_conversation_finalization_client_terminal('success', job)
return JSONResponse(status_code=200, content={'status': 'done'})
except asyncio.CancelledError:
release_lock = False
logger.warning('listen finalization handler cancelled job=%s; preserving run lock until TTL', job_id)
raise
except Exception:
if claimed_lease_epoch is not None:
try:
terminal = await _retry_or_dead_letter(
job_id,
dispatch_generation,
claimed_lease_epoch,
task_retry_count,
'worker_failed',
)
except Exception:
logger.error('listen finalization recovery update failed job=%s failure=worker_failed', job_id)
else:
if terminal:
logger.error('listen finalization final attempt failed job=%s failure=worker_failed', job_id)
return JSONResponse(status_code=200, content={'status': 'dead_letter'})
return JSONResponse(status_code=500, content={'status': 'retry'})
logger.error('listen finalization handler failed job=%s failure=worker_failed', job_id)
return JSONResponse(status_code=500, content={'status': 'retry'})
finally:
if release_lock:
await run_blocking(db_executor, release_job_run_lock, lock_key, lock_token)