forked from BasedHardware/omi
-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathtest_conversation_finalization.py
More file actions
175 lines (134 loc) · 7.6 KB
/
Copy pathtest_conversation_finalization.py
File metadata and controls
175 lines (134 loc) · 7.6 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
import pytest
from unittest import mock
from services import conversation_finalization
from services.conversation_finalization import reconcile_listen_finalization_jobs
from services.conversation_finalization import reconcile_meeting_receipts
@pytest.fixture
def mock_dependencies(monkeypatch):
mocks = {
"is_enabled": mock.Mock(return_value=True),
"publish_metrics": mock.Mock(),
"get_stale_after": mock.Mock(return_value="stale_after"),
"get_candidates": mock.Mock(return_value=[]),
"claim_replay": mock.Mock(),
"enqueue_job": mock.Mock(),
"record_reconciliation": mock.Mock(),
"record_fallback": mock.Mock(),
"inc_retries": mock.Mock(),
}
monkeypatch.setattr(conversation_finalization, "is_listen_finalization_dispatch_enabled", mocks["is_enabled"])
monkeypatch.setattr(conversation_finalization, "_publish_job_metrics", mocks["publish_metrics"])
monkeypatch.setattr(
conversation_finalization.jobs_db, "get_finalization_reconcile_stale_after", mocks["get_stale_after"]
)
monkeypatch.setattr(
conversation_finalization.jobs_db, "get_finalization_replay_candidates", mocks["get_candidates"]
)
monkeypatch.setattr(conversation_finalization.jobs_db, "claim_finalization_replay", mocks["claim_replay"])
monkeypatch.setattr(conversation_finalization, "enqueue_listen_finalization_job", mocks["enqueue_job"])
monkeypatch.setattr(
conversation_finalization, "record_capture_finalization_reconciliation", mocks["record_reconciliation"]
)
monkeypatch.setattr(conversation_finalization, "record_fallback", mocks["record_fallback"])
monkeypatch.setattr(conversation_finalization.LISTEN_FINALIZATION_RETRIES_TOTAL, "inc", mocks["inc_retries"])
return mocks
def test_reconcile_listen_finalization_jobs_disabled(mock_dependencies):
mock_dependencies["is_enabled"].return_value = False
result = reconcile_listen_finalization_jobs()
assert result == {'requeued': 0, 'skipped': 0, 'enqueue_failed': 0}
mock_dependencies["publish_metrics"].assert_called_once()
mock_dependencies["get_candidates"].assert_not_called()
def test_reconcile_listen_finalization_jobs_query_fails(mock_dependencies):
mock_dependencies["get_candidates"].side_effect = Exception("DB error")
result = reconcile_listen_finalization_jobs()
assert result == {'requeued': 0, 'skipped': 0, 'enqueue_failed': 0, 'error': 1}
mock_dependencies["publish_metrics"].assert_called_once()
def test_reconcile_listen_finalization_jobs_skips_invalid_job_id(mock_dependencies):
mock_dependencies["get_candidates"].return_value = [{"job_id": None}, {"job_id": 123}, {}]
result = reconcile_listen_finalization_jobs()
assert result == {'requeued': 0, 'skipped': 3, 'enqueue_failed': 0}
mock_dependencies["claim_replay"].assert_not_called()
mock_dependencies["publish_metrics"].assert_called_once()
def test_reconcile_listen_finalization_jobs_claim_fails(mock_dependencies):
mock_dependencies["get_candidates"].return_value = [{"job_id": "job1"}]
mock_dependencies["claim_replay"].side_effect = Exception("Claim error")
result = reconcile_listen_finalization_jobs()
assert result == {'requeued': 0, 'skipped': 1, 'enqueue_failed': 0}
mock_dependencies["publish_metrics"].assert_called_once()
def test_reconcile_listen_finalization_jobs_claim_not_queued(mock_dependencies):
mock_dependencies["get_candidates"].return_value = [{"job_id": "job1"}, {"job_id": "job2"}]
mock_dependencies["claim_replay"].side_effect = [
{"status": "processing", "dispatch_generation": 1},
{"status": "queued", "dispatch_generation": None},
]
result = reconcile_listen_finalization_jobs()
assert result == {'requeued': 0, 'skipped': 2, 'enqueue_failed': 0}
assert mock_dependencies["claim_replay"].call_count == 2
mock_dependencies["enqueue_job"].assert_not_called()
mock_dependencies["publish_metrics"].assert_called_once()
def test_reconcile_listen_finalization_jobs_enqueue_fails(mock_dependencies):
mock_dependencies["get_candidates"].return_value = [{"job_id": "job1"}]
mock_dependencies["claim_replay"].return_value = {"status": "queued", "dispatch_generation": 1}
mock_dependencies["enqueue_job"].side_effect = Exception("Enqueue error")
result = reconcile_listen_finalization_jobs()
assert result == {'requeued': 0, 'skipped': 0, 'enqueue_failed': 1}
mock_dependencies["record_reconciliation"].assert_called_once_with('enqueue_failed')
mock_dependencies["record_fallback"].assert_called_once()
mock_dependencies["publish_metrics"].assert_called_once()
def test_reconcile_listen_finalization_jobs_success(mock_dependencies):
mock_dependencies["get_candidates"].return_value = [{"job_id": "job1"}]
mock_dependencies["claim_replay"].return_value = {"status": "queued", "dispatch_generation": 1}
result = reconcile_listen_finalization_jobs()
assert result == {'requeued': 1, 'skipped': 0, 'enqueue_failed': 0}
mock_dependencies["claim_replay"].assert_called_once_with("job1", stale_after="stale_after", firestore_client=None)
mock_dependencies["enqueue_job"].assert_called_once_with("job1", 1)
mock_dependencies["record_reconciliation"].assert_called_once_with('requeued')
mock_dependencies["inc_retries"].assert_called_once()
mock_dependencies["publish_metrics"].assert_called_once()
def _stub_meeting_backfill(monkeypatch, candidates=None):
monkeypatch.setattr(
conversation_finalization.jobs_db,
'get_meeting_receipt_backfill_cursor',
lambda **kwargs: {'resume_after_path': None, 'generation': 0},
)
monkeypatch.setattr(
conversation_finalization.jobs_db,
'get_meeting_receipt_backfill_candidates',
lambda **kwargs: {'candidates': candidates or [], 'resume_after_path': None, 'exhausted': True},
)
monkeypatch.setattr(
conversation_finalization.jobs_db,
'advance_meeting_receipt_backfill_cursor',
lambda *args, **kwargs: True,
)
def test_meeting_receipt_reconciler_redrives_one_missing_intent(monkeypatch):
candidate = {'job_id': 'job-1', 'uid': 'uid-1', 'conversation_id': 'conversation-1'}
monkeypatch.setattr(conversation_finalization, 'is_meeting_receipt_reconciler_enabled', lambda: True)
monkeypatch.setattr(
conversation_finalization.jobs_db,
'get_meeting_receipt_reconcile_candidates',
lambda **kwargs: [candidate],
)
repair = mock.Mock(return_value=True)
monkeypatch.setattr(conversation_finalization, 'repair_meeting_receipt_intent', repair)
_stub_meeting_backfill(monkeypatch)
result = reconcile_meeting_receipts()
assert result == {'repaired': 1, 'backfilled': 0, 'skipped': 0, 'error': 0}
repair.assert_called_once_with(candidate)
def test_meeting_receipt_backfill_repairs_two_2026_08_19_shaped_rows(monkeypatch):
candidates = [
{'uid': 'uid-1', 'conversation': {'id': 'meeting-1'}},
{'uid': 'uid-1', 'conversation': {'id': 'meeting-2'}},
]
monkeypatch.setattr(conversation_finalization, 'is_meeting_receipt_reconciler_enabled', lambda: True)
monkeypatch.setattr(
conversation_finalization.jobs_db,
'get_meeting_receipt_reconcile_candidates',
lambda **kwargs: [],
)
_stub_meeting_backfill(monkeypatch, candidates)
record = mock.Mock(return_value={'status': 'recorded'})
monkeypatch.setattr(conversation_finalization, 'record_and_persist_finalized_meeting_receipt', record)
result = reconcile_meeting_receipts()
assert result == {'repaired': 0, 'backfilled': 2, 'skipped': 0, 'error': 0}
assert record.call_count == 2