添加警告相关接口
This commit is contained in:
parent
f634449be8
commit
d1563215d7
@ -3,7 +3,7 @@ from fastapi import APIRouter, Depends, HTTPException, Query, Body
|
||||
from sqlalchemy.ext.asyncio import AsyncSession
|
||||
from app.core.database import get_db
|
||||
from app.crud.event import event
|
||||
from app.schemas.event import EventList, EventDetail, EventUpdate, EventQuery, BackStageEvent, BackStageEventDto, BackStageEventDetail, EditTemperatureDto
|
||||
from app.schemas.event import EventList, EventDetail, EventUpdate, EventQuery, BackStageEvent, BackStageEventDto, BackStageEventDetail, EditTemperatureDto, OcrAlertMessage, OcrAlertMessageDto
|
||||
from app.util.httpResponse import BaseResponse
|
||||
|
||||
import datetime
|
||||
@ -102,4 +102,48 @@ async def delete_event(
|
||||
event_obj = await event.delete_event(db, event_id=event_id)
|
||||
if not event_obj:
|
||||
return BaseResponse(code=404,msg="事件不存在")
|
||||
return BaseResponse(code=200, msg="success", data=EventDetail.model_validate(event_obj))
|
||||
return BaseResponse(code=200, msg="success", data=EventDetail.model_validate(event_obj))
|
||||
|
||||
@router.get("/events/messages", response_model=BaseResponse[List[OcrAlertMessage]])
|
||||
async def get_messages(
|
||||
db: AsyncSession = Depends(get_db)
|
||||
):
|
||||
"""_summary_
|
||||
获取告警消息列表
|
||||
"""
|
||||
message = await event.get_messages(db)
|
||||
return BaseResponse(code=200, msg="success", data=message)
|
||||
|
||||
# 批量处理告警数据
|
||||
@router.post("/events/handleOcrAlerts", response_model=BaseResponse)
|
||||
async def handle_ocr_alerts(
|
||||
db: AsyncSession = Depends(get_db),
|
||||
messageIdList: List[int] = Body(...)
|
||||
):
|
||||
"""_summary_
|
||||
一建处理告警
|
||||
"""
|
||||
|
||||
flag = await event.handle_ocr_alerts(db, messageIdList=messageIdList)
|
||||
if flag:
|
||||
return BaseResponse(code=200, msg="success")
|
||||
return BaseResponse(code=500, msg="fail to update data")
|
||||
|
||||
# 处理单个告警数据
|
||||
@router.post("/events/handleOcrAlert", response_model=BaseResponse)
|
||||
async def handle_ocr_alert(
|
||||
db: AsyncSession = Depends(get_db),
|
||||
ocrAlertMessageDto: OcrAlertMessageDto = Body(...)
|
||||
):
|
||||
"""
|
||||
处理单个警告
|
||||
"""
|
||||
flag = await event.handle_ocr_alert(db, ocrAlertMessageDto=ocrAlertMessageDto)
|
||||
if flag:
|
||||
return BaseResponse(code=200, msg="success")
|
||||
return BaseResponse(code=500, msg="fail to update data")
|
||||
|
||||
|
||||
|
||||
|
||||
|
||||
|
||||
@ -1,12 +1,13 @@
|
||||
from typing import List, Optional, Dict, Any
|
||||
from sqlalchemy import select, and_, or_, update
|
||||
from sqlalchemy import select, and_, or_, update, bindparam
|
||||
from sqlalchemy.ext.asyncio import AsyncSession
|
||||
from sqlalchemy.orm import selectinload
|
||||
from app.crud.base import CRUDBase
|
||||
from app.models.models import Event, Image, Temperature
|
||||
from app.schemas.event import EventUpdate, EventQuery, BackStageEvent, BackStageEventDto, BackStageEventDetail, EditTemperatureDto
|
||||
from app.models.models import Event, Image, Temperature, Message
|
||||
from app.schemas.event import EventUpdate, EventQuery, BackStageEvent, BackStageEventDto, BackStageEventDetail, EditTemperatureDto, OcrAlertMessage, OcrAlertMessageDto
|
||||
|
||||
from pydantic import VERSION as PYDANTIC_VERSION
|
||||
from datetime import datetime
|
||||
|
||||
class CRUDEvent(CRUDBase[Event, EventUpdate, EventUpdate]):
|
||||
async def get_by_id(self, db: AsyncSession, *, event_id: str) -> Optional[Event]:
|
||||
@ -185,5 +186,92 @@ class CRUDEvent(CRUDBase[Event, EventUpdate, EventUpdate]):
|
||||
await db.delete(event)
|
||||
await db.commit()
|
||||
return event
|
||||
|
||||
async def get_messages(
|
||||
self,
|
||||
db: AsyncSession,
|
||||
) -> Optional[List[OcrAlertMessage]]:
|
||||
query_stmt = (
|
||||
select(Message.messageId,Message.eventId ,Message.messageType, Message.eventType, Message.handle, Message.remark, Message.createTime,
|
||||
Event.number,Event.name,
|
||||
Image.imageUrl, Image.localPath,
|
||||
Temperature.temperature)
|
||||
.select_from(Message)
|
||||
.outerjoin(Event, Message.eventId == Event.eventId)
|
||||
.outerjoin(Image, Message.eventId == Image.eventId)
|
||||
.outerjoin(Temperature, Message.eventId == Temperature.eventId)
|
||||
)
|
||||
|
||||
messages = await db.execute(query_stmt)
|
||||
messages = messages.mappings().all()
|
||||
return [OcrAlertMessage(**t) for t in messages]
|
||||
|
||||
async def handle_ocr_alerts(
|
||||
self,
|
||||
db:AsyncSession,
|
||||
messageIdList: List[int]
|
||||
):
|
||||
"""处理OCR告警消息
|
||||
|
||||
Args:
|
||||
db: 数据库会话
|
||||
messageIdList: 要处理的messageId列表
|
||||
|
||||
Returns:
|
||||
bool: 更新是否成功
|
||||
"""
|
||||
try:
|
||||
# 获取所有需要更新的消息ID
|
||||
# 构建更新语句
|
||||
update_stmt = (
|
||||
update(Message)
|
||||
.where(Message.messageId.in_(messageIdList))
|
||||
.values(
|
||||
handle="1",
|
||||
updateTime=datetime.now()
|
||||
)
|
||||
)
|
||||
|
||||
# 准备更新数据
|
||||
|
||||
# 执行更新
|
||||
await db.execute(
|
||||
update_stmt,
|
||||
# update_data,
|
||||
execution_options={"synchronize_session": False}
|
||||
)
|
||||
await db.commit()
|
||||
return True
|
||||
|
||||
except Exception as e:
|
||||
print(f"更新OCR告警消息失败: {str(e)}")
|
||||
await db.rollback()
|
||||
return False
|
||||
|
||||
async def handle_ocr_alert(
|
||||
self,
|
||||
db: AsyncSession,
|
||||
ocrAlertMessageDto: OcrAlertMessageDto
|
||||
):
|
||||
try:
|
||||
update_stmt = (
|
||||
update(Message)
|
||||
.where(Message.messageId == ocrAlertMessageDto.messageId)
|
||||
.values(
|
||||
remark = ocrAlertMessageDto.remark,
|
||||
updateTime=datetime.now()
|
||||
)
|
||||
)
|
||||
await db.execute(update_stmt)
|
||||
await db.commit()
|
||||
return True
|
||||
|
||||
except Exception as e:
|
||||
print(f"处理单个消息失败: {str(e)}")
|
||||
await db.rollback()
|
||||
return False
|
||||
|
||||
|
||||
|
||||
|
||||
event = CRUDEvent(Event)
|
||||
@ -145,6 +145,30 @@ class BackStageEventDetail(BaseModel):
|
||||
status: Optional[str] = None
|
||||
reason: Optional[str] = None
|
||||
remark: Optional[str] = None
|
||||
|
||||
# 获取报警信息列表
|
||||
class OcrAlertMessage(BaseModel):
|
||||
messageId: int
|
||||
eventId: str
|
||||
messageType: Optional[str] = None
|
||||
eventType: Optional[str] = None
|
||||
createTime: Optional[datetime] = None
|
||||
handle: Optional[str] = None
|
||||
remark: Optional[str] = None
|
||||
number: Optional[str] = None
|
||||
name: Optional[str] = None
|
||||
imageUrl : Optional[str] = None
|
||||
localPath: Optional[str] = None
|
||||
temperature : Optional[str] = None
|
||||
|
||||
model_config = ConfigDict(from_attributes=True)
|
||||
|
||||
# 处理消息DTO
|
||||
class OcrAlertMessageDto(BaseModel):
|
||||
messageId: int
|
||||
remark: Optional[str] = None
|
||||
|
||||
|
||||
|
||||
|
||||
|
||||
|
||||
@ -7,6 +7,7 @@ from app.core.database import async_session
|
||||
from app.models.models import Event, Image, Temperature, Message
|
||||
from app.util.kangda import Kangda
|
||||
from app.util.baiduOcr import BadiduOcr
|
||||
from app.util.status import EventType
|
||||
|
||||
|
||||
class EventSyncService:
|
||||
@ -195,7 +196,8 @@ class EventSyncService:
|
||||
if status.value != 0:
|
||||
message = Message(
|
||||
eventId = image.eventId,
|
||||
messageType = status.value
|
||||
messageType = status.value,
|
||||
eventType = EventType.HOT.value
|
||||
)
|
||||
session.add(message)
|
||||
try:
|
||||
|
||||
@ -9,5 +9,9 @@ class TemperatureStatus(Enum):
|
||||
ALDATE = 2 # AL53数据
|
||||
BLURRY = 3 # 图片模糊--读数异常
|
||||
|
||||
@unique
|
||||
class EventType(Enum):
|
||||
HOT = 0 # 温度异常报警
|
||||
|
||||
|
||||
|
||||
@ -8,4 +8,4 @@ from app.services.event_sync_service import run_sync, run_sync_event
|
||||
if __name__ == "__main__":
|
||||
print("启动事件同步服务...")
|
||||
# asyncio.run(run_sync())
|
||||
asyncio.run(run_sync_event("3d685e6927724b9db47e52b544bf6103"))
|
||||
asyncio.run(run_sync_event("95279ebe957f4a2bb919883f9e2fad24"))
|
||||
Loading…
Reference in New Issue
Block a user