标准Task模型
This commit is contained in:
82
main.py
82
main.py
@@ -1,4 +1,5 @@
|
||||
import uvicorn
|
||||
import sys
|
||||
from fastapi import FastAPI, WebSocket, WebSocketDisconnect, Request, Depends, HTTPException
|
||||
from fastapi.responses import HTMLResponse
|
||||
from ObjectManager import ConnectionManager, TaskManager, ServerManager, Task
|
||||
@@ -26,10 +27,23 @@ html = """
|
||||
</form>
|
||||
<ul id='messages'>
|
||||
</ul>
|
||||
<script>
|
||||
var client_id = Date.now()
|
||||
document.querySelector("#ws-id").textContent = client_id;
|
||||
var ws = new WebSocket(`ws://localhost:8000/ws/${client_id}`);
|
||||
<script type="module">
|
||||
var task = await fetch('/tasks', {
|
||||
method: 'POST',
|
||||
headers: {
|
||||
'Content-Type': 'application/json'
|
||||
},
|
||||
body: JSON.stringify({
|
||||
"name": "test",
|
||||
"description": "test",
|
||||
"status": "running",
|
||||
"created_at": "2020-01-01 00:00:00",
|
||||
"updated_at": "2020-01-01 00:00:00"
|
||||
})
|
||||
}).then(response => response.json())
|
||||
console.log(task)
|
||||
document.querySelector("#ws-id").textContent = task.id;
|
||||
var ws = new WebSocket(`ws://localhost:8000/tasks/${task.id}`);
|
||||
ws.onmessage = function(event) {
|
||||
var messages = document.getElementById('messages')
|
||||
var message = document.createElement('li')
|
||||
@@ -37,6 +51,9 @@ html = """
|
||||
message.appendChild(content)
|
||||
messages.appendChild(message)
|
||||
};
|
||||
ws.onclose = function(event) {
|
||||
console.log('Socket is closed. Reconnect will be attempted in 1 second.', event.reason);
|
||||
};
|
||||
function sendMessage(event) {
|
||||
var input = document.getElementById("messageText")
|
||||
ws.send(input.value)
|
||||
@@ -59,26 +76,6 @@ async def get():
|
||||
connection_manager = ConnectionManager()
|
||||
task_manager = TaskManager()
|
||||
|
||||
# 接收客户端的websocket连接
|
||||
@app.websocket("/ws/{client_id}")
|
||||
async def websocket_endpoint(websocket: WebSocket, client_id: int):
|
||||
# TODO: 验证客户端的身份(使用TOKEN)
|
||||
# 获取TOKEN (从数据库中获取用户的信息)
|
||||
# token = request.headers.get('Authorization')
|
||||
# if token is None:
|
||||
# raise HTTPException(status_code=401, detail="Unauthorized")
|
||||
|
||||
await connection_manager.connect(websocket=websocket, client_id=client_id)
|
||||
try:
|
||||
while True:
|
||||
data = await websocket.receive_text()
|
||||
await connection_manager.send_personal_message(f"You wrote: {data}", client_id=client_id)
|
||||
await connection_manager.broadcast(f"Client #{client_id} says: {data}")
|
||||
# TODO: 处理客户端的请求变化(理论上并没有)
|
||||
except WebSocketDisconnect:
|
||||
connection_manager.disconnect(client_id=client_id)
|
||||
await connection_manager.broadcast(f"Client #{client_id} left the chat")
|
||||
|
||||
|
||||
# 通知所有的ws客户端
|
||||
@app.post("/notify")
|
||||
@@ -95,5 +92,38 @@ async def get_tasks():
|
||||
async def create_task(task: Task):
|
||||
return task_manager.add(task)
|
||||
|
||||
# 维护一个任务队列, 任务队列中的任务会被分发给worker节点
|
||||
# 任务状态变化时通知对应的客户端
|
||||
|
||||
'''
|
||||
监听任务进度
|
||||
可能有多个客户端监听同一个任务(向任务的观察者列表中添加websocket连接)
|
||||
可能有多个任务被同一客户端监听(向客户端的观察目标列表中添加任务)
|
||||
应检查目标任务是否存在
|
||||
'''
|
||||
|
||||
|
||||
@app.websocket("/tasks/{task_id}")
|
||||
async def task_endpoint(websocket: WebSocket, task_id: str):
|
||||
await websocket.accept()
|
||||
if not task_manager.has_task(task_id):
|
||||
await websocket.close()
|
||||
print(f"close websocket: {task_id}")
|
||||
return
|
||||
task_manager.add_observer(task_id, websocket)
|
||||
try:
|
||||
while True:
|
||||
data = await websocket.receive_text()
|
||||
print(f"Client #says: {data}")
|
||||
except WebSocketDisconnect:
|
||||
task_manager.remove_observer(task_id, websocket)
|
||||
print(f"close websocket: {task_id}")
|
||||
|
||||
|
||||
'''
|
||||
维护一个任务队列, 任务队列中的任务会被分发给worker节点
|
||||
任务状态变化时通知对应的客户端
|
||||
'''
|
||||
|
||||
# 启动服务
|
||||
if __name__ == '__main__':
|
||||
port = 8000 if len(sys.argv) < 2 else int(sys.argv[1])
|
||||
uvicorn.run(app='main:app', host='0.0.0.0', port=port, reload=True, workers=1)
|
||||
|
Reference in New Issue
Block a user