forked from BasedHardware/omi
-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathfinalization_decision.py
More file actions
133 lines (103 loc) · 4.52 KB
/
Copy pathfinalization_decision.py
File metadata and controls
133 lines (103 loc) · 4.52 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
"""Pure, deterministic decisions for a conversation finalization generation.
The lifecycle service is the runtime owner. This reducer deliberately has no
Firestore, clock, or queue dependency so its ordering contract can be fuzzed
before the runtime callers are migrated.
"""
from __future__ import annotations
from dataclasses import dataclass, field
from enum import Enum
class LifecyclePhase(str, Enum):
IN_PROGRESS = 'in_progress'
PROCESSING = 'processing'
MERGING = 'merging'
COMPLETED = 'completed'
FAILED = 'failed'
DISCARDED = 'discarded'
class FinalizationEvent(str, Enum):
DISCONNECT = 'disconnect'
FINALIZE = 'finalize'
REPROCESS = 'reprocess'
MERGE = 'merge'
MERGE_COMPLETED = 'merge_completed'
PROCESSING_COMPLETED = 'processing_completed'
DISCARD = 'discard'
RESTART = 'restart'
TERMINAL_PHASES = frozenset({LifecyclePhase.COMPLETED, LifecyclePhase.FAILED, LifecyclePhase.DISCARDED})
@dataclass(frozen=True)
class FinalizationDecisionState:
"""The complete reducer state for one immutable lifecycle generation."""
phase: LifecyclePhase = LifecyclePhase.IN_PROGRESS
terminal_outcome: LifecyclePhase | None = None
emitted_fanout_keys: frozenset[str] = field(default_factory=frozenset)
@property
def is_terminal(self) -> bool:
return self.phase in TERMINAL_PHASES
@dataclass(frozen=True)
class FinalizationDecision:
state: FinalizationDecisionState
fanout_key: str | None = None
reason: str = 'accepted'
def _terminal(state: FinalizationDecisionState, phase: LifecyclePhase) -> FinalizationDecisionState:
if state.terminal_outcome is not None:
return state
return FinalizationDecisionState(
phase=phase,
terminal_outcome=phase,
emitted_fanout_keys=state.emitted_fanout_keys,
)
def _admit_finalization(
state: FinalizationDecisionState,
conversation_id: str,
fanout_key: str | None,
) -> FinalizationDecision:
key = fanout_key or f'conversation:{conversation_id}:finalization'
if key in state.emitted_fanout_keys:
return FinalizationDecision(state=state, reason='duplicate_finalization')
return FinalizationDecision(
state=FinalizationDecisionState(
phase=LifecyclePhase.PROCESSING,
emitted_fanout_keys=state.emitted_fanout_keys | {key},
),
fanout_key=key,
)
def decide_finalization(
state: FinalizationDecisionState,
event: FinalizationEvent,
*,
conversation_id: str,
fanout_key: str | None = None,
) -> FinalizationDecision:
"""Return the one allowed transition and, at most, one durable fanout key.
A ``FinalizationDecisionState`` represents exactly one lifecycle
generation. Reprocessing therefore needs a new generation rather than an
escape hatch out of a terminal state. This is what lets duplicate events
and reclaimed workers share a single idempotency key.
"""
if state.is_terminal:
return FinalizationDecision(state=state, reason='terminal')
if event in {FinalizationEvent.DISCONNECT, FinalizationEvent.FINALIZE}:
if state.phase == LifecyclePhase.IN_PROGRESS:
return _admit_finalization(state, conversation_id, fanout_key)
return FinalizationDecision(state=state, reason='already_finalizing')
if event == FinalizationEvent.PROCESSING_COMPLETED and state.phase == LifecyclePhase.PROCESSING:
return FinalizationDecision(state=_terminal(state, LifecyclePhase.COMPLETED))
if event == FinalizationEvent.DISCARD and state.phase in {
LifecyclePhase.IN_PROGRESS,
LifecyclePhase.PROCESSING,
LifecyclePhase.MERGING,
}:
return FinalizationDecision(state=_terminal(state, LifecyclePhase.DISCARDED))
if event == FinalizationEvent.MERGE and state.phase == LifecyclePhase.IN_PROGRESS:
return FinalizationDecision(
state=FinalizationDecisionState(
phase=LifecyclePhase.MERGING,
emitted_fanout_keys=state.emitted_fanout_keys,
)
)
if event == FinalizationEvent.MERGE_COMPLETED and state.phase == LifecyclePhase.MERGING:
return FinalizationDecision(state=_terminal(state, LifecyclePhase.COMPLETED))
if event == FinalizationEvent.RESTART:
return FinalizationDecision(state=state, reason='restart_replays_state')
if event == FinalizationEvent.REPROCESS:
return FinalizationDecision(state=state, reason='reprocess_requires_new_generation')
return FinalizationDecision(state=state, reason='invalid_transition')