forked from ZhuLinsen/daily_stock_analysis
-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathanalysis.py
More file actions
344 lines (297 loc) · 11 KB
/
Copy pathanalysis.py
File metadata and controls
344 lines (297 loc) · 11 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
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
# -*- coding: utf-8 -*-
"""
===================================
分析相关模型
===================================
职责:
1. 定义分析请求和响应模型
2. 定义任务状态模型
3. 定义异步任务队列相关模型
"""
from typing import Optional, List, Any
from enum import Enum
from pydantic import BaseModel, Field
from src.utils.analysis_metadata import SELECTION_SOURCE_PATTERN
class TaskStatusEnum(str, Enum):
"""任务状态枚举"""
PENDING = "pending"
PROCESSING = "processing"
COMPLETED = "completed"
FAILED = "failed"
class AnalyzeRequest(BaseModel):
"""Analysis request parameters"""
asset_type: str = Field(
"stock",
description="资产类型:stock(股票) / futures(国内期货主力合约)",
pattern="^(stock|futures)$",
example="stock",
)
stock_code: Optional[str] = Field(
None,
description="单只股票代码",
example="600519"
)
stock_codes: Optional[List[str]] = Field(
None,
description="多只股票代码(与 stock_code 二选一)",
example=["600519", "000858"]
)
report_type: str = Field(
"detailed",
description="报告类型:simple(精简) / detailed(完整) / full(完整) / brief(简洁)",
pattern="^(simple|detailed|full|brief)$",
)
force_refresh: bool = Field(
False,
description="是否强制刷新(忽略缓存)"
)
async_mode: bool = Field(
False,
description="是否使用异步模式"
)
stock_name: Optional[str] = Field(
None,
description="用户选中的股票名称(自动补全时提供)",
example="贵州茅台"
)
original_query: Optional[str] = Field(
None,
description="用户原始输入(如茅台、gzmt、600519)",
example="茅台"
)
selection_source: Optional[str] = Field(
None,
description="股票选择来源:manual(手动输入) | autocomplete(自动补全) | import(导入) | image(图片识别)",
pattern=SELECTION_SOURCE_PATTERN,
example="autocomplete"
)
notify: bool = Field(
True,
description="是否发送推送通知(Telegram/企业微信等)"
)
class Config:
json_schema_extra = {
"example": {
"stock_code": "600519",
"report_type": "detailed",
"force_refresh": False,
"async_mode": False,
"stock_name": "贵州茅台",
"original_query": "茅台",
"selection_source": "autocomplete",
"notify": True
}
}
class AnalysisResultResponse(BaseModel):
"""分析结果响应模型"""
query_id: str = Field(..., description="分析记录唯一标识")
stock_code: str = Field(..., description="股票代码")
stock_name: Optional[str] = Field(None, description="股票名称")
report: Optional[Any] = Field(None, description="分析报告")
created_at: str = Field(..., description="创建时间")
class Config:
json_schema_extra = {
"example": {
"query_id": "abc123def456",
"stock_code": "600519",
"stock_name": "贵州茅台",
"report": {
"summary": {
"sentiment_score": 75,
"operation_advice": "持有"
}
},
"created_at": "2024-01-01T12:00:00"
}
}
class TaskAccepted(BaseModel):
"""异步任务接受响应"""
task_id: str = Field(..., description="任务 ID,用于查询状态")
status: str = Field(
...,
description="任务状态",
pattern="^(pending|processing)$"
)
message: Optional[str] = Field(None, description="提示信息")
class Config:
json_schema_extra = {
"example": {
"task_id": "task_abc123",
"status": "pending",
"message": "Analysis task accepted"
}
}
class BatchTaskAcceptedItem(BaseModel):
"""批量异步任务中的单个成功提交项。"""
task_id: str = Field(..., description="任务 ID,用于查询状态")
stock_code: str = Field(..., description="股票代码")
status: str = Field(
...,
description="任务状态",
pattern="^(pending|processing)$"
)
message: Optional[str] = Field(None, description="提示信息")
class Config:
json_schema_extra = {
"example": {
"task_id": "task_abc123",
"stock_code": "600519",
"status": "pending",
"message": "分析任务已加入队列: 600519"
}
}
class BatchDuplicateTaskItem(BaseModel):
"""批量异步任务中的重复提交项。"""
stock_code: str = Field(..., description="股票代码")
existing_task_id: str = Field(..., description="已存在的任务 ID")
message: str = Field(..., description="错误信息")
class Config:
json_schema_extra = {
"example": {
"stock_code": "600519",
"existing_task_id": "task_existing_123",
"message": "股票 600519 正在分析中 (task_id: task_existing_123)"
}
}
class BatchTaskAcceptedResponse(BaseModel):
"""批量异步任务接受响应。"""
accepted: List[BatchTaskAcceptedItem] = Field(default_factory=list, description="成功提交的任务列表")
duplicates: List[BatchDuplicateTaskItem] = Field(default_factory=list, description="重复而跳过的任务列表")
message: str = Field(..., description="汇总信息")
class Config:
json_schema_extra = {
"example": {
"accepted": [
{
"task_id": "task_abc123",
"stock_code": "600519",
"status": "pending",
"message": "分析任务已加入队列: 600519"
}
],
"duplicates": [
{
"stock_code": "000858",
"existing_task_id": "task_existing_456",
"message": "股票 000858 正在分析中 (task_id: task_existing_456)"
}
],
"message": "已提交 1 个任务,1 个重复跳过"
}
}
class TaskStatus(BaseModel):
"""Task status model"""
task_id: str = Field(..., description="任务 ID")
status: str = Field(
...,
description="任务状态",
pattern="^(pending|processing|completed|failed)$"
)
progress: Optional[int] = Field(
None,
description="进度百分比 (0-100)",
ge=0,
le=100
)
result: Optional[AnalysisResultResponse] = Field(
None,
description="分析结果(仅在 completed 时存在)"
)
error: Optional[str] = Field(
None,
description="错误信息(仅在 failed 时存在)"
)
asset_type: str = Field("stock", description="资产类型:stock/futures")
stock_name: Optional[str] = Field(None, description="股票名称")
original_query: Optional[str] = Field(None, description="用户原始输入")
selection_source: Optional[str] = Field(
None,
description="选择来源",
pattern=SELECTION_SOURCE_PATTERN,
)
class Config:
json_schema_extra = {
"example": {
"task_id": "task_abc123",
"status": "completed",
"progress": 100,
"result": None,
"error": None,
"asset_type": "stock",
"stock_name": "贵州茅台",
"original_query": "茅台",
"selection_source": "autocomplete"
}
}
class TaskInfo(BaseModel):
"""
Task details model
Used for task list and SSE event delivery
"""
task_id: str = Field(..., description="任务 ID")
stock_code: str = Field(..., description="股票代码")
stock_name: Optional[str] = Field(None, description="股票名称")
asset_type: str = Field("stock", description="资产类型:stock/futures")
status: TaskStatusEnum = Field(..., description="任务状态")
progress: int = Field(0, description="进度百分比 (0-100)", ge=0, le=100)
message: Optional[str] = Field(None, description="状态消息")
report_type: str = Field("detailed", description="报告类型")
created_at: str = Field(..., description="创建时间")
started_at: Optional[str] = Field(None, description="开始执行时间")
completed_at: Optional[str] = Field(None, description="完成时间")
error: Optional[str] = Field(None, description="错误信息(仅在 failed 时存在)")
original_query: Optional[str] = Field(None, description="用户原始输入")
selection_source: Optional[str] = Field(
None,
description="选择来源",
pattern=SELECTION_SOURCE_PATTERN,
)
class Config:
json_schema_extra = {
"example": {
"task_id": "abc123def456",
"stock_code": "600519",
"stock_name": "贵州茅台",
"asset_type": "stock",
"status": "processing",
"progress": 50,
"message": "正在分析中...",
"report_type": "detailed",
"created_at": "2026-02-05T10:30:00",
"started_at": "2026-02-05T10:30:01",
"completed_at": None,
"error": None,
"original_query": "茅台",
"selection_source": "autocomplete"
}
}
class TaskListResponse(BaseModel):
"""任务列表响应模型"""
total: int = Field(..., description="任务总数")
pending: int = Field(..., description="等待中的任务数")
processing: int = Field(..., description="处理中的任务数")
tasks: List[TaskInfo] = Field(..., description="任务列表")
class Config:
json_schema_extra = {
"example": {
"total": 3,
"pending": 1,
"processing": 2,
"tasks": []
}
}
class DuplicateTaskErrorResponse(BaseModel):
"""重复任务错误响应模型"""
error: str = Field("duplicate_task", description="错误类型")
message: str = Field(..., description="错误信息")
stock_code: str = Field(..., description="股票代码")
existing_task_id: str = Field(..., description="已存在的任务 ID")
class Config:
json_schema_extra = {
"example": {
"error": "duplicate_task",
"message": "股票 600519 正在分析中",
"stock_code": "600519",
"existing_task_id": "abc123def456"
}
}