-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathflow.py
More file actions
225 lines (193 loc) · 7.09 KB
/
Copy pathflow.py
File metadata and controls
225 lines (193 loc) · 7.09 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
"""
舆情分析智能体 - Flow 编排定义
系统采用线性主链模式,通过 start_stage 选择入口。
"""
from pocketflow import AsyncFlow
from nodes import (
# 管线状态节点
TerminalNode,
Stage1CompletionNode,
Stage2CompletionNode,
Stage3CompletionNode,
# Stage1 节点
DataLoadNode,
SaveEnhancedDataNode,
DataValidationAndOverviewNode,
NLPEnrichmentNode,
AsyncSentimentPolarityAnalysisBatchNode,
AsyncSentimentAttributeAnalysisBatchNode,
AsyncTwoLevelTopicAnalysisBatchNode,
AsyncPublisherObjectAnalysisBatchNode,
AsyncBeliefSystemAnalysisBatchNode,
AsyncPublisherDecisionAnalysisBatchNode,
# Stage2 节点
LoadEnhancedDataNode,
DataSummaryNode,
ClearReportDirNode,
SaveAnalysisResultsNode,
ChartAnalysisNode,
LLMInsightNode,
create_query_search_flow,
create_parallel_agent_flow,
ForumHostNode,
SupplementDataNode,
SupplementSearchNode,
VisualAnalysisNode,
MergeResultsNode,
# Stage3 节点
ClearStage3OutputsNode,
LoadAnalysisResultsNode,
PlanOutlineNode,
GenerateChaptersBatchNode,
ReviewChaptersNode,
IRRendererNode,
InjectTraceNode,
MethodologyAppendixNode,
FormatReportNode,
RenderHTMLNode,
SaveReportNode,
)
def _create_async_enhancement_flow(
concurrent_num: int,
max_retries: int,
wait_time: int,
) -> AsyncFlow:
"""Create Stage1 async enhancement flow."""
data_load_node = DataLoadNode()
nlp_enrichment_node = NLPEnrichmentNode()
sentiment_polarity_node = AsyncSentimentPolarityAnalysisBatchNode(
max_retries=max_retries,
wait=wait_time,
max_concurrent=concurrent_num,
)
sentiment_attribute_node = AsyncSentimentAttributeAnalysisBatchNode(
max_retries=max_retries,
wait=wait_time,
max_concurrent=concurrent_num,
)
topic_analysis_node = AsyncTwoLevelTopicAnalysisBatchNode(
max_retries=max_retries,
wait=wait_time,
max_concurrent=concurrent_num,
)
publisher_analysis_node = AsyncPublisherObjectAnalysisBatchNode(
max_retries=max_retries,
wait=wait_time,
max_concurrent=concurrent_num,
)
belief_analysis_node = AsyncBeliefSystemAnalysisBatchNode(
max_retries=max_retries,
wait=wait_time,
max_concurrent=concurrent_num,
)
publisher_decision_node = AsyncPublisherDecisionAnalysisBatchNode(
max_retries=max_retries,
wait=wait_time,
max_concurrent=concurrent_num,
)
save_data_node = SaveEnhancedDataNode()
validation_node = DataValidationAndOverviewNode()
completion_node = Stage1CompletionNode()
data_load_node >> nlp_enrichment_node
nlp_enrichment_node >> sentiment_polarity_node
sentiment_polarity_node >> sentiment_attribute_node
sentiment_attribute_node >> topic_analysis_node
topic_analysis_node >> publisher_analysis_node
publisher_analysis_node >> belief_analysis_node
belief_analysis_node >> publisher_decision_node
publisher_decision_node >> save_data_node
save_data_node >> validation_node
validation_node >> completion_node
return AsyncFlow(start=data_load_node)
def _create_agent_analysis_flow() -> AsyncFlow:
"""Create Stage2 dual-source analysis flow."""
clear_report_node = ClearReportDirNode()
load_data_node = LoadEnhancedDataNode()
data_summary_node = DataSummaryNode()
query_search_flow = create_query_search_flow()
parallel_agent_flow = create_parallel_agent_flow()
forum_host_node = ForumHostNode()
supplement_data_node = SupplementDataNode()
supplement_search_node = SupplementSearchNode()
visual_analysis_node = VisualAnalysisNode()
merge_node = MergeResultsNode()
chart_analysis_node = ChartAnalysisNode(max_retries=2, wait=3)
llm_insight_node = LLMInsightNode()
save_results_node = SaveAnalysisResultsNode()
completion_node = Stage2CompletionNode()
clear_report_node >> load_data_node
load_data_node >> data_summary_node
data_summary_node >> query_search_flow
query_search_flow >> parallel_agent_flow
parallel_agent_flow >> forum_host_node
forum_host_node - "supplement_data" >> supplement_data_node
forum_host_node - "supplement_search" >> supplement_search_node
forum_host_node - "supplement_visual" >> visual_analysis_node
forum_host_node - "sufficient" >> merge_node
supplement_data_node >> forum_host_node
supplement_search_node >> forum_host_node
visual_analysis_node >> forum_host_node
merge_node >> chart_analysis_node
chart_analysis_node >> llm_insight_node
llm_insight_node >> save_results_node
save_results_node >> completion_node
return AsyncFlow(start=clear_report_node)
def _create_unified_report_flow() -> AsyncFlow:
"""Create unified Stage3 report flow."""
clear_outputs_node = ClearStage3OutputsNode()
load_results_node = LoadAnalysisResultsNode()
outline_node = PlanOutlineNode()
generate_chapters_node = GenerateChaptersBatchNode(max_concurrent=3)
review_chapters_node = ReviewChaptersNode()
ir_renderer_node = IRRendererNode()
inject_trace_node = InjectTraceNode()
methodology_node = MethodologyAppendixNode()
format_node = FormatReportNode()
render_html_node = RenderHTMLNode()
save_node = SaveReportNode()
completion_node = Stage3CompletionNode()
clear_outputs_node >> load_results_node
load_results_node >> outline_node
outline_node >> generate_chapters_node
generate_chapters_node >> review_chapters_node
review_chapters_node - "needs_revision" >> generate_chapters_node
review_chapters_node - "satisfied" >> ir_renderer_node
ir_renderer_node >> inject_trace_node
inject_trace_node >> methodology_node
methodology_node >> format_node
format_node >> render_html_node
render_html_node >> save_node
save_node >> completion_node
return AsyncFlow(start=clear_outputs_node)
def create_main_flow(
start_stage: int = 1,
concurrent_num: int = 60,
max_retries: int = 3,
wait_time: int = 8,
) -> AsyncFlow:
"""Create linear Stage1->Stage2->Stage3 flow with selectable start stage."""
if start_stage not in {1, 2, 3}:
raise ValueError(f"start_stage must be one of [1,2,3], got {start_stage}")
terminal = TerminalNode()
async_enhancement_flow = _create_async_enhancement_flow(
concurrent_num=concurrent_num,
max_retries=max_retries,
wait_time=wait_time,
)
agent_analysis_flow = _create_agent_analysis_flow()
unified_report_flow = _create_unified_report_flow()
async_enhancement_flow >> agent_analysis_flow
agent_analysis_flow >> unified_report_flow
unified_report_flow >> terminal
entry_by_stage = {
1: async_enhancement_flow,
2: agent_analysis_flow,
3: unified_report_flow,
}
return AsyncFlow(start=entry_by_stage[start_stage])
def create_stage2_only_flow() -> AsyncFlow:
"""Create Stage2-only flow."""
return _create_agent_analysis_flow()
def create_stage3_only_flow() -> AsyncFlow:
"""Create Stage3-only unified flow."""
return _create_unified_report_flow()