-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathmain.py
More file actions
372 lines (329 loc) · 12.9 KB
/
Copy pathmain.py
File metadata and controls
372 lines (329 loc) · 12.9 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
import yaml # type: ignore
with open("config.yaml", "r") as f:
config = yaml.safe_load(f)
import time
import numpy as np
import uuid
import datetime
import traceback
import uvicorn
from typing_extensions import Annotated
from fastapi import FastAPI, File, UploadFile, Form, BackgroundTasks
from router_schemas import ( # type: ignore
EmbeddingRequest, EmbeddingResponse,
RerankRequest, RerankResponse,
KnowledgeRequest, KnowledgeResponse,
DocumentRequest, DocumentResponse,
RAGRequest, RAGResponse
)
from rag_api import RAG
from db_api import ( # type: ignore
KnowledgeDocument, KnowledgeDatabase,
Session
)
app = FastAPI()
# 知识库查询API接口,具有重试机制和错误处理,使用SQLite数据库存储知识库元信息,符合RESTful API设计原则
@app.get("/v1/knowledge_base")
def get_knowledge_base(knowledge_id: int, token: str) -> KnowledgeResponse:
start_time = time.time()
try:
for retry_time in range(10):
with Session() as session:
record = session.query(KnowledgeDatabase).filter(KnowledgeDatabase.knowledge_id == knowledge_id).first()
if record is not None:
return KnowledgeResponse( # type: ignore
request_id=str(uuid.uuid4()),
knowledge_id=knowledge_id,
title=str(record.title),
category=str(record.category),
response_code=200,
response_msg="知识库查询成功",
process_status="completed",
processing_time=time.time() - start_time
)
except Exception as e:
# TODO 打印日志
pass
return KnowledgeResponse( # type: ignore
request_id=str(uuid.uuid4()),
knowledge_id=knowledge_id,
category="",
title="",
response_code=404,
response_msg="知识库不存在",
process_status="completed",
processing_time=time.time() - start_time
)
# 知识库删除API接口,具有重试机制和错误处理,使用SQLite数据库进行知识库元信息的删除操作,符合RESTful API设计原则
@app.delete("/v1/knowledge_base")
def delete_knowledge_base(knowledge_id: int, token: str) -> KnowledgeResponse:
start_time = time.time()
try:
for retry_time in range(10):
with Session() as session:
record = session.query(KnowledgeDatabase).filter(KnowledgeDatabase.knowledge_id == knowledge_id).first()
if record is None:
break
session.delete(record)
session.commit()
return KnowledgeResponse( # type: ignore
request_id=str(uuid.uuid4()),
knowledge_id=knowledge_id,
category=str(record.category),
title=str(record.title),
response_code=200,
response_msg="知识库删除成功",
process_status="completed",
processing_time=time.time() - start_time
)
except Exception as e:
# TODO 打印日志
pass
return KnowledgeResponse( # type: ignore
request_id=str(uuid.uuid4()),
knowledge_id=knowledge_id,
category="",
title="",
response_code=404,
response_msg="知识库不存在",
process_status="completed",
processing_time=time.time() - start_time
)
# 知识库创建API接口,具有重试机制和错误处理,使用SQLite数据库进行知识库元信息的创建操作,符合RESTful API设计原则。当用户发起POST请求到/v1/knowledge_base时,会创建一个新的知识库记录并返回新纪录的ID
@app.post("/v1/knowledge_base")
def add_knowledge_base(req: KnowledgeRequest) -> KnowledgeResponse:
start_time = time.time()
try:
for retry_time in range(10):
with Session() as session:
record = KnowledgeDatabase(
title=req.title,
category=req.category,
create_dt=datetime.datetime.now(),
update_dt=datetime.datetime.now(),
)
session.add(record)
session.flush() # Flushes changes to generate primary key if using autoincrement
knowledge_id = record.knowledge_id
session.commit()
return KnowledgeResponse( # type: ignore
request_id=str(uuid.uuid4()),
knowledge_id=knowledge_id,
category=req.category,
title=req.title,
response_code=200,
response_msg="知识库插入成功",
process_status="completed",
processing_time=time.time() - start_time
)
except Exception as e:
print(traceback.format_exc())
# TODO 打印日志
pass
return KnowledgeResponse( # type: ignore
request_id=str(uuid.uuid4()),
knowledge_id=0,
category="",
title="",
response_code=504,
response_msg="知识库插入失败",
process_status="completed",
processing_time=time.time() - start_time
)
# 文档查询API接口
@app.get("/v1/document")
def get_document(document_id: int, token: str) -> DocumentResponse:
start_time = time.time()
try:
for retry_time in range(10):
with Session() as session:
record = session.query(KnowledgeDocument) \
.filter(KnowledgeDocument.document_id == document_id).first()
if record is not None:
return DocumentResponse( # type: ignore
request_id=str(uuid.uuid4()),
document_id=document_id,
category=record.category,
title=record.title,
knowledge_id=record.knowledge_id,
file_type=record.file_type,
response_code=200,
response_msg="文档查询成功",
process_status="completed",
processing_time=time.time() - start_time
)
break
except Exception as e:
print(traceback.format_exc())
# TODO 打印日志
pass
return DocumentResponse( # type: ignore
request_id=str(uuid.uuid4()),
document_id=document_id,
category="",
title="",
knowledge_id=0,
file_type="",
response_code=404,
response_msg="文档不存在",
process_status="completed",
processing_time=time.time() - start_time
)
# 文档删除API接口
@app.delete("/v1/document")
def delete_document(document_id: int, token: str) -> DocumentResponse:
start_time = time.time()
try:
for retry_time in range(10):
with Session() as session:
record = session.query(KnowledgeDocument).filter(KnowledgeDocument.document_id == document_id).first()
if record is None:
break
session.delete(record)
session.commit()
return KnowledgeResponse( # type: ignore
request_id=str(uuid.uuid4()),
document_id=document_id,
knowledge_id=record.knowledge_id,
category=record.category,
title=record.title,
file_type=record.file_type,
response_code=200,
response_msg="文档删除成功",
process_status="completed",
processing_time=time.time() - start_time
)
except Exception as e:
print(traceback.format_exc())
pass
return DocumentResponse( # type: ignore
request_id=str(uuid.uuid4()),
document_id=document_id,
category="",
title="",
knowledge_id=0,
file_type="",
response_code=404,
response_msg="文档不存在",
process_status="completed",
processing_time=time.time() - start_time
)
# 添加文档:判断知识库是否存在、解析上传的文件保存到本地、后台解析上传文件的内容
# BackgroundTasks 后台任务执行
@app.post("/v1/document")
async def add_document(
knowledge_id: int = Annotated[str, Form()],
title: str = Annotated[str, Form()],
category: str = Annotated[str, Form()],
file: UploadFile = Annotated[str, File(...)],
background_tasks: BackgroundTasks = Annotated[BackgroundTasks, Form()]
) -> DocumentResponse:
start_time = time.time()
response_msg = "新增文档失败"
try:
for retry_time in range(10):
# 上传的文档,记录在关系型数据库中, orm 添加记录
with Session() as session:
record = session.query(KnowledgeDatabase).filter(KnowledgeDatabase.knowledge_id == knowledge_id).first()
if record is None:
response_msg = "知识库不存在,请提前创建"
break
record = KnowledgeDocument(
title=title,
category=category,
knowledge_id=knowledge_id,
file_path="",
file_type=file.content_type,
create_dt=datetime.datetime.now(),
update_dt=datetime.datetime.now(),
)
session.add(record)
session.flush() # Flushes changes to generate primary key if using autoincrement
document_id = record.document_id
record
session.commit()
# 存储数据到文件
file_path = f"upload_files/document_id_{document_id}_" + file.filename
with open(file_path, "wb") as buffer:
buffer.write(file.file.read())
record = session.query(KnowledgeDocument).filter(KnowledgeDocument.document_id == document_id).first()
record.file_path = file_path
session.commit()
# 文档内容解析,后台执行,后台提取数据
background_tasks.add_task(
RAG().extract_content, # 后台运行的函数名
knowledge_id=knowledge_id,
document_id=document_id,
title=title,
file_type=file.content_type,
file_path=file_path
)
return DocumentResponse(
request_id=str(uuid.uuid4()),
document_id=document_id,
category=category,
title=title,
knowledge_id=knowledge_id,
file_type=file.content_type,
response_code=200,
response_msg="文档添加成功",
process_status="completed",
processing_time=time.time() - start_time
)
except Exception as e:
print(traceback.format_exc())
pass
return DocumentResponse( # type: ignore
request_id=str(uuid.uuid4()),
document_id=0,
category="",
title="",
knowledge_id=0,
file_type="",
response_code=404,
response_msg=response_msg,
process_status="completed",
processing_time=time.time() - start_time
)
@app.post("/v1/embedding")
async def semantic_embedding(req: EmbeddingRequest) -> EmbeddingResponse:
start_time = time.time()
if not isinstance(req.text, list):
text = [req.text]
else:
text = req.text
vector: np.ndarray = RAG().get_embedding(text)
return EmbeddingResponse(
request_id=str(uuid.uuid4()),
vector=vector.astype(float).tolist(),
response_code=200,
response_msg="ok",
process_status="completed",
processing_time=time.time() - start_time
)
@app.post("/v1/rerank")
async def semantic_rerank(req: RerankRequest) -> RerankResponse:
start_time = time.time()
vector: np.ndarray = RAG().get_rank(req.text_pair)
return RerankResponse(
request_id=str(uuid.uuid4()),
vector=vector.astype(float).tolist(),
response_code=200,
response_msg="ok",
process_status="completed",
processing_time=time.time() - start_time
)
@app.post("/chat")
def chat(req: RAGRequest) -> RAGResponse:
start_time = time.time()
message = RAG().chat_with_rag(req.knowledge_id, req.message)
return RAGResponse(
request_id=str(uuid.uuid4()),
message=message,
response_code=200,
response_msg="ok",
process_status="completed",
processing_time=time.time() - start_time
)
if __name__ == "__main__":
uvicorn.run(app, host="0.0.0.0", port=config["rag"]["port"], workers=1)