-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathbuilder.py
More file actions
285 lines (244 loc) · 11.4 KB
/
Copy pathbuilder.py
File metadata and controls
285 lines (244 loc) · 11.4 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
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
"""
Graph Builder — Wires all nodes + conditional edges into a LangGraph.
=========================================================================
Uses LangGraph's native parallel execution for Discovery and Strategy pods.
Flow:
Input → Node 1 → Node 2 ─┬─ Clean ──→ Node 3 → Node 4 ─┬─→ Forensic Detective ─┐
│ └─→ Pattern Matcher ────┤
└─ Score < T → Retry → Node 1 ↓
Diagnosis Merge
(+ Skeptic)
↓
Node 6 (Hypo. Validate)
├─ Verified → Node 7 → Node 8 ─┬─→ Unit Economist ──┐
└─ Weak → loop back ├─→ JTBD Specialist ─┤
└─→ Growth Hacker ──┤
↓
Strategy Merge
↓
Node 10 → Node 11 (Critic)
├ Approved → Node 12 → END
├ Low Lift → loop back (exec)
└ Fail 3+ → loop back (disc)
"""
from __future__ import annotations
import asyncio
import inspect
import sys
import time
import traceback
from langgraph.graph import StateGraph, END
from app.graph.state import RetentionGraphState
from app.graph.conditions import (
route_after_data_audit,
route_after_retry,
route_after_hypothesis_validation,
route_after_strategy_critic,
)
from app.shared import is_cancelled, JobCancelled
from app.graph.nodes import (
input_ingest_node,
data_audit_node,
feature_engineering_node,
behavioral_map_node,
hypothesis_validation_node,
constraint_add_node,
adaptive_hitl_node,
simulation_node,
strategy_critic_node,
execution_architect_node,
retry_handler_node,
# Discovery parallel nodes
forensic_detective_node,
pattern_matcher_node,
competitor_research_node,
diagnosis_merge_node,
# Execution parallel nodes
unit_economist_node,
jtbd_specialist_node,
growth_hacker_node,
strategy_merge_node,
strategy_skeptic_node,
evidence_dossier_node,
)
try:
import resource as _resource
def _rss_mb() -> float:
# Linux: ru_maxrss is in kilobytes. macOS: bytes. Render is Linux.
val = _resource.getrusage(_resource.RUSAGE_SELF).ru_maxrss
return val / 1024.0 if sys.platform.startswith("linux") else val / (1024.0 * 1024.0)
except Exception:
def _rss_mb() -> float:
return -1.0
def _wrap_node(name: str, fn):
"""Log entry, exit duration, RSS, and any exception for a node.
Mirrors the underlying node's arity so LangGraph signature-inspection still
decides correctly whether to pass `config`. A `(state, *args, **kwargs)`
wrapper would look like a 1-arg node to LangGraph and config would be
dropped, breaking nodes that declare `(state, config)`.
"""
is_async = inspect.iscoroutinefunction(fn)
try:
params = list(inspect.signature(fn).parameters.values())
accepts_config = len(params) >= 2
except (TypeError, ValueError):
accepts_config = False
def _check_cancel(state):
job_id = state.get("job_id") if isinstance(state, dict) else None
if is_cancelled(job_id):
raise JobCancelled(job_id)
def _log_enter():
t0 = time.time()
print(f"[NODE→] {name} start | rss={_rss_mb():.0f}MB", flush=True)
return t0
def _log_exit(t0):
print(f"[NODE✓] {name} done in {time.time() - t0:.1f}s | rss={_rss_mb():.0f}MB", flush=True)
def _log_fail(t0, e):
print(f"[NODE✗] {name} FAILED after {time.time() - t0:.1f}s: {type(e).__name__}: {e}", flush=True, file=sys.stderr)
traceback.print_exc()
if is_async and accepts_config:
async def _w(state, config):
_check_cancel(state)
t0 = _log_enter()
try:
result = await fn(state, config)
_log_exit(t0); return result
except Exception as e:
_log_fail(t0, e); raise
return _w
if is_async and not accepts_config:
async def _w(state):
_check_cancel(state)
t0 = _log_enter()
try:
result = await fn(state)
_log_exit(t0); return result
except Exception as e:
_log_fail(t0, e); raise
return _w
if accepts_config:
def _w(state, config):
_check_cancel(state)
t0 = _log_enter()
try:
result = fn(state, config)
_log_exit(t0); return result
except Exception as e:
_log_fail(t0, e); raise
return _w
def _w(state):
_check_cancel(state)
t0 = _log_enter()
try:
result = fn(state)
_log_exit(t0); return result
except Exception as e:
_log_fail(t0, e); raise
return _w
def build_retention_graph() -> StateGraph:
"""
Construct and compile the full retention analysis graph.
Uses LangGraph native fan-out/fan-in for parallel agent execution:
- Discovery Pod: forensic_detective + pattern_matcher run in parallel → diagnosis_merge
- Strategy Pod: unit_economist + jtbd_specialist + growth_hacker run in parallel → strategy_merge
"""
graph = StateGraph(RetentionGraphState)
# ── Register all nodes (wrapped with entry/exit/RSS logging) ─────
def _add(name, fn):
graph.add_node(name, _wrap_node(name, fn))
_add("input_ingest", input_ingest_node)
_add("data_audit", data_audit_node)
_add("retry_handler", retry_handler_node)
_add("feature_engineering", feature_engineering_node)
_add("behavioral_map", behavioral_map_node)
# Discovery Agent nodes (parallel)
_add("forensic_detective", forensic_detective_node)
_add("pattern_matcher", pattern_matcher_node)
_add("competitor_research", competitor_research_node)
_add("diagnosis_merge", diagnosis_merge_node)
_add("hypothesis_validation", hypothesis_validation_node)
_add("constraint_add", constraint_add_node)
_add("adaptive_hitl", adaptive_hitl_node)
# Execution Agent nodes (parallel)
_add("unit_economist", unit_economist_node)
_add("jtbd_specialist", jtbd_specialist_node)
_add("growth_hacker", growth_hacker_node)
_add("strategy_merge", strategy_merge_node)
_add("strategy_skeptic", strategy_skeptic_node)
_add("simulation", simulation_node)
_add("strategy_critic", strategy_critic_node)
_add("evidence_dossier", evidence_dossier_node)
_add("execution_architect", execution_architect_node)
# ── Entry point ──────────────────────────────────────────────────
graph.set_entry_point("input_ingest")
# ── Linear edges ─────────────────────────────────────────────────
graph.add_edge("input_ingest", "data_audit")
# Node 2 → conditional: clean vs retry
graph.add_conditional_edges(
"data_audit",
route_after_data_audit,
{
"feature_engineering": "feature_engineering",
"retry_handler": "retry_handler",
},
)
# Retry loops back or ends
graph.add_conditional_edges(
"retry_handler",
route_after_retry,
{
"input_ingest": "input_ingest",
"feature_engineering": "feature_engineering",
},
)
# Node 3 → Node 4
graph.add_edge("feature_engineering", "behavioral_map")
# ── Discovery Pod: Fan-out (parallel) ────────────────────────────
# behavioral_map fans out to forensic_detective, pattern_matcher, competitor_research
graph.add_edge("behavioral_map", "forensic_detective")
graph.add_edge("behavioral_map", "pattern_matcher")
graph.add_edge("behavioral_map", "competitor_research")
# All three fan-in to diagnosis_merge
graph.add_edge("forensic_detective", "diagnosis_merge")
graph.add_edge("pattern_matcher", "diagnosis_merge")
graph.add_edge("competitor_research", "diagnosis_merge")
# Diagnosis merge → hypothesis validation
graph.add_edge("diagnosis_merge", "hypothesis_validation")
# Node 6 → conditional: verified vs weak proof
graph.add_conditional_edges(
"hypothesis_validation",
route_after_hypothesis_validation,
{
"constraint_add": "constraint_add",
"behavioral_map": "behavioral_map",
},
)
# Node 7 → Node 8
graph.add_edge("constraint_add", "adaptive_hitl")
# ── Strategy Pod: Fan-out (parallel) ─────────────────────────────
# adaptive_hitl fans out to all three execution agents
graph.add_edge("adaptive_hitl", "unit_economist")
graph.add_edge("adaptive_hitl", "jtbd_specialist")
graph.add_edge("adaptive_hitl", "growth_hacker")
# All three fan-in to strategy_merge
graph.add_edge("unit_economist", "strategy_merge")
graph.add_edge("jtbd_specialist", "strategy_merge")
graph.add_edge("growth_hacker", "strategy_merge")
# Strategy merge → skeptic → simulation → critic
graph.add_edge("strategy_merge", "strategy_skeptic")
graph.add_edge("strategy_skeptic", "simulation")
graph.add_edge("simulation", "strategy_critic")
# Node 11 → conditional: approved/exhausted → evidence_dossier; else loop back
graph.add_conditional_edges(
"strategy_critic",
route_after_strategy_critic,
{
"evidence_dossier": "evidence_dossier",
"adaptive_hitl": "adaptive_hitl",
},
)
# Node 11b → Node 12 → END
graph.add_edge("evidence_dossier", "execution_architect")
graph.add_edge("execution_architect", END)
# ── Compile ──────────────────────────────────────────────────────
return graph.compile()