forked from Templeton1129/qronos
-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathmain.py
More file actions
1803 lines (1453 loc) · 70.5 KB
/
Copy pathmain.py
File metadata and controls
1803 lines (1453 loc) · 70.5 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
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
645
646
647
648
649
650
651
652
653
654
655
656
657
658
659
660
661
662
663
664
665
666
667
668
669
670
671
672
673
674
675
676
677
678
679
680
681
682
683
684
685
686
687
688
689
690
691
692
693
694
695
696
697
698
699
700
701
702
703
704
705
706
707
708
709
710
711
712
713
714
715
716
717
718
719
720
721
722
723
724
725
726
727
728
729
730
731
732
733
734
735
736
737
738
739
740
741
742
743
744
745
746
747
748
749
750
751
752
753
754
755
756
757
758
759
760
761
762
763
764
765
766
767
768
769
770
771
772
773
774
775
776
777
778
779
780
781
782
783
784
785
786
787
788
789
790
791
792
793
794
795
796
797
798
799
800
801
802
803
804
805
806
807
808
809
810
811
812
813
814
815
816
817
818
819
820
821
822
823
824
825
826
827
828
829
830
831
832
833
834
835
836
837
838
839
840
841
842
843
844
845
846
847
848
849
850
851
852
853
854
855
856
857
858
859
860
861
862
863
864
865
866
867
868
869
870
871
872
873
874
875
876
877
878
879
880
881
882
883
884
885
886
887
888
889
890
891
892
893
894
895
896
897
898
899
900
901
902
903
904
905
906
907
908
909
910
911
912
913
914
915
916
917
918
919
920
921
922
923
924
925
926
927
928
929
930
931
932
933
934
935
936
937
938
939
940
941
942
943
944
945
946
947
948
949
950
951
952
953
954
955
956
957
958
959
960
961
962
963
964
965
966
967
968
969
970
971
972
973
974
975
976
977
978
979
980
981
982
983
984
985
986
987
988
989
990
991
992
993
994
995
996
997
998
999
1000
"""
量化交易框架管理系统 - FastAPI主应用
该模块是量化交易框架管理系统的核心FastAPI应用,提供完整的Web API服务。
主要功能包括:
1. 用户认证管理
- Google Authenticator 2FA登录
- JWT token管理和自动刷新
- 用户会话管理
2. 框架管理
- 基础代码版本获取和下载
- 框架状态监控和管理
- PM2进程管理集成
3. 数据中心配置
- 数据中心参数配置
- 市值数据下载管理
- 实时数据路径配置
4. 账户管理
- 交易账户配置
- 策略绑定和配置
- 账户文件生成
5. 文件管理
- 因子文件上传(时序/截面)
- 仓管策略上传
- 文件列表查询
技术特性:
- FastAPI框架,支持自动API文档生成
- 异步处理和后台任务
- 动态CORS配置
- 统一的响应模型
- 完善的错误处理和日志记录
- SQLite数据库持久化
"""
import json
import shutil
import subprocess
import time
from pathlib import Path
from typing import Optional
from fastapi import (
FastAPI, HTTPException, Request, Response, BackgroundTasks, UploadFile, File
)
from starlette.middleware.cors import CORSMiddleware
from db.db import init_db
from db.db_ops import (
get_framework_status, get_all_framework_status, delete_framework_status, get_finished_data_center_status,
del_user_token, get_user, save_google_secret, get_all_finished_framework_status
)
from model.enum_kit import StatusEnum, UploadFolderEnum
from model.model import (
LoginRequest, ResponseModel, DataCenterCfgModel, BasicCodeOperateModel, AccountModel, FrameworkCfgModel,
ApiKeySecretModel
)
from service.basic_code import (
generate_account_py_file_from_config, extract_variables_from_py,
generate_account_py_file_from_json, process_framework_account_statistics,
migrate_framework_data
)
from service.command import (
get_pm2_list, del_pm2, get_pm2_env
)
from service.xbx_api import XbxAPI, TokenExpiredException
from utils.auth import google_login, AuthMiddleware
from utils.constant import DATA_CENTER_ID, PREFIX, CACHE_CODE_FILE, LOCAL_CODE_FILE, SELECT_COIN_ID
from utils.log_kit import get_logger
from service.log_parser import parse_data_center_logs
# 初始化日志记录器
logger = get_logger()
# 创建FastAPI应用实例
app = FastAPI(
title="交易框架管理系统",
description="提供量化交易框架的完整管理功能,包括用户认证、框架下载、配置管理等",
version="0.0.1"
)
# 配置认证中间件 - 统一处理JWT token验证和刷新
app.add_middleware(AuthMiddleware)
# 配置CORS中间件 - 允许跨域请求
app.add_middleware(
CORSMiddleware,
allow_origins=[], # 动态配置,允许任意origin
allow_credentials=True,
allow_methods=["*"],
allow_headers=["*"],
expose_headers=["X-Refresh-Token"], # 暴露token刷新头
)
@app.get(f"/{PREFIX}/declaration")
def declaration(code: str):
"""
系统声明代码验证接口
验证客户端提供的声明代码是否与系统预设的声明代码一致。
成功验证后会将代码缓存到data/code.txt文件中,供后续接口使用。
用于系统身份验证或特定功能的准入控制。
:param code: 客户端提供的声明代码
:type code: str
:return: 验证结果
:rtype: ResponseModel
Returns:
ResponseModel:
- data: bool - True表示代码匹配,False表示代码不匹配
Process:
1. 从code.txt文件读取系统预设的声明代码
2. 比较客户端代码与系统代码是否一致
3. 验证成功时缓存代码到data/code.txt文件
4. 返回验证结果
"""
logger.info(f"收到系统声明代码验证请求,客户端代码: {code}")
try:
# 读取系统预设的声明代码
with open(LOCAL_CODE_FILE, "r", encoding="utf-8") as f:
local_code = f.read().strip() # 去除可能的换行符
logger.debug(f"系统声明代码: {local_code}")
# 验证代码是否匹配
is_match = local_code == code
if is_match:
# 验证成功,缓存代码到指定文件
CACHE_CODE_FILE.write_text(code, encoding="utf-8")
logger.info(f"系统声明代码验证成功,代码匹配,已缓存到: {CACHE_CODE_FILE}")
logger.info(f"缓存文件路径: {CACHE_CODE_FILE.absolute()}")
else:
logger.warning(f"系统声明代码验证失败,代码不匹配")
logger.debug(f"期望代码内容: {local_code}")
logger.debug(f"实际代码内容: {code}")
return ResponseModel.ok(data=is_match)
except FileNotFoundError:
logger.error("系统声明代码文件不存在: code.txt")
logger.error("请确保项目根目录存在code.txt文件")
return ResponseModel.error(msg="系统配置异常,声明代码文件不存在")
except Exception as e:
logger.error(f"系统声明代码验证过程中发生异常: {e}")
return ResponseModel.error(msg=f"验证过程中发生异常: {str(e)}")
@app.get(f"/{PREFIX}/first")
def first():
"""
检查系统初始化状态和声明代码状态
检查系统是否为首次使用,并对比系统声明代码与缓存代码的一致性。
用于前端判断是否需要显示初始化向导和验证用户权限状态。
:return: 包含系统状态信息的响应
:rtype: ResponseModel
Returns:
ResponseModel:
- data: dict - 包含系统状态信息
- is_first_use: bool - True表示首次使用,False表示已初始化
- is_declaration: bool - True表示声明代码已确认,False表示需要验证
Process:
1. 检查数据库是否有用户记录,判断是否首次使用
2. 读取系统预设声明代码(code.txt)
3. 读取缓存的确认声明代码(data/code.txt)
4. 对比两个代码是否一致
5. 返回系统状态和声明验证状态
"""
logger.info("收到系统初始化状态检查请求")
try:
# 检查是否首次使用
is_first_use = get_user() is None
logger.info(f"首次使用状态检查: {is_first_use} (基于数据库用户记录)")
# 读取系统预设的声明代码
with open(LOCAL_CODE_FILE, "r", encoding="utf-8") as f:
local_code = f.read().strip() # 去除可能的换行符
# 读取缓存的确认声明代码
if CACHE_CODE_FILE.exists():
with open(CACHE_CODE_FILE, "r", encoding="utf-8") as f:
cache_code = f.read().strip() # 去除可能的换行符
logger.debug(f"缓存声明代码读取成功: {cache_code}")
else:
# 缓存文件不存在,说明用户从未成功验证过声明代码
cache_code = ""
logger.info(f"缓存声明代码文件不存在: {CACHE_CODE_FILE},用户尚未验证声明代码")
# 对比声明代码是否一致
is_declaration = local_code == cache_code
logger.info(f"声明代码对比结果: {is_declaration}")
if is_declaration:
logger.info("声明代码验证状态: 已确认 ✓")
else:
logger.warning("声明代码验证状态: 需要验证 ✗")
logger.debug(f"系统代码: {local_code}, 缓存代码: {cache_code}")
result_data = {
"is_first_use": is_first_use,
"is_declaration": is_declaration,
}
logger.info(f"系统状态检查完成: {result_data}")
return ResponseModel.ok(data=result_data)
except Exception as e:
logger.error(f"系统初始化状态检查失败: {e}")
return ResponseModel.error(msg=f"系统状态检查失败: {str(e)}")
@app.post(f"/{PREFIX}/login")
def login(body: LoginRequest, response: Response):
"""
用户登录接口
使用Google Authenticator进行2FA认证登录。
支持首次登录时绑定Google Secret Key。
:param body: 登录请求数据
:type body: LoginRequest
:param response: HTTP响应对象
:type response: Response
:return: 登录结果和JWT token
:rtype: ResponseModel
Process:
1. 验证Google Authenticator代码
2. 生成JWT访问token
3. 保存用户认证信息
4. 添加wx_token到响应头
5. 返回token和用户信息
"""
logger.info(f"用户登录请求,参数: {body}")
try:
# 执行Google登录验证
data = google_login(getattr(body, 'google_secret_key', None), getattr(body, 'code', None))
logger.info("Google认证验证成功")
# 保存Google Secret Key到数据库
success = save_google_secret(body.google_secret_key, data.get('access_token'))
if not success:
logger.warning("Google Secret Key已存在,拒绝重复绑定")
return ResponseModel.error(msg="已经绑定过 secret,请勿重复绑定")
is_bind = False
# 获取用户信息并添加wx_token到响应头
user = get_user()
if user:
# 没有 apikey, uuid, token,需要重新扫码
if not (user.apikey and user.apikey and user.xbx_token):
is_bind = False
else:
try:
api = XbxAPI.get_instance()
api._ensure_token()
is_bind = True
except Exception as e:
is_bind = False
logger.info("用户登录成功,token已生成")
return ResponseModel.ok(data={**data, **{'is_bind': is_bind}})
except Exception as e:
logger.error(f"用户登录失败: {e}")
return ResponseModel.error(msg=f"登录失败: {str(e)}")
@app.post(f"/{PREFIX}/logout")
def logout():
"""
用户登出接口
清除用户的认证token,结束当前会话。
:return: 登出成功响应
:rtype: ResponseModel
"""
logger.info("用户登出请求")
try:
del_user_token()
logger.info("用户登出成功,token已清除")
return ResponseModel.ok(msg="Logged out")
except Exception as e:
logger.error(f"用户登出失败: {e}")
return ResponseModel.error(msg=f"登出失败: {str(e)}")
@app.post(f"/{PREFIX}/user/info")
def user_info(request: Request, background_tasks: BackgroundTasks):
"""
获取用户信息接口
通过XBX授权token获取用户详细信息,并自动触发数据中心下载。
:param request: HTTP请求对象
:type request: Request
:param background_tasks: 后台任务管理器
:type background_tasks: BackgroundTasks
:return: 用户信息数据
:rtype: ResponseModel
Process:
1. 从请求头获取XBX授权token
2. 调用XBX API获取用户信息
3. 设置用户凭据并登录XBX系统
4. 后台任务下载最新数据中心代码
5. 返回用户信息
"""
authorization = request.headers.get("xbx-Authorization", None)
logger.info(f"获取用户信息请求,token前缀: {authorization[:20] if authorization else 'None'}...")
try:
api = XbxAPI.get_instance()
data = api.get_user_info(authorization)
if data:
logger.info(f"成功获取用户信息,UUID: {data.get('uuid')}")
# 设置用户凭据并自动登录
api.set_credentials(data.get("uuid"), data.get("apiKey"))
if not api.login():
logger.error("XBX系统登录失败,uuid或apikey错误")
return ResponseModel.error(code=444, msg="系统认证失败,请重新扫描二维码绑定用户")
logger.info("XBX系统登录成功,启动数据中心下载任务")
# 后台任务:下载最新数据中心代码
background_tasks.add_task(api.download_data_center_latest)
return ResponseModel.ok(data=data)
else:
logger.error("获取用户信息失败:XBX API返回空数据")
return ResponseModel.error(code=444, msg="获取用户信息失败,请重新扫描二维码绑定用户")
except TokenExpiredException as e:
logger.error(f"Token已过期,需要重新认证: {e}")
return ResponseModel.error(code=444, msg="Token已过期,请重新扫描二维码登录")
except Exception as e:
logger.error(f"获取用户信息异常: {e}")
return ResponseModel.error(code=500, msg=f"获取用户信息异常: {str(e)}")
@app.get(f"/{PREFIX}/basic_code/list")
def get_basic_code():
"""
获取基础代码版本列表
从XBX服务器获取所有可用的基础代码框架版本信息。
自动过滤掉数据中心框架,仅返回业务框架。
同时过滤版本列表,只保留时间大于2025-06-01的版本。
:return: 基础代码版本列表
:rtype: ResponseModel
Returns:
ResponseModel:
- data: list - 框架版本信息列表
- 每个框架包含:id, name, versions等信息
- versions中只包含time > "2025-06-01"的版本
"""
logger.info("获取基础代码版本列表")
try:
api = XbxAPI.get_instance()
result = api.get_basic_code_version()
if result.get("error") == "token_invalid":
logger.error("获取基础代码版本失败:XBX token无效")
raise HTTPException(status_code=401, detail="三方token失效,请重新登录")
# 过滤掉数据中心框架
data = result.get('data', [])
filtered_data = [item for item in data if not item.get('id') in [DATA_CENTER_ID, SELECT_COIN_ID]]
# 过滤版本列表,只保留time大于2025-06-01的版本(6月份更新的代码,配合当前框架可以使用)
# 0.2.0版本,更新了账户统计接口,需要限制仓管框架版本必须要 1.3.4 版本, 2025-07-14
time_threshold = "2025-07-14"
for framework in filtered_data:
versions = framework.get('versions', [])
# 过滤版本:只保留time大于threshold的版本
filtered_versions = []
for version in versions:
# 这里直接使用字符串比较
version_time = version.get('time', '')
if version_time > time_threshold:
filtered_versions.append(version)
framework['versions'] = filtered_versions
# 统计过滤后的版本数量
total_versions = sum(len(framework.get('versions', [])) for framework in filtered_data)
logger.info(
f"成功获取基础代码版本列表,共{len(filtered_data)}个框架,{total_versions}个版本(时间>{time_threshold})")
return ResponseModel.ok(data=filtered_data)
except TokenExpiredException as e:
logger.error(f"Token已过期,需要重新认证: {e}")
return ResponseModel.error(code=444, msg="Token已过期,请重新扫描二维码登录")
except Exception as e:
logger.error(f"获取基础代码版本列表异常: {e}")
return ResponseModel.error(msg=f"获取版本列表失败: {str(e)}")
@app.post(f"/{PREFIX}/save_config/data_center")
def save_config_data_center(data_center_cfg: DataCenterCfgModel):
"""
保存数据中心配置
保存数据中心的配置参数,包括API配置、数据源配置等。
如果启用了市值数据,会自动下载历史市值数据。
:param data_center_cfg: 数据中心配置数据
:type data_center_cfg: DataCenterCfgModel
:return: 保存结果
:rtype: ResponseModel
Process:
1. 验证数据中心下载状态
2. 下载市值数据(如果启用)
3. 保存配置到数据库
4. 生成config.json配置文件
"""
logger.info(f"保存数据中心配置请求: {data_center_cfg.id}")
try:
api = XbxAPI.get_instance()
# 设置API凭据信息
data_center_cfg.data_api_key = api.apikey
data_center_cfg.data_api_uuid = api.uuid
data_center_cfg.is_first = False
# 检查数据中心框架状态
framework_status = get_framework_status(data_center_cfg.id)
if not framework_status or framework_status.status != StatusEnum.FINISHED or not framework_status.path:
logger.warning(f"数据中心未下载完成,状态: {framework_status.status if framework_status else 'None'}")
return ResponseModel.ok(msg='数据中心还没有下载完毕')
logger.info(f"数据中心框架路径: {framework_status.path}")
# 下载市值数据(如果启用)
if data_center_cfg.use_api.coin_cap:
logger.info("开始下载市值数据...")
coin_cap_path = (Path(framework_status.path) / 'data' / 'coin_cap')
if api.download_coin_cap_hist(coin_cap_path):
logger.info('市值数据下载成功')
else:
logger.warning('市值数据下载失败')
# 生成配置文件
config_file_path = Path(framework_status.path) / 'config.json'
config_file_path.write_text(
json.dumps(data_center_cfg.model_dump(), ensure_ascii=False, indent=2))
logger.info(f"配置文件已生成: {config_file_path}")
return ResponseModel.ok()
except TokenExpiredException as e:
logger.error(f"Token已过期,需要重新认证: {e}")
return ResponseModel.error(code=444, msg="Token已过期,请重新扫描二维码登录")
except Exception as e:
logger.error(f"保存数据中心配置失败: {e}")
return ResponseModel.error(msg=f"保存配置失败: {str(e)}")
@app.put(f"/{PREFIX}/save_config/data_center")
def update_config_data_center(data_center_cfg: DataCenterCfgModel):
"""
更新数据中心配置
更新已存在的数据中心配置参数。
:param data_center_cfg: 数据中心配置数据
:type data_center_cfg: DataCenterCfgModel
:return: 更新结果
:rtype: ResponseModel
"""
logger.info(f"更新数据中心配置请求: {data_center_cfg.id}")
try:
api = XbxAPI.get_instance()
# 设置API凭据信息
data_center_cfg.data_api_key = api.apikey
data_center_cfg.data_api_uuid = api.uuid
data_center_cfg.is_first = False
# 更新配置文件
framework_status = get_framework_status(data_center_cfg.id)
if framework_status and framework_status.path:
config_file_path = Path(framework_status.path) / 'config.json'
config_file_path.write_text(
json.dumps(data_center_cfg.model_dump(), ensure_ascii=False, indent=2))
logger.info(f"配置文件已更新: {config_file_path}")
# 下载市值数据(如果启用)
if data_center_cfg.use_api.coin_cap:
logger.info("开始下载市值数据...")
coin_cap_path = (Path(framework_status.path) / 'data' / 'coin_cap')
if api.download_coin_cap_hist(coin_cap_path):
logger.info('市值数据下载成功')
else:
logger.warning('市值数据下载失败')
return ResponseModel.ok()
except TokenExpiredException as e:
logger.error(f"Token已过期,需要重新认证: {e}")
return ResponseModel.error(code=444, msg="Token已过期,请重新扫描二维码登录")
except Exception as e:
logger.error(f"更新数据中心配置失败: {e}")
return ResponseModel.error(msg=f"更新配置失败: {str(e)}")
@app.get(f"/{PREFIX}/basic_code/query_config")
def basic_code_query_config(framework_id: str):
"""
查询框架配置
获取指定框架配置信息。
:param framework_id: 框架ID
:type framework_id: str
:return: 配置数据
:rtype: ResponseModel
"""
logger.info(f"查询框架配置: {framework_id}")
try:
# 验证框架下载状态
framework_status = get_framework_status(framework_id)
if not framework_status:
logger.error(f"框架未下载完成: {framework_id}")
return ResponseModel.error(msg=f'框架未下载完成')
config_json_path = Path(framework_status.path) / 'config.json'
if config_json_path.exists():
config_json = json.loads(config_json_path.read_text(encoding='utf-8'))
return ResponseModel.ok(data=config_json)
return ResponseModel.ok()
except Exception as e:
logger.error(f"查询框架配置失败: {e}")
return ResponseModel.error(msg=f"查询框架配置失败: {str(e)}")
@app.get(f"/{PREFIX}/basic_code/download")
def basic_code_download(framework_id: str, background_tasks: BackgroundTasks):
"""
启动框架下载任务
异步下载指定的基础代码框架。
:param framework_id: 要下载的框架ID
:type framework_id: str
:param background_tasks: 后台任务管理器
:type background_tasks: BackgroundTasks
:return: 任务启动结果
:rtype: ResponseModel
"""
logger.info(f"启动框架下载任务: {framework_id}")
try:
api = XbxAPI.get_instance()
background_tasks.add_task(api.download_basic_code_for_id, framework_id)
logger.info(f"框架下载任务已添加到后台队列: {framework_id}")
return ResponseModel.ok()
except TokenExpiredException as e:
logger.error(f"Token已过期,需要重新认证: {e}")
return ResponseModel.error(code=444, msg="Token已过期,请重新扫描二维码登录")
except Exception as e:
logger.error(f"启动框架下载任务失败: {e}")
return ResponseModel.error(msg=f"启动下载任务失败: {str(e)}")
@app.get(f"/{PREFIX}/basic_code/download/status")
def basic_code_download_status():
"""
获取框架下载状态
查询所有框架的下载状态信息。
:return: 框架状态列表
:rtype: ResponseModel
"""
logger.info("查询框架下载状态")
try:
data = get_all_framework_status()
logger.info(f"成功获取框架状态,共{len(data)}个框架")
return ResponseModel.ok(data=data)
except Exception as e:
logger.error(f"获取框架下载状态失败: {e}")
return ResponseModel.error(msg=f"获取状态失败: {str(e)}")
@app.delete(f"/{PREFIX}/basic_code")
def basic_code_delete(framework_id: str):
"""
删除框架
删除指定的框架,包括停止PM2进程、删除文件和数据库记录。
:param framework_id: 要删除的框架ID
:type framework_id: str
:return: 删除结果
:rtype: ResponseModel
Process:
1. 检查框架状态
2. 停止PM2进程
3. 删除数据库记录
4. 删除本地文件
"""
logger.info(f"删除框架请求: {framework_id}")
try:
framework_status = get_framework_status(framework_id)
if not framework_status:
logger.warning(f"框架不存在或未下载完成: {framework_id}")
return ResponseModel.error(msg=f'框架未下载完成')
logger.info(f"开始删除框架,路径: {framework_status.path}")
# 停止并删除PM2进程
del_pm2(framework_id)
logger.info(f"PM2进程已停止: {framework_id}")
# 删除数据库记录
delete_framework_status(framework_id)
logger.info(f"数据库记录已删除: {framework_id}")
# 只删数据库,不删磁盘文件
# 删除本地文件
# if framework_status.path:
# shutil.rmtree(framework_status.path, ignore_errors=True)
# logger.info(f"本地文件已删除: {framework_status.path}")
logger.info(f"框架删除完成: {framework_id}")
return ResponseModel.ok()
except Exception as e:
logger.error(f"删除框架失败: {e}")
return ResponseModel.error(msg=f"删除框架失败: {str(e)}")
# ========== 框架启停/日志 ==========
@app.post(f"/{PREFIX}/basic_code/operate")
def basic_code_operate(operate: BasicCodeOperateModel):
"""
框架操作接口
对框架进行启动、停止、重启或获取日志等操作。
支持PM2进程管理集成。
:param operate: 操作请求数据
:type operate: BasicCodeOperateModel
:return: 操作结果
:rtype: ResponseModel
支持的操作类型:
- start: 启动框架
- stop: 停止框架
- restart: 重启框架
- log: 获取框架日志
"""
logger.info(f"框架操作请求: {operate.framework_id}, 操作类型: {operate.type}")
try:
if operate.type in ["start", "stop", "restart"]:
logger.info(f"执行PM2操作: {operate.type}")
framework_status = get_framework_status(operate.framework_id)
if not framework_status:
logger.error(f"框架未下载完成: {operate.framework_id}")
return ResponseModel.error(msg=f'框架未下载完成')
config_json = Path(framework_status.path) / 'config.json'
if not config_json.exists():
return ResponseModel.error(msg=f'当前框架未导入策略配置,禁止实盘启停操作')
# 检查PM2进程列表
data = get_pm2_list()
if not any([item['framework_id'] == operate.framework_id for item in data]):
logger.info(f"PM2进程不存在,需要先启动: {operate.framework_id}")
# 启动PM2进程
startup_config = Path(framework_status.path) / 'startup.json'
logger.info(f"使用配置文件启动PM2: {startup_config}")
try:
result = subprocess.run(f"pm2 start {startup_config}", env=get_pm2_env(),
shell=True, capture_output=True, text=True)
logger.info(f'PM2启动结果: {result.stdout}')
if result.stderr:
logger.warning(f'PM2启动警告: {result.stderr}')
# 启动后直接保存并返回,不需要再执行额外操作
subprocess.Popen(f"pm2 save -f", env=get_pm2_env(), shell=True)
return ResponseModel.ok(data=f"框架已启动并使用namespace配置")
except Exception as e:
logger.error(f'PM2启动异常: {e}')
return ResponseModel.error(msg=f"PM2启动失败: {str(e)}")
else:
# 执行对namespace的操作(支持PM2 namespace功能)
command = f"pm2 {operate.type} {operate.framework_id}"
logger.info(f"执行PM2命令: {command}")
subprocess.Popen(command, env=get_pm2_env(), shell=True)
logger.info(f"PM2操作已执行: {operate.type}")
subprocess.Popen(f"pm2 save -f", env=get_pm2_env(), shell=True)
return ResponseModel.ok(data=f"{operate.type} 命令已执行")
elif operate.type == "log":
logger.info(f"获取框架日志: {operate.framework_id}, 行数: {operate.lines}")
try:
log_command = f"pm2 logs {operate.framework_id} --lines {operate.lines} --nostream"
result = subprocess.run(log_command, env=get_pm2_env(), shell=True,
capture_output=True, text=True, timeout=30)
logger.info(f"成功获取框架日志,输出长度: {len(result.stdout)}")
return ResponseModel.ok(data=result.stdout)
except subprocess.TimeoutExpired:
logger.error("获取日志超时")
return ResponseModel.error(msg="日志获取超时")
except Exception as e:
logger.error(f"获取日志异常: {e}")
return ResponseModel.error(msg=f"日志获取失败: {e}")
else:
logger.warning(f"不支持的操作类型: {operate.type}")
return ResponseModel.error(msg="不支持的操作类型")
except Exception as e:
logger.error(f"框架操作失败: {e}")
return ResponseModel.error(msg=f"命令执行失败: {e}")
# ========== 框架运行状态 ==========
@app.get(f"/{PREFIX}/basic_code/status")
def basic_code_status():
"""
获取框架运行状态
查询所有框架的PM2进程运行状态。
:return: 框架运行状态列表
:rtype: ResponseModel
Returns:
ResponseModel:
- data: list - PM2进程状态信息列表
- 包含进程ID、状态、CPU、内存等信息
"""
logger.info("查询框架运行状态")
try:
data = get_pm2_list()
logger.info(f"成功获取框架运行状态,共{len(data)}个进程")
return ResponseModel.ok(data=data)
except Exception as e:
logger.error(f"获取框架运行状态失败: {e}")
return ResponseModel.error(msg=f'获取框架运行状态失败, {e}')
@app.get(f"/{PREFIX}/basic_code/cfg/overview")
def basic_code_detail(framework_id: str):
"""
获取框架配置概览
获取指定框架的详细配置信息。
:param framework_id: 框架ID
:type framework_id: str
:return: 框架配置信息
:rtype: ResponseModel
Note:
当前为占位实现,后续可扩展具体配置信息
"""
logger.info(f"获取框架配置概览: {framework_id}")
# TODO: 实现具体的配置概览逻辑
return ResponseModel.ok()
# ========== 上传文件(时序因子/截面因子/仓管策略) ==========
@app.post(f"/{PREFIX}/basic_code/upload/file")
def basic_code_upload_file(framework_id: str, upload_folder: UploadFolderEnum, files: list[UploadFile] = File(...)):
"""
上传文件到框架
上传时序因子、截面因子或仓管策略文件到指定框架。
:param framework_id: 目标框架ID
:type framework_id: str
:param upload_folder: 上传文件夹类型
:type upload_folder: UploadFolderEnum
:param files: 要上传的文件列表
:type files: list[UploadFile]
:return: 上传结果
:rtype: ResponseModel
支持的文件夹类型:
- factors: 时序因子
- sections: 截面因子
- positions: 仓管策略
"""
logger.info(f"文件上传请求: 框架={framework_id}, 文件夹={upload_folder.value}, 文件数={len(files)}")
try:
framework_status = get_framework_status(framework_id)
if not framework_status:
logger.error(f"框架未下载完成: {framework_id}")
return ResponseModel.error(msg=f'框架未下载完成')
target_dir = Path(framework_status.path) / upload_folder.value
logger.info(f"目标上传目录: {target_dir}")
saved_files = []
for file in files:
# 处理子目录情况,提取文件名
filename = file.filename.split('/')[-1]
logger.debug(f"处理文件: {file.filename} -> {filename}")
# 跳过__init__.py文件 和 非py文件
file_path = target_dir / filename
if file_path.stem == '__init__' or file_path.suffix != '.py':
logger.debug(f"跳过__init__.py文件 和 不是 py 的脚本文件: {filename}")
continue
# 确保目录存在
file_path.parent.mkdir(parents=True, exist_ok=True)
# 保存文件
with open(file_path, "wb") as f:
content = file.file.read()
f.write(content)
logger.info(f"文件保存成功: {file_path}")
saved_files.append(file_path.stem)
logger.info(f"文件上传完成,成功保存{len(saved_files)}个文件")
return ResponseModel.ok(data={"saved_files": saved_files})
except Exception as e:
logger.error(f"文件上传失败: {e}")
return ResponseModel.error(msg=f"文件上传失败: {str(e)}")
# ========== 获取框架文件列表(时序因子/截面因子/仓管策略) ==========
@app.get(f"/{PREFIX}/basic_code/file/list")
def basic_code_file_factor(framework_id: str, upload_folder: UploadFolderEnum):
"""
获取框架文件列表
获取指定框架中特定文件夹的Python文件列表。
:param framework_id: 框架ID
:type framework_id: str
:param upload_folder: 文件夹类型
:type upload_folder: UploadFolderEnum
:return: 文件名列表
:rtype: ResponseModel
"""
logger.info(f"获取文件列表: 框架={framework_id}, 文件夹={upload_folder.value}")
try:
framework_status = get_framework_status(framework_id)
if not framework_status:
logger.error(f"框架未下载完成: {framework_id}")
return ResponseModel.error(msg=f'框架未下载完成')
target_dir = Path(framework_status.path) / upload_folder.value
if not target_dir.exists():
logger.warning(f"目标目录不存在: {target_dir}")
return ResponseModel.error(msg=f'{target_dir} 路径不存在')
# 获取Python文件列表(排除__init__.py)
file_names = [file.name for file in target_dir.iterdir()
if file.is_file() and file.suffix == ".py" and file.name != '__init__.py']
logger.info(f"成功获取文件列表,共{len(file_names)}个文件")
return ResponseModel.ok(data=file_names)
except Exception as e:
logger.error(f"获取文件列表失败: {e}")
return ResponseModel.error(msg=f"获取文件列表失败: {str(e)}")
@app.post(f"/{PREFIX}/basic_code/global_config")
def basic_code_global_config(framework_cfg: FrameworkCfgModel):
"""
保存框架全局配置
保存指定框架的全局配置参数,包括数据路径、调试模式、错误通知等。
自动关联数据中心路径,生成框架运行所需的全局配置文件。
:param framework_cfg: 框架全局配置数据
:type framework_cfg: FrameworkCfgModel
:return: 保存结果
:rtype: ResponseModel
Process:
1. 验证框架下载状态
2. 验证数据中心状态
3. 自动配置实时数据路径
4. 生成config.json配置文件
Configuration Fields:
- framework_id: 框架唯一标识
- realtime_data_path: 实时数据存储路径(自动设置)
- is_debug: 是否启用调试模式
- error_webhook_url: 错误通知webhook地址
"""
logger.info(f"保存框架全局配置: 框架={framework_cfg.framework_id}, 调试模式={framework_cfg.is_debug}")
try:
# 验证框架下载状态
framework_status = get_framework_status(framework_cfg.framework_id)
if not framework_status:
logger.error(f"框架未下载完成: {framework_cfg.framework_id}")
return ResponseModel.error(msg=f'框架未下载完成')
if not framework_status.path:
logger.error(f"框架路径为空: {framework_cfg.framework_id}")
return ResponseModel.error(msg=f"磁盘上未存储框架")
# 验证数据中心状态
data_center_status = get_finished_data_center_status()
if not data_center_status:
logger.error("数据中心未下载完成")
return ResponseModel.error(msg="数据中心未下载完成")
if not data_center_status.path:
logger.error("数据中心路径为空")
return ResponseModel.error(msg="数据中心路径异常")
# 自动配置数据中心存储数据路径
framework_cfg.realtime_data_path = str(Path(data_center_status.path) / 'data')
logger.info(f"自动配置实时数据路径: {framework_cfg.realtime_data_path}")
# 保存JSON配置文件
config_json_path = Path(framework_status.path) / 'config.json'
config_data = framework_cfg.model_dump()
config_json_path.write_text(
json.dumps(config_data, ensure_ascii=False, indent=2),
encoding='utf-8'
)
logger.info(f"框架全局配置保存成功")
logger.info(f"配置文件路径: {config_json_path}")
logger.info(f"配置内容: {config_data}")
return ResponseModel.ok(msg="全局配置保存成功")
except Exception as e:
logger.error(f"保存框架全局配置失败: {e}")
return ResponseModel.error(msg=f"保存全局配置失败: {str(e)}")
@app.post(f"/{PREFIX}/basic_code/account")
def basic_code_account(account_cfg: AccountModel):
"""
保存账户配置
保存交易账户的配置信息,包括API密钥、杠杆、黑白名单等。
同时生成对应的Python配置文件。
:param account_cfg: 账户配置数据
:type account_cfg: AccountModel
:return: 保存结果
:rtype: ResponseModel
Process:
1. 验证框架状态
2. 保存JSON配置文件
3. 生成Python配置文件
"""
logger.info(f"保存账户配置: 框架={account_cfg.framework_id}, 账户={account_cfg.account_name}")