forked from BasedHardware/omi
-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathcandidate_service.py
More file actions
267 lines (238 loc) · 9.24 KB
/
Copy pathcandidate_service.py
File metadata and controls
267 lines (238 loc) · 9.24 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
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
"""Candidate lifecycle orchestration and post-commit integration policy."""
import asyncio
from typing import Optional, Protocol
import database.action_items as action_items_db
import database.candidate_integration_outbox as integration_outbox_db
import database.candidates as candidates_db
import database.workstreams as workstreams_db
from database.durable_queue import OutcomeKind, ProcessOutcome, drain_isolated
from models.candidate import (
CandidateAction,
CandidateCreate,
CandidateRecord,
CandidateResolutionReceipt,
CandidateStatus,
CandidateSubjectKind,
)
from utils.executors import postprocess_executor, submit_with_context
from utils.observability.fallback import record_fallback
from utils.task_sync import auto_sync_action_item
from utils.task_intelligence import task_links
from utils.task_intelligence.workstream_index import refresh_workstream_association_index
class WorkstreamCandidateResolver(Protocol):
def __call__(self, uid: str, candidate: CandidateRecord, account_generation: int) -> CandidateResolutionReceipt: ...
_workstream_resolver: Optional[WorkstreamCandidateResolver] = workstreams_db.resolve_workstream_candidate
def register_workstream_candidate_resolver(resolver: WorkstreamCandidateResolver) -> None:
global _workstream_resolver
_workstream_resolver = resolver
def clear_workstream_candidate_resolver() -> None:
global _workstream_resolver
_workstream_resolver = None
def create_candidate(
uid: str,
proposal: CandidateCreate,
*,
idempotency_key: str,
account_generation: int,
) -> CandidateRecord:
return candidates_db.create_candidate(
uid,
proposal,
idempotency_key=idempotency_key,
account_generation=account_generation,
)
def _dispatch_task_integration(uid: str, candidate_id: str, task_id: str, *, account_generation: int) -> bool:
lease_token = integration_outbox_db.claim_candidate_integration_dispatch(
uid,
candidate_id,
account_generation=account_generation,
)
if lease_token is None:
return False
task = action_items_db.get_action_item(uid, task_id)
if not task:
integration_outbox_db.complete_candidate_integration_dispatch(
uid,
candidate_id,
account_generation=account_generation,
lease_token=lease_token,
succeeded=False,
error_text='task_not_found',
)
record_fallback(
component='other',
from_mode='candidate_integration',
to_mode='retry_queue',
reason='other',
outcome='degraded',
)
return False
def run_sync() -> None:
try:
result = asyncio.run(auto_sync_action_item(uid, task, skip_apple_reminders=False))
except Exception:
integration_outbox_db.complete_candidate_integration_dispatch(
uid,
candidate_id,
account_generation=account_generation,
lease_token=lease_token,
succeeded=False,
error_text='integration_exception',
)
record_fallback(
component='other',
from_mode='candidate_integration',
to_mode='retry_queue',
reason='other',
outcome='degraded',
)
raise
terminal_noop = result.get('reason') in {
'no_default_integration',
'integration_not_found',
'integration_not_connected',
'client_handles_sync',
}
succeeded = bool(result.get('synced')) or terminal_noop
integration_outbox_db.complete_candidate_integration_dispatch(
uid,
candidate_id,
account_generation=account_generation,
lease_token=lease_token,
succeeded=succeeded,
error_text=None if succeeded else 'integration_not_synced',
)
if not succeeded:
record_fallback(
component='other',
from_mode='candidate_integration',
to_mode='retry_queue',
reason='other',
outcome='degraded',
)
submit_with_context(postprocess_executor, run_sync)
return True
def drain_candidate_integrations(uid: str, *, account_generation: int, limit: int = 100) -> int:
items = integration_outbox_db.list_candidate_integration_dispatches(
uid,
account_generation=account_generation,
limit=limit,
)
def process_one(item: dict) -> ProcessOutcome:
candidate_id = item.get('candidate_id')
task_id = item.get('task_id')
if not isinstance(candidate_id, str) or not isinstance(task_id, str):
if isinstance(candidate_id, str):
integration_outbox_db.dead_letter_malformed_candidate_integration(
uid,
candidate_id,
account_generation=account_generation,
error_text='malformed integration outbox item',
)
return ProcessOutcome.reject('malformed integration outbox item', reason='malformed')
scheduled = _dispatch_task_integration(
uid,
candidate_id,
task_id,
account_generation=account_generation,
)
if not scheduled:
return ProcessOutcome.retry('not_scheduled', reason='not_scheduled')
return ProcessOutcome.ack()
results = drain_isolated(items, process_one)
return sum(1 for result in results if result.outcome.kind == OutcomeKind.ACK)
def accept_candidate(uid: str, candidate_id: str, *, account_generation: int) -> CandidateResolutionReceipt:
candidate = candidates_db.get_candidate(uid, candidate_id)
if candidate is None:
raise candidates_db.CandidateNotFoundError(candidate_id)
if candidate.subject_kind == CandidateSubjectKind.workstream:
if _workstream_resolver is None:
raise candidates_db.WorkstreamCandidateResolverUnavailableError(
'Ticket 04 workstream resolver is not registered'
)
try:
receipt = _workstream_resolver(uid, candidate, account_generation)
except workstreams_db.WorkstreamNotFoundError as exc:
raise candidates_db.CandidateNotFoundError(candidate_id) from exc
except workstreams_db.WorkstreamGenerationMismatchError as exc:
raise candidates_db.CandidateGenerationMismatchError(candidate_id) from exc
except workstreams_db.WorkstreamConflictError as exc:
raise candidates_db.CandidateConflictError(str(exc)) from exc
if receipt.workstream_id:
refresh_workstream_association_index(uid, receipt.workstream_id)
if receipt.task_id:
_dispatch_task_integration(
uid,
candidate_id,
receipt.task_id,
account_generation=account_generation,
)
return receipt
expected_task_links = None
final_goal_id = candidate.goal_id
final_workstream_id = candidate.workstream_id
if candidate.proposed_action != CandidateAction.create:
task_id = candidate.task_id
if task_id is None:
raise candidates_db.CandidateConflictError('task mutation Candidate is missing task_id')
task = action_items_db.get_action_item(uid, task_id)
if task is None:
raise candidates_db.CandidateNotFoundError(f'task:{task_id}')
expected_task_links = (task.get('goal_id'), task.get('workstream_id'))
if final_goal_id is None:
final_goal_id = expected_task_links[0]
if final_workstream_id is None:
final_workstream_id = expected_task_links[1]
task_links.validate_task_links(uid, goal_id=final_goal_id, workstream_id=final_workstream_id)
receipt = candidates_db.resolve_task_candidate(
uid,
candidate_id,
account_generation=account_generation,
expected_task_links=expected_task_links,
)
if candidate.proposed_action == CandidateAction.create and receipt.task_id:
_dispatch_task_integration(
uid,
candidate_id,
receipt.task_id,
account_generation=account_generation,
)
return receipt
def reject_candidate(
uid: str,
candidate_id: str,
*,
reason: Optional[str],
account_generation: int,
) -> CandidateResolutionReceipt:
return candidates_db.resolve_candidate_without_mutation(
uid,
candidate_id,
status=CandidateStatus.rejected,
reason=reason,
account_generation=account_generation,
)
def expire_candidate(
uid: str,
candidate_id: str,
*,
reason: Optional[str],
account_generation: int,
) -> CandidateResolutionReceipt:
return candidates_db.resolve_candidate_without_mutation(
uid,
candidate_id,
status=CandidateStatus.expired,
reason=reason,
account_generation=account_generation,
)
__all__ = [
'WorkstreamCandidateResolver',
'accept_candidate',
'clear_workstream_candidate_resolver',
'create_candidate',
'drain_candidate_integrations',
'expire_candidate',
'register_workstream_candidate_resolver',
'reject_candidate',
]