forked from BasedHardware/omi
-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathdurable_queue.py
More file actions
58 lines (50 loc) · 1.45 KB
/
Copy pathdurable_queue.py
File metadata and controls
58 lines (50 loc) · 1.45 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
"""Firestore-touching durable-queue drain helpers and redrive.
Pure policy (outcomes, decide_attempt, drain_isolated, age helpers) lives in
``utils.durable_queue_policy`` so async callers can import it without the
async-blocker treating ``database.*`` as I/O. This module re-exports those
names for database-layer callers and owns the additive redrive patch.
Store-wide oldest-ready sampling lives in ``database.durable_queue_age``.
"""
from __future__ import annotations
from datetime import datetime
from utils.durable_queue_policy import (
AttemptDecision,
EnqueueDecision,
IsolatedResult,
OutcomeKind,
ProcessOutcome,
QueuePolicy,
adopt_on_identity,
backoff_seconds,
decide_attempt,
drain_isolated,
drain_isolated_async,
oldest_ready_age_seconds,
ready_sort_key,
)
def redrive_patch(*, now: datetime) -> dict:
"""Additive fields to move a dead letter back to ready by identity."""
return {
'status': 'pending',
'attempt_count': 0,
'available_at': now,
'dead_letter_reason': None,
'last_error_text': None,
'updated_at': now,
}
__all__ = [
'AttemptDecision',
'EnqueueDecision',
'IsolatedResult',
'OutcomeKind',
'ProcessOutcome',
'QueuePolicy',
'adopt_on_identity',
'backoff_seconds',
'decide_attempt',
'drain_isolated',
'drain_isolated_async',
'oldest_ready_age_seconds',
'ready_sort_key',
'redrive_patch',
]