From d1563215d7300360d9c29cfb82519c0afa72d1b8 Mon Sep 17 00:00:00 2001 From: haotian <2421912570@qq.com> Date: Tue, 27 May 2025 17:52:45 +0800 Subject: [PATCH] =?UTF-8?q?=E6=B7=BB=E5=8A=A0=E8=AD=A6=E5=91=8A=E7=9B=B8?= =?UTF-8?q?=E5=85=B3=E6=8E=A5=E5=8F=A3?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- app/api/v1/endpoints/events.py | 48 ++++++++++++++- app/crud/event.py | 94 +++++++++++++++++++++++++++++- app/schemas/event.py | 24 ++++++++ app/services/event_sync_service.py | 4 +- app/util/status.py | 4 ++ run_sync.py | 2 +- 6 files changed, 169 insertions(+), 7 deletions(-) diff --git a/app/api/v1/endpoints/events.py b/app/api/v1/endpoints/events.py index afbd651..ad509bf 100644 --- a/app/api/v1/endpoints/events.py +++ b/app/api/v1/endpoints/events.py @@ -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)) \ No newline at end of file + 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") + + + + + diff --git a/app/crud/event.py b/app/crud/event.py index 85e8caa..9e95b75 100644 --- a/app/crud/event.py +++ b/app/crud/event.py @@ -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) \ No newline at end of file diff --git a/app/schemas/event.py b/app/schemas/event.py index fe0475b..fe388fb 100644 --- a/app/schemas/event.py +++ b/app/schemas/event.py @@ -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 + + diff --git a/app/services/event_sync_service.py b/app/services/event_sync_service.py index 91ad5bd..5183809 100644 --- a/app/services/event_sync_service.py +++ b/app/services/event_sync_service.py @@ -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: diff --git a/app/util/status.py b/app/util/status.py index 7860b66..c7e5fc2 100644 --- a/app/util/status.py +++ b/app/util/status.py @@ -9,5 +9,9 @@ class TemperatureStatus(Enum): ALDATE = 2 # AL53数据 BLURRY = 3 # 图片模糊--读数异常 +@unique +class EventType(Enum): + HOT = 0 # 温度异常报警 + \ No newline at end of file diff --git a/run_sync.py b/run_sync.py index 0829e80..a59a305 100644 --- a/run_sync.py +++ b/run_sync.py @@ -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")) \ No newline at end of file + asyncio.run(run_sync_event("95279ebe957f4a2bb919883f9e2fad24")) \ No newline at end of file